1use std::sync::{Arc, Mutex};
39
40use axum::Router;
41use axum::body::Bytes;
42use axum::extract::{Path, Request, State};
43use axum::http::{HeaderMap, HeaderValue, StatusCode, header};
44use axum::middleware::Next;
45use axum::response::{IntoResponse, Json, Response};
46use axum::routing::{MethodFilter, MethodRouter, any, get, on};
47use serde::{Deserialize, Serialize};
48use serde_json::{Value, json};
49use tower::ServiceExt as _;
50
51use super::openapi::Method;
52use super::rebinding::{Rebinding, refuse_rebinding};
53use super::{Api, ApiError, Authenticator, Session};
54use crate::core::{
55 ACTION_ADMIT, ACTION_PERFORM, ACTION_RELEASE, Digest, PolicyBundleIdentity, PolicyDecision,
56 PolicyEngine, PolicyRequest,
57};
58use crate::runtime::Runtime;
59
60pub mod action {
62 pub const RUN: &str = "api:dev.run";
64 pub const REPLAY: &str = "api:dev.replay";
66 pub const EXPORT: &str = "api:dev.export";
68 pub const MANIFEST: &str = "api:dev.manifest";
70 pub const HISTORY: &str = "api:dev.history";
72 pub const RUNS: &str = "api:dev.runs";
74
75 pub const ALL: &[&str] = &[RUN, REPLAY, EXPORT, MANIFEST, HISTORY, RUNS];
77}
78
79pub const TENANT: &str = "dev";
81
82pub const CONTENT_SECURITY_POLICY: &str = "default-src 'none'; script-src 'self'; \
84 style-src 'self'; connect-src 'self'; img-src 'self'; base-uri 'none'; \
85 form-action 'none'; frame-ancestors 'none'; require-trusted-types-for 'script'; \
86 trusted-types 'none'";
87
88pub const SHELL: &[(&str, &str, &str)] = &[
90 (
91 "/",
92 "text/html; charset=utf-8",
93 include_str!("dev/index.html"),
94 ),
95 (
96 "/assets/app.js",
97 "text/javascript; charset=utf-8",
98 include_str!("dev/app.js"),
99 ),
100 (
101 "/assets/app.css",
102 "text/css; charset=utf-8",
103 include_str!("dev/app.css"),
104 ),
105];
106
107#[async_trait::async_trait]
113pub trait Workbench: Send + Sync + 'static {
114 async fn plane(&self) -> Arc<Runtime>;
116
117 async fn declaration(&self) -> Declaration;
120
121 async fn start(&self, request: StartRequest) -> Result<Started, String>;
127
128 fn streams(&self) -> Arc<StreamHub>;
130
131 async fn replay(&self, run: Option<crate::core::RunId>) -> Result<Vec<Replayed>, String>;
138}
139
140#[derive(Debug)]
146pub struct StreamHub {
147 sender: tokio::sync::broadcast::Sender<(String, Value)>,
148 closed: tokio::sync::watch::Sender<bool>,
151}
152
153impl StreamHub {
154 const BACKLOG: usize = 1024;
156
157 #[must_use]
158 pub fn new() -> Arc<Self> {
159 Arc::new(Self {
160 sender: tokio::sync::broadcast::channel(Self::BACKLOG).0,
161 closed: tokio::sync::watch::channel(false).0,
162 })
163 }
164
165 pub fn close(&self) {
167 self.closed.send_replace(true);
168 }
169}
170
171impl crate::runtime::RunStreamObserver for StreamHub {
172 fn event(
173 &self,
174 run: crate::core::RunId,
175 event: crate::core::Tainted<crate::model::ModelStreamEvent>,
176 ) {
177 let trust = event.label().trust;
178 let line = match event.into_unlabelled() {
179 crate::model::ModelStreamEvent::TextDelta(text) => {
180 json!({ "type": "text_delta", "value": crate::core::visible::escape(&text).0, "trust": trust })
181 }
182 crate::model::ModelStreamEvent::Usage(usage) => {
183 json!({ "type": "usage", "value": usage, "trust": trust })
184 }
185 };
186 let _ = self.sender.send((run.to_string(), line));
188 }
189}
190
191#[derive(Debug, Clone, Default, Serialize, Deserialize)]
193pub struct Declaration {
194 pub file: String,
196 pub agents: Vec<DeclaredAgent>,
198 pub refused: Option<String>,
201 pub live: Vec<String>,
203}
204
205#[derive(Debug, Clone, Default, Serialize, Deserialize)]
207pub struct DeclaredAgent {
208 pub name: String,
209 pub version: String,
210 pub digest: String,
211 pub bound: Vec<String>,
213 #[serde(default)]
215 pub provides: Vec<String>,
216 #[serde(default)]
218 pub input_schema: Option<Value>,
219}
220
221#[derive(Debug, Clone, Deserialize)]
223#[serde(deny_unknown_fields)]
224pub struct StartRequest {
225 pub input: Value,
227 #[serde(default)]
229 pub correlate: Vec<String>,
230 #[serde(default)]
232 pub capability: Option<String>,
233}
234
235#[derive(Debug, Clone, Serialize, Deserialize)]
237pub struct Started {
238 pub run: String,
239 pub status: String,
240}
241
242#[derive(Debug, Clone, Serialize, Deserialize)]
244pub struct Replayed {
245 pub run: String,
246 pub verdict: String,
248 pub detail: String,
250}
251
252#[derive(Debug, Clone, Default, Deserialize)]
254#[serde(deny_unknown_fields)]
255pub struct ReplayRequest {
256 #[serde(default)]
258 pub run: Option<String>,
259}
260
261#[derive(Debug, Clone, Copy)]
263pub struct DevRoute {
264 pub method: Method,
265 pub path: &'static str,
266 pub action: &'static str,
268 serve: fn(MethodFilter) -> MethodRouter<Arc<Surface>>,
269}
270
271pub const ROUTES: &[DevRoute] = &[
273 DevRoute {
274 method: Method::Get,
275 path: "/dev/manifest",
276 action: action::MANIFEST,
277 serve: |m| on(m, declaration),
278 },
279 DevRoute {
280 method: Method::Post,
281 path: "/dev/runs",
282 action: action::RUN,
283 serve: |m| on(m, start),
284 },
285 DevRoute {
286 method: Method::Get,
287 path: "/dev/runs",
288 action: action::RUNS,
289 serve: |m| on(m, runs),
290 },
291 DevRoute {
292 method: Method::Post,
293 path: "/dev/replay",
294 action: action::REPLAY,
295 serve: |m| on(m, replay),
296 },
297 DevRoute {
298 method: Method::Get,
299 path: "/dev/export",
300 action: action::EXPORT,
301 serve: |m| on(m, export),
302 },
303 DevRoute {
304 method: Method::Get,
305 path: "/dev/runs/{run}",
306 action: action::HISTORY,
307 serve: |m| on(m, run_view),
308 },
309 DevRoute {
310 method: Method::Get,
311 path: "/dev/stream",
312 action: action::RUNS,
313 serve: |m| on(m, stream),
314 },
315 DevRoute {
316 method: Method::Get,
317 path: "/dev/runs/{run}/history",
318 action: action::HISTORY,
319 serve: |m| on(m, history),
320 },
321];
322
323const EXPORT_LIMIT: usize = 10_000;
325
326#[derive(Debug, Clone)]
333pub struct DevPolicy {
334 actor: String,
335}
336
337impl DevPolicy {
338 #[must_use]
340 pub fn new(actor: impl Into<String>) -> Self {
341 Self {
342 actor: actor.into(),
343 }
344 }
345}
346
347impl PolicyEngine for DevPolicy {
348 fn authorize(&self, request: &PolicyRequest<'_>) -> PolicyDecision {
349 if request.action.starts_with("api:") {
350 let tenant = request.context.get("tenant").and_then(Value::as_str);
351 if request.principal == self.actor && tenant == Some(TENANT) {
352 return PolicyDecision::Permit;
353 }
354 return PolicyDecision::deny("only this dev session's own token reaches its plane");
355 }
356 match request.action {
357 ACTION_ADMIT | ACTION_PERFORM => PolicyDecision::Permit,
358 ACTION_RELEASE => PolicyDecision::deny(
359 "`agentplane dev` permits no release, as `agentplane run` permits none: a \
360 release lowers a label only when a rule permits `data:release`",
361 ),
362 _ => PolicyDecision::deny("a dev plane permits admission and effects only"),
363 }
364 }
365
366 fn bundle(&self) -> PolicyBundleIdentity {
367 PolicyBundleIdentity::new(
368 Digest::of(b"agentplane dev policy"),
369 "agentplane/dev-policy-v1",
370 )
371 }
372}
373
374struct Surface {
376 bench: Arc<dyn Workbench>,
377 auth: Arc<dyn Authenticator>,
378 operator: Mutex<Option<(Arc<Runtime>, Api)>>,
380}
381
382impl std::fmt::Debug for Surface {
383 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
384 f.debug_struct("Surface").finish_non_exhaustive()
385 }
386}
387
388impl Surface {
389 async fn api(&self) -> Result<Api, ApiError> {
391 let plane = self.bench.plane().await;
392 let mut held = self
393 .operator
394 .lock()
395 .unwrap_or_else(std::sync::PoisonError::into_inner);
396 if let Some((built_for, api)) = held.as_ref()
397 && Arc::ptr_eq(built_for, &plane)
398 {
399 return Ok(api.clone());
400 }
401 let api = Api::new(Arc::clone(&plane), Arc::clone(&self.auth))
402 .map_err(|e| ApiError(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
403 *held = Some((plane, api.clone()));
404 Ok(api)
405 }
406
407 async fn gate(
410 &self,
411 headers: &HeaderMap,
412 action: &str,
413 resource: &str,
414 ) -> Result<Session, ApiError> {
415 let api = self.api().await?;
416 let caller = self.auth.authenticate(headers).await?;
417 api.authorize(caller, action, resource)
418 }
419}
420
421pub fn router(bench: Arc<dyn Workbench>, auth: Arc<dyn Authenticator>, port: u16) -> Router {
423 let surface = Arc::new(Surface {
424 bench,
425 auth,
426 operator: Mutex::new(None),
427 });
428 let authorities = [
429 format!("127.0.0.1:{port}"),
430 format!("localhost:{port}"),
431 format!("[::1]:{port}"),
432 ];
433 let rebinding = Rebinding::new(
434 authorities.to_vec(),
435 authorities.iter().map(|a| format!("http://{a}")).collect(),
436 )
437 .origin_on_writes();
438
439 let mut routes = Router::new();
440 for route in ROUTES {
441 let filter = match route.method {
442 Method::Get => MethodFilter::GET,
443 Method::Post => MethodFilter::POST,
444 };
445 routes = routes.route(route.path, (route.serve)(filter));
446 }
447 for (path, content_type, body) in SHELL {
448 routes = routes.route(
449 path,
450 get(move || async move { ([(header::CONTENT_TYPE, *content_type)], *body) }),
451 );
452 }
453 routes
454 .route("/api/{*rest}", any(operator))
455 .fallback(not_found)
456 .with_state(surface)
457 .layer(axum::middleware::from_fn_with_state(
458 rebinding,
459 refuse_rebinding,
460 ))
461 .layer(axum::middleware::from_fn(security_headers))
462}
463
464async fn security_headers(request: Request, next: Next) -> Response {
466 let mut response = next.run(request).await;
467 let headers = response.headers_mut();
468 headers.insert(
469 header::CONTENT_SECURITY_POLICY,
470 HeaderValue::from_static(CONTENT_SECURITY_POLICY),
471 );
472 headers.insert(
473 header::X_CONTENT_TYPE_OPTIONS,
474 HeaderValue::from_static("nosniff"),
475 );
476 headers.insert(
477 header::REFERRER_POLICY,
478 HeaderValue::from_static("no-referrer"),
479 );
480 headers.insert(header::CACHE_CONTROL, HeaderValue::from_static("no-store"));
481 response
482}
483
484async fn not_found() -> Response {
485 ApiError(StatusCode::NOT_FOUND, "no such page".to_owned()).into_response()
486}
487
488async fn operator(State(surface): State<Arc<Surface>>, request: Request) -> Response {
490 let api = match surface.api().await {
491 Ok(api) => api,
492 Err(refused) => return refused.into_response(),
493 };
494 let (parts, body) = request.into_parts();
498 let path = parts
499 .uri
500 .path_and_query()
501 .map_or("/", axum::http::uri::PathAndQuery::as_str);
502 let Ok(uri) = path
503 .strip_prefix("/api")
504 .unwrap_or(path)
505 .parse::<axum::http::Uri>()
506 else {
507 return not_found().await;
508 };
509 let mut inner = Request::new(body);
510 *inner.method_mut() = parts.method;
511 *inner.uri_mut() = uri;
512 *inner.version_mut() = parts.version;
513 *inner.headers_mut() = parts.headers;
514 match api.router().oneshot(inner).await {
515 Ok(response) => response,
516 Err(never) => match never {},
517 }
518}
519
520fn unreadable(e: &serde_json::Error) -> ApiError {
522 ApiError(StatusCode::UNPROCESSABLE_ENTITY, e.to_string())
523}
524
525async fn declaration(
526 State(surface): State<Arc<Surface>>,
527 headers: HeaderMap,
528) -> Result<Json<Declaration>, ApiError> {
529 surface.gate(&headers, action::MANIFEST, "manifest").await?;
530 Ok(Json(surface.bench.declaration().await))
531}
532
533async fn start(
534 State(surface): State<Arc<Surface>>,
535 headers: HeaderMap,
536 body: Bytes,
537) -> Result<Json<Started>, ApiError> {
538 surface.gate(&headers, action::RUN, "run").await?;
539 let request: StartRequest = serde_json::from_slice(&body).map_err(|e| unreadable(&e))?;
540 surface
541 .bench
542 .start(request)
543 .await
544 .map(Json)
545 .map_err(|why| ApiError(StatusCode::UNPROCESSABLE_ENTITY, why))
546}
547
548async fn replay(
549 State(surface): State<Arc<Surface>>,
550 headers: HeaderMap,
551 body: Bytes,
552) -> Result<Json<Vec<Replayed>>, ApiError> {
553 surface.gate(&headers, action::REPLAY, "run").await?;
554 let request: ReplayRequest = if body.is_empty() {
555 ReplayRequest::default()
556 } else {
557 serde_json::from_slice(&body).map_err(|e| unreadable(&e))?
558 };
559 let run = request
560 .run
561 .as_deref()
562 .map(crate::core::RunId::parse)
563 .transpose()
564 .map_err(|_| super::bad("run"))?;
565 surface
566 .bench
567 .replay(run)
568 .await
569 .map(Json)
570 .map_err(|why| ApiError(StatusCode::UNPROCESSABLE_ENTITY, why))
571}
572
573async fn runs(
576 State(surface): State<Arc<Surface>>,
577 headers: HeaderMap,
578) -> Result<Json<Value>, ApiError> {
579 let s = surface.gate(&headers, action::RUNS, "store").await?;
580 let outcomes: Vec<String> = crate::runtime::OUTCOMES_OF_RECORD
581 .iter()
582 .map(|o| (*o).to_owned())
583 .collect();
584 let found = crate::export::runs_to_read(s.plane.journal(), &outcomes, true, EXPORT_LIMIT)
585 .await
586 .map_err(|_| super::store_failed())?;
587 let runs: Vec<String> = found.runs.iter().map(ToString::to_string).collect();
588 Ok(Json(
589 json!({ "runs": runs, "partial": !found.reached.is_empty() }),
590 ))
591}
592
593async fn export(
596 State(surface): State<Arc<Surface>>,
597 headers: HeaderMap,
598) -> Result<Json<Value>, ApiError> {
599 let s = surface.gate(&headers, action::EXPORT, "store").await?;
600 let journal = Arc::clone(s.plane.journal());
601 let cases = s
602 .plane
603 .cases()
604 .cloned()
605 .ok_or_else(|| super::unavailable("case store"))?;
606 let outcomes: Vec<String> = crate::runtime::OUTCOMES_OF_RECORD
607 .iter()
608 .map(|o| (*o).to_owned())
609 .collect();
610 let found = crate::export::runs_to_read(&journal, &outcomes, true, EXPORT_LIMIT)
611 .await
612 .map_err(|_| super::store_failed())?;
613 let mut bytes = Vec::new();
614 crate::export::to_jsonl(&journal, &cases, &found.runs, &mut bytes)
615 .await
616 .map_err(|_| super::store_failed())?;
617 let report = crate::export::verify(&bytes[..], None, &[]).map_err(|_| super::store_failed())?;
618 Ok(Json(json!({
619 "export": String::from_utf8_lossy(&bytes),
620 "partial": !found.reached.is_empty(),
621 "report": report,
622 "verify": [
623 "agentplane verify export.jsonl",
624 "python3 tools/verify_export.py export.jsonl",
625 ],
626 })))
627}
628
629async fn history(
634 State(surface): State<Arc<Surface>>,
635 Path(run): Path<String>,
636 request: Request,
637) -> Response {
638 let (parts, _) = request.into_parts();
639 let query = parts
640 .uri
641 .query()
642 .map_or_else(String::new, |q| format!("?{q}"));
643 escaped_read(
644 &surface,
645 &parts.headers,
646 &run,
647 &format!("/history{query}"),
648 "records",
649 )
650 .await
651}
652
653async fn stream(State(surface): State<Arc<Surface>>, headers: HeaderMap) -> Response {
661 if let Err(refused) = surface.gate(&headers, action::RUNS, "store").await {
662 return refused.into_response();
663 }
664 let hub = surface.bench.streams();
665 let state = (hub.sender.subscribe(), hub.closed.subscribe());
666 let lines = futures_util::stream::unfold(state, |(mut receiver, mut closed)| async move {
667 if *closed.borrow() {
668 return None;
669 }
670 let line = tokio::select! {
671 _ = closed.changed() => return None,
672 next = receiver.recv() => match next {
673 Ok((run, mut line)) => {
674 line["run"] = Value::String(run);
675 line
676 }
677 Err(tokio::sync::broadcast::error::RecvError::Lagged(missed)) => {
678 json!({ "type": "lagged", "value": missed })
679 }
680 Err(tokio::sync::broadcast::error::RecvError::Closed) => return None,
681 },
682 };
683 let mut bytes = line.to_string().into_bytes();
684 bytes.push(b'\n');
685 Some((
686 Ok::<_, std::convert::Infallible>(Bytes::from(bytes)),
687 (receiver, closed),
688 ))
689 });
690 (
691 [(header::CONTENT_TYPE, "application/x-ndjson")],
692 axum::body::Body::from_stream(lines),
693 )
694 .into_response()
695}
696
697async fn run_view(
699 State(surface): State<Arc<Surface>>,
700 Path(run): Path<String>,
701 headers: HeaderMap,
702) -> Response {
703 escaped_read(&surface, &headers, &run, "", "").await
704}
705
706async fn escaped_read(
714 surface: &Surface,
715 headers: &HeaderMap,
716 run: &str,
717 tail: &str,
718 member: &str,
719) -> Response {
720 if let Err(refused) = surface.gate(headers, action::HISTORY, run).await {
721 return refused.into_response();
722 }
723 let Ok(run) = crate::core::RunId::parse(run) else {
724 return super::bad("run").into_response();
725 };
726 let api = match surface.api().await {
727 Ok(api) => api,
728 Err(refused) => return refused.into_response(),
729 };
730 let Ok(uri) = format!("/runs/{run}{tail}").parse::<axum::http::Uri>() else {
731 return super::bad("from").into_response();
732 };
733 let mut inner = Request::new(axum::body::Body::empty());
734 *inner.uri_mut() = uri;
735 *inner.headers_mut() = headers.clone();
736 let answer = match api.router().oneshot(inner).await {
737 Ok(answer) => answer,
738 Err(never) => match never {},
739 };
740 if answer.status() != StatusCode::OK {
741 return answer;
742 }
743 let Ok(bytes) = axum::body::to_bytes(answer.into_body(), usize::MAX).await else {
744 return super::store_failed().into_response();
745 };
746 let Ok(mut page) = serde_json::from_slice::<Value>(&bytes) else {
747 return super::store_failed().into_response();
748 };
749 let mut escaped = false;
750 if member.is_empty() {
751 page = escape_strings(page, &mut escaped);
752 } else if let Some(inner) = page.get_mut(member) {
753 *inner = escape_strings(inner.take(), &mut escaped);
754 }
755 page["escaped"] = Value::Bool(escaped);
756 Json(page).into_response()
757}
758
759fn escape_strings(value: Value, escaped: &mut bool) -> Value {
762 match value {
763 Value::String(text) => {
764 let (text, hit) = crate::core::visible::escape(&text);
765 *escaped |= hit;
766 Value::String(text)
767 }
768 Value::Array(items) => Value::Array(
769 items
770 .into_iter()
771 .map(|item| escape_strings(item, escaped))
772 .collect(),
773 ),
774 Value::Object(fields) => Value::Object(
775 fields
776 .into_iter()
777 .map(|(key, item)| {
778 let (key, hit) = crate::core::visible::escape(&key);
779 *escaped |= hit;
780 (key, escape_strings(item, escaped))
781 })
782 .collect(),
783 ),
784 other => other,
785 }
786}