use std::{path::PathBuf, sync::Arc, time::Duration};
use anyhow::{Context, Result, bail};
use proofborne_runtime::SessionStore;
use serde_json::{Value, json};
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
net::{TcpListener, TcpStream},
};
use uuid::Uuid;
use crate::daemon;
const MAX_REQUEST_LINE: usize = 4096;
const MAX_HEADER_BYTES: usize = 16 * 1024;
const MAX_BODY_BYTES: usize = 16 * 1024 * 1024;
const SESSION_LIST_LIMIT: usize = 50;
const TASK_RUN_TIMEOUT: Duration = Duration::from_secs(30 * 60);
pub async fn serve(
workspace: PathBuf,
listener: TcpListener,
daemon_port: Option<u16>,
) -> Result<()> {
let store = SessionStore::open(
proofborne_runtime::RuntimePaths::for_workspace(&workspace)?
.workspace_data
.join("sessions.sqlite3"),
)
.with_context(|| format!("web store open failed for {}", workspace.display()))?;
let store = Arc::new(store);
let daemon_addr = daemon_port.map(|port| format!("127.0.0.1:{port}"));
loop {
let (stream, _) = listener.accept().await?;
let store = store.clone();
let daemon_addr = daemon_addr.clone();
tokio::spawn(async move {
if let Err(error) = handle_connection(stream, store, daemon_addr.as_deref()).await {
eprintln!("web connection error: {error:#}");
}
});
}
}
async fn handle_connection(
mut stream: TcpStream,
store: Arc<SessionStore>,
daemon_addr: Option<&str>,
) -> Result<()> {
let request_line = read_request_line(&mut stream).await?;
let parts: Vec<&str> = request_line.split_whitespace().collect();
if parts.len() != 3 {
write_response(&mut stream, 400, "text/plain", b"bad request line").await?;
return Ok(());
}
let (method, path) = (parts[0], parts[1]);
let content_length = match drain_headers(&mut stream).await? {
Some(length) if length > MAX_BODY_BYTES => {
write_response(&mut stream, 413, "text/plain", b"body too large").await?;
return Ok(());
}
other => other,
};
if let Some(addr) = daemon_addr {
handle_daemon_mode(&mut stream, method, path, content_length, addr).await
} else {
handle_store_mode(&mut stream, method, path, content_length, &store).await
}
}
async fn handle_store_mode(
stream: &mut TcpStream,
method: &str,
path: &str,
content_length: Option<usize>,
store: &SessionStore,
) -> Result<()> {
if method != "GET" {
if let Some(length) = content_length {
let _ = read_body(stream, length).await?;
}
write_response(stream, 405, "text/plain", b"method not allowed").await?;
return Ok(());
}
match route_store(path, store) {
Ok(Some((content_type, body))) => {
write_response(stream, 200, content_type, &body).await?;
}
Ok(None) => {
write_response(stream, 404, "text/plain", b"not found").await?;
}
Err(_) => {
write_response(stream, 400, "text/plain", b"bad request").await?;
}
}
Ok(())
}
async fn handle_daemon_mode(
stream: &mut TcpStream,
method: &str,
path: &str,
content_length: Option<usize>,
daemon_addr: &str,
) -> Result<()> {
match (method, path) {
("GET", "/") => {
write_response(
stream,
200,
"text/html; charset=utf-8",
client_html().as_bytes(),
)
.await?;
}
("GET", "/api/status") => {
let result = daemon::request(daemon_addr, "daemon.status", json!({})).await;
respond_json(stream, result).await?;
}
("GET", "/api/sessions") => {
let result = daemon::request(daemon_addr, "daemon.sessions", json!({})).await;
respond_json(stream, result).await?;
}
("GET", session_path) if path.starts_with("/api/sessions/") => {
let id = session_path
.strip_prefix("/api/sessions/")
.unwrap_or_default();
let result = daemon::request(daemon_addr, "daemon.session", json!({"id": id})).await;
respond_json(stream, result).await?;
}
("POST", "/api/tasks") => {
let Some(length) = content_length else {
write_response(stream, 400, "text/plain", b"missing Content-Length").await?;
return Ok(());
};
let body = read_body(stream, length).await?;
let Ok(params) = serde_json::from_slice::<Value>(&body) else {
write_response(stream, 400, "text/plain", b"invalid JSON body").await?;
return Ok(());
};
let result =
daemon::request_with_timeout(daemon_addr, "task.run", params, TASK_RUN_TIMEOUT)
.await;
respond_json(stream, result).await?;
}
_ => {
write_response(stream, 404, "text/plain", b"not found").await?;
}
}
Ok(())
}
async fn respond_json(stream: &mut TcpStream, result: Result<Value>) -> Result<()> {
match result {
Ok(value) => {
let body = serde_json::to_vec(&value)?;
write_response(stream, 200, "application/json", &body).await?;
}
Err(error) => {
let body = serde_json::to_vec(&json!({"error": format!("{error:#}")}))?;
write_response(stream, 502, "application/json", &body).await?;
}
}
Ok(())
}
fn route_store(path: &str, store: &SessionStore) -> Result<Option<(&'static str, Vec<u8>)>> {
match path {
"/" => Ok(Some((
"text/html; charset=utf-8",
dashboard_html().into_bytes(),
))),
"/api/sessions" => {
let sessions = store
.list_sessions(SESSION_LIST_LIMIT)
.context("list sessions")?;
Ok(Some((
"application/json",
serde_json::to_vec(&json!({ "sessions": sessions }))?,
)))
}
_ => {
let session_path = path.strip_prefix("/api/sessions/");
let Some(id_text) = session_path else {
return Ok(None);
};
let id = Uuid::parse_str(id_text).context("invalid session id")?;
let record = store.load_session(id).context("load session")?;
let events = store.load_events(id).context("load events")?;
let final_event = events
.iter()
.rev()
.find(|event| event.kind == "proof.finalized");
let summary = json!({
"id": record.id,
"workspace": record.workspace,
"provider": record.provider,
"model": record.model,
"status": record.status,
"createdAt": record.created_at,
"updatedAt": record.updated_at,
"eventCount": events.len(),
"finalEvent": final_event.map(|event| {
json!({
"outcome": event.payload.get("outcome"),
"proofHash": event.payload.get("proofHash"),
})
}),
});
Ok(Some(("application/json", serde_json::to_vec(&summary)?)))
}
}
}
async fn read_request_line(stream: &mut TcpStream) -> Result<String> {
let mut buffer = Vec::new();
let mut byte = [0u8; 1];
loop {
if stream.read(&mut byte).await? == 0 {
bail!("connection closed before request line");
}
buffer.push(byte[0]);
if buffer.ends_with(b"\r\n") {
break;
}
if buffer.len() > MAX_REQUEST_LINE {
bail!("request line too long");
}
}
let line = String::from_utf8_lossy(&buffer);
Ok(line.trim_end_matches(['\r', '\n']).to_owned())
}
async fn drain_headers(stream: &mut TcpStream) -> Result<Option<usize>> {
let mut buffer = Vec::new();
let mut byte = [0u8; 1];
let mut content_length = None;
loop {
if stream.read(&mut byte).await? == 0 {
bail!("connection closed before headers");
}
buffer.push(byte[0]);
if buffer.ends_with(b"\r\n\r\n") {
break;
}
if buffer.len() > MAX_HEADER_BYTES {
bail!("headers too large");
}
}
let header_text = String::from_utf8_lossy(&buffer);
for line in header_text.lines() {
if let Some((name, value)) = line.split_once(':')
&& name.eq_ignore_ascii_case("content-length")
{
content_length = Some(value.trim().parse::<usize>()?);
}
}
Ok(content_length)
}
async fn read_body(stream: &mut TcpStream, length: usize) -> Result<Vec<u8>> {
let mut body = vec![0u8; length];
stream.read_exact(&mut body).await?;
Ok(body)
}
async fn write_response(
stream: &mut TcpStream,
status: u16,
content_type: &str,
body: &[u8],
) -> Result<()> {
let reason = match status {
200 => "OK",
400 => "Bad Request",
404 => "Not Found",
405 => "Method Not Allowed",
413 => "Payload Too Large",
502 => "Bad Gateway",
_ => "Internal Server Error",
};
stream
.write_all(
format!(
"HTTP/1.1 {status} {reason}\r\nContent-Type: {content_type}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
body.len()
)
.as_bytes(),
)
.await?;
stream.write_all(body).await?;
stream.flush().await?;
Ok(())
}
fn dashboard_html() -> String {
r#"<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8">
<title>Proofborne dashboard</title>
<style>
body{font-family:system-ui,sans-serif;margin:2rem;color:#1f2933}
table{border-collapse:collapse;width:100%;margin-top:1rem}
th,td{border:1px solid #d3d9e0;padding:.5rem;text-align:left;font-size:.9rem}
th{background:#f4f6f8}
.verified{color:#0a7d33;font-weight:600}
</style>
</head>
<body>
<h1>Proofborne dashboard</h1>
<p id="status">Loading sessions…</p>
<table id="sessions"><thead><tr><th>Session</th><th>Status</th><th>Provider</th><th>Model</th><th>Created</th></tr></thead><tbody></tbody></table>
<script>
fetch('/api/sessions').then(r => r.json()).then(data => {
const rows = (data.sessions || []).map(s => `
<tr>
<td><a href="/api/sessions/${s.id}">${s.id}</a></td>
<td class="${s.status === 'verified' ? 'verified' : ''}">${s.status}</td>
<td>${s.provider}</td>
<td>${s.model}</td>
<td>${s.createdAt}</td>
</tr>`).join('');
document.querySelector('#sessions tbody').innerHTML = rows;
document.querySelector('#status').textContent = data.sessions.length + ' sessions';
}).catch(() => { document.querySelector('#status').textContent = 'dashboard unavailable'; });
</script>
</body>
</html>"#
.to_owned()
}
fn client_html() -> String {
r#"<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width, initial-scale=1">
<title>Proofborne client</title>
<style>
body{font-family:system-ui,sans-serif;margin:0;color:#1f2933;background:#f7f9fb}
header{background:#0f2a43;color:#fff;padding:1rem}
header h1{margin:0;font-size:1.2rem}
main{padding:1rem;max-width:48rem;margin:0 auto}
section{background:#fff;border:1px solid #d3d9e0;border-radius:.5rem;padding:1rem;margin-bottom:1rem}
h2{margin:0 0 .75rem;font-size:1rem}
input{width:100%;box-sizing:border-box;padding:.6rem;margin:.25rem 0;font-size:1rem;border:1px solid #c3ccd5;border-radius:.3rem}
button{width:100%;padding:.7rem;font-size:1rem;border:0;border-radius:.3rem;background:#0f6db5;color:#fff;margin-top:.5rem}
button:disabled{background:#9db4c7}
li{padding:.5rem;border-bottom:1px solid #eef1f4;list-style:none;font-size:.95rem}
li:last-child{border-bottom:0}
ul{padding:0;margin:0}
.verified{color:#0a7d33;font-weight:600}
.failed{color:#b51d1d}
.meta{color:#5b6b7a;font-size:.85rem}
.error{color:#b51d1d}
pre{white-space:pre-wrap;word-break:break-all;font-size:.85rem;background:#f4f6f8;padding:.5rem;border-radius:.3rem}
</style>
</head>
<body>
<header><h1>Proofborne client</h1></header>
<main>
<section id="status-card"><h2>Daemon status</h2><p class="meta">loading…</p></section>
<section id="run-card">
<h2>Run proof-gated task</h2>
<input id="contract" placeholder="contract path (inside the daemon workspace)" autocomplete="off">
<input id="policy" placeholder="policy (blank = daemon default: workspace | review | ci)" autocomplete="off">
<label><input type="checkbox" id="trust"> trust workspace</label>
<button id="run">Run task</button>
<p id="run-result" class="meta"></p>
</section>
<section id="sessions-card"><h2>Sessions</h2><p class="meta">loading…</p></section>
</main>
<script>
const $ = id => document.getElementById(id);
function showStatus(data) {
$('status-card').innerHTML = '<h2>Daemon status</h2><p class="meta">version ' +
(data.version || '?') + ' · pid ' + (data.pid || '?') + '</p>' +
'<p class="meta">workspace: ' + (data.workspace || '?') + '</p>' +
(data.pluginToolCount !== undefined ? '<p class="meta">plugins: ' + data.pluginToolCount + ' tool(s)</p>' : '') +
(data.pluginError ? '<p class="error">plugin error: ' + data.pluginError + '</p>' : '');
}
function showSessions(data) {
const sessions = data.sessions || [];
const items = sessions.map(s => {
const cls = s.status === 'verified' ? 'verified' : (s.status === 'failed' ? 'failed' : '');
return '<li><a href="/api/sessions/' + s.id + '">' + s.id + '</a> ' +
'<span class="' + cls + '">' + s.status + '</span>' +
'<div class="meta">' + (s.provider || '') + ' · ' + (s.model || '') + ' · ' + (s.createdAt || '') + '</div></li>';
}).join('');
$('sessions-card').innerHTML = '<h2>Sessions</h2><ul>' + items + '</ul>';
}
fetch('/api/status').then(r => r.json()).then(showStatus).catch(() => {
$('status-card').innerHTML = '<h2>Daemon status</h2><p class="error">daemon unreachable</p>';
});
fetch('/api/sessions').then(r => r.json()).then(showSessions).catch(() => {
$('sessions-card').innerHTML = '<h2>Sessions</h2><p class="error">cannot list sessions</p>';
});
$('run').addEventListener('click', async () => {
const params = {
contract: $('contract').value.trim(),
trustWorkspace: $('trust').checked
};
const policy = $('policy').value.trim();
if (policy) params.policy = policy;
if (!params.contract) { $('run-result').textContent = 'contract path is required'; return; }
$('run').disabled = true;
$('run-result').textContent = 'running proof-gated task…';
try {
const response = await fetch('/api/tasks', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify(params)
});
const data = await response.json();
if (!response.ok) throw new Error(data.error || response.status);
$('run-result').innerHTML = '<span class="verified">' + data.outcome + '</span> session ' +
data.sessionId + (data.proofHash ? '<br><span class="meta">proofHash: ' + data.proofHash + '</span>' : '');
} catch (error) {
$('run-result').className = 'error';
$('run-result').textContent = 'task failed: ' + error.message;
} finally {
$('run').disabled = false;
}
});
</script>
</body>
</html>"#
.to_owned()
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::Value;
use std::net::SocketAddr;
use std::path::Path;
use tokio::io::{AsyncBufReadExt, BufReader};
async fn request(addr: SocketAddr, raw: &str) -> (u16, String) {
let mut stream = TcpStream::connect(addr).await.unwrap();
stream.write_all(raw.as_bytes()).await.unwrap();
let mut reader = BufReader::new(stream);
let mut status_line = String::new();
reader.read_line(&mut status_line).await.unwrap();
let status = status_line
.split_whitespace()
.nth(1)
.unwrap()
.parse()
.unwrap();
let mut rest = String::new();
reader.read_to_string(&mut rest).await.unwrap();
let body = rest
.split_once("\r\n\r\n")
.map(|(_, body)| body.to_owned())
.unwrap_or(rest);
(status, body)
}
fn temp_store(workspace: &Path) -> (uuid::Uuid, std::path::PathBuf) {
let paths = proofborne_runtime::RuntimePaths::for_workspace(workspace).unwrap();
paths.ensure_data_directories().unwrap();
let store = SessionStore::open(paths.workspace_data.join("sessions.sqlite3")).unwrap();
let id = uuid::Uuid::now_v7();
store
.create_session(id, workspace, "mock", "deterministic-v1")
.unwrap();
drop(store);
(id, paths.workspace_data.join("sessions.sqlite3"))
}
#[tokio::test]
async fn dashboard_serves_page_and_session_api() {
let workspace = tempfile::tempdir().unwrap();
let (id, _) = temp_store(workspace.path());
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let workspace = workspace.path().to_path_buf();
let server = tokio::spawn(serve(workspace, listener, None));
let (status, page) = request(addr, "GET / HTTP/1.1\r\nHost: localhost\r\n\r\n").await;
assert_eq!(status, 200);
assert!(page.contains("Proofborne dashboard"));
let (status, body) = request(
addr,
"GET /api/sessions HTTP/1.1\r\nHost: localhost\r\n\r\n",
)
.await;
assert_eq!(status, 200);
let sessions: Value = serde_json::from_str(&body).unwrap();
assert_eq!(sessions["sessions"].as_array().unwrap().len(), 1);
assert!(body.contains(&id.to_string()));
let (status, body) = request(
addr,
&format!("GET /api/sessions/{id} HTTP/1.1\r\nHost: localhost\r\n\r\n"),
)
.await;
assert_eq!(status, 200);
let detail: Value = serde_json::from_str(&body).unwrap();
assert_eq!(detail["status"], "running");
let (status, _) = request(addr, "GET /nope HTTP/1.1\r\nHost: localhost\r\n\r\n").await;
assert_eq!(status, 404);
let (status, _) = request(
addr,
"POST /api/tasks HTTP/1.1\r\nHost: localhost\r\nContent-Length: 2\r\n\r\n{}",
)
.await;
assert_eq!(status, 405);
server.abort();
}
#[tokio::test]
async fn client_mode_serves_daemon_status_sessions_and_task_run() {
let workspace = tempfile::tempdir().unwrap();
let config_dir = workspace.path().join(".proofborne");
std::fs::create_dir_all(&config_dir).unwrap();
std::fs::write(
config_dir.join("config.toml"),
"default_profile = \"mock\"\n\n[providers.mock]\nprovider = \"mock\"\nmodel = \"deterministic-v1\"\n",
)
.unwrap();
let contract = workspace.path().join("contract.json");
std::fs::write(
&contract,
serde_json::to_string_pretty(&json!({
"schemaVersion": "proofborne.v1",
"id": "019fbf00-0000-7000-8000-000000000002",
"goal": "client mode task",
"claimScope": "task",
"criteria": [{
"id": "runtime_completed",
"description": "runtime completed",
"required": true,
"minimumEvidence": 1,
"evidenceRequirement": {
"allowedKinds": ["runtime"],
"allowedProducers": ["proofborne.runtime"],
"minimumObservations": 1,
"freshness": "final_workspace_state",
"requireArtifacts": false,
"minimumIndependentProducers": 1,
"minimumAssurance": "observed"
},
"state": "pending",
"evidenceIds": []
},
{
"id": "tests_pass",
"description": "verification passes",
"required": true,
"minimumEvidence": 1,
"evidenceRequirement": {
"allowedKinds": ["process"],
"allowedProducers": ["verify.exec"],
"minimumObservations": 1,
"freshness": "final_workspace_state",
"requireArtifacts": false,
"minimumIndependentProducers": 1,
"minimumAssurance": "observed"
},
"state": "pending",
"evidenceIds": []
}],
"createdAt": "2026-08-02T00:00:00Z",
"confirmed": true
}))
.unwrap(),
)
.unwrap();
let verification = workspace.path().join("verification.json");
std::fs::write(
&verification,
serde_json::to_string(&json!([{
"criterionId": "tests_pass",
"program": std::env::current_exe().unwrap().to_string_lossy(),
"args": ["--help"],
"timeoutSeconds": 30
}]))
.unwrap(),
)
.unwrap();
let daemon_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let daemon_addr = daemon_listener.local_addr().unwrap();
let daemon_workspace = workspace.path().to_path_buf();
let daemon_server = tokio::spawn(crate::daemon::serve(daemon_workspace, daemon_listener));
let web_listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let web_addr = web_listener.local_addr().unwrap();
let web_workspace = workspace.path().to_path_buf();
let web_server = tokio::spawn(serve(web_workspace, web_listener, Some(daemon_addr.port())));
let (status, body) = request(
web_addr,
"GET /api/status HTTP/1.1\r\nHost: localhost\r\n\r\n",
)
.await;
assert_eq!(status, 200);
let status_value: Value = serde_json::from_str(&body).unwrap();
assert_eq!(status_value["version"], env!("CARGO_PKG_VERSION"));
let (status, page) = request(web_addr, "GET / HTTP/1.1\r\nHost: localhost\r\n\r\n").await;
assert_eq!(status, 200);
assert!(page.contains("Proofborne client"));
let contract_text = contract.to_string_lossy();
let body = format!(
"{{\"contract\":\"{}\",\"verification\":\"{}\",\"policy\":\"workspace\",\"trustWorkspace\":true}}",
contract_text.replace('\\', "\\\\"),
verification.to_string_lossy().replace('\\', "\\\\")
);
let (status, body) = request(
web_addr,
&format!(
"POST /api/tasks HTTP/1.1\r\nHost: localhost\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{}",
body.len(),
body
),
)
.await;
assert_eq!(status, 200, "unexpected body: {body}");
let result: Value = serde_json::from_str(&body).unwrap();
assert_eq!(result["outcome"], "Verified");
assert!(result["proofHash"].as_str().is_some());
let (status, body) = request(
web_addr,
"GET /api/sessions HTTP/1.1\r\nHost: localhost\r\n\r\n",
)
.await;
assert_eq!(status, 200);
let sessions: Value = serde_json::from_str(&body).unwrap();
assert!(!sessions["sessions"].as_array().unwrap().is_empty());
let _ = crate::daemon::request(
&format!("127.0.0.1:{}", daemon_addr.port()),
"daemon.shutdown",
json!({}),
)
.await;
daemon_server.await.unwrap().unwrap();
web_server.abort();
}
}