use std::sync::{Arc, Mutex};
use axum::Router;
use axum::body::Bytes;
use axum::extract::{Path, Request, State};
use axum::http::{HeaderMap, HeaderValue, StatusCode, header};
use axum::middleware::Next;
use axum::response::{IntoResponse, Json, Response};
use axum::routing::{MethodFilter, MethodRouter, any, get, on};
use serde::{Deserialize, Serialize};
use serde_json::{Value, json};
use tower::ServiceExt as _;
use super::openapi::Method;
use super::rebinding::{Rebinding, refuse_rebinding};
use super::{Api, ApiError, Authenticator, Session};
use crate::core::{
ACTION_ADMIT, ACTION_PERFORM, ACTION_RELEASE, Digest, PolicyBundleIdentity, PolicyDecision,
PolicyEngine, PolicyRequest,
};
use crate::runtime::Runtime;
pub mod action {
pub const RUN: &str = "api:dev.run";
pub const REPLAY: &str = "api:dev.replay";
pub const EXPORT: &str = "api:dev.export";
pub const MANIFEST: &str = "api:dev.manifest";
pub const HISTORY: &str = "api:dev.history";
pub const RUNS: &str = "api:dev.runs";
pub const ALL: &[&str] = &[RUN, REPLAY, EXPORT, MANIFEST, HISTORY, RUNS];
}
pub const TENANT: &str = "dev";
pub const CONTENT_SECURITY_POLICY: &str = "default-src 'none'; script-src 'self'; \
style-src 'self'; connect-src 'self'; img-src 'self'; base-uri 'none'; \
form-action 'none'; frame-ancestors 'none'; require-trusted-types-for 'script'; \
trusted-types 'none'";
pub const SHELL: &[(&str, &str, &str)] = &[
(
"/",
"text/html; charset=utf-8",
include_str!("dev/index.html"),
),
(
"/assets/app.js",
"text/javascript; charset=utf-8",
include_str!("dev/app.js"),
),
(
"/assets/app.css",
"text/css; charset=utf-8",
include_str!("dev/app.css"),
),
];
#[async_trait::async_trait]
pub trait Workbench: Send + Sync + 'static {
async fn plane(&self) -> Arc<Runtime>;
async fn declaration(&self) -> Declaration;
async fn start(&self, request: StartRequest) -> Result<Started, String>;
fn streams(&self) -> Arc<StreamHub>;
async fn replay(&self, run: Option<crate::core::RunId>) -> Result<Vec<Replayed>, String>;
}
#[derive(Debug)]
pub struct StreamHub {
sender: tokio::sync::broadcast::Sender<(String, Value)>,
closed: tokio::sync::watch::Sender<bool>,
}
impl StreamHub {
const BACKLOG: usize = 1024;
#[must_use]
pub fn new() -> Arc<Self> {
Arc::new(Self {
sender: tokio::sync::broadcast::channel(Self::BACKLOG).0,
closed: tokio::sync::watch::channel(false).0,
})
}
pub fn close(&self) {
self.closed.send_replace(true);
}
}
impl crate::runtime::RunStreamObserver for StreamHub {
fn event(
&self,
run: crate::core::RunId,
event: crate::core::Tainted<crate::model::ModelStreamEvent>,
) {
let trust = event.label().trust;
let line = match event.into_unlabelled() {
crate::model::ModelStreamEvent::TextDelta(text) => {
json!({ "type": "text_delta", "value": crate::core::visible::escape(&text).0, "trust": trust })
}
crate::model::ModelStreamEvent::Usage(usage) => {
json!({ "type": "usage", "value": usage, "trust": trust })
}
};
let _ = self.sender.send((run.to_string(), line));
}
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct Declaration {
pub file: String,
pub agents: Vec<DeclaredAgent>,
pub refused: Option<String>,
pub live: Vec<String>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct DeclaredAgent {
pub name: String,
pub version: String,
pub digest: String,
pub bound: Vec<String>,
#[serde(default)]
pub provides: Vec<String>,
#[serde(default)]
pub input_schema: Option<Value>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct StartRequest {
pub input: Value,
#[serde(default)]
pub correlate: Vec<String>,
#[serde(default)]
pub capability: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Started {
pub run: String,
pub status: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Replayed {
pub run: String,
pub verdict: String,
pub detail: String,
}
#[derive(Debug, Clone, Default, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ReplayRequest {
#[serde(default)]
pub run: Option<String>,
}
#[derive(Debug, Clone, Copy)]
pub struct DevRoute {
pub method: Method,
pub path: &'static str,
pub action: &'static str,
serve: fn(MethodFilter) -> MethodRouter<Arc<Surface>>,
}
pub const ROUTES: &[DevRoute] = &[
DevRoute {
method: Method::Get,
path: "/dev/manifest",
action: action::MANIFEST,
serve: |m| on(m, declaration),
},
DevRoute {
method: Method::Post,
path: "/dev/runs",
action: action::RUN,
serve: |m| on(m, start),
},
DevRoute {
method: Method::Get,
path: "/dev/runs",
action: action::RUNS,
serve: |m| on(m, runs),
},
DevRoute {
method: Method::Post,
path: "/dev/replay",
action: action::REPLAY,
serve: |m| on(m, replay),
},
DevRoute {
method: Method::Get,
path: "/dev/export",
action: action::EXPORT,
serve: |m| on(m, export),
},
DevRoute {
method: Method::Get,
path: "/dev/runs/{run}",
action: action::HISTORY,
serve: |m| on(m, run_view),
},
DevRoute {
method: Method::Get,
path: "/dev/stream",
action: action::RUNS,
serve: |m| on(m, stream),
},
DevRoute {
method: Method::Get,
path: "/dev/runs/{run}/history",
action: action::HISTORY,
serve: |m| on(m, history),
},
];
const EXPORT_LIMIT: usize = 10_000;
#[derive(Debug, Clone)]
pub struct DevPolicy {
actor: String,
}
impl DevPolicy {
#[must_use]
pub fn new(actor: impl Into<String>) -> Self {
Self {
actor: actor.into(),
}
}
}
impl PolicyEngine for DevPolicy {
fn authorize(&self, request: &PolicyRequest<'_>) -> PolicyDecision {
if request.action.starts_with("api:") {
let tenant = request.context.get("tenant").and_then(Value::as_str);
if request.principal == self.actor && tenant == Some(TENANT) {
return PolicyDecision::Permit;
}
return PolicyDecision::deny("only this dev session's own token reaches its plane");
}
match request.action {
ACTION_ADMIT | ACTION_PERFORM => PolicyDecision::Permit,
ACTION_RELEASE => PolicyDecision::deny(
"`agentplane dev` permits no release, as `agentplane run` permits none: a \
release lowers a label only when a rule permits `data:release`",
),
_ => PolicyDecision::deny("a dev plane permits admission and effects only"),
}
}
fn bundle(&self) -> PolicyBundleIdentity {
PolicyBundleIdentity::new(
Digest::of(b"agentplane dev policy"),
"agentplane/dev-policy-v1",
)
}
}
struct Surface {
bench: Arc<dyn Workbench>,
auth: Arc<dyn Authenticator>,
operator: Mutex<Option<(Arc<Runtime>, Api)>>,
}
impl std::fmt::Debug for Surface {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Surface").finish_non_exhaustive()
}
}
impl Surface {
async fn api(&self) -> Result<Api, ApiError> {
let plane = self.bench.plane().await;
let mut held = self
.operator
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some((built_for, api)) = held.as_ref()
&& Arc::ptr_eq(built_for, &plane)
{
return Ok(api.clone());
}
let api = Api::new(Arc::clone(&plane), Arc::clone(&self.auth))
.map_err(|e| ApiError(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
*held = Some((plane, api.clone()));
Ok(api)
}
async fn gate(
&self,
headers: &HeaderMap,
action: &str,
resource: &str,
) -> Result<Session, ApiError> {
let api = self.api().await?;
let caller = self.auth.authenticate(headers).await?;
api.authorize(caller, action, resource)
}
}
pub fn router(bench: Arc<dyn Workbench>, auth: Arc<dyn Authenticator>, port: u16) -> Router {
let surface = Arc::new(Surface {
bench,
auth,
operator: Mutex::new(None),
});
let authorities = [
format!("127.0.0.1:{port}"),
format!("localhost:{port}"),
format!("[::1]:{port}"),
];
let rebinding = Rebinding::new(
authorities.to_vec(),
authorities.iter().map(|a| format!("http://{a}")).collect(),
)
.origin_on_writes();
let mut routes = Router::new();
for route in ROUTES {
let filter = match route.method {
Method::Get => MethodFilter::GET,
Method::Post => MethodFilter::POST,
};
routes = routes.route(route.path, (route.serve)(filter));
}
for (path, content_type, body) in SHELL {
routes = routes.route(
path,
get(move || async move { ([(header::CONTENT_TYPE, *content_type)], *body) }),
);
}
routes
.route("/api/{*rest}", any(operator))
.fallback(not_found)
.with_state(surface)
.layer(axum::middleware::from_fn_with_state(
rebinding,
refuse_rebinding,
))
.layer(axum::middleware::from_fn(security_headers))
}
async fn security_headers(request: Request, next: Next) -> Response {
let mut response = next.run(request).await;
let headers = response.headers_mut();
headers.insert(
header::CONTENT_SECURITY_POLICY,
HeaderValue::from_static(CONTENT_SECURITY_POLICY),
);
headers.insert(
header::X_CONTENT_TYPE_OPTIONS,
HeaderValue::from_static("nosniff"),
);
headers.insert(
header::REFERRER_POLICY,
HeaderValue::from_static("no-referrer"),
);
headers.insert(header::CACHE_CONTROL, HeaderValue::from_static("no-store"));
response
}
async fn not_found() -> Response {
ApiError(StatusCode::NOT_FOUND, "no such page".to_owned()).into_response()
}
async fn operator(State(surface): State<Arc<Surface>>, request: Request) -> Response {
let api = match surface.api().await {
Ok(api) => api,
Err(refused) => return refused.into_response(),
};
let (parts, body) = request.into_parts();
let path = parts
.uri
.path_and_query()
.map_or("/", axum::http::uri::PathAndQuery::as_str);
let Ok(uri) = path
.strip_prefix("/api")
.unwrap_or(path)
.parse::<axum::http::Uri>()
else {
return not_found().await;
};
let mut inner = Request::new(body);
*inner.method_mut() = parts.method;
*inner.uri_mut() = uri;
*inner.version_mut() = parts.version;
*inner.headers_mut() = parts.headers;
match api.router().oneshot(inner).await {
Ok(response) => response,
Err(never) => match never {},
}
}
fn unreadable(e: &serde_json::Error) -> ApiError {
ApiError(StatusCode::UNPROCESSABLE_ENTITY, e.to_string())
}
async fn declaration(
State(surface): State<Arc<Surface>>,
headers: HeaderMap,
) -> Result<Json<Declaration>, ApiError> {
surface.gate(&headers, action::MANIFEST, "manifest").await?;
Ok(Json(surface.bench.declaration().await))
}
async fn start(
State(surface): State<Arc<Surface>>,
headers: HeaderMap,
body: Bytes,
) -> Result<Json<Started>, ApiError> {
surface.gate(&headers, action::RUN, "run").await?;
let request: StartRequest = serde_json::from_slice(&body).map_err(|e| unreadable(&e))?;
surface
.bench
.start(request)
.await
.map(Json)
.map_err(|why| ApiError(StatusCode::UNPROCESSABLE_ENTITY, why))
}
async fn replay(
State(surface): State<Arc<Surface>>,
headers: HeaderMap,
body: Bytes,
) -> Result<Json<Vec<Replayed>>, ApiError> {
surface.gate(&headers, action::REPLAY, "run").await?;
let request: ReplayRequest = if body.is_empty() {
ReplayRequest::default()
} else {
serde_json::from_slice(&body).map_err(|e| unreadable(&e))?
};
let run = request
.run
.as_deref()
.map(crate::core::RunId::parse)
.transpose()
.map_err(|_| super::bad("run"))?;
surface
.bench
.replay(run)
.await
.map(Json)
.map_err(|why| ApiError(StatusCode::UNPROCESSABLE_ENTITY, why))
}
async fn runs(
State(surface): State<Arc<Surface>>,
headers: HeaderMap,
) -> Result<Json<Value>, ApiError> {
let s = surface.gate(&headers, action::RUNS, "store").await?;
let outcomes: Vec<String> = crate::runtime::OUTCOMES_OF_RECORD
.iter()
.map(|o| (*o).to_owned())
.collect();
let found = crate::export::runs_to_read(s.plane.journal(), &outcomes, true, EXPORT_LIMIT)
.await
.map_err(|_| super::store_failed())?;
let runs: Vec<String> = found.runs.iter().map(ToString::to_string).collect();
Ok(Json(
json!({ "runs": runs, "partial": !found.reached.is_empty() }),
))
}
async fn export(
State(surface): State<Arc<Surface>>,
headers: HeaderMap,
) -> Result<Json<Value>, ApiError> {
let s = surface.gate(&headers, action::EXPORT, "store").await?;
let journal = Arc::clone(s.plane.journal());
let cases = s
.plane
.cases()
.cloned()
.ok_or_else(|| super::unavailable("case store"))?;
let outcomes: Vec<String> = crate::runtime::OUTCOMES_OF_RECORD
.iter()
.map(|o| (*o).to_owned())
.collect();
let found = crate::export::runs_to_read(&journal, &outcomes, true, EXPORT_LIMIT)
.await
.map_err(|_| super::store_failed())?;
let mut bytes = Vec::new();
crate::export::to_jsonl(&journal, &cases, &found.runs, &mut bytes)
.await
.map_err(|_| super::store_failed())?;
let report = crate::export::verify(&bytes[..], None, &[]).map_err(|_| super::store_failed())?;
Ok(Json(json!({
"export": String::from_utf8_lossy(&bytes),
"partial": !found.reached.is_empty(),
"report": report,
"verify": [
"agentplane verify export.jsonl",
"python3 tools/verify_export.py export.jsonl",
],
})))
}
async fn history(
State(surface): State<Arc<Surface>>,
Path(run): Path<String>,
request: Request,
) -> Response {
let (parts, _) = request.into_parts();
let query = parts
.uri
.query()
.map_or_else(String::new, |q| format!("?{q}"));
escaped_read(
&surface,
&parts.headers,
&run,
&format!("/history{query}"),
"records",
)
.await
}
async fn stream(State(surface): State<Arc<Surface>>, headers: HeaderMap) -> Response {
if let Err(refused) = surface.gate(&headers, action::RUNS, "store").await {
return refused.into_response();
}
let hub = surface.bench.streams();
let state = (hub.sender.subscribe(), hub.closed.subscribe());
let lines = futures_util::stream::unfold(state, |(mut receiver, mut closed)| async move {
if *closed.borrow() {
return None;
}
let line = tokio::select! {
_ = closed.changed() => return None,
next = receiver.recv() => match next {
Ok((run, mut line)) => {
line["run"] = Value::String(run);
line
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(missed)) => {
json!({ "type": "lagged", "value": missed })
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => return None,
},
};
let mut bytes = line.to_string().into_bytes();
bytes.push(b'\n');
Some((
Ok::<_, std::convert::Infallible>(Bytes::from(bytes)),
(receiver, closed),
))
});
(
[(header::CONTENT_TYPE, "application/x-ndjson")],
axum::body::Body::from_stream(lines),
)
.into_response()
}
async fn run_view(
State(surface): State<Arc<Surface>>,
Path(run): Path<String>,
headers: HeaderMap,
) -> Response {
escaped_read(&surface, &headers, &run, "", "").await
}
async fn escaped_read(
surface: &Surface,
headers: &HeaderMap,
run: &str,
tail: &str,
member: &str,
) -> Response {
if let Err(refused) = surface.gate(headers, action::HISTORY, run).await {
return refused.into_response();
}
let Ok(run) = crate::core::RunId::parse(run) else {
return super::bad("run").into_response();
};
let api = match surface.api().await {
Ok(api) => api,
Err(refused) => return refused.into_response(),
};
let Ok(uri) = format!("/runs/{run}{tail}").parse::<axum::http::Uri>() else {
return super::bad("from").into_response();
};
let mut inner = Request::new(axum::body::Body::empty());
*inner.uri_mut() = uri;
*inner.headers_mut() = headers.clone();
let answer = match api.router().oneshot(inner).await {
Ok(answer) => answer,
Err(never) => match never {},
};
if answer.status() != StatusCode::OK {
return answer;
}
let Ok(bytes) = axum::body::to_bytes(answer.into_body(), usize::MAX).await else {
return super::store_failed().into_response();
};
let Ok(mut page) = serde_json::from_slice::<Value>(&bytes) else {
return super::store_failed().into_response();
};
let mut escaped = false;
if member.is_empty() {
page = escape_strings(page, &mut escaped);
} else if let Some(inner) = page.get_mut(member) {
*inner = escape_strings(inner.take(), &mut escaped);
}
page["escaped"] = Value::Bool(escaped);
Json(page).into_response()
}
fn escape_strings(value: Value, escaped: &mut bool) -> Value {
match value {
Value::String(text) => {
let (text, hit) = crate::core::visible::escape(&text);
*escaped |= hit;
Value::String(text)
}
Value::Array(items) => Value::Array(
items
.into_iter()
.map(|item| escape_strings(item, escaped))
.collect(),
),
Value::Object(fields) => Value::Object(
fields
.into_iter()
.map(|(key, item)| {
let (key, hit) = crate::core::visible::escape(&key);
*escaped |= hit;
(key, escape_strings(item, escaped))
})
.collect(),
),
other => other,
}
}