anvilsign in

collin/anvil

1//! Web surface for agent sessions: the session list and page, the websocket
2//! that bridges a browser terminal to the container's tmux, and the vendored
3//! terminal emulator assets.
4//!
5//! The wire format on the websocket is ours rather than wterm's built-in
6//! transport, which is a raw byte pass-through with no control channel and so
7//! no way to carry a resize:
8//!
9//! - **binary frames** are raw terminal bytes, in both directions;
10//! - **text frames** are JSON control messages — `{"t":"resize",…}` from the
11//! browser, `{"t":"exit",…}` back.
12
13use anvil_agent::StartError;
14use anvil_core::{
15 App,
16 User,
17 access,
18 agent::{
19 self,
20 kind,
21 status,
22 },
23};
24use axum::{
25 Router,
26 extract::{
27 Path,
28 State,
29 ws::{
30 Message,
31 WebSocket,
32 WebSocketUpgrade,
33 },
34 },
35 http::header,
36 response::{
37 IntoResponse,
38 Redirect,
39 Response,
40 },
41 routing::{
42 get,
43 post,
44 },
45};
46use futures_util::{
47 SinkExt,
48 StreamExt,
49};
50use maud::{
51 Markup,
52 PreEscaped,
53 html,
54};
55use tokio::io::AsyncWriteExt;
56
57use crate::{
58 auth::{
59 Csrf,
60 CsrfForm,
61 CurrentUser,
62 verify_csrf,
63 },
64 ui::{
65 self,
66 forbidden,
67 layout,
68 not_found,
69 server_error,
70 },
71};
72
73pub fn routes(router: Router<App>) -> Router<App> {
74 router
75 .route("/{owner}/{repo}/-/agent", get(list).post(create))
76 .route("/{owner}/{repo}/-/agent/{id}", get(show))
77 .route("/{owner}/{repo}/-/agent/{id}/stop", post(stop))
78 .route("/{owner}/{repo}/-/agent/{id}/ws", get(socket))
79 .route("/-/static/wterm/{*path}", get(wterm_asset))
80}
81
82// --- vendored terminal emulator --------------------------------------------
83
84/// wterm (Apache-2.0), vendored as published ESM and embedded in the binary,
85/// exactly as htmx is. No bundler is involved: the tree has a single bare
86/// specifier (`@wterm/core`), which the page resolves with an import map.
87///
88/// The built-in core's WASM is inlined as base64 inside `wasm-inline.js`, so
89/// there is no separate binary to fetch and no `import.meta.url` resolution to
90/// get wrong.
91const WTERM_ASSETS: &[(&str, &str, &str)] = &[
92 (
93 "core/index.js",
94 include_str!("../assets/wterm/core/index.js"),
95 JS,
96 ),
97 (
98 "core/terminal-core.js",
99 include_str!("../assets/wterm/core/terminal-core.js"),
100 JS,
101 ),
102 (
103 "core/transport.js",
104 include_str!("../assets/wterm/core/transport.js"),
105 JS,
106 ),
107 (
108 "core/wasm-bridge.js",
109 include_str!("../assets/wterm/core/wasm-bridge.js"),
110 JS,
111 ),
112 (
113 "core/wasm-inline.js",
114 include_str!("../assets/wterm/core/wasm-inline.js"),
115 JS,
116 ),
117 (
118 "dom/index.js",
119 include_str!("../assets/wterm/dom/index.js"),
120 JS,
121 ),
122 (
123 "dom/debug.js",
124 include_str!("../assets/wterm/dom/debug.js"),
125 JS,
126 ),
127 (
128 "dom/hyperlink.js",
129 include_str!("../assets/wterm/dom/hyperlink.js"),
130 JS,
131 ),
132 (
133 "dom/input.js",
134 include_str!("../assets/wterm/dom/input.js"),
135 JS,
136 ),
137 (
138 "dom/renderer.js",
139 include_str!("../assets/wterm/dom/renderer.js"),
140 JS,
141 ),
142 (
143 "dom/wterm.js",
144 include_str!("../assets/wterm/dom/wterm.js"),
145 JS,
146 ),
147 (
148 "terminal.css",
149 include_str!("../assets/wterm/terminal.css"),
150 CSS,
151 ),
152];
153
154const JS: &str = "application/javascript; charset=utf-8";
155const CSS: &str = "text/css; charset=utf-8";
156
157/// Serve one vendored wterm file. A `match` over embedded constants rather than
158/// a route each, since it is a dozen files that only ever change together.
159async fn wterm_asset(Path(path): Path<String>) -> Response {
160 let Some((_, body, content_type)) = WTERM_ASSETS.iter().find(|(name, _, _)| *name == path)
161 else {
162 return not_found("no such asset");
163 };
164 (
165 [
166 (header::CONTENT_TYPE, *content_type),
167 // Vendored at a fixed version and only replaced by a redeploy.
168 (header::CACHE_CONTROL, "public, max-age=31536000, immutable"),
169 ],
170 *body,
171 )
172 .into_response()
173}
174
175// --- pages ------------------------------------------------------------------
176
177/// Resolve the repo and require write access: starting a session runs code
178/// against the repository and (from M2) pushes to it, so it is an owner action,
179/// not a reader one.
180async fn resolve_writable(
181 app: &App,
182 viewer: Option<&User>,
183 owner: &str,
184 name: &str,
185) -> Result<anvil_core::Repository, Response> {
186 let (_, repo) = ui::resolve_repo(app, viewer, owner, name).await?;
187 if !access::can_write(&repo, viewer) {
188 return Err(forbidden());
189 }
190 Ok(repo)
191}
192
193/// `GET /{owner}/{repo}/-/agent` — this repository's sessions.
194async fn list(
195 State(app): State<App>,
196 CurrentUser(user): CurrentUser,
197 Path((owner, name)): Path<(String, String)>,
198) -> Result<Markup, Response> {
199 let repo = resolve_writable(&app, user.as_ref(), &owner, &name).await?;
200 let sessions = agent::list_by_repo(&app.db, repo.id, 50)
201 .await
202 .map_err(server_error)?;
203 let enabled = app.config.agent.enabled;
204 let csrf = crate::auth::current_csrf();
205
206 Ok(layout(
207 &format!("{owner}/{name} · agent"),
208 user.as_ref(),
209 html! {
210 h1 { "Agent sessions" }
211 @if !enabled {
212 p.notice {
213 "Agent sessions are disabled. Set "
214 code { "agent.enabled" }
215 " in anvil.toml to turn them on."
216 }
217 } @else {
218 form method="post" action={"/" (owner) "/" (name) "/-/agent"} {
219 (ui::csrf_input(&csrf))
220 label {
221 "Start from branch "
222 input type="text" name="base_ref" value="main" required;
223 }
224 label {
225 "Opening prompt (optional — leave empty to drive it yourself)"
226 textarea name="prompt" rows="3" {}
227 }
228 button type="submit" { "Start session" }
229 }
230 }
231 @if sessions.is_empty() {
232 p { "No sessions yet." }
233 } @else {
234 ul.sessions {
235 @for session in &sessions {
236 li {
237 a href={"/" (owner) "/" (name) "/-/agent/" (session.id)} {
238 "#" (session.id)
239 }
240 " " span.status { (session.status) }
241 " " span.muted { (session.base_ref) " @ "
242 (short_commit(&session.base_commit)) }
243 " " span.muted { (ui::fmt_relative(session.created_at)) }
244 }
245 }
246 }
247 }
248 },
249 ))
250}
251
252fn short_commit(commit: &str) -> &str {
253 &commit[..commit.len().min(12)]
254}
255
256/// `POST /{owner}/{repo}/-/agent` — start a session.
257async fn create(
258 State(app): State<App>,
259 CurrentUser(user): CurrentUser,
260 csrf: Csrf,
261 Path((owner, name)): Path<(String, String)>,
262 axum::Form(form): axum::Form<CreateForm>,
263) -> Result<Response, Response> {
264 verify_csrf(&csrf, &form.csrf)?;
265 let repo = resolve_writable(&app, user.as_ref(), &owner, &name).await?;
266 let Some(user) = user else {
267 return Err(forbidden());
268 };
269
270 let base_ref = if form.base_ref.trim().is_empty() {
271 repo.default_branch.clone()
272 } else {
273 form.base_ref.trim().to_string()
274 };
275 let prompt = form.prompt.trim();
276 let session_kind = if prompt.is_empty() {
277 kind::INTERACTIVE
278 } else {
279 kind::AUTONOMOUS
280 };
281
282 match anvil_agent::start(&app, repo.id, user.id, session_kind, &base_ref, prompt).await {
283 Ok(id) => Ok(Redirect::to(&format!("/{owner}/{name}/-/agent/{id}")).into_response()),
284 // Being at capacity, or switched off, is a normal answer rather than a
285 // fault — "something went wrong" would send someone hunting a bug that
286 // isn't there.
287 Err(e @ (StartError::Disabled | StartError::AtCapacity(_))) => Err(unavailable(&e)),
288 Err(e) => Err(server_error(e)),
289 }
290}
291
292/// A session could not be started for a reason the caller can act on.
293fn unavailable(reason: &StartError) -> Response {
294 (
295 axum::http::StatusCode::SERVICE_UNAVAILABLE,
296 layout(
297 "Cannot start a session",
298 None,
299 html! {
300 h1 { "Cannot start a session" }
301 p.muted { (reason.to_string()) }
302 @if matches!(reason, StartError::AtCapacity(_)) {
303 p { "Stop a running session, or raise " code { "agent.max_concurrent" } "." }
304 }
305 },
306 ),
307 )
308 .into_response()
309}
310
311#[derive(serde::Deserialize)]
312struct CreateForm {
313 #[serde(default)]
314 csrf: String,
315 #[serde(default)]
316 base_ref: String,
317 #[serde(default)]
318 prompt: String,
319}
320
321/// `POST /{owner}/{repo}/-/agent/{id}/stop`
322async fn stop(
323 State(app): State<App>,
324 CurrentUser(user): CurrentUser,
325 csrf: Csrf,
326 Path((owner, name, id)): Path<(String, String, i64)>,
327 axum::Form(form): axum::Form<CsrfForm>,
328) -> Result<Response, Response> {
329 verify_csrf(&csrf, &form.csrf)?;
330 let repo = resolve_writable(&app, user.as_ref(), &owner, &name).await?;
331 let session = session_of(&app, repo.id, id).await?;
332 anvil_agent::stop(&app, session.id).await;
333 Ok(Redirect::to(&format!("/{owner}/{name}/-/agent/{id}")).into_response())
334}
335
336/// Load a session, checking it really belongs to this repository — otherwise
337/// the id alone would read across repos.
338async fn session_of(
339 app: &App,
340 repo_id: i64,
341 id: i64,
342) -> Result<anvil_core::AgentSession, Response> {
343 let session = agent::get(&app.db, id)
344 .await
345 .map_err(server_error)?
346 .ok_or_else(|| not_found("no such session"))?;
347 if session.repo_id != repo_id {
348 return Err(not_found("no such session"));
349 }
350 Ok(session)
351}
352
353/// `GET /{owner}/{repo}/-/agent/{id}` — the terminal.
354async fn show(
355 State(app): State<App>,
356 CurrentUser(user): CurrentUser,
357 Path((owner, name, id)): Path<(String, String, i64)>,
358) -> Result<Markup, Response> {
359 let repo = resolve_writable(&app, user.as_ref(), &owner, &name).await?;
360 let session = session_of(&app, repo.id, id).await?;
361 let live = status::is_live(&session.status);
362 let csrf = crate::auth::current_csrf();
363 let ws_path = format!("/{owner}/{name}/-/agent/{id}/ws");
364
365 Ok(layout(
366 &format!("{owner}/{name} · agent #{id}"),
367 user.as_ref(),
368 html! {
369 h1 { "Agent session #" (id) }
370 p.muted {
371 (session.status) " · " (session.base_ref) " @ "
372 (short_commit(&session.base_commit))
373 @if !session.error.is_empty() { " · " (session.error) }
374 }
375 @if live {
376 form method="post" action={(ws_path.trim_end_matches("/ws")) "/stop"} {
377 (ui::csrf_input(&csrf))
378 button type="submit" { "Stop session" }
379 }
380 link rel="stylesheet" href="/-/static/wterm/terminal.css";
381 div #terminal style="height:70vh" {}
382 script type="importmap" {
383 (PreEscaped(IMPORT_MAP))
384 }
385 script type="module" {
386 (PreEscaped(format!("const WS_PATH = {};\n{}",
387 serde_json::to_string(&ws_path).unwrap_or_else(|_| "\"\"".into()),
388 TERMINAL_JS)))
389 }
390 } @else {
391 p { "This session has ended." }
392 }
393 },
394 ))
395}
396
397/// Resolves wterm's single bare specifier. Everything else in the vendored
398/// tree imports by relative path.
399const IMPORT_MAP: &str = r#"{"imports":{"@wterm/core":"/-/static/wterm/core/index.js"}}"#;
400
401/// Browser half of the terminal.
402///
403/// Note the `htmx:beforeSwap` teardown: the layout sets `hx-boost` on `<body>`,
404/// so navigating away swaps the DOM without a page load. Without this the
405/// websocket and the WASM instance would leak on every navigation.
406const TERMINAL_JS: &str = r#"
407import { WTerm } from "/-/static/wterm/dom/index.js";
408
409const url = new URL(WS_PATH, location.href);
410url.protocol = location.protocol === "https:" ? "wss:" : "ws:";
411const ws = new WebSocket(url);
412ws.binaryType = "arraybuffer";
413
414const encode = new TextEncoder();
415let torndown = false;
416
417const sendResize = () => {
418 if (ws.readyState === WebSocket.OPEN) {
419 ws.send(JSON.stringify({ t: "resize", cols: term.cols, rows: term.rows }));
420 }
421};
422
423// Callbacks are read off the instance at call time, so setting them here (in
424// the constructor options) is the same as setting them later — this is just
425// the clearer place. `wasmUrl` is left unset on purpose: the built-in core's
426// WASM is inlined as base64, so there is no second request to make.
427const term = new WTerm(document.getElementById("terminal"), {
428 autoResize: true,
429 cursorBlink: true,
430 // Keystrokes out. onData hands us a string; the wire carries bytes.
431 onData: (data) => {
432 if (ws.readyState === WebSocket.OPEN) ws.send(encode.encode(data));
433 },
434 // Fires on the ResizeObserver as well as an explicit resize(), so the
435 // container's pty follows the browser window.
436 onResize: sendResize,
437});
438await term.init();
439
440ws.onopen = sendResize;
441
442ws.onmessage = (event) => {
443 // Binary is terminal bytes; text is our control channel.
444 if (event.data instanceof ArrayBuffer) {
445 term.write(new Uint8Array(event.data));
446 return;
447 }
448 let msg;
449 try {
450 msg = JSON.parse(event.data);
451 } catch {
452 return;
453 }
454 if (msg.t === "exit") {
455 term.write("\r\n\x1b[2m[session ended, status " + msg.code + "]\x1b[0m\r\n");
456 } else if (msg.t === "error") {
457 term.write("\r\n\x1b[31m[" + msg.message + "]\x1b[0m\r\n");
458 }
459};
460
461ws.onclose = () => {
462 if (!torndown) term.write("\r\n\x1b[2m[disconnected]\x1b[0m\r\n");
463};
464
465// The layout sets hx-boost on <body>, so navigating away swaps the DOM without
466// a page load. Without this teardown the socket and the WASM instance would
467// leak on every navigation.
468const teardown = () => {
469 if (torndown) return;
470 torndown = true;
471 try { ws.close(); } catch {}
472 try { term.destroy(); } catch {}
473};
474document.body.addEventListener("htmx:beforeSwap", teardown, { once: true });
475window.addEventListener("pagehide", teardown, { once: true });
476
477term.focus();
478"#;
479
480// --- the websocket ----------------------------------------------------------
481
482/// `GET /{owner}/{repo}/-/agent/{id}/ws` — bridge the browser to the
483/// container's tmux.
484async fn socket(
485 State(app): State<App>,
486 CurrentUser(user): CurrentUser,
487 Path((owner, name, id)): Path<(String, String, i64)>,
488 upgrade: WebSocketUpgrade,
489) -> Result<Response, Response> {
490 let repo = resolve_writable(&app, user.as_ref(), &owner, &name).await?;
491 let session = session_of(&app, repo.id, id).await?;
492 if !status::is_live(&session.status) {
493 return Err(not_found("session is not running"));
494 }
495 Ok(upgrade.on_upgrade(move |socket| bridge(app, session.id, socket)))
496}
497
498/// Pump bytes between one websocket and one tmux client.
499///
500/// Each browser gets its own `docker exec … tmux attach`, so a dropped socket
501/// takes down that client and nothing else — the agent is a process inside
502/// tmux, not a child of the exec. tmux repaints a newly attached client, which
503/// is what makes reconnect show the current screen rather than a blank one.
504async fn bridge(app: App, session_id: i64, socket: WebSocket) {
505 let (mut sink, mut stream) = socket.split();
506
507 let (terminal, docker) = match anvil_agent::attach_stream(&app, session_id, 80, 24).await {
508 Ok(attached) => attached,
509 Err(e) => {
510 let _ = sink
511 .send(Message::Text(
512 format!(r#"{{"t":"error","message":{}}}"#, json_string(&e)).into(),
513 ))
514 .await;
515 return;
516 }
517 };
518 let anvil_agent::container::Terminal {
519 exec_id,
520 mut output,
521 mut input,
522 } = terminal;
523
524 // Container -> browser.
525 let mut to_browser = tokio::spawn(async move {
526 while let Some(chunk) = output.next().await {
527 let Ok(log) = chunk else { break };
528 let bytes = log.into_bytes();
529 if bytes.is_empty() {
530 continue;
531 }
532 if sink
533 .send(Message::Binary(bytes.to_vec().into()))
534 .await
535 .is_err()
536 {
537 break;
538 }
539 }
540 let _ = sink.close().await;
541 });
542
543 // Browser -> container, plus the control channel.
544 let control_docker = docker.clone();
545 let control_exec = exec_id.clone();
546 let mut to_container = tokio::spawn(async move {
547 while let Some(message) = stream.next().await {
548 match message {
549 Ok(Message::Binary(data)) => {
550 if input.write_all(&data).await.is_err() {
551 break;
552 }
553 let _ = input.flush().await;
554 }
555 Ok(Message::Text(text)) => {
556 if let Some((cols, rows)) = parse_resize(&text) {
557 anvil_agent::container::resize(&control_docker, &control_exec, cols, rows)
558 .await;
559 }
560 }
561 Ok(Message::Close(_)) | Err(_) => break,
562 _ => {}
563 }
564 }
565 // Detach this tmux client explicitly.
566 //
567 // Docker has no "kill exec" call, and dropping bollard's streams does
568 // not reliably tear the exec down — without this, every page view left
569 // a tmux client attached forever, counting against the container's pids
570 // limit and taking part in tmux's window sizing.
571 //
572 // Writing the prefix (C-b) followed by `d` into the exec's own stdin
573 // detaches precisely the client on the other end of it, with no need to
574 // work out which of several clients is ours. `tmux.conf` keeps the
575 // default prefix, so this stays in step with it.
576 let _ = input.write_all(b"\x02d").await;
577 let _ = input.flush().await;
578 });
579
580 // Either direction ending means this viewer is gone. Abort BOTH halves
581 // rather than letting the survivor linger: each holds one end of bollard's
582 // hijacked connection, and while either is alive the `docker exec` — and so
583 // the tmux client behind it — stays up. Leaving them accumulated a stale
584 // client per page view, which counts against the container's pids limit and
585 // takes part in tmux's window sizing.
586 tokio::select! {
587 _ = &mut to_browser => {}
588 _ = &mut to_container => {}
589 }
590 to_browser.abort();
591 to_container.abort();
592 anvil_agent::detached(&app, session_id);
593}
594
595/// Parse `{"t":"resize","cols":N,"rows":M}`.
596fn parse_resize(text: &str) -> Option<(u16, u16)> {
597 let value: serde_json::Value = serde_json::from_str(text).ok()?;
598 if value.get("t")?.as_str()? != "resize" {
599 return None;
600 }
601 let cols = value.get("cols")?.as_u64()?.try_into().ok()?;
602 let rows = value.get("rows")?.as_u64()?.try_into().ok()?;
603 Some((cols, rows))
604}
605
606/// JSON-encode a string, for the small hand-built control messages.
607fn json_string(s: &str) -> String {
608 serde_json::to_string(s).unwrap_or_else(|_| "\"\"".to_string())
609}