Skip to main content

supercode_harness/
server.rs

1//! §2 module 31 `server` (COMPOSABLE-HARNESS-DESIGN.md, D7 "full
2//! programmatic RPC/HTTP server", D8 "remote attach", D10 "daemon"; §1.9
3//! Obligation 9's out-of-process half — the in-process SDK already meets
4//! the core commitment via [`crate::EventSink`]).
5//!
6//! The embedding ladder this module builds:
7//!
8//! 1. **`--output-format stream-json`** (the CLI's existing rung, UX-23) —
9//!    already ships a JSONL [`crate::AgentEvent`] stream over stdout. This
10//!    unit completes it: [`crate::AgentEvent::to_json`] is now the single
11//!    canonical projection both that sink AND this module's notifications
12//!    share, and it covers the FULL event set (previously
13//!    [`crate::AgentEvent::BackgroundOutput`] fell into a generic
14//!    "unknown" catch-all).
15//! 2. **JSONL-RPC over stdio** ([`run_stdio`]) — the SDK-out-of-process
16//!    surface: a parent process drives this agent's loop over stdin/stdout
17//!    with `{"id","method","params"}` request lines, getting back
18//!    `{"id","result"|"error"}` responses interleaved with
19//!    `{"event":...}` notifications. Parent-process-trusted (same trust
20//!    model as [`crate::mcp::serve_stdio`]) — no auth token.
21//! 3. **The same RPC surface over HTTP** ([`run_http`], D8 "remote
22//!    attach") — `POST /rpc` for request/response, `GET /events` for the
23//!    event stream (SSE-shaped: `data: <json>\n\n` per line). Unlike
24//!    stdio, a network client is UNTRUSTED by default, so every request
25//!    must carry the bearer token (`check_auth`).
26//!
27//! **Security posture (this is a listener — the highest-risk module
28//! class):**
29//! - `[capabilities.server]` is project-forbidden (D-10) — see
30//!   `crates/cli/src/userconfig.rs`'s `PROJECT_FORBIDDEN_CAPABILITY_TABLES`
31//!   and this crate's `configfile::PROJECT_FORBIDDEN_CAPABILITY_TABLES`
32//!   (both already listed `"server"` before this unit landed; this module
33//!   is what makes the listener the strip was already guarding against
34//!   real).
35//! - Default-off: nothing in this module is ever reached unless a caller
36//!   explicitly invokes [`run_stdio`]/[`run_http`] AND the CLI's own
37//!   gate (`capabilities.server.enabled == Some(true)`, checked before
38//!   either is called) passed.
39//! - Loopback-only HTTP bind by default — enforced by the CALLER (the
40//!   CLI's `serve` command defaults `bind` to `127.0.0.1:0` and only binds
41//!   elsewhere on an explicit `bind`/`--bind` override, with a printed
42//!   exposure warning); [`run_http`] itself binds whatever address it's
43//!   given, since the loopback POLICY decision belongs to the config/CLI
44//!   layer, not the transport.
45//! - **No permission/sandbox bypass.** [`RpcEngine::new`] takes an already
46//!   fully-constructed [`crate::Agent`] — the SAME `Agent` a local
47//!   `run`/`chat` session would build (same `Config`, same permission
48//!   rules, same sandbox). This module installs NO approval handler of its
49//!   own and provides no channel for a remote/RPC caller to answer an
50//!   approval prompt; combined with `Agent`'s existing fail-closed rule
51//!   ("absent handler denies" — `crates/harness/src/agent.rs`'s
52//!   `prepare_tool_call`), any tool call that would need interactive
53//!   approval is DENIED, never silently approved, when driven through this
54//!   module. See `crates/harness/tests/server_engine.rs` for a fail-on-revert
55//!   proof.
56//! - Bounded buffering throughout ([`SERVER_MAX_LINE_BYTES`],
57//!   [`SERVER_EVENT_CHANNEL_CAPACITY`]) — same 16MiB-class discipline P5-2
58//!   established for `crate::mcp`'s SSE reader, reused here rather than
59//!   re-derived.
60//! - Graceful shutdown: the `shutdown` RPC method stops the stdio loop and
61//!   the HTTP accept loop alike (both select on the same
62//!   [`RpcEngine::wait_for_shutdown`]) — no orphaned listener/accept task
63//!   survives a `shutdown` call, mirroring P5-3/P5-6's drop-abort
64//!   discipline for background work.
65
66use std::collections::{BTreeMap, HashMap, VecDeque};
67#[cfg(feature = "adapter-api")]
68use std::net::SocketAddr;
69use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
70use std::sync::{Arc, Mutex as StdMutex};
71
72use async_trait::async_trait;
73#[cfg(feature = "adapter-api")]
74use futures::{SinkExt, StreamExt};
75use serde_json::{json, Value};
76use tokio::io::AsyncBufRead;
77#[cfg(feature = "adapter-api")]
78use tokio::io::{AsyncRead, AsyncReadExt};
79#[cfg(feature = "adapter-api")]
80use tokio::io::{AsyncWrite, AsyncWriteExt};
81#[cfg(feature = "adapter-api")]
82use tokio::net::TcpListener;
83#[cfg(feature = "adapter-api")]
84use tokio::sync::mpsc;
85use tokio::sync::{broadcast, Mutex, Notify, RwLock};
86
87/// How often an idle event stream carries an SSE comment, well inside the body
88/// timeouts of common HTTP clients (Node's fetch: 300 s).
89const EVENT_STREAM_KEEPALIVE: std::time::Duration = std::time::Duration::from_secs(15);
90
91use crate::agent::SteerInbox;
92use crate::frontend::{
93    FrontendActions, FrontendApprovalDecision, FrontendAttachSnapshot, FrontendAttachment,
94    FrontendCommandDescriptor, FrontendConnectionState, FrontendDisplayCapabilities, FrontendEvent,
95    FrontendOperationDescriptor, FrontendOperationInvocation, FrontendOperationKind,
96    FrontendOperationResult, FrontendProjectionState, FrontendRequest, FrontendRequestKind,
97    FrontendResponse, FrontendRuntime, FrontendRuntimeDescriptor, FrontendRuntimeError,
98    FrontendRuntimeMetadata, FrontendTurnState, FRONTEND_EVENT_SCHEMA_VERSION,
99    FRONTEND_REPLAY_CAPACITY, FRONTEND_RUNTIME_SCHEMA_VERSION,
100};
101use crate::mcp::{
102    ElicitationAction, ElicitationRequest, ElicitationResponse, McpElicitationHandler,
103};
104use crate::permissions::{ApprovalOutcome, ApprovalRequest, PermissionsApprovalHandler};
105pub use crate::sdk::RuntimeSubmitError;
106use crate::sdk::SdkAgent;
107#[cfg(feature = "adapter-api")]
108use crate::{CoordinatedRuntime, CoordinatedRuntimeClient, RuntimeAuthorization, RuntimeClientId};
109use supercode_interchange::ChatMessage;
110
111/// Bounded broadcast capacity for the event-notification channel — mirrors
112/// `crate::mcp::MCP_SSE_CHANNEL_CAPACITY`'s bounded-buffering discipline
113/// (P5-2): a slow/absent subscriber can never make the sender block or
114/// grow memory unboundedly; a lagging receiver just misses old events
115/// (`broadcast::error::RecvError::Lagged`) rather than stalling the agent
116/// loop or accumulating unbounded backlog.
117pub const SERVER_EVENT_CHANNEL_CAPACITY: usize = 1024;
118
119/// Maximum canonical messages retained for late-attaching frontends. This
120/// matches the `history` RPC limit and prevents the lock-independent snapshot
121/// from duplicating an arbitrarily large agent transcript.
122pub(crate) const SERVER_HISTORY_CAPACITY: usize = 200;
123
124/// Maximum accepted line/body length (bytes) for both the stdio JSONL-RPC
125/// reader and the HTTP transport's request line/headers/body — the same
126/// 16MiB-class cap P5-2 established for `crate::mcp`'s SSE frame reader
127/// (`MCP_MAX_SSE_FRAME_BYTES`), reused here so an adversarial or simply
128/// broken client can never make either transport buffer an unbounded
129/// amount of data in memory.
130pub const SERVER_MAX_LINE_BYTES: usize = 16 * 1024 * 1024;
131
132/// Cap on HTTP header line COUNT (independent of [`SERVER_MAX_LINE_BYTES`],
133/// which only bounds any one line's length) — without this, a client could
134/// send an unbounded NUMBER of small, individually-under-cap header lines
135/// and still exhaust memory over one connection.
136#[cfg(feature = "adapter-api")]
137const MAX_HEADER_LINES: usize = 200;
138
139/// One JSONL-RPC request line a client sends: `{"id", "method", "params"}`.
140/// `params` defaults to `null` when omitted (a method that takes no
141/// arguments, e.g. `status`/`shutdown`, never requires callers to spell out
142/// `"params": null}` explicitly).
143#[derive(Debug, Clone, serde::Deserialize)]
144pub struct RpcRequest {
145    /// Caller-chosen correlation id, echoed back verbatim on the matching
146    /// response — never interpreted, so any JSON value (string, number,
147    /// null) a caller likes works.
148    pub id: Value,
149    /// The method name (`submit` | `interrupt` | `status` | `shutdown`).
150    pub method: String,
151    /// Method-specific arguments; `submit` reads `params.prompt`.
152    #[serde(default)]
153    pub params: Value,
154}
155
156/// Build a `{"id", "result"}` response line.
157fn rpc_ok(id: Value, result: Value) -> Value {
158    json!({"id": id, "result": result})
159}
160
161/// Build a `{"id", "error": {"code","message"}}` response line. `code`
162/// follows JSON-RPC 2.0's reserved-range convention where a natural fit
163/// exists (`-32700` parse error, `-32601` method not found, `-32602`
164/// invalid params) purely as a familiar, self-documenting convention — this
165/// protocol does not otherwise claim JSON-RPC 2.0 compliance (no
166/// `"jsonrpc":"2.0"` envelope; see the module doc's minimal wire shape).
167fn rpc_error(id: Value, code: i32, message: impl Into<String>) -> Value {
168    json!({"id": id, "error": {"code": code, "message": message.into()}})
169}
170
171fn sdk_runtime_rpc_error(id: Value, code: i32, error: &FrontendRuntimeError) -> Value {
172    let code = match error.code() {
173        crate::SdkErrorCode::Unauthenticated => -32030,
174        crate::SdkErrorCode::Unauthorized => -32031,
175        crate::SdkErrorCode::ControllerRequired => -32032,
176        crate::SdkErrorCode::LeaseExpired => -32033,
177        _ => code,
178    };
179    let mut envelope = json!({
180        "id": id,
181        "error": {
182            "code": code,
183            "name": error.code(),
184            "operation": error.operation(),
185            "message": error.to_string(),
186        }
187    });
188    if let Some(detail) = envelope.get_mut("error").and_then(Value::as_object_mut) {
189        match error {
190            FrontendRuntimeError::Unauthorized { permission } => {
191                detail.insert("permission".into(), Value::String(permission.clone()));
192            }
193            FrontendRuntimeError::ControllerRequired {
194                holder,
195                expires_at_ms,
196            } => {
197                if let Some(holder) = holder {
198                    detail.insert("holder".into(), Value::String(holder.clone()));
199                }
200                if let Some(expires_at_ms) = expires_at_ms {
201                    detail.insert("expiresAtMs".into(), json!(expires_at_ms));
202                }
203            }
204            _ => {}
205        }
206    }
207    envelope
208}
209
210/// Read one line (trailing `\n`/`\r\n` stripped) from `reader`, capped at
211/// `cap` bytes — mirrors `crate::mcp`'s `SseLineAccumulator` bounded-
212/// buffering discipline (P5-2). Returns `Ok(None)` at a clean EOF with no
213/// partial line pending. On an over-cap line, the REST of that oversized
214/// line is drained and discarded (up to the next `\n`) so the stream
215/// resyncs at the next real line boundary instead of desyncing forever,
216/// and `Err` is returned naming the cap.
217async fn read_bounded_line<R>(reader: &mut R, cap: usize) -> std::io::Result<Option<String>>
218where
219    R: AsyncBufRead + Unpin,
220{
221    use tokio::io::AsyncBufReadExt;
222    let mut out: Vec<u8> = Vec::new();
223    loop {
224        let buf = reader.fill_buf().await?;
225        if buf.is_empty() {
226            return Ok(if out.is_empty() {
227                None
228            } else {
229                Some(strip_crlf(out))
230            });
231        }
232        if let Some(pos) = buf.iter().position(|&b| b == b'\n') {
233            if out.len() + pos > cap {
234                reader.consume(pos + 1);
235                return Err(std::io::Error::new(
236                    std::io::ErrorKind::InvalidData,
237                    format!("line exceeded {cap} byte cap"),
238                ));
239            }
240            out.extend_from_slice(&buf[..pos]);
241            reader.consume(pos + 1);
242            return Ok(Some(strip_crlf(out)));
243        }
244        let take = buf.len();
245        if out.len() + take > cap {
246            reader.consume(take);
247            // Drain/discard the rest of this oversized line so a future
248            // read starts at the next real line boundary.
249            loop {
250                let b = reader.fill_buf().await?;
251                if b.is_empty() {
252                    break;
253                }
254                if let Some(p) = b.iter().position(|&x| x == b'\n') {
255                    reader.consume(p + 1);
256                    break;
257                }
258                let n = b.len();
259                reader.consume(n);
260            }
261            return Err(std::io::Error::new(
262                std::io::ErrorKind::InvalidData,
263                format!("line exceeded {cap} byte cap"),
264            ));
265        }
266        out.extend_from_slice(buf);
267        reader.consume(take);
268    }
269}
270
271fn strip_crlf(mut v: Vec<u8>) -> String {
272    if v.last() == Some(&b'\r') {
273        v.pop();
274    }
275    String::from_utf8_lossy(&v).into_owned()
276}
277
278/// Constant-time byte-slice comparison (avoids leaking the bearer token
279/// through a timing side-channel on `==`) — small, self-contained, no new
280/// dependency for one comparison.
281#[cfg(feature = "adapter-api")]
282fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
283    if a.len() != b.len() {
284        return false;
285    }
286    let mut diff = 0u8;
287    for (x, y) in a.iter().zip(b.iter()) {
288        diff |= x ^ y;
289    }
290    diff == 0
291}
292
293/// Mint a random per-session bearer token (32 bytes, hex-encoded) for the
294/// HTTP transport, when the operator hasn't configured a fixed
295/// `capabilities.server.token`. Uses `getrandom` (already resolved
296/// transitively via `reqwest`'s rustls/ring stack; promoted to a direct
297/// dependency here so this crate can call it directly, rather than rolling
298/// a hand-written PRNG for a value that must actually be unguessable).
299pub fn generate_token() -> String {
300    let mut bytes = [0u8; 32];
301    // `getrandom::getrandom` only fails if the OS entropy source itself is
302    // unavailable/misconfigured — effectively never on a real target this
303    // crate supports. Falling back to a process-time/PID-derived value
304    // would be a WORSE, easily-guessable token, so a failure here is
305    // treated as fatal (panic) rather than silently minting a weak secret.
306    getrandom::getrandom(&mut bytes).expect("OS entropy source for the server bearer token");
307    bytes.iter().map(|b| format!("{b:02x}")).collect()
308}
309
310/// The `on_turn_complete` hook's type — factored out (clippy
311/// `type_complexity`) rather than spelled out at both
312/// [`RpcEngine`]'s field and [`RpcEngine::new`]'s parameter.
313type TurnCompleteHook = Box<dyn Fn(&SdkAgent) + Send + Sync>;
314
315/// Protocol-neutral snapshot of one SDK-owned agent runtime.
316///
317/// Transport adapters (CLI/RPC/HTTP/ACP) project this value into their own
318/// wire shapes instead of each inventing a separate definition of "busy".
319#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
320pub struct RuntimeStatus {
321    /// Stable identity chosen by the embedding surface/session store.
322    pub session_id: String,
323    /// Model currently backing the agent.
324    pub model: String,
325    /// Whether a user or scheduler turn currently owns the agent loop.
326    pub busy: bool,
327    /// Whether graceful shutdown has been requested.
328    pub shutting_down: bool,
329}
330
331type PendingFrontendResponses = StdMutex<
332    HashMap<
333        u64,
334        (
335            FrontendRequestKind,
336            std::sync::mpsc::Sender<AcceptedFrontendResponse>,
337        ),
338    >,
339>;
340
341struct AcceptedFrontendResponse {
342    response: FrontendResponse,
343    /// The blocked handler may resume execution only after the canonical
344    /// resolution event has entered the sequenced frontend projection.
345    published: std::sync::mpsc::Receiver<()>,
346}
347
348/// Runtime-owned interactive request broker. It publishes complete request
349/// payloads into the same sequenced event stream and resolves each id once.
350struct FrontendRequestBroker {
351    next_id: std::sync::atomic::AtomicU64,
352    pending: PendingFrontendResponses,
353    transport: StdMutex<Option<FrontendRequestTransport>>,
354}
355
356#[derive(Clone)]
357struct FrontendRequestTransport {
358    events: broadcast::Sender<FrontendEvent>,
359    state: Arc<StdMutex<FrontendProjectionState>>,
360}
361
362impl FrontendRequestBroker {
363    fn new() -> Arc<Self> {
364        Arc::new(Self {
365            next_id: std::sync::atomic::AtomicU64::new(1),
366            pending: StdMutex::new(HashMap::new()),
367            transport: StdMutex::new(None),
368        })
369    }
370
371    fn bind(
372        &self,
373        events: broadcast::Sender<FrontendEvent>,
374        state: Arc<StdMutex<FrontendProjectionState>>,
375    ) {
376        *self
377            .transport
378            .lock()
379            .unwrap_or_else(std::sync::PoisonError::into_inner) =
380            Some(FrontendRequestTransport { events, state });
381    }
382
383    fn transport(&self) -> Option<FrontendRequestTransport> {
384        self.transport
385            .lock()
386            .unwrap_or_else(std::sync::PoisonError::into_inner)
387            .clone()
388    }
389
390    fn publish(&self, request: &FrontendRequest) -> bool {
391        self.publish_payload(json!({"type": "request", "request": request}))
392    }
393
394    fn publish_payload(&self, payload: Value) -> bool {
395        let Some(transport) = self.transport() else {
396            return false;
397        };
398        let event = {
399            let mut state = transport
400                .state
401                .lock()
402                .unwrap_or_else(std::sync::PoisonError::into_inner);
403            let event = FrontendEvent::new(state.next_sequence, payload);
404            state.next_sequence = state.next_sequence.saturating_add(1);
405            state.replay.push_back(event.clone());
406            while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
407                state.replay.pop_front();
408            }
409            event
410        };
411        transport.events.send(event).is_ok()
412    }
413
414    fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
415        let request_id = response.request_id();
416        let response_kind = match &response {
417            FrontendResponse::Approval { .. } => FrontendRequestKind::Approval,
418            FrontendResponse::Elicitation { .. } => FrontendRequestKind::Elicitation,
419            FrontendResponse::Other { .. } => FrontendRequestKind::Other,
420        };
421        let mut pending = self
422            .pending
423            .lock()
424            .unwrap_or_else(std::sync::PoisonError::into_inner);
425        let expected = pending
426            .get(&request_id)
427            .map(|(kind, _)| *kind)
428            .ok_or(FrontendRuntimeError::UnknownRequest(request_id))?;
429        if expected != response_kind {
430            return Err(FrontendRuntimeError::InvalidResponse(format!(
431                "request {request_id} expects {expected:?}, got {response_kind:?}"
432            )));
433        }
434        let (_, sender) = pending
435            .remove(&request_id)
436            .ok_or(FrontendRuntimeError::UnknownRequest(request_id))?;
437        drop(pending);
438        let payload = json!({
439            "type": "request_resolved",
440            "request_id": request_id,
441            "response": &response,
442        });
443        let (published_tx, published_rx) = std::sync::mpsc::channel();
444        sender
445            .send(AcceptedFrontendResponse {
446                response,
447                published: published_rx,
448            })
449            .map_err(|_| FrontendRuntimeError::UnknownRequest(request_id))?;
450        self.publish_payload(payload);
451        let _ = published_tx.send(());
452        Ok(())
453    }
454
455    fn ask_approval(
456        &self,
457        req: &ApprovalRequest<'_>,
458        child: Option<(&str, &Arc<StdMutex<Vec<crate::subagents::QueuedApproval>>>)>,
459    ) -> ApprovalOutcome {
460        // Record with no outcome first: while this call is in flight the
461        // entry IS a pending request, and the decision is written back onto
462        // the same record below so a reader never has to guess.
463        let queued = child.and_then(|(child_agent_id, queue)| {
464            crate::subagents::queue_approval(
465                queue,
466                crate::subagents::QueuedApproval {
467                    child_agent_id: child_agent_id.to_string(),
468                    tool: req.tool.to_string(),
469                    subject: req.subject.map(String::from),
470                    queued_at_ms: std::time::SystemTime::now()
471                        .duration_since(std::time::UNIX_EPOCH)
472                        .map(|duration| duration.as_millis() as i64)
473                        .unwrap_or_default(),
474                    outcome: None,
475                },
476            )
477            .map(|index| (queue.clone(), index))
478        });
479        let outcome = self.decide_approval(req, child.map(|(id, _)| id));
480        if let Some((queue, index)) = queued {
481            crate::subagents::record_queued_outcome(&queue, index, outcome.into());
482        }
483        outcome
484    }
485
486    fn decide_approval(
487        &self,
488        req: &ApprovalRequest<'_>,
489        child_agent_id: Option<&str>,
490    ) -> ApprovalOutcome {
491        // No observer means no human can answer. Preserve the server's
492        // existing fail-closed, non-blocking headless behavior.
493        let Some(transport) = self.transport() else {
494            return ApprovalOutcome::Deny;
495        };
496        if transport.events.receiver_count() == 0 {
497            return ApprovalOutcome::Deny;
498        }
499        let id = self.next_id.fetch_add(1, Ordering::SeqCst);
500        let mut payload = json!({
501            "tool": req.tool,
502            "subject": req.subject,
503            "raw_args": req.raw_args,
504        });
505        if let Some(child_agent_id) = child_agent_id {
506            payload["child_agent_id"] = Value::String(child_agent_id.to_string());
507        }
508        let request = FrontendRequest {
509            id,
510            kind: FrontendRequestKind::Approval,
511            payload,
512        };
513        let (tx, rx) = std::sync::mpsc::channel();
514        self.pending
515            .lock()
516            .unwrap_or_else(std::sync::PoisonError::into_inner)
517            .insert(id, (FrontendRequestKind::Approval, tx));
518        if !self.publish(&request) {
519            self.pending
520                .lock()
521                .unwrap_or_else(std::sync::PoisonError::into_inner)
522                .remove(&id);
523            return ApprovalOutcome::Deny;
524        }
525        let wait_for_response = || loop {
526            match rx.recv_timeout(std::time::Duration::from_millis(100)) {
527                Ok(accepted) => {
528                    let _ = accepted.published.recv();
529                    let FrontendResponse::Approval { decision, .. } = accepted.response else {
530                        return ApprovalOutcome::Deny;
531                    };
532                    return match decision {
533                        FrontendApprovalDecision::Deny => ApprovalOutcome::Deny,
534                        FrontendApprovalDecision::Allow => ApprovalOutcome::Allow,
535                        FrontendApprovalDecision::AllowForSession => {
536                            ApprovalOutcome::AllowForSession
537                        }
538                    };
539                }
540                Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
541                    return ApprovalOutcome::Deny;
542                }
543                Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
544                    if self
545                        .transport()
546                        .map(|transport| transport.events.receiver_count() == 0)
547                        .unwrap_or(true)
548                    {
549                        self.pending
550                            .lock()
551                            .unwrap_or_else(std::sync::PoisonError::into_inner)
552                            .remove(&id);
553                        return ApprovalOutcome::Deny;
554                    }
555                }
556            }
557        };
558        if tokio::runtime::Handle::try_current()
559            .map(|handle| handle.runtime_flavor() == tokio::runtime::RuntimeFlavor::MultiThread)
560            .unwrap_or(false)
561        {
562            tokio::task::block_in_place(wait_for_response)
563        } else {
564            wait_for_response()
565        }
566    }
567
568    async fn ask_elicitation(self: Arc<Self>, req: &ElicitationRequest) -> ElicitationResponse {
569        let cancel = || ElicitationResponse {
570            action: ElicitationAction::Cancel,
571            content: None,
572        };
573        let Some(transport) = self.transport() else {
574            return cancel();
575        };
576        if transport.events.receiver_count() == 0 {
577            return cancel();
578        }
579        let id = self.next_id.fetch_add(1, Ordering::SeqCst);
580        let request = FrontendRequest {
581            id,
582            kind: FrontendRequestKind::Elicitation,
583            payload: json!({
584                "message": req.message,
585                "requested_schema": req.requested_schema,
586            }),
587        };
588        let (tx, rx) = std::sync::mpsc::channel();
589        self.pending
590            .lock()
591            .unwrap_or_else(std::sync::PoisonError::into_inner)
592            .insert(id, (FrontendRequestKind::Elicitation, tx));
593        if !self.publish(&request) {
594            self.pending
595                .lock()
596                .unwrap_or_else(std::sync::PoisonError::into_inner)
597                .remove(&id);
598            return cancel();
599        }
600        let broker = self.clone();
601        tokio::task::spawn_blocking(move || loop {
602            match rx.recv_timeout(std::time::Duration::from_millis(100)) {
603                Ok(accepted) => {
604                    let _ = accepted.published.recv();
605                    let FrontendResponse::Elicitation {
606                        action, content, ..
607                    } = accepted.response
608                    else {
609                        return cancel();
610                    };
611                    return ElicitationResponse {
612                        action: match action {
613                            crate::frontend::FrontendElicitationAction::Accept => {
614                                ElicitationAction::Accept
615                            }
616                            crate::frontend::FrontendElicitationAction::Decline => {
617                                ElicitationAction::Decline
618                            }
619                            crate::frontend::FrontendElicitationAction::Cancel => {
620                                ElicitationAction::Cancel
621                            }
622                        },
623                        content,
624                    };
625                }
626                Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => return cancel(),
627                Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
628                    if broker
629                        .transport()
630                        .map(|transport| transport.events.receiver_count() == 0)
631                        .unwrap_or(true)
632                    {
633                        broker
634                            .pending
635                            .lock()
636                            .unwrap_or_else(std::sync::PoisonError::into_inner)
637                            .remove(&id);
638                        return cancel();
639                    }
640                }
641            }
642        })
643        .await
644        .unwrap_or_else(|_| cancel())
645    }
646}
647
648struct FrontendApprovalHandler(Arc<FrontendRequestBroker>);
649
650impl PermissionsApprovalHandler for FrontendApprovalHandler {
651    fn ask(&self, req: &ApprovalRequest<'_>) -> ApprovalOutcome {
652        self.0.ask_approval(req, None)
653    }
654}
655
656struct FrontendChildApprovalHandler {
657    broker: Arc<FrontendRequestBroker>,
658    child_agent_id: String,
659    queue: Arc<StdMutex<Vec<crate::subagents::QueuedApproval>>>,
660}
661
662impl PermissionsApprovalHandler for FrontendChildApprovalHandler {
663    fn ask(&self, req: &ApprovalRequest<'_>) -> ApprovalOutcome {
664        self.broker
665            .ask_approval(req, Some((&self.child_agent_id, &self.queue)))
666    }
667}
668
669/// Pre-runtime bridge for MCP clients that must receive their elicitation
670/// handler before they are consumed into agent tool registration.
671#[derive(Clone)]
672pub struct FrontendRequestBridge {
673    broker: Arc<FrontendRequestBroker>,
674}
675
676impl FrontendRequestBridge {
677    /// Create an unbound bridge. Pass it to
678    /// [`RpcEngine::new_named_with_frontend_bridge`] after MCP registration.
679    pub fn new() -> Self {
680        Self {
681            broker: FrontendRequestBroker::new(),
682        }
683    }
684
685    /// Handler installed on every interactive MCP client before registration.
686    pub fn elicitation_handler(&self) -> Arc<dyn McpElicitationHandler> {
687        Arc::new(FrontendElicitationHandler(self.broker.clone()))
688    }
689}
690
691impl Default for FrontendRequestBridge {
692    fn default() -> Self {
693        Self::new()
694    }
695}
696
697struct FrontendElicitationHandler(Arc<FrontendRequestBroker>);
698
699#[async_trait]
700impl McpElicitationHandler for FrontendElicitationHandler {
701    async fn handle(&self, request: &ElicitationRequest) -> ElicitationResponse {
702        self.0.clone().ask_elicitation(request).await
703    }
704}
705
706/// The out-of-process RPC driver: wraps one already-constructed
707/// [`crate::Agent`] with the `submit`/`interrupt`/`status`/`shutdown`
708/// method set (§ module doc). Shared by both transports ([`run_stdio`],
709/// [`run_http`]) so the method semantics — including the fail-closed
710/// permission behavior — can never drift between them.
711pub struct RpcEngine {
712    agent: Mutex<SdkAgent>,
713    /// Last canonical transcript observed at a turn boundary. Frontends must
714    /// be able to attach and replay prior history while `submit` holds the
715    /// agent lock for an active turn, so reads use this independent snapshot.
716    history_snapshot: RwLock<Vec<ChatMessage>>,
717    session_id: String,
718    /// The model label, captured once at construction so `status` never
719    /// needs to lock `agent` (which `submit` holds for the WHOLE turn) —
720    /// `status` must stay answerable while a turn is in flight.
721    model: String,
722    busy: Arc<AtomicBool>,
723    /// The in-flight turn's cancellation handle, if any — see
724    /// [`Self::handle_submit`]/[`Self::handle_interrupt`]'s doc comments
725    /// for why this is a fresh [`Notify`] per turn rather than one shared
726    /// instance (`notify_one`'s stored-permit semantics only give the
727    /// correctness guarantee this needs when each turn gets a clean one).
728    current_cancel: Arc<StdMutex<Option<Arc<Notify>>>>,
729    /// Signals that the terminal event for an interrupted/completed turn is
730    /// published and its canonical history snapshot is stable.
731    turn_finished: Arc<Notify>,
732    /// Shared control queue owned by `Agent` but writable without waiting for
733    /// the active turn's long-held async agent lock.
734    steer_queue: Arc<StdMutex<SteerInbox>>,
735    events: broadcast::Sender<Value>,
736    /// Sequenced frontend events used for atomic replay/live attachment.
737    frontend_events: broadcast::Sender<FrontendEvent>,
738    /// Short-held projection lock. It is never held across model/tool I/O.
739    frontend_state: Arc<StdMutex<FrontendProjectionState>>,
740    frontend_metadata: FrontendRuntimeMetadata,
741    frontend_active_modules: Vec<String>,
742    frontend_commands: Vec<FrontendCommandDescriptor>,
743    frontend_operations: Vec<FrontendOperationDescriptor>,
744    frontend_requests: Option<Arc<FrontendRequestBroker>>,
745    shutdown: Notify,
746    shutting_down: AtomicBool,
747    /// Explicit owner shutdown seals new model-loop claims before it
748    /// interrupts the active one. Plain transport EOF only raises
749    /// `shutting_down` so already-buffered stdio requests can still flush.
750    accepting_submits: AtomicBool,
751    /// Serializes explicit shutdown barriers so every concurrent caller
752    /// returns only after the same admitted turn.
753    shutdown_barrier: Mutex<()>,
754    /// Fires after every SUCCESSFUL `submit` (never on an errored/
755    /// interrupted turn — see the call site), with the agent still locked
756    /// so the hook sees fully-consistent state (e.g. `agent.history()`).
757    /// The CLI installs session persistence/auto-titling here — this
758    /// module itself has no opinion on session storage.
759    on_turn_complete: Option<TurnCompleteHook>,
760}
761
762/// Owns every externally visible piece of an SDK submit claim from the
763/// instant the claim succeeds until the submit reaches a terminal boundary.
764/// The agent loop closes the ordinary final-answer steering boundary
765/// atomically with its last drain; this outer guard also restores steering,
766/// cancellation, busy state, and lifecycle waiters when the public submit
767/// future is dropped or fails before `Agent::run_loop` is entered.
768struct SdkSubmitClaim {
769    inbox: Arc<StdMutex<SteerInbox>>,
770    busy: Arc<AtomicBool>,
771    cancel: Arc<Notify>,
772    current_cancel: Arc<StdMutex<Option<Arc<Notify>>>>,
773    turn_finished: Arc<Notify>,
774    frontend_events: broadcast::Sender<FrontendEvent>,
775    frontend_state: Arc<StdMutex<FrontendProjectionState>>,
776    lifecycle_started: bool,
777}
778
779impl SdkSubmitClaim {
780    fn mark_lifecycle_started(&mut self) {
781        self.lifecycle_started = true;
782    }
783
784    fn mark_lifecycle_finished(&mut self) {
785        self.lifecycle_started = false;
786    }
787}
788
789impl Drop for SdkSubmitClaim {
790    fn drop(&mut self) {
791        // A caller may cancel the public `submit` future after the SDK has
792        // exposed `turn_started` but before `submit_claimed` can publish its
793        // ordinary terminal event. Close that exact lifecycle before making
794        // the runtime idle so replay and live frontends cannot remain busy on
795        // an abandoned turn. There is no await between the normal terminal
796        // publication and disarming this fallback, so exactly one terminal
797        // event is observable for every started claim.
798        if self.lifecycle_started {
799            let event = {
800                let mut state = self
801                    .frontend_state
802                    .lock()
803                    .unwrap_or_else(std::sync::PoisonError::into_inner);
804                let event = FrontendEvent::new(
805                    state.next_sequence,
806                    json!({
807                        "type": "turn_interrupted",
808                        "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
809                    }),
810                );
811                state.next_sequence = state.next_sequence.saturating_add(1);
812                state.replay.push_back(event.clone());
813                while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
814                    state.replay.pop_front();
815                }
816                event
817            };
818            let _ = self.frontend_events.send(event);
819        }
820        self.inbox
821            .lock()
822            .unwrap_or_else(std::sync::PoisonError::into_inner)
823            .close();
824        *self
825            .current_cancel
826            .lock()
827            .unwrap_or_else(std::sync::PoisonError::into_inner) = None;
828        self.busy.store(false, Ordering::SeqCst);
829        self.turn_finished.notify_waiters();
830        self.turn_finished.notify_one();
831    }
832}
833
834impl RpcEngine {
835    /// Wrap `agent` (already fully built by the caller — same `Config`,
836    /// same permission/sandbox posture as a local session) for out-of-
837    /// process driving. Installs its OWN event sink via
838    /// `Agent::set_event_sink`, overwriting whatever sink `agent` may
839    /// already have had wired (callers of this module drive an agent
840    /// exclusively through the RPC surface, so there is never a second,
841    /// competing consumer of its events).
842    pub fn new(
843        agent: impl Into<SdkAgent>,
844        on_turn_complete: Option<TurnCompleteHook>,
845    ) -> Arc<Self> {
846        let agent = agent.into();
847        let session_id = agent
848            .session_name()
849            .map(str::to_owned)
850            .unwrap_or_else(|| format!("supercode-{}", std::process::id()));
851        Self::new_named(agent, session_id, on_turn_complete)
852    }
853
854    /// Construct the canonical SDK runtime with an explicit durable session
855    /// identity.  Every frontend must use this identity when referring to the
856    /// same live agent; transport-local connection ids are not session ids.
857    pub fn new_named(
858        agent: impl Into<SdkAgent>,
859        session_id: impl Into<String>,
860        on_turn_complete: Option<TurnCompleteHook>,
861    ) -> Arc<Self> {
862        Self::new_named_with_frontend_metadata(
863            agent.into(),
864            session_id,
865            FrontendRuntimeMetadata::default(),
866            on_turn_complete,
867        )
868    }
869
870    /// Construct the canonical runtime with explicit source-harness and
871    /// emulation-profile identity for every attached frontend.
872    pub fn new_named_with_frontend_metadata(
873        agent: impl Into<SdkAgent>,
874        session_id: impl Into<String>,
875        frontend_metadata: FrontendRuntimeMetadata,
876        on_turn_complete: Option<TurnCompleteHook>,
877    ) -> Arc<Self> {
878        Self::build(
879            agent.into(),
880            session_id.into(),
881            frontend_metadata,
882            None,
883            on_turn_complete,
884        )
885    }
886
887    /// Construct a canonical runtime whose attached frontends may answer
888    /// policy-authorized approval requests. Existing constructors retain the
889    /// historical fail-closed headless behavior and report `respond=false`.
890    pub fn new_named_with_frontend_requests(
891        agent: impl Into<SdkAgent>,
892        session_id: impl Into<String>,
893        frontend_metadata: FrontendRuntimeMetadata,
894        on_turn_complete: Option<TurnCompleteHook>,
895    ) -> Arc<Self> {
896        let bridge = FrontendRequestBridge::new();
897        Self::new_named_with_frontend_bridge(
898            agent.into(),
899            session_id,
900            frontend_metadata,
901            bridge,
902            on_turn_complete,
903        )
904    }
905
906    /// Bind a pre-created request bridge after its elicitation handler has
907    /// been installed on MCP clients.
908    pub fn new_named_with_frontend_bridge(
909        agent: impl Into<SdkAgent>,
910        session_id: impl Into<String>,
911        frontend_metadata: FrontendRuntimeMetadata,
912        bridge: FrontendRequestBridge,
913        on_turn_complete: Option<TurnCompleteHook>,
914    ) -> Arc<Self> {
915        Self::build(
916            agent.into(),
917            session_id.into(),
918            frontend_metadata,
919            Some(bridge.broker),
920            on_turn_complete,
921        )
922    }
923
924    fn build(
925        mut agent: SdkAgent,
926        session_id: String,
927        frontend_metadata: FrontendRuntimeMetadata,
928        frontend_requests: Option<Arc<FrontendRequestBroker>>,
929        on_turn_complete: Option<TurnCompleteHook>,
930    ) -> Arc<Self> {
931        let (tx, _rx) = broadcast::channel(SERVER_EVENT_CHANNEL_CAPACITY);
932        let events_tx = tx.clone();
933        let (frontend_tx, _frontend_rx) = broadcast::channel(SERVER_EVENT_CHANNEL_CAPACITY);
934        let frontend_events_tx = frontend_tx.clone();
935        let model = agent.config().model.clone();
936        let steer_queue = agent.inner().steer_queue_handle();
937        let history_snapshot = bounded_history_snapshot(agent.history());
938        let frontend_state = Arc::new(StdMutex::new(FrontendProjectionState {
939            history: history_snapshot.clone(),
940            history_cursor: 0,
941            next_sequence: 1,
942            replay: VecDeque::new(),
943        }));
944        if let Some(broker) = &frontend_requests {
945            broker.bind(frontend_tx.clone(), frontend_state.clone());
946            let legacy_broker = broker.clone();
947            agent
948                .inner_mut()
949                .set_legacy_approval_handler(Box::new(move |call| {
950                    let Ok(raw_args) = call.function.parsed_arguments() else {
951                        return false;
952                    };
953                    let subject = raw_args
954                        .get("command")
955                        .or_else(|| raw_args.get("path"))
956                        .or_else(|| raw_args.get("file_path"))
957                        .or_else(|| raw_args.get("patch"))
958                        .and_then(Value::as_str);
959                    matches!(
960                        legacy_broker.ask_approval(
961                            &ApprovalRequest {
962                                tool: &call.function.name,
963                                subject,
964                                raw_args: &raw_args,
965                            },
966                            None,
967                        ),
968                        ApprovalOutcome::Allow | ApprovalOutcome::AllowForSession
969                    )
970                }));
971            agent
972                .inner_mut()
973                .set_permissions_approval_handler(FrontendApprovalHandler(broker.clone()));
974            // BP-3 (§2 module 6 `tools.question`): `ask_user` asks through
975            // the SAME broker — the request is published into the sequenced
976            // frontend event stream and the turn blocks on it until
977            // `respond` answers, exactly like the approval above. Without
978            // this the tool would be deny-default even on an attached
979            // frontend.
980            agent
981                .inner_mut()
982                .set_user_question_handler(Arc::new(FrontendElicitationHandler(broker.clone())));
983            let broker = broker.clone();
984            agent
985                .inner_mut()
986                .set_child_approval_handler_factory(move |child_agent_id, queue| {
987                    Arc::new(FrontendChildApprovalHandler {
988                        broker: broker.clone(),
989                        child_agent_id,
990                        queue,
991                    }) as Arc<dyn PermissionsApprovalHandler>
992                });
993        }
994        let event_frontend_state = frontend_state.clone();
995        let frontend_active_modules = agent
996            .config()
997            .module_activation
998            .iter()
999            .map(ToString::to_string)
1000            .collect();
1001        let mut frontend_operations = agent
1002            .config()
1003            .prompts
1004            .keys()
1005            .filter(|name| valid_frontend_command_name(name))
1006            .map(|name| FrontendOperationDescriptor {
1007                id: format!("prompt:{name}"),
1008                kind: FrontendOperationKind::Prompt,
1009                command: Some(FrontendCommandDescriptor {
1010                    name: name.clone(),
1011                    description: None,
1012                    argument_hint: Some("[arguments]".into()),
1013                }),
1014            })
1015            .collect::<Vec<_>>();
1016        frontend_operations.sort_by(|left, right| left.id.cmp(&right.id));
1017        // BP-4 (catalog:109): the one operation every runtime can always
1018        // answer — it reads the accounting it already computes for the
1019        // context guard, submits nothing, and mutates nothing. Advertised
1020        // unconditionally for that reason (it is not inferred from a module
1021        // set, which is what the "phantom operation" gate forbids), and
1022        // pushed AFTER the id sort so the prompt catalog's own ordering is
1023        // untouched.
1024        // BP-13 (catalog D9 "Mid-session model switching"): the model
1025        // control, advertised only when this runtime's config actually
1026        // allows a switch — the same "never advertise a phantom operation"
1027        // rule the context row above obeys, read from the one gate
1028        // (`[core.model_switch] allow_switch`) rather than inferred.
1029        if agent.config().model_switch_allow_switch {
1030            frontend_operations.push(FrontendOperationDescriptor {
1031                id: "model:switch".into(),
1032                kind: FrontendOperationKind::Model,
1033                command: Some(FrontendCommandDescriptor {
1034                    name: "model".into(),
1035                    description: Some("show or switch this session's model".into()),
1036                    argument_hint: Some("[model]".into()),
1037                }),
1038            });
1039        }
1040        frontend_operations.push(FrontendOperationDescriptor {
1041            id: "context:usage".into(),
1042            kind: FrontendOperationKind::Context,
1043            command: Some(FrontendCommandDescriptor {
1044                name: "context".into(),
1045                description: Some("context-window usage for this session".into()),
1046                argument_hint: None,
1047            }),
1048        });
1049        // Schema-v1 compatibility projection. New frontends use only the
1050        // typed operation catalog and never fall back to this list.
1051        let frontend_commands = frontend_operations
1052            .iter()
1053            .filter_map(|operation| operation.command.as_ref())
1054            .map(|command| FrontendCommandDescriptor {
1055                name: command.name.clone(),
1056                description: command.description.clone(),
1057                argument_hint: None,
1058            })
1059            .collect();
1060        agent.inner_mut().set_event_sink(Box::new(move |event| {
1061            // A `send` error here only means "no subscriber is currently
1062            // listening" (every receiver dropped) — never a reason to fail
1063            // the turn itself, so it's intentionally discarded.
1064            let payload = event.to_json();
1065            let _ = events_tx.send(payload.clone());
1066            let sequenced = {
1067                let mut state = event_frontend_state
1068                    .lock()
1069                    .unwrap_or_else(std::sync::PoisonError::into_inner);
1070                let event = FrontendEvent::new(state.next_sequence, payload);
1071                state.next_sequence = state.next_sequence.saturating_add(1);
1072                state.replay.push_back(event.clone());
1073                while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
1074                    state.replay.pop_front();
1075                }
1076                event
1077            };
1078            let _ = frontend_events_tx.send(sequenced);
1079        }));
1080        Arc::new(RpcEngine {
1081            agent: Mutex::new(agent),
1082            history_snapshot: RwLock::new(history_snapshot),
1083            session_id,
1084            model,
1085            busy: Arc::new(AtomicBool::new(false)),
1086            current_cancel: Arc::new(StdMutex::new(None)),
1087            turn_finished: Arc::new(Notify::new()),
1088            steer_queue,
1089            events: tx,
1090            frontend_events: frontend_tx,
1091            frontend_state,
1092            frontend_metadata,
1093            frontend_active_modules,
1094            frontend_commands,
1095            frontend_operations,
1096            frontend_requests,
1097            shutdown: Notify::new(),
1098            shutting_down: AtomicBool::new(false),
1099            accepting_submits: AtomicBool::new(true),
1100            shutdown_barrier: Mutex::new(()),
1101            on_turn_complete,
1102        })
1103    }
1104
1105    /// Subscribe to this engine's event-notification stream (already
1106    /// `AgentEvent::to_json`-projected) — each subscriber gets every event
1107    /// emitted from this point on, independent of any other subscriber.
1108    pub fn subscribe(&self) -> broadcast::Receiver<Value> {
1109        self.events.subscribe()
1110    }
1111
1112    /// Describe the SDK-owned runtime without locking the active agent turn.
1113    pub fn frontend_descriptor(&self) -> FrontendRuntimeDescriptor {
1114        FrontendRuntimeDescriptor {
1115            schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
1116            session_id: self.session_id.clone(),
1117            source_harness: self.frontend_metadata.source_harness.clone(),
1118            emulation_profile: self.frontend_metadata.emulation_profile.clone(),
1119            active_modules: self.frontend_active_modules.clone(),
1120            commands: self.frontend_commands.clone(),
1121            operations: self.frontend_operations.clone(),
1122            actions: FrontendActions {
1123                submit: true,
1124                interrupt: true,
1125                steer: true,
1126                respond: self.frontend_requests.is_some(),
1127                detach: true,
1128                // The canonical runtime supports close. Coordinated client
1129                // projections mask this unless the authenticated grant owns
1130                // the independent terminate capability.
1131                close: true,
1132            },
1133            display: FrontendDisplayCapabilities {
1134                event_kinds: vec![
1135                    "user_message".into(),
1136                    "turn_started".into(),
1137                    "turn_succeeded".into(),
1138                    "turn_interrupted".into(),
1139                    "turn_failed".into(),
1140                    "text_delta".into(),
1141                    "turn_completed".into(),
1142                    "tool_call_started".into(),
1143                    "tool_call_completed".into(),
1144                    "cache_warning".into(),
1145                    "usage".into(),
1146                    "background_output".into(),
1147                    "request".into(),
1148                    "request_resolved".into(),
1149                    "scheduled_prompt_started".into(),
1150                    "scheduled_prompt_deferred".into(),
1151                    "scheduled_prompt_completed".into(),
1152                    "scheduler_error".into(),
1153                ],
1154                opaque_fallback: true,
1155            },
1156            model: self.model.clone(),
1157            turn_state: if self.busy.load(Ordering::SeqCst) {
1158                FrontendTurnState::Busy
1159            } else {
1160                FrontendTurnState::Idle
1161            },
1162            connection_state: if self.is_shutting_down() {
1163                FrontendConnectionState::ShuttingDown
1164            } else {
1165                FrontendConnectionState::Connected
1166            },
1167            extensions: Default::default(),
1168        }
1169    }
1170
1171    /// Attach to one atomic history/replay/live boundary. The live receiver
1172    /// is created before the projection snapshot is locked; events racing the
1173    /// snapshot therefore appear either in replay or in the receiver, and
1174    /// [`FrontendAttachment::next_event`] removes any overlap by sequence.
1175    pub fn frontend_attach(
1176        &self,
1177        history_limit: usize,
1178    ) -> Result<FrontendAttachment, FrontendRuntimeError> {
1179        let live = self.frontend_subscribe();
1180        let snapshot = self.frontend_snapshot(history_limit)?;
1181        Ok(FrontendAttachment::new(
1182            snapshot.descriptor,
1183            snapshot.history,
1184            snapshot.history_cursor,
1185            snapshot.replay,
1186            live,
1187            None,
1188        ))
1189    }
1190
1191    /// Subscribe to sequenced frontend events. Transport adapters subscribe
1192    /// before taking [`Self::frontend_snapshot`] so boundary events cannot be
1193    /// missed.
1194    pub fn frontend_subscribe(&self) -> broadcast::Receiver<FrontendEvent> {
1195        self.frontend_events.subscribe()
1196    }
1197
1198    /// Capture the serializable history/replay half of a frontend attachment.
1199    pub fn frontend_snapshot(
1200        &self,
1201        history_limit: usize,
1202    ) -> Result<FrontendAttachSnapshot, FrontendRuntimeError> {
1203        let state = self
1204            .frontend_state
1205            .lock()
1206            .unwrap_or_else(std::sync::PoisonError::into_inner);
1207        let limit = history_limit.min(SERVER_HISTORY_CAPACITY);
1208        let start = state.history.len().saturating_sub(limit);
1209        let replay = state
1210            .replay
1211            .iter()
1212            .filter(|event| event.sequence > state.history_cursor)
1213            .cloned()
1214            .collect::<VecDeque<_>>();
1215        if let Some(first) = replay.front() {
1216            let expected = state.history_cursor.saturating_add(1);
1217            if first.sequence > expected {
1218                return Err(FrontendRuntimeError::ReplayGap(first.sequence - expected));
1219            }
1220        }
1221        Ok(FrontendAttachSnapshot {
1222            descriptor: self.frontend_descriptor(),
1223            history: state.history[start..].to_vec(),
1224            history_cursor: state.history_cursor,
1225            replay,
1226        })
1227    }
1228
1229    fn publish_frontend_payload(&self, payload: Value) {
1230        let event = {
1231            let mut state = self
1232                .frontend_state
1233                .lock()
1234                .unwrap_or_else(std::sync::PoisonError::into_inner);
1235            let event = FrontendEvent::new(state.next_sequence, payload);
1236            state.next_sequence = state.next_sequence.saturating_add(1);
1237            state.replay.push_back(event.clone());
1238            while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
1239                state.replay.pop_front();
1240            }
1241            event
1242        };
1243        let _ = self.frontend_events.send(event);
1244    }
1245
1246    /// Stable SDK session identity shared by every frontend.
1247    pub fn session_id(&self) -> &str {
1248        &self.session_id
1249    }
1250
1251    fn claim_submit(&self) -> Result<SdkSubmitClaim, RuntimeSubmitError> {
1252        // Hold the cancellation slot across admission. Shutdown seals first,
1253        // then takes this same lock through `interrupt`: it therefore either
1254        // wins before `busy` is claimed or observes the admitted turn's
1255        // installed token. There is no busy-without-cancel interval.
1256        let mut current_cancel = self
1257            .current_cancel
1258            .lock()
1259            .unwrap_or_else(std::sync::PoisonError::into_inner);
1260        if !self.accepting_submits.load(Ordering::SeqCst) {
1261            return Err(RuntimeSubmitError::Interrupted);
1262        }
1263        if self.busy.swap(true, Ordering::SeqCst) {
1264            return Err(RuntimeSubmitError::Busy);
1265        }
1266        // Close the race with a concurrent shutdown between the first state
1267        // check and ownership of `busy`; relinquish this claim without
1268        // exposing a turn when the shutdown seal won.
1269        if !self.accepting_submits.load(Ordering::SeqCst) {
1270            self.busy.store(false, Ordering::SeqCst);
1271            self.turn_finished.notify_waiters();
1272            return Err(RuntimeSubmitError::Interrupted);
1273        }
1274        let cancel = Arc::new(Notify::new());
1275        *current_cancel = Some(cancel.clone());
1276        drop(current_cancel);
1277        self.steer_queue
1278            .lock()
1279            .unwrap_or_else(std::sync::PoisonError::into_inner)
1280            .open();
1281        Ok(SdkSubmitClaim {
1282            inbox: self.steer_queue.clone(),
1283            busy: self.busy.clone(),
1284            cancel,
1285            current_cancel: self.current_cancel.clone(),
1286            turn_finished: self.turn_finished.clone(),
1287            frontend_events: self.frontend_events.clone(),
1288            frontend_state: self.frontend_state.clone(),
1289            lifecycle_started: false,
1290        })
1291    }
1292
1293    async fn submit_claimed(
1294        &self,
1295        prompt: String,
1296        image_urls: Vec<String>,
1297        mut submit_claim: SdkSubmitClaim,
1298    ) -> Result<String, RuntimeSubmitError> {
1299        let cancel = submit_claim.cancel.clone();
1300        self.publish_frontend_payload(json!({"type": "user_message", "text": &prompt}));
1301        self.publish_frontend_payload(json!({
1302            "type": "turn_started",
1303            "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
1304        }));
1305        submit_claim.mark_lifecycle_started();
1306        let outcome = {
1307            let mut agent = self.agent.lock().await;
1308            let result = tokio::select! {
1309                biased;
1310                _ = cancel.notified() => Err(RuntimeSubmitError::Interrupted),
1311                result = async {
1312                    if image_urls.is_empty() {
1313                        agent.inner_mut().send(&prompt).await
1314                    } else {
1315                        agent.inner_mut().send_with_images(&prompt, &image_urls).await
1316                    }
1317                } => result.map_err(|error| RuntimeSubmitError::Agent(error.to_string())),
1318            };
1319            if result.is_ok() {
1320                if let Some(hook) = &self.on_turn_complete {
1321                    hook(&agent);
1322                }
1323            }
1324            // `Agent::send` is cancellation-safe at await boundaries. Publish
1325            // its latest well-formed history on success, failure, or
1326            // interruption without making attach readers wait on `agent`.
1327            let history = bounded_history_snapshot(agent.history());
1328            *self.history_snapshot.write().await = history.clone();
1329            let mut state = self
1330                .frontend_state
1331                .lock()
1332                .unwrap_or_else(std::sync::PoisonError::into_inner);
1333            state.history = history;
1334            state.history_cursor = state.next_sequence.saturating_sub(1);
1335            // Canonical ChatMessage history represents user/model/tool
1336            // content, but not interactive frontend requests or the human's
1337            // typed decision. Re-sequence those semantic events immediately
1338            // after the history boundary so later attachments retain a
1339            // resolved transcript without replaying ordinary turn events
1340            // already represented by `history`.
1341            let request_history = compact_frontend_request_history(&state.replay);
1342            state.replay.clear();
1343            for payload in request_history {
1344                let event = FrontendEvent::new(state.next_sequence, payload);
1345                state.next_sequence = state.next_sequence.saturating_add(1);
1346                state.replay.push_back(event);
1347                while state.replay.len() > FRONTEND_REPLAY_CAPACITY {
1348                    state.replay.pop_front();
1349                }
1350            }
1351            result
1352        };
1353        *self
1354            .current_cancel
1355            .lock()
1356            .unwrap_or_else(std::sync::PoisonError::into_inner) = None;
1357        let lifecycle = match &outcome {
1358            Ok(reply) => json!({
1359                "type": "turn_succeeded",
1360                "schema_version": FRONTEND_EVENT_SCHEMA_VERSION,
1361                "reply": reply,
1362            }),
1363            Err(RuntimeSubmitError::Interrupted) => json!({
1364                "type": "turn_interrupted",
1365                "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
1366            }),
1367            Err(error) => json!({
1368                "type": "turn_failed",
1369                "schema_version": FRONTEND_EVENT_SCHEMA_VERSION,
1370                "message": error.to_string()
1371            }),
1372        };
1373        self.publish_frontend_payload(lifecycle);
1374        submit_claim.mark_lifecycle_finished();
1375        // Close before publishing idle. This is redundant with the agent's
1376        // normal final-boundary close, but is authoritative for pre-loop
1377        // validation/record failures and biased immediate interruption.
1378        drop(submit_claim);
1379        outcome
1380    }
1381
1382    /// Submit one prompt through the canonical runtime and wait for its reply.
1383    pub async fn submit(&self, prompt: impl Into<String>) -> Result<String, RuntimeSubmitError> {
1384        let submit_claim = self.claim_submit()?;
1385        self.submit_claimed(prompt.into(), Vec::new(), submit_claim)
1386            .await
1387    }
1388
1389    /// Submit one prompt with runtime-owned multimodal image inputs.
1390    pub async fn submit_with_images(
1391        &self,
1392        prompt: impl Into<String>,
1393        image_urls: Vec<String>,
1394    ) -> Result<String, RuntimeSubmitError> {
1395        let submit_claim = self.claim_submit()?;
1396        self.submit_claimed(prompt.into(), image_urls, submit_claim)
1397            .await
1398    }
1399
1400    /// Atomically claim one prompt, then run it on the SDK owner while the
1401    /// caller consumes the canonical event stream.
1402    pub fn send_input(self: &Arc<Self>, prompt: String) -> Result<(), RuntimeSubmitError> {
1403        self.send_input_with_images(prompt, Vec::new())
1404    }
1405
1406    /// Atomically claim one multimodal prompt, then run it while callers
1407    /// consume the canonical event stream.
1408    pub fn send_input_with_images(
1409        self: &Arc<Self>,
1410        prompt: String,
1411        image_urls: Vec<String>,
1412    ) -> Result<(), RuntimeSubmitError> {
1413        let submit_claim = self.claim_submit()?;
1414        let runtime = self.clone();
1415        tokio::spawn(async move {
1416            let _ = runtime
1417                .submit_claimed(prompt, image_urls, submit_claim)
1418                .await;
1419        });
1420        Ok(())
1421    }
1422
1423    /// Queue steering for the active agent loop without waiting for its
1424    /// long-held async lock. The agent consumes it at the next model-loop
1425    /// boundary according to the configured steering mode.
1426    pub fn steer(&self, prompt: impl Into<String>) -> Result<(), FrontendRuntimeError> {
1427        let accepted = self
1428            .steer_queue
1429            .lock()
1430            .unwrap_or_else(std::sync::PoisonError::into_inner)
1431            .enqueue(prompt.into());
1432        if accepted {
1433            Ok(())
1434        } else {
1435            Err(FrontendRuntimeError::UnsupportedAction("steer"))
1436        }
1437    }
1438
1439    /// Resolve one pending interactive request exactly once.
1440    pub fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
1441        self.frontend_requests
1442            .as_ref()
1443            .ok_or(FrontendRuntimeError::UnsupportedAction("respond"))?
1444            .respond(response)
1445    }
1446
1447    /// Invoke one operation after resolving its opaque identifier solely
1448    /// against the trusted catalog captured at runtime construction.
1449    pub async fn invoke(
1450        &self,
1451        operation: FrontendOperationInvocation,
1452    ) -> Result<FrontendOperationResult, FrontendRuntimeError> {
1453        match operation {
1454            FrontendOperationInvocation::Prompt {
1455                operation_id,
1456                arguments,
1457            } => {
1458                let prompt_name = self
1459                    .frontend_operations
1460                    .iter()
1461                    .find(|descriptor| {
1462                        descriptor.id == operation_id
1463                            && descriptor.kind == FrontendOperationKind::Prompt
1464                    })
1465                    .and_then(|descriptor| descriptor.command.as_ref())
1466                    .map(|command| command.name.as_str())
1467                    .ok_or_else(|| {
1468                        FrontendRuntimeError::UnsupportedOperation(operation_id.clone())
1469                    })?;
1470                let prompt = if arguments.is_empty() {
1471                    format!("/{prompt_name}")
1472                } else {
1473                    format!("/{prompt_name} {arguments}")
1474                };
1475                let reply = self.submit(prompt).await?;
1476                Ok(FrontendOperationResult::Prompt { reply })
1477            }
1478            FrontendOperationInvocation::Context { operation_id } => {
1479                if !self.frontend_operations.iter().any(|descriptor| {
1480                    descriptor.id == operation_id
1481                        && descriptor.kind == FrontendOperationKind::Context
1482                }) {
1483                    return Err(FrontendRuntimeError::UnsupportedOperation(operation_id));
1484                }
1485                // Short-held read lock: `context_usage` is pure and never
1486                // touches the provider, so this can never be held across
1487                // model I/O.
1488                let usage = self.agent.lock().await.context_usage();
1489                Ok(FrontendOperationResult::Context { usage })
1490            }
1491            FrontendOperationInvocation::Model {
1492                operation_id,
1493                model,
1494            } => {
1495                if !self.frontend_operations.iter().any(|descriptor| {
1496                    descriptor.id == operation_id && descriptor.kind == FrontendOperationKind::Model
1497                }) {
1498                    return Err(FrontendRuntimeError::UnsupportedOperation(operation_id));
1499                }
1500                let mut agent = self.agent.lock().await;
1501                let previous = agent.model().to_string();
1502                if model.trim().is_empty() {
1503                    return Ok(FrontendOperationResult::Model {
1504                        model: previous.clone(),
1505                        previous,
1506                    });
1507                }
1508                // BP-13: the ONE resolution path — the same routing table
1509                // `--model` and the request build read, so a friendly name
1510                // typed here means exactly what it means everywhere else.
1511                let resolved = agent
1512                    .config()
1513                    .model_routing
1514                    .resolve_alias(model.trim())
1515                    .to_string();
1516                if let Some(refusal) = agent.config().model_routing.refusal(&resolved) {
1517                    return Err(FrontendRuntimeError::UnsupportedOperation(refusal));
1518                }
1519                // `switch_model`, never `set_model`: the governed switch
1520                // filters model-A reasoning artifacts out of the live
1521                // history (dep 8) and records the change in the session
1522                // journal before model B ever builds a request from it.
1523                agent.switch_model(resolved.clone());
1524                Ok(FrontendOperationResult::Model {
1525                    model: resolved,
1526                    previous,
1527                })
1528            }
1529        }
1530    }
1531
1532    /// Cancel the current turn through the shared runtime handle.
1533    pub async fn interrupt(&self) -> bool {
1534        let cancel = self
1535            .current_cancel
1536            .lock()
1537            .unwrap_or_else(std::sync::PoisonError::into_inner)
1538            .clone();
1539        match cancel {
1540            Some(cancel) => {
1541                cancel.notify_one();
1542                true
1543            }
1544            None => false,
1545        }
1546    }
1547
1548    /// Return a lock-free runtime snapshot, including while a turn is active.
1549    pub fn status(&self) -> RuntimeStatus {
1550        RuntimeStatus {
1551            session_id: self.session_id.clone(),
1552            model: self.model.clone(),
1553            busy: self.busy.load(Ordering::SeqCst),
1554            shutting_down: self.is_shutting_down(),
1555        }
1556    }
1557
1558    /// Return the tail of the canonical conversation for newly attached
1559    /// frontends. This is SDK state, not a transport-local replay buffer, and
1560    /// remains answerable while a turn owns the agent lock.
1561    pub async fn history(&self, limit: usize) -> Vec<ChatMessage> {
1562        let history = self.history_snapshot.read().await;
1563        let start = history.len().saturating_sub(limit);
1564        history[start..].to_vec()
1565    }
1566
1567    /// Run one caller-owned finalization projection while holding the agent
1568    /// at a quiescent boundary. This is the persistence/inspection seam for
1569    /// local frontends that transfer `Agent` ownership into the SDK runtime;
1570    /// it does not expose a second way to drive the model loop.
1571    pub async fn finalize_with<R>(&self, finalize: impl FnOnce(&SdkAgent) -> R) -> R {
1572        let agent = self.agent.lock().await;
1573        finalize(&agent)
1574    }
1575
1576    /// Request graceful runtime shutdown.  All connected frontends observe
1577    /// the same transition and any in-flight turn is interrupted.
1578    pub async fn shutdown(&self) {
1579        let _barrier = self.shutdown_barrier.lock().await;
1580        // Seal first so no frontend can claim replacement work while the
1581        // active claim is interrupted. Signal transports only after its
1582        // terminal lifecycle is published, so SSE observers receive that
1583        // boundary before their connections close.
1584        self.accepting_submits.store(false, Ordering::SeqCst);
1585        self.interrupt().await;
1586        loop {
1587            let finished = self.turn_finished.notified();
1588            if !self.busy.load(Ordering::SeqCst) {
1589                break;
1590            }
1591            finished.await;
1592        }
1593        self.signal_shutdown();
1594    }
1595
1596    /// Whether `shutdown` has been requested — callers use this to stop
1597    /// accepting new work/connections.
1598    pub fn is_shutting_down(&self) -> bool {
1599        self.shutting_down.load(Ordering::SeqCst)
1600    }
1601
1602    /// Resolves once `shutdown` has been requested. Cheap to call
1603    /// repeatedly/concurrently — every waiter is woken.
1604    pub async fn wait_for_shutdown(&self) {
1605        // `shutting_down` may already be `true` by the time a caller starts
1606        // waiting (e.g. a connection accepted right after `shutdown` fired)
1607        // — check first so this never blocks forever on a signal that
1608        // already happened.
1609        if self.is_shutting_down() {
1610            return;
1611        }
1612        self.shutdown.notified().await;
1613    }
1614
1615    /// Flip `shutting_down` and wake every [`Self::wait_for_shutdown`]
1616    /// waiter — factored out of [`Self::handle_shutdown`] so [`run_stdio`]
1617    /// can raise the EXACT same signal on its OTHER termination path
1618    /// (`reader` hitting EOF with no explicit `shutdown` RPC) without
1619    /// duplicating the store-then-notify sequence. Deliberately does NOT
1620    /// touch `current_cancel` (unlike [`Self::handle_shutdown`], which also
1621    /// interrupts an in-flight turn) — a plain stdio EOF should let an
1622    /// already-accepted turn run to completion and flush its reply, not cut
1623    /// it off.
1624    fn signal_shutdown(&self) {
1625        self.shutting_down.store(true, Ordering::SeqCst);
1626        self.shutdown.notify_waiters();
1627    }
1628
1629    /// Dispatch one already-parsed [`RpcRequest`] to the right method
1630    /// handler. Every recognized method is fully wired to real `Agent`
1631    /// behavior — there is no method that parses but no-ops.
1632    pub async fn handle_request(self: &Arc<Self>, req: RpcRequest) -> Value {
1633        match req.method.as_str() {
1634            "submit" => self.handle_submit(req).await,
1635            "frontend.send_input" => self.handle_frontend_send_input(req),
1636            "interrupt" => self.handle_interrupt(req).await,
1637            "steer" => self.handle_steer(req),
1638            "respond" => self.handle_respond(req),
1639            "status" => self.handle_status(req).await,
1640            "history" => self.handle_history(req).await,
1641            "frontend.describe" => self.handle_frontend_describe(req),
1642            "frontend.attach" => self.handle_frontend_attach(req),
1643            "frontend.invoke" => self.handle_frontend_invoke(req).await,
1644            "shutdown" => self.handle_shutdown(req).await,
1645            other => rpc_error(req.id, -32601, format!("unknown method `{other}`")),
1646        }
1647    }
1648
1649    fn handle_frontend_send_input(self: &Arc<Self>, req: RpcRequest) -> Value {
1650        let Some(prompt) = req.params.get("prompt").and_then(Value::as_str) else {
1651            return rpc_error(
1652                req.id,
1653                -32602,
1654                "frontend.send_input requires a string `params.prompt`",
1655            );
1656        };
1657        let image_urls = match parse_image_urls(&req.params, "frontend.send_input") {
1658            Ok(image_urls) => image_urls,
1659            Err(message) => return rpc_error(req.id, -32602, message),
1660        };
1661        match self.send_input_with_images(prompt.to_string(), image_urls) {
1662            Ok(()) => rpc_ok(req.id, json!({"accepted": true})),
1663            Err(RuntimeSubmitError::Busy) => {
1664                rpc_error(req.id, -32000, "a turn is already in progress")
1665            }
1666            Err(RuntimeSubmitError::Interrupted) => rpc_error(req.id, -32001, "turn interrupted"),
1667            Err(RuntimeSubmitError::Agent(error)) => rpc_error(req.id, -32002, error),
1668        }
1669    }
1670
1671    /// `submit`: drive one turn (`Agent::send`) with `params.prompt`.
1672    /// Refuses (fail, never queues) a second `submit` while one is already
1673    /// in flight — "no method that no-ops": a caller either gets a real
1674    /// answer or an explicit "busy" error, never a silently-dropped
1675    /// request. Races the turn against this call's own fresh cancellation
1676    /// handle so a concurrent `interrupt` can drop it mid-flight — the
1677    /// SAME cancel-safety the CLI's own `race_ctrl_c` relies on
1678    /// (`Agent::send`/`run_loop` only ever mutate `history`/the sidecar
1679    /// BETWEEN `.await` points, never during one, so dropping the future
1680    /// mid-poll always lands in a well-formed place).
1681    async fn handle_submit(&self, req: RpcRequest) -> Value {
1682        let Some(prompt) = req.params.get("prompt").and_then(|v| v.as_str()) else {
1683            return rpc_error(req.id, -32602, "submit requires a string `params.prompt`");
1684        };
1685        match self.submit(prompt).await {
1686            Ok(reply) => rpc_ok(req.id, json!({"reply": reply})),
1687            Err(RuntimeSubmitError::Busy) => rpc_error(
1688                req.id,
1689                -32000,
1690                "a turn is already in progress; `interrupt` it or wait for its response before submitting another",
1691            ),
1692            Err(RuntimeSubmitError::Interrupted) => {
1693                rpc_error(req.id, -32001, "turn interrupted")
1694            }
1695            Err(RuntimeSubmitError::Agent(error)) => rpc_error(req.id, -32002, error),
1696        }
1697    }
1698
1699    /// `interrupt`: cancel the in-flight turn, if any. A no-op-but-honest
1700    /// `{"interrupted": false}` (never an error) when nothing is running —
1701    /// calling `interrupt` with no turn in flight is a normal, harmless
1702    /// race a client can't always avoid (it may not know yet that the
1703    /// previous `submit` just finished).
1704    async fn handle_interrupt(&self, req: RpcRequest) -> Value {
1705        if self.interrupt().await {
1706            rpc_ok(req.id, json!({"interrupted": true}))
1707        } else {
1708            rpc_ok(
1709                req.id,
1710                json!({"interrupted": false, "reason": "no turn in progress"}),
1711            )
1712        }
1713    }
1714
1715    /// `steer`: enqueue an instruction for the next model-loop boundary.
1716    /// This remains responsive while `submit` owns the active agent lock.
1717    fn handle_steer(&self, req: RpcRequest) -> Value {
1718        let Some(prompt) = req.params.get("prompt").and_then(Value::as_str) else {
1719            return rpc_error(req.id, -32602, "steer requires a string `params.prompt`");
1720        };
1721        match self.steer(prompt) {
1722            Ok(()) => rpc_ok(req.id, json!({"queued": true})),
1723            Err(error) => rpc_error(req.id, -32020, error.to_string()),
1724        }
1725    }
1726
1727    fn handle_respond(&self, req: RpcRequest) -> Value {
1728        let response = match req.params.get("response").cloned() {
1729            Some(value) => match serde_json::from_value::<FrontendResponse>(value) {
1730                Ok(response) => response,
1731                Err(error) => return rpc_error(req.id, -32602, error.to_string()),
1732            },
1733            None => return rpc_error(req.id, -32602, "respond requires `params.response`"),
1734        };
1735        match self.respond(response) {
1736            Ok(()) => rpc_ok(req.id, json!({"accepted": true})),
1737            Err(FrontendRuntimeError::UnsupportedAction(_)) => {
1738                rpc_error(req.id, -32020, "frontend respond is not enabled")
1739            }
1740            Err(FrontendRuntimeError::UnknownRequest(id)) => rpc_error(
1741                req.id,
1742                -32021,
1743                format!("frontend request {id} is not pending"),
1744            ),
1745            Err(error) => rpc_error(req.id, -32022, error.to_string()),
1746        }
1747    }
1748
1749    /// `status`: current busy/idle state + the model label. Deliberately
1750    /// never locks `agent` (see [`Self::model`]'s doc comment) — answerable
1751    /// even while a turn is running.
1752    async fn handle_status(&self, req: RpcRequest) -> Value {
1753        rpc_ok(
1754            req.id,
1755            serde_json::to_value(self.status()).unwrap_or_default(),
1756        )
1757    }
1758
1759    /// `history`: bounded canonical transcript replay for a frontend that
1760    /// attached after earlier events were emitted.
1761    async fn handle_history(&self, req: RpcRequest) -> Value {
1762        let limit = req
1763            .params
1764            .get("limit")
1765            .and_then(Value::as_u64)
1766            .unwrap_or(50)
1767            .clamp(1, SERVER_HISTORY_CAPACITY as u64) as usize;
1768        rpc_ok(req.id, json!({"messages": self.history(limit).await}))
1769    }
1770
1771    fn handle_frontend_describe(&self, req: RpcRequest) -> Value {
1772        rpc_ok(
1773            req.id,
1774            serde_json::to_value(self.frontend_descriptor()).unwrap_or_default(),
1775        )
1776    }
1777
1778    fn handle_frontend_attach(&self, req: RpcRequest) -> Value {
1779        let limit = req
1780            .params
1781            .get("limit")
1782            .and_then(Value::as_u64)
1783            .unwrap_or(50)
1784            .clamp(1, SERVER_HISTORY_CAPACITY as u64) as usize;
1785        match self.frontend_snapshot(limit) {
1786            Ok(snapshot) => rpc_ok(req.id, serde_json::to_value(snapshot).unwrap_or_default()),
1787            Err(error) => rpc_error(req.id, -32010, error.to_string()),
1788        }
1789    }
1790
1791    async fn handle_frontend_invoke(&self, req: RpcRequest) -> Value {
1792        let operation = match req.params.get("operation").cloned() {
1793            Some(value) => match serde_json::from_value::<FrontendOperationInvocation>(value) {
1794                Ok(operation) => operation,
1795                Err(error) => return rpc_error(req.id, -32602, error.to_string()),
1796            },
1797            None => {
1798                return rpc_error(
1799                    req.id,
1800                    -32602,
1801                    "frontend.invoke requires `params.operation`",
1802                )
1803            }
1804        };
1805        match self.invoke(operation).await {
1806            Ok(result) => rpc_ok(req.id, serde_json::to_value(result).unwrap_or_default()),
1807            Err(FrontendRuntimeError::UnsupportedOperation(id)) => rpc_error(
1808                req.id,
1809                -32023,
1810                FrontendRuntimeError::UnsupportedOperation(id).to_string(),
1811            ),
1812            Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
1813                rpc_error(req.id, -32000, "a turn is already in progress")
1814            }
1815            Err(error) => rpc_error(req.id, -32022, error.to_string()),
1816        }
1817    }
1818
1819    /// `shutdown`: request graceful teardown — interrupts any in-flight
1820    /// turn (never leaves a caller hanging on a `submit` that will now
1821    /// never get a transport to answer on) and wakes every
1822    /// [`Self::wait_for_shutdown`] waiter (both transports' accept/read
1823    /// loops select on it, so neither survives as an orphan).
1824    async fn handle_shutdown(&self, req: RpcRequest) -> Value {
1825        self.shutdown().await;
1826        rpc_ok(req.id, json!({"shutting_down": true}))
1827    }
1828}
1829
1830/// Retain one canonical request and at most one resolution for every runtime
1831/// request id. Request history is copied across ordinary ChatMessage snapshot
1832/// boundaries, so blindly copying the prior replay would duplicate it on
1833/// every turn and could eventually leave a later attachment on a stale
1834/// duplicate request. Stable id order and request-before-resolution ordering
1835/// make the reconstructed semantic transcript deterministic.
1836fn compact_frontend_request_history(replay: &VecDeque<FrontendEvent>) -> Vec<Value> {
1837    let mut by_id: BTreeMap<u64, (Option<Value>, Option<Value>)> = BTreeMap::new();
1838    for event in replay {
1839        let (request_id, resolved) = match event.kind.as_str() {
1840            "request" => (
1841                event.payload.pointer("/request/id").and_then(Value::as_u64),
1842                false,
1843            ),
1844            "request_resolved" => (
1845                event.payload.get("request_id").and_then(Value::as_u64),
1846                true,
1847            ),
1848            _ => continue,
1849        };
1850        let Some(request_id) = request_id else {
1851            continue;
1852        };
1853        let entry = by_id.entry(request_id).or_default();
1854        let slot = if resolved { &mut entry.1 } else { &mut entry.0 };
1855        slot.get_or_insert_with(|| event.payload.clone());
1856    }
1857    by_id
1858        .into_values()
1859        .flat_map(|(request, resolution)| request.into_iter().chain(resolution))
1860        .collect()
1861}
1862
1863fn valid_frontend_command_name(name: &str) -> bool {
1864    !name.is_empty()
1865        && !name.starts_with('/')
1866        && name
1867            .chars()
1868            .all(|character| !character.is_whitespace() && !character.is_control())
1869}
1870
1871fn bounded_history_snapshot(history: &[ChatMessage]) -> Vec<ChatMessage> {
1872    let start = history.len().saturating_sub(SERVER_HISTORY_CAPACITY);
1873    history[start..].to_vec()
1874}
1875
1876/// Drive the JSONL-RPC protocol over `reader`/`writer` (the stdio rung —
1877/// parent-process-trusted, no auth token; see the module doc). Each
1878/// request line is dispatched on its OWN spawned task so a `submit`
1879/// in-flight never blocks the reader from picking up a subsequent
1880/// `interrupt`/`status` line — every outgoing line (a response OR an event
1881/// notification) is funneled through one mpsc channel into a single writer
1882/// task, so two concurrent handlers can never interleave a line's bytes.
1883/// Returns once `reader` hits EOF or a `shutdown` request lands.
1884///
1885/// Both termination paths raise `RpcEngine::signal_shutdown` (EOF does it
1886/// directly here; the `shutdown` RPC does it inside
1887/// `RpcEngine::handle_shutdown`), and `event_task` below SELECTS against
1888/// [`RpcEngine::wait_for_shutdown`] rather than merely looping on
1889/// `events.recv()`. This is deliberate: `events.recv()` alone only ends via
1890/// `RecvError::Closed`, which fires only once EVERY clone of
1891/// `engine.events` (the broadcast `Sender`) has dropped — and `engine`
1892/// itself, which keeps that `Sender` alive, is owned by THIS function for
1893/// its whole body. Waiting on `events.recv()` to close would therefore mean
1894/// waiting on `engine` to drop, which can't happen until `event_task`
1895/// itself finishes — a circular wait that never resolves (the bug this fn
1896/// exists to fix). Selecting on the shutdown signal instead lets
1897/// `event_task` end WITHOUT needing `engine`'s refcount to reach zero, so
1898/// there is no orphaned task and no leaked `engine`/`writer_task` blocking
1899/// on it in turn.
1900#[cfg(feature = "adapter-api")]
1901pub async fn run_stdio<R, W>(engine: Arc<RpcEngine>, reader: R, writer: W) -> std::io::Result<()>
1902where
1903    R: AsyncBufRead + Unpin + Send + 'static,
1904    W: AsyncWrite + Unpin + Send + 'static,
1905{
1906    let (out_tx, mut out_rx) = mpsc::unbounded_channel::<Value>();
1907
1908    let writer_task = tokio::spawn(async move {
1909        let mut writer = writer;
1910        while let Some(v) = out_rx.recv().await {
1911            let line = format!("{v}\n");
1912            if writer.write_all(line.as_bytes()).await.is_err() {
1913                break;
1914            }
1915            if writer.flush().await.is_err() {
1916                break;
1917            }
1918        }
1919    });
1920
1921    let mut events = engine.subscribe();
1922    let evt_tx = out_tx.clone();
1923    let evt_engine = engine.clone();
1924    let event_task = tokio::spawn(async move {
1925        loop {
1926            tokio::select! {
1927                biased;
1928                _ = evt_engine.wait_for_shutdown() => break,
1929                recv = events.recv() => {
1930                    match recv {
1931                        Ok(v) => {
1932                            if evt_tx.send(json!({"event": v})).is_err() {
1933                                break;
1934                            }
1935                        }
1936                        Err(broadcast::error::RecvError::Lagged(_)) => continue,
1937                        Err(broadcast::error::RecvError::Closed) => break,
1938                    }
1939                }
1940            }
1941        }
1942    });
1943
1944    let mut reader = reader;
1945    loop {
1946        if engine.is_shutting_down() {
1947            break;
1948        }
1949        tokio::select! {
1950            biased;
1951            _ = engine.wait_for_shutdown() => break,
1952            line = read_bounded_line(&mut reader, SERVER_MAX_LINE_BYTES) => {
1953                match line {
1954                    Ok(None) => {
1955                        // EOF: no explicit `shutdown` RPC landed, but stdin
1956                        // closing is this fn's OTHER documented
1957                        // termination signal — raise the same shutdown
1958                        // signal `event_task` (and any other
1959                        // `wait_for_shutdown` caller) already knows how to
1960                        // watch for, so teardown below actually completes
1961                        // instead of blocking forever on `event_task`.
1962                        engine.signal_shutdown();
1963                        break;
1964                    }
1965                    Ok(Some(text)) => {
1966                        let text = text.trim();
1967                        if text.is_empty() {
1968                            continue;
1969                        }
1970                        match serde_json::from_str::<RpcRequest>(text) {
1971                            Ok(req) => {
1972                                let engine = engine.clone();
1973                                let out_tx = out_tx.clone();
1974                                tokio::spawn(async move {
1975                                    let resp = engine.handle_request(req).await;
1976                                    let _ = out_tx.send(resp);
1977                                });
1978                            }
1979                            Err(e) => {
1980                                let _ = out_tx.send(rpc_error(Value::Null, -32700, format!("parse error: {e}")));
1981                            }
1982                        }
1983                    }
1984                    Err(e) => {
1985                        let _ = out_tx.send(rpc_error(Value::Null, -32700, format!("{e}")));
1986                    }
1987                }
1988            }
1989        }
1990    }
1991    // `event_task` now ends promptly (it's shutdown-signalled above, on
1992    // EITHER termination path) rather than waiting on `engine`'s broadcast
1993    // `Sender` to drop — so awaiting it here no longer deadlocks. Dropping
1994    // this fn's own `out_tx` clone (plus `event_task`'s, once it exits)
1995    // lets `writer_task` see `out_rx.recv()` return `None` once every
1996    // OTHER in-flight per-request task (spawned above, each holding its own
1997    // `out_tx` clone) has also sent its reply and dropped its clone — so
1998    // any reply already accepted before shutdown is still flushed to
1999    // `writer` before this fn returns.
2000    drop(out_tx);
2001    let _ = event_task.await;
2002    let _ = writer_task.await;
2003    Ok(())
2004}
2005
2006/// One parsed HTTP/1.1 request (the minimal subset this module's two
2007/// routes need — no keep-alive, no chunked request bodies).
2008#[cfg(feature = "adapter-api")]
2009pub(crate) struct HttpRequest {
2010    pub(crate) method: String,
2011    /// Path WITHOUT the query string (see `query` for that).
2012    pub(crate) path: String,
2013    pub(crate) query: String,
2014    pub(crate) headers: HashMap<String, String>,
2015    pub(crate) body: Vec<u8>,
2016}
2017
2018/// Read and parse one HTTP/1.1 request from `reader`. `Ok(None)` at a
2019/// clean EOF before any bytes arrive (an idle keep-alive-less connection
2020/// closing). Bounded throughout: the request line and each header line go
2021/// through [`read_bounded_line`] (8KiB — generous for a request
2022/// line/header, far below [`SERVER_MAX_LINE_BYTES`]), the header COUNT is
2023/// capped at [`MAX_HEADER_LINES`], and the body is capped at
2024/// [`SERVER_MAX_LINE_BYTES`].
2025#[cfg(feature = "adapter-api")]
2026pub(crate) async fn read_http_request<R>(reader: &mut R) -> std::io::Result<Option<HttpRequest>>
2027where
2028    R: AsyncBufRead + AsyncRead + Unpin,
2029{
2030    const HEAD_LINE_CAP: usize = 8 * 1024;
2031    let Some(request_line) = read_bounded_line(reader, HEAD_LINE_CAP).await? else {
2032        return Ok(None);
2033    };
2034    let mut parts = request_line.split_whitespace();
2035    let method = parts.next().unwrap_or("").to_string();
2036    let target = parts.next().unwrap_or("").to_string();
2037    if method.is_empty() || target.is_empty() {
2038        return Err(std::io::Error::new(
2039            std::io::ErrorKind::InvalidData,
2040            "malformed request line",
2041        ));
2042    }
2043    let (path, query) = match target.split_once('?') {
2044        Some((p, q)) => (p.to_string(), q.to_string()),
2045        None => (target, String::new()),
2046    };
2047
2048    let mut headers = HashMap::new();
2049    let mut content_length: usize = 0;
2050    for _ in 0..MAX_HEADER_LINES {
2051        let Some(line) = read_bounded_line(reader, HEAD_LINE_CAP).await? else {
2052            return Ok(None);
2053        };
2054        if line.is_empty() {
2055            break;
2056        }
2057        if let Some((k, v)) = line.split_once(':') {
2058            let k = k.trim().to_ascii_lowercase();
2059            let v = v.trim().to_string();
2060            if k == "content-length" {
2061                content_length = v.parse().unwrap_or(0);
2062            }
2063            headers.insert(k, v);
2064        }
2065    }
2066    if content_length > SERVER_MAX_LINE_BYTES {
2067        return Err(std::io::Error::new(
2068            std::io::ErrorKind::InvalidData,
2069            format!("request body exceeded {SERVER_MAX_LINE_BYTES} byte cap"),
2070        ));
2071    }
2072    let mut body = vec![0u8; content_length];
2073    if content_length > 0 {
2074        reader.read_exact(&mut body).await?;
2075    }
2076    Ok(Some(HttpRequest {
2077        method,
2078        path,
2079        query,
2080        headers,
2081        body,
2082    }))
2083}
2084
2085/// How long a kept-alive frontend RPC connection may sit idle before this side closes it.
2086#[cfg(feature = "adapter-api")]
2087const FRONTEND_HTTP_IDLE: std::time::Duration = std::time::Duration::from_secs(300);
2088
2089/// A 200 JSON answer that keeps the connection for the client's next request.
2090#[cfg(feature = "adapter-api")]
2091async fn write_http_response_keep_alive<W: AsyncWrite + Unpin>(
2092    writer: &mut W,
2093    body: &[u8],
2094) -> std::io::Result<()> {
2095    let head = format!(
2096        "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: keep-alive\r\n\r\n",
2097        body.len()
2098    );
2099    writer.write_all(head.as_bytes()).await?;
2100    writer.write_all(body).await?;
2101    writer.flush().await
2102}
2103
2104#[cfg(feature = "adapter-api")]
2105pub(crate) async fn write_http_response<W: AsyncWrite + Unpin>(
2106    writer: &mut W,
2107    status: u16,
2108    reason: &str,
2109    content_type: &str,
2110    body: &[u8],
2111) -> std::io::Result<()> {
2112    let head = format!(
2113        "HTTP/1.1 {status} {reason}\r\nContent-Type: {content_type}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
2114        body.len()
2115    );
2116    writer.write_all(head.as_bytes()).await?;
2117    writer.write_all(body).await?;
2118    writer.flush().await
2119}
2120
2121#[cfg(feature = "adapter-api")]
2122fn browser_observer_asset(path: &str) -> Option<(&'static str, &'static [u8])> {
2123    match path {
2124        "/observer" | "/observer/" => Some((
2125            "text/html; charset=utf-8",
2126            include_bytes!("../embedded/frontend-browser/index.html"),
2127        )),
2128        "/observer/app.mjs" => Some((
2129            "text/javascript; charset=utf-8",
2130            include_bytes!("../embedded/frontend-browser/app.mjs"),
2131        )),
2132        "/observer/client.mjs" => Some((
2133            "text/javascript; charset=utf-8",
2134            include_bytes!("../embedded/frontend-browser/client.mjs"),
2135        )),
2136        "/observer/view.mjs" => Some((
2137            "text/javascript; charset=utf-8",
2138            include_bytes!("../embedded/frontend-browser/view.mjs"),
2139        )),
2140        "/observer/style.css" => Some((
2141            "text/css; charset=utf-8",
2142            include_bytes!("../embedded/frontend-browser/style.css"),
2143        )),
2144        "/frontend/client.mjs" => Some((
2145            "text/javascript; charset=utf-8",
2146            include_bytes!("../embedded/frontend/client.mjs"),
2147        )),
2148        "/frontend/generated-client.mjs" => Some((
2149            "text/javascript; charset=utf-8",
2150            include_bytes!("../embedded/frontend/generated-client.mjs"),
2151        )),
2152        "/frontend/generated.mjs" => Some((
2153            "text/javascript; charset=utf-8",
2154            include_bytes!("../embedded/frontend/generated.mjs"),
2155        )),
2156        _ => None,
2157    }
2158}
2159
2160#[cfg(feature = "adapter-api")]
2161async fn write_browser_observer_asset<W: AsyncWrite + Unpin>(
2162    writer: &mut W,
2163    content_type: &str,
2164    body: &[u8],
2165) -> std::io::Result<()> {
2166    let head = format!(
2167        "HTTP/1.1 200 OK\r\n\
2168         Content-Type: {content_type}\r\n\
2169         Content-Length: {}\r\n\
2170         Cache-Control: no-store\r\n\
2171         Content-Security-Policy: default-src 'none'; script-src 'self'; style-src 'self'; connect-src 'self'; img-src 'self' https://brand.volter.ai; base-uri 'none'; form-action 'self'; frame-ancestors 'none'\r\n\
2172         Referrer-Policy: no-referrer\r\n\
2173         X-Content-Type-Options: nosniff\r\n\
2174         Connection: close\r\n\r\n",
2175        body.len()
2176    );
2177    writer.write_all(head.as_bytes()).await?;
2178    writer.write_all(body).await?;
2179    writer.flush().await
2180}
2181
2182/// One bearer credential and its exact SDK authorization grant.
2183#[cfg(feature = "adapter-api")]
2184#[derive(Clone)]
2185pub struct RuntimeHttpCredential {
2186    token: Arc<str>,
2187    authorization: RuntimeAuthorization,
2188    client_id: Option<crate::RuntimeClientId>,
2189    bootstrap: bool,
2190    runtime_id: Option<Arc<str>>,
2191    generation: Option<[u8; 16]>,
2192    revocation: Option<Arc<RuntimeCredentialRevocation>>,
2193}
2194
2195#[cfg(feature = "adapter-api")]
2196impl RuntimeHttpCredential {
2197    /// Create a scoped credential. Tokens are deliberately private and never
2198    /// implement `Debug` or serialization.
2199    pub fn new(token: impl Into<Arc<str>>, authorization: RuntimeAuthorization) -> Self {
2200        Self {
2201            token: token.into(),
2202            authorization,
2203            client_id: None,
2204            bootstrap: false,
2205            runtime_id: None,
2206            generation: None,
2207            revocation: None,
2208        }
2209    }
2210
2211    /// Full owner credential preserving the historical `run_http` contract.
2212    ///
2213    /// New frontend-host paths must keep this bootstrap credential private and
2214    /// exchange it through the local mint endpoint for a bound frontend grant.
2215    pub fn owner(token: impl Into<Arc<str>>) -> Self {
2216        Self {
2217            token: token.into(),
2218            authorization: RuntimeAuthorization::owner(),
2219            client_id: None,
2220            bootstrap: true,
2221            runtime_id: None,
2222            generation: None,
2223            revocation: None,
2224        }
2225    }
2226
2227    /// Read-only observer credential.
2228    pub fn observer(token: impl Into<Arc<str>>) -> Self {
2229        Self::new(token, RuntimeAuthorization::observer())
2230    }
2231
2232    fn frontend(
2233        token: impl Into<Arc<str>>,
2234        client_id: crate::RuntimeClientId,
2235        authorization: RuntimeAuthorization,
2236        runtime_id: impl Into<Arc<str>>,
2237        generation: [u8; 16],
2238    ) -> Self {
2239        Self {
2240            token: token.into(),
2241            authorization,
2242            client_id: Some(client_id),
2243            bootstrap: false,
2244            runtime_id: Some(runtime_id.into()),
2245            generation: Some(generation),
2246            revocation: Some(Arc::new(RuntimeCredentialRevocation::new())),
2247        }
2248    }
2249}
2250
2251#[cfg(feature = "adapter-api")]
2252struct AuthenticatedRuntimeHttpCredential {
2253    authorization: RuntimeAuthorization,
2254    client_id: Option<crate::RuntimeClientId>,
2255    bootstrap: bool,
2256    revocation: Option<tokio::sync::watch::Receiver<bool>>,
2257    attachment: Option<RuntimeCredentialAttachment>,
2258    via_bearer_header: bool,
2259}
2260
2261#[cfg(feature = "adapter-api")]
2262struct RuntimeCredentialRevocation {
2263    signal: tokio::sync::watch::Sender<bool>,
2264    active_attachments: AtomicUsize,
2265    drained: tokio::sync::Notify,
2266}
2267
2268#[cfg(feature = "adapter-api")]
2269impl RuntimeCredentialRevocation {
2270    fn new() -> Self {
2271        let (signal, _) = tokio::sync::watch::channel(false);
2272        Self {
2273            signal,
2274            active_attachments: AtomicUsize::new(0),
2275            drained: tokio::sync::Notify::new(),
2276        }
2277    }
2278
2279    fn register(self: &Arc<Self>) -> RuntimeCredentialAttachment {
2280        self.active_attachments.fetch_add(1, Ordering::AcqRel);
2281        RuntimeCredentialAttachment {
2282            revocation: self.clone(),
2283        }
2284    }
2285
2286    async fn revoke_and_wait(&self) {
2287        let _ = self.signal.send(true);
2288        loop {
2289            let drained = self.drained.notified();
2290            if self.active_attachments.load(Ordering::Acquire) == 0 {
2291                return;
2292            }
2293            drained.await;
2294        }
2295    }
2296}
2297
2298#[cfg(feature = "adapter-api")]
2299struct RuntimeCredentialAttachment {
2300    revocation: Arc<RuntimeCredentialRevocation>,
2301}
2302
2303#[cfg(feature = "adapter-api")]
2304impl Drop for RuntimeCredentialAttachment {
2305    fn drop(&mut self) {
2306        if self
2307            .revocation
2308            .active_attachments
2309            .fetch_sub(1, Ordering::AcqRel)
2310            == 1
2311        {
2312            // Registry removal makes the first revoke the sole waiter for a
2313            // credential; concurrent revokes find no credential. `notify_one`
2314            // stores a permit if this lands between the counter check and the
2315            // first poll of `notified()`, preventing a lost wakeup.
2316            self.revocation.drained.notify_one();
2317        }
2318    }
2319}
2320
2321#[cfg(feature = "adapter-api")]
2322struct IssuedRuntimeHttpCredential(Arc<str>);
2323
2324#[cfg(feature = "adapter-api")]
2325impl IssuedRuntimeHttpCredential {
2326    fn as_bytes(&self) -> &[u8] {
2327        self.0.as_bytes()
2328    }
2329}
2330
2331#[cfg(feature = "adapter-api")]
2332struct RuntimeHttpCredentialRegistry {
2333    credentials: StdMutex<Vec<RuntimeHttpCredential>>,
2334    runtime_id: String,
2335    generation: [u8; 16],
2336}
2337
2338#[cfg(feature = "adapter-api")]
2339impl RuntimeHttpCredentialRegistry {
2340    fn new(
2341        runtime_id: impl Into<String>,
2342        credentials: Vec<RuntimeHttpCredential>,
2343    ) -> std::io::Result<Arc<Self>> {
2344        let mut generation = [0_u8; 16];
2345        getrandom::getrandom(&mut generation).map_err(|error| {
2346            std::io::Error::other(format!(
2347                "cannot create runtime credential generation: {error}"
2348            ))
2349        })?;
2350        Ok(Arc::new(Self {
2351            credentials: StdMutex::new(credentials),
2352            runtime_id: runtime_id.into(),
2353            generation,
2354        }))
2355    }
2356
2357    fn authenticate(&self, request: &HttpRequest) -> Option<AuthenticatedRuntimeHttpCredential> {
2358        let credentials = self
2359            .credentials
2360            .lock()
2361            .unwrap_or_else(std::sync::PoisonError::into_inner);
2362        debug_assert!(credentials.iter().all(|credential| {
2363            credential.client_id.is_none()
2364                || (credential.runtime_id.as_deref() == Some(self.runtime_id.as_str())
2365                    && credential.generation == Some(self.generation))
2366        }));
2367        check_auth(request, &credentials)
2368    }
2369
2370    fn issue_frontend(
2371        &self,
2372        client_id: crate::RuntimeClientId,
2373        observer: bool,
2374    ) -> std::io::Result<IssuedRuntimeHttpCredential> {
2375        let authorization = if observer {
2376            RuntimeAuthorization::observer()
2377        } else {
2378            RuntimeAuthorization::interactive()
2379        };
2380        for _ in 0..3 {
2381            let mut secret = [0_u8; 32];
2382            getrandom::getrandom(&mut secret).map_err(|error| {
2383                std::io::Error::other(format!("cannot mint frontend credential: {error}"))
2384            })?;
2385            let token: Arc<str> = encode_credential(&secret).into();
2386            secret.fill(0);
2387            let mut credentials = self
2388                .credentials
2389                .lock()
2390                .unwrap_or_else(std::sync::PoisonError::into_inner);
2391            if credentials
2392                .iter()
2393                .any(|credential| constant_time_eq(token.as_bytes(), credential.token.as_bytes()))
2394            {
2395                continue;
2396            }
2397            credentials.push(RuntimeHttpCredential::frontend(
2398                token.clone(),
2399                client_id,
2400                authorization,
2401                self.runtime_id.clone(),
2402                self.generation,
2403            ));
2404            return Ok(IssuedRuntimeHttpCredential(token));
2405        }
2406        Err(std::io::Error::new(
2407            std::io::ErrorKind::AlreadyExists,
2408            "frontend credential collision limit exceeded",
2409        ))
2410    }
2411
2412    async fn revoke_client(&self, client_id: &crate::RuntimeClientId) -> bool {
2413        // Authentication removal and attachment registration share this lock,
2414        // so no credential-owned channel can appear after the removal point.
2415        let revocations = {
2416            let mut credentials = self
2417                .credentials
2418                .lock()
2419                .unwrap_or_else(std::sync::PoisonError::into_inner);
2420            let mut revocations = Vec::new();
2421            credentials.retain(|credential| {
2422                if credential.client_id.as_ref() == Some(client_id) {
2423                    if let Some(revocation) = &credential.revocation {
2424                        revocations.push(revocation.clone());
2425                    }
2426                    false
2427                } else {
2428                    true
2429                }
2430            });
2431            revocations
2432        };
2433        let revoked = !revocations.is_empty();
2434        for revocation in revocations {
2435            revocation.revoke_and_wait().await;
2436        }
2437        revoked
2438    }
2439}
2440
2441#[cfg(feature = "adapter-api")]
2442fn encode_credential(secret: &[u8; 32]) -> String {
2443    const HEX: &[u8; 16] = b"0123456789abcdef";
2444    let mut encoded = String::with_capacity(64);
2445    for byte in secret {
2446        encoded.push(HEX[(byte >> 4) as usize] as char);
2447        encoded.push(HEX[(byte & 0x0f) as usize] as char);
2448    }
2449    encoded
2450}
2451
2452/// Does `req` carry a recognized bearer credential? Checked two ways: the
2453/// standard `Authorization: Bearer <token>` header, or, for unbound legacy
2454/// credentials only, a `?token=` query-string parameter (kept for
2455/// `GET /events`, since browser `EventSource` cannot set custom headers).
2456/// Scoped frontend credentials are header-only. Compared with
2457/// [`constant_time_eq`].
2458#[cfg(feature = "adapter-api")]
2459fn check_auth(
2460    req: &HttpRequest,
2461    credentials: &[RuntimeHttpCredential],
2462) -> Option<AuthenticatedRuntimeHttpCredential> {
2463    if let Some(auth) = req.headers.get("authorization") {
2464        if let Some(t) = auth.strip_prefix("Bearer ") {
2465            for credential in credentials {
2466                if constant_time_eq(t.as_bytes(), credential.token.as_bytes()) {
2467                    return Some(AuthenticatedRuntimeHttpCredential {
2468                        authorization: credential.authorization.clone(),
2469                        client_id: credential.client_id.clone(),
2470                        bootstrap: credential.bootstrap,
2471                        revocation: credential
2472                            .revocation
2473                            .as_ref()
2474                            .map(|revocation| revocation.signal.subscribe()),
2475                        attachment: credential.revocation.as_ref().and_then(|revocation| {
2476                            matches!(req.path.as_str(), "/events" | "/frontend/events")
2477                                .then(|| revocation.register())
2478                        }),
2479                        via_bearer_header: true,
2480                    });
2481                }
2482            }
2483        }
2484    }
2485    for pair in req.query.split('&') {
2486        if let Some((k, v)) = pair.split_once('=') {
2487            if k == "token" {
2488                for credential in credentials
2489                    .iter()
2490                    .filter(|credential| credential.client_id.is_none())
2491                {
2492                    if constant_time_eq(v.as_bytes(), credential.token.as_bytes()) {
2493                        return Some(AuthenticatedRuntimeHttpCredential {
2494                            authorization: credential.authorization.clone(),
2495                            client_id: credential.client_id.clone(),
2496                            bootstrap: credential.bootstrap,
2497                            revocation: None,
2498                            attachment: None,
2499                            via_bearer_header: false,
2500                        });
2501                    }
2502                }
2503            }
2504        }
2505    }
2506    None
2507}
2508
2509#[cfg(feature = "adapter-api")]
2510fn coordinated_http_client(
2511    request: &HttpRequest,
2512    coordinator: &Arc<CoordinatedRuntime>,
2513    credential: AuthenticatedRuntimeHttpCredential,
2514) -> Result<Arc<CoordinatedRuntimeClient>, crate::RuntimeLeaseError> {
2515    // Old authenticated API clients predate explicit client ids. Preserve
2516    // them as one named compatibility controller; current SDK clients always
2517    // send a random stable id and therefore coordinate independently.
2518    let supplied_client_id = request
2519        .headers
2520        .get("x-supercode-client-id")
2521        .map(String::as_str);
2522    let client_id = match credential.client_id.as_ref() {
2523        Some(bound) if supplied_client_id == Some(bound.as_str()) => bound.as_str(),
2524        Some(_) => return Err(crate::RuntimeLeaseError::InvalidClientId),
2525        None => supplied_client_id.unwrap_or("legacy-owner"),
2526    };
2527    let mut authorization = credential.authorization;
2528    if let Some(requested) = request.headers.get("x-supercode-permissions") {
2529        authorization = authorization.restrict_to(&RuntimeAuthorization::parse_header(requested)?);
2530    }
2531    Ok(coordinator.client(RuntimeClientId::parse(client_id)?, authorization))
2532}
2533
2534#[cfg(feature = "adapter-api")]
2535async fn coordinated_runtime_rpc(
2536    client: Arc<CoordinatedRuntimeClient>,
2537    request: RpcRequest,
2538) -> Value {
2539    let id = request.id.clone();
2540    let method = crate::FrontendFacadeMethod::from_wire_name(&request.method);
2541    let result = match method {
2542        Some(crate::FrontendFacadeMethod::AcquireControl) => client
2543            .acquire_control()
2544            .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2545        Some(crate::FrontendFacadeMethod::TakeControl) => client
2546            .take_control()
2547            .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2548        Some(crate::FrontendFacadeMethod::Heartbeat) => client
2549            .heartbeat()
2550            .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2551        Some(crate::FrontendFacadeMethod::Lease) => client
2552            .lease_snapshot()
2553            .and_then(|snapshot| serde_json::to_value(snapshot).map_err(json_sdk_error)),
2554        Some(crate::FrontendFacadeMethod::Detach) => {
2555            serde_json::to_value(client.detach()).map_err(json_sdk_error)
2556        }
2557        Some(crate::FrontendFacadeMethod::Close) => match client.close().await {
2558            Ok(()) => Ok(json!({"closed":true})),
2559            Err(error) => Err(error),
2560        },
2561        None if request.method == "shutdown" => match client.close().await {
2562            Ok(()) => Ok(json!({"shutting_down":true})),
2563            Err(error) => Err(error),
2564        },
2565        _ => return frontend_http_rpc(client, request).await,
2566    };
2567    match result {
2568        Ok(value) => rpc_ok(id, value),
2569        Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
2570    }
2571}
2572
2573#[cfg(feature = "adapter-api")]
2574fn json_sdk_error(error: serde_json::Error) -> FrontendRuntimeError {
2575    FrontendRuntimeError::Transport(error.to_string())
2576}
2577
2578/// The runtime's own loopback credential door: `POST` mint and revoke.
2579///
2580/// It lives in ONE place and both HTTP conns call it, because the live-runtime
2581/// receipt names a single `base_url` and `sdk/typescript/live-runtime.mjs`
2582/// mints against exactly that URL — "the runtime's own loopback mint door" —
2583/// while warning that two copies of a security-relevant door drift. A hosted
2584/// runtime registers its receipt with the FRONTEND conn's address
2585/// (`harness_service.rs`'s `insert_hosted_runtime`), so that conn has to serve
2586/// this door or the receipt points at a 404.
2587///
2588/// Returns `None` when the request is not one of these two paths.
2589#[cfg(feature = "adapter-api")]
2590async fn serve_frontend_credential_door<W: AsyncWrite + Unpin>(
2591    req: &HttpRequest,
2592    write_half: &mut W,
2593    peer_is_loopback: bool,
2594    credential: &AuthenticatedRuntimeHttpCredential,
2595    credentials: &RuntimeHttpCredentialRegistry,
2596    coordinator: &Arc<CoordinatedRuntime>,
2597) -> Option<std::io::Result<()>> {
2598    if matches!(
2599        req.path.as_str(),
2600        "/_supercode/frontend-credentials/mint" | "/_supercode/frontend-credentials/revoke"
2601    ) && !credential.via_bearer_header
2602    {
2603        let body =
2604            sdk_runtime_rpc_error(Value::Null, -32030, &FrontendRuntimeError::Unauthenticated)
2605                .to_string();
2606        return Some(
2607            write_http_response(
2608                write_half,
2609                401,
2610                "Unauthorized",
2611                "application/json",
2612                body.as_bytes(),
2613            )
2614            .await,
2615        );
2616    }
2617    if req.path == "/_supercode/frontend-credentials/mint" {
2618        if req.method != "POST" || !peer_is_loopback || !credential.bootstrap {
2619            let body = sdk_runtime_rpc_error(
2620                Value::Null,
2621                -32031,
2622                &FrontendRuntimeError::Unauthorized {
2623                    permission: "bootstrap".into(),
2624                },
2625            )
2626            .to_string();
2627            return Some(
2628                write_http_response(
2629                    write_half,
2630                    403,
2631                    "Forbidden",
2632                    "application/json",
2633                    body.as_bytes(),
2634                )
2635                .await,
2636            );
2637        }
2638        let request: Value = match serde_json::from_slice(&req.body) {
2639            Ok(request) => request,
2640            Err(error) => {
2641                return Some(
2642                    write_http_response(
2643                        write_half,
2644                        400,
2645                        "Bad Request",
2646                        "application/json",
2647                        format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2648                    )
2649                    .await,
2650                );
2651            }
2652        };
2653        let Some(client_id) = request.get("clientId").and_then(Value::as_str) else {
2654            return Some(
2655                write_http_response(
2656                    write_half,
2657                    400,
2658                    "Bad Request",
2659                    "application/json",
2660                    b"{\"error\":\"mint request omitted clientId\"}",
2661                )
2662                .await,
2663            );
2664        };
2665        let client_id = match crate::RuntimeClientId::parse(client_id) {
2666            Ok(client_id) => client_id,
2667            Err(error) => {
2668                return Some(
2669                    write_http_response(
2670                        write_half,
2671                        400,
2672                        "Bad Request",
2673                        "application/json",
2674                        format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2675                    )
2676                    .await,
2677                );
2678            }
2679        };
2680        let observer = match request.get("grant").and_then(Value::as_str) {
2681            Some("observer") => true,
2682            Some("interactive") => false,
2683            _ => {
2684                return Some(
2685                    write_http_response(
2686                        write_half,
2687                        400,
2688                        "Bad Request",
2689                        "application/json",
2690                        b"{\"error\":\"grant must be observer or interactive\"}",
2691                    )
2692                    .await,
2693                );
2694            }
2695        };
2696        let token = match credentials.issue_frontend(client_id, observer) {
2697            Ok(token) => token,
2698            Err(error) => return Some(Err(error)),
2699        };
2700        return Some(
2701            write_http_response(
2702                write_half,
2703                200,
2704                "OK",
2705                "application/octet-stream",
2706                token.as_bytes(),
2707            )
2708            .await,
2709        );
2710    }
2711
2712    if req.path == "/_supercode/frontend-credentials/revoke" {
2713        if req.method != "POST" || !peer_is_loopback || !credential.bootstrap {
2714            let body = sdk_runtime_rpc_error(
2715                Value::Null,
2716                -32031,
2717                &FrontendRuntimeError::Unauthorized {
2718                    permission: "bootstrap".into(),
2719                },
2720            )
2721            .to_string();
2722            return Some(
2723                write_http_response(
2724                    write_half,
2725                    403,
2726                    "Forbidden",
2727                    "application/json",
2728                    body.as_bytes(),
2729                )
2730                .await,
2731            );
2732        }
2733        let request: Value = match serde_json::from_slice(&req.body) {
2734            Ok(request) => request,
2735            Err(error) => {
2736                return Some(
2737                    write_http_response(
2738                        write_half,
2739                        400,
2740                        "Bad Request",
2741                        "application/json",
2742                        format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2743                    )
2744                    .await,
2745                );
2746            }
2747        };
2748        let Some(client_id) = request.get("clientId").and_then(Value::as_str) else {
2749            return Some(
2750                write_http_response(
2751                    write_half,
2752                    400,
2753                    "Bad Request",
2754                    "application/json",
2755                    b"{\"error\":\"revoke request omitted clientId\"}",
2756                )
2757                .await,
2758            );
2759        };
2760        let client_id = match crate::RuntimeClientId::parse(client_id) {
2761            Ok(client_id) => client_id,
2762            Err(error) => {
2763                return Some(
2764                    write_http_response(
2765                        write_half,
2766                        400,
2767                        "Bad Request",
2768                        "application/json",
2769                        format!("{{\"error\":{}}}", json!(error.to_string())).as_bytes(),
2770                    )
2771                    .await,
2772                );
2773            }
2774        };
2775        let revoked = credentials.revoke_client(&client_id).await;
2776        if revoked {
2777            coordinator
2778                .client(client_id, RuntimeAuthorization::observer())
2779                .detach();
2780        }
2781        return Some(
2782            write_http_response(
2783                write_half,
2784                200,
2785                "OK",
2786                "application/json",
2787                if revoked {
2788                    b"{\"revoked\":true}"
2789                } else {
2790                    b"{\"revoked\":false}"
2791                },
2792            )
2793            .await,
2794        );
2795    }
2796
2797    None
2798}
2799
2800#[cfg(feature = "adapter-api")]
2801async fn handle_http_conn(
2802    stream: tokio::net::TcpStream,
2803    engine: Arc<RpcEngine>,
2804    coordinator: Arc<CoordinatedRuntime>,
2805    credentials: Arc<RuntimeHttpCredentialRegistry>,
2806) -> std::io::Result<()> {
2807    let peer_is_loopback = stream.peer_addr()?.ip().is_loopback();
2808    let (read_half, mut write_half) = stream.into_split();
2809    let mut reader = tokio::io::BufReader::new(read_half);
2810    let Some(req) = read_http_request(&mut reader).await? else {
2811        return Ok(());
2812    };
2813
2814    // Static observer assets contain no runtime descriptor or session content.
2815    // Serve them before authentication so a user can load the credential form;
2816    // every SDK request made by that page remains bearer-authenticated below.
2817    if req.method == "GET" {
2818        if let Some((content_type, body)) = browser_observer_asset(&req.path) {
2819            return write_browser_observer_asset(&mut write_half, content_type, body).await;
2820        }
2821    }
2822
2823    let Some(credential) = credentials.authenticate(&req) else {
2824        let body =
2825            sdk_runtime_rpc_error(Value::Null, -32030, &FrontendRuntimeError::Unauthenticated)
2826                .to_string();
2827        return write_http_response(
2828            &mut write_half,
2829            401,
2830            "Unauthorized",
2831            "application/json",
2832            body.as_bytes(),
2833        )
2834        .await;
2835    };
2836
2837    if let Some(result) = serve_frontend_credential_door(
2838        &req,
2839        &mut write_half,
2840        peer_is_loopback,
2841        &credential,
2842        &credentials,
2843        &coordinator,
2844    )
2845    .await
2846    {
2847        return result;
2848    }
2849
2850    let mut revocation = credential.revocation.clone();
2851    let mut attachment = credential.attachment;
2852    let credential = AuthenticatedRuntimeHttpCredential {
2853        authorization: credential.authorization,
2854        client_id: credential.client_id,
2855        bootstrap: credential.bootstrap,
2856        revocation: None,
2857        attachment: None,
2858        via_bearer_header: credential.via_bearer_header,
2859    };
2860    let client = match coordinated_http_client(&req, &coordinator, credential) {
2861        Ok(client) => client,
2862        Err(error) => {
2863            let permission = match error {
2864                crate::RuntimeLeaseError::InvalidClientId => "client_id",
2865                crate::RuntimeLeaseError::InvalidAuthorization => "authorization",
2866                _ => "runtime",
2867            };
2868            let body = sdk_runtime_rpc_error(
2869                Value::Null,
2870                -32031,
2871                &FrontendRuntimeError::Unauthorized {
2872                    permission: permission.into(),
2873                },
2874            )
2875            .to_string();
2876            return write_http_response(
2877                &mut write_half,
2878                403,
2879                "Forbidden",
2880                "application/json",
2881                body.as_bytes(),
2882            )
2883            .await;
2884        }
2885    };
2886
2887    match (req.method.as_str(), req.path.as_str()) {
2888        ("POST", "/rpc") => {
2889            let body_text = String::from_utf8_lossy(&req.body);
2890            let resp = match serde_json::from_str::<RpcRequest>(&body_text) {
2891                Ok(rpc_req) if matches!(rpc_req.method.as_str(), "status" | "history") => {
2892                    engine.handle_request(rpc_req).await
2893                }
2894                Ok(rpc_req) => coordinated_runtime_rpc(client.clone(), rpc_req).await,
2895                Err(e) => rpc_error(Value::Null, -32700, format!("parse error: {e}")),
2896            };
2897            let body = resp.to_string();
2898            write_http_response(
2899                &mut write_half,
2900                200,
2901                "OK",
2902                "application/json",
2903                body.as_bytes(),
2904            )
2905            .await
2906        }
2907        ("GET", "/events") => {
2908            if let Err(error) = client.observe() {
2909                let body = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
2910                return write_http_response(
2911                    &mut write_half,
2912                    403,
2913                    "Forbidden",
2914                    "application/json",
2915                    body.as_bytes(),
2916                )
2917                .await;
2918            }
2919            // Arm the subscriber before acknowledging SSE readiness. Otherwise a
2920            // client can receive 200, immediately submit a turn on /rpc, and lose
2921            // every event emitted before this branch reaches `subscribe()`.
2922            let mut events = engine.subscribe();
2923            let head = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: close\r\n\r\n";
2924            if write_half.write_all(head.as_bytes()).await.is_err() {
2925                client.detach();
2926                return Ok(());
2927            }
2928            let _ = write_half.flush().await;
2929            // A comment every 15 s while no event flows: an HTTP client's body
2930            // timeout (Node's fetch drops a body after 300 s without data) ended
2931            // an idle frontend's stream, and it never saw the turn that followed.
2932            let mut keepalive = tokio::time::interval(EVENT_STREAM_KEEPALIVE);
2933            keepalive.tick().await;
2934            loop {
2935                tokio::select! {
2936                    biased;
2937                    recv = events.recv() => {
2938                        match recv {
2939                            Ok(v) => {
2940                                let line = format!("data: {v}\n\n");
2941                                if write_half.write_all(line.as_bytes()).await.is_err() {
2942                                    break;
2943                                }
2944                                if write_half.flush().await.is_err() {
2945                                    break;
2946                                }
2947                            }
2948                            Err(broadcast::error::RecvError::Lagged(_)) => continue,
2949                            Err(broadcast::error::RecvError::Closed) => break,
2950                        }
2951                    }
2952                    // `RpcEngine::shutdown` publishes the turn terminal
2953                    // event before signaling shutdown. Drain that already-
2954                    // buffered event first so a close racing an active turn
2955                    // cannot make observers miss its final state.
2956                    _ = engine.wait_for_shutdown() => break,
2957                    _ = wait_for_credential_revocation(&mut revocation) => break,
2958                    _ = reader.read_u8() => break,
2959                    _ = keepalive.tick() => {
2960                        if write_half.write_all(b": keep-alive\n\n").await.is_err() || write_half.flush().await.is_err() {
2961                            break;
2962                        }
2963                    }
2964                }
2965            }
2966            let _ = write_half.shutdown().await;
2967            client.detach();
2968            drop(attachment.take());
2969            Ok(())
2970        }
2971        ("GET", "/frontend/events") => {
2972            if let Err(error) = client.observe() {
2973                let body = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
2974                return write_http_response(
2975                    &mut write_half,
2976                    403,
2977                    "Forbidden",
2978                    "application/json",
2979                    body.as_bytes(),
2980                )
2981                .await;
2982            }
2983            // Subscribe before acknowledging the stream. Once the client has
2984            // received the 200 response, every later frontend event is either
2985            // buffered here or delivered live; there is no header/subscription
2986            // race at connection startup.
2987            let mut events = engine.frontend_subscribe();
2988            let head = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: close\r\n\r\n";
2989            if write_half.write_all(head.as_bytes()).await.is_err() {
2990                client.detach();
2991                return Ok(());
2992            }
2993            let _ = write_half.flush().await;
2994            // A comment every 15 s while no event flows: an HTTP client's body
2995            // timeout (Node's fetch drops a body after 300 s without data) ended
2996            // an idle frontend's stream, and it never saw the turn that followed.
2997            let mut keepalive = tokio::time::interval(EVENT_STREAM_KEEPALIVE);
2998            keepalive.tick().await;
2999            loop {
3000                tokio::select! {
3001                    biased;
3002                    recv = events.recv() => {
3003                        match recv {
3004                            Ok(event) => {
3005                                let value = serde_json::to_string(&event).unwrap_or_default();
3006                                let line = format!("data: {value}\n\n");
3007                                if write_half.write_all(line.as_bytes()).await.is_err() {
3008                                    break;
3009                                }
3010                                if write_half.flush().await.is_err() {
3011                                    break;
3012                                }
3013                            }
3014                            Err(broadcast::error::RecvError::Lagged(_)) => break,
3015                            Err(broadcast::error::RecvError::Closed) => break,
3016                        }
3017                    }
3018                    // Shutdown is signaled only after the terminal frontend
3019                    // event is published. Prefer the receiver when both are
3020                    // ready so every observer sees that final sequence.
3021                    _ = engine.wait_for_shutdown() => break,
3022                    _ = wait_for_credential_revocation(&mut revocation) => break,
3023                    _ = reader.read_u8() => break,
3024                    _ = keepalive.tick() => {
3025                        if write_half.write_all(b": keep-alive\n\n").await.is_err() || write_half.flush().await.is_err() {
3026                            break;
3027                        }
3028                    }
3029                }
3030            }
3031            let _ = write_half.shutdown().await;
3032            client.detach();
3033            drop(attachment.take());
3034            Ok(())
3035        }
3036        _ => {
3037            write_http_response(
3038                &mut write_half,
3039                404,
3040                "Not Found",
3041                "application/json",
3042                b"{\"error\":\"not found\"}",
3043            )
3044            .await
3045        }
3046    }
3047}
3048
3049#[cfg(feature = "adapter-api")]
3050async fn wait_for_credential_revocation(receiver: &mut Option<tokio::sync::watch::Receiver<bool>>) {
3051    let Some(receiver) = receiver else {
3052        std::future::pending::<()>().await;
3053        return;
3054    };
3055    if *receiver.borrow() {
3056        return;
3057    }
3058    while receiver.changed().await.is_ok() {
3059        if *receiver.borrow() {
3060            return;
3061        }
3062    }
3063}
3064
3065#[cfg(feature = "adapter-api")]
3066async fn handle_frontend_http_conn(
3067    stream: tokio::net::TcpStream,
3068    coordinator: Arc<CoordinatedRuntime>,
3069    events: broadcast::Sender<FrontendEvent>,
3070    credentials: Arc<RuntimeHttpCredentialRegistry>,
3071) -> std::io::Result<()> {
3072    let peer_is_loopback = stream.peer_addr()?.ip().is_loopback();
3073    let (read_half, mut write_half) = stream.into_split();
3074    let mut reader = tokio::io::BufReader::new(read_half);
3075    // An RPC client keeps its connection: a frontend or the mail watcher asking every hosted runtime's
3076    // turn state every two seconds opened (and this side closed) one connection per request, until the
3077    // machine ran out of ephemeral ports in TIME_WAIT (t_f3f102db t_4f29b43e). Every other answer closes.
3078    loop {
3079        let Ok(read) =
3080            tokio::time::timeout(FRONTEND_HTTP_IDLE, read_http_request(&mut reader)).await
3081        else {
3082            return Ok(());
3083        };
3084        let Some(req) = read? else {
3085            return Ok(());
3086        };
3087        let keep_alive = req.method == "POST"
3088            && req.path == "/rpc"
3089            && !req
3090                .headers
3091                .get("connection")
3092                .is_some_and(|value| value.eq_ignore_ascii_case("close"));
3093        return {
3094            if req.method == "GET" {
3095                if let Some((content_type, body)) = browser_observer_asset(&req.path) {
3096                    return write_browser_observer_asset(&mut write_half, content_type, body).await;
3097                }
3098            }
3099            let Some(credential) = credentials.authenticate(&req) else {
3100                return write_http_response(
3101                    &mut write_half,
3102                    401,
3103                    "Unauthorized",
3104                    "application/json",
3105                    b"{\"error\":\"missing or invalid bearer token\"}",
3106                )
3107                .await;
3108            };
3109            // The receipt a hosted runtime registers names THIS listener, so the mint
3110            // door a frontend is told to use has to answer here.
3111            if let Some(result) = serve_frontend_credential_door(
3112                &req,
3113                &mut write_half,
3114                peer_is_loopback,
3115                &credential,
3116                &credentials,
3117                &coordinator,
3118            )
3119            .await
3120            {
3121                return result;
3122            }
3123            let client = match coordinated_http_client(&req, &coordinator, credential) {
3124                Ok(client) => client,
3125                Err(error) => {
3126                    let body = json!({"error":error.to_string()}).to_string();
3127                    return write_http_response(
3128                        &mut write_half,
3129                        400,
3130                        "Bad Request",
3131                        "application/json",
3132                        body.as_bytes(),
3133                    )
3134                    .await;
3135                }
3136            };
3137            match (req.method.as_str(), req.path.as_str()) {
3138                ("POST", "/rpc") => {
3139                    let body_text = String::from_utf8_lossy(&req.body);
3140                    let response = match serde_json::from_str::<RpcRequest>(&body_text) {
3141                        Ok(request) => coordinated_runtime_rpc(client.clone(), request).await,
3142                        Err(error) => {
3143                            rpc_error(Value::Null, -32700, format!("parse error: {error}"))
3144                        }
3145                    };
3146                    let body = response.to_string();
3147                    if keep_alive {
3148                        write_http_response_keep_alive(&mut write_half, body.as_bytes()).await?;
3149                        continue;
3150                    }
3151                    write_http_response(
3152                        &mut write_half,
3153                        200,
3154                        "OK",
3155                        "application/json",
3156                        body.as_bytes(),
3157                    )
3158                    .await
3159                }
3160                ("GET", "/frontend/events") => {
3161                    if let Err(error) = client.observe() {
3162                        let body = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
3163                        return write_http_response(
3164                            &mut write_half,
3165                            403,
3166                            "Forbidden",
3167                            "application/json",
3168                            body.as_bytes(),
3169                        )
3170                        .await;
3171                    }
3172                    let mut receiver = events.subscribe();
3173                    let head = "HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nConnection: close\r\n\r\n";
3174                    if write_half.write_all(head.as_bytes()).await.is_err() {
3175                        client.detach();
3176                        return Ok(());
3177                    }
3178                    let _ = write_half.flush().await;
3179                    let mut keepalive = tokio::time::interval(EVENT_STREAM_KEEPALIVE);
3180                    keepalive.tick().await;
3181                    loop {
3182                        tokio::select! {
3183                            _ = reader.read_u8() => break,
3184                            _ = keepalive.tick() => {
3185                                if write_half.write_all(b": keep-alive\n\n").await.is_err() || write_half.flush().await.is_err() {
3186                                    break;
3187                                }
3188                            }
3189                            event = receiver.recv() => match event {
3190                                Ok(event) => {
3191                                    let value = serde_json::to_string(&event).unwrap_or_default();
3192                                    let line = format!("data: {value}\n\n");
3193                                    if write_half.write_all(line.as_bytes()).await.is_err() || write_half.flush().await.is_err() {
3194                                        break;
3195                                    }
3196                                }
3197                                Err(broadcast::error::RecvError::Lagged(_)) => break,
3198                                Err(broadcast::error::RecvError::Closed) => break,
3199                            }
3200                        }
3201                    }
3202                    client.detach();
3203                    Ok(())
3204                }
3205                _ => {
3206                    write_http_response(
3207                        &mut write_half,
3208                        404,
3209                        "Not Found",
3210                        "application/json",
3211                        b"{\"error\":\"not found\"}",
3212                    )
3213                    .await
3214                }
3215            }
3216        };
3217    }
3218}
3219
3220#[cfg(feature = "adapter-api")]
3221async fn frontend_http_rpc(runtime: Arc<dyn FrontendRuntime>, request: RpcRequest) -> Value {
3222    let id = request.id;
3223    let Some(method) = crate::FrontendFacadeMethod::from_wire_name(&request.method) else {
3224        return rpc_error(id, -32601, format!("unknown method `{}`", request.method));
3225    };
3226    match method {
3227        crate::FrontendFacadeMethod::Describe => match runtime.describe().await {
3228            Ok(descriptor) => rpc_ok(id, serde_json::to_value(descriptor).unwrap_or_default()),
3229            Err(error) => sdk_runtime_rpc_error(id, -32010, &error),
3230        },
3231        crate::FrontendFacadeMethod::Attach => {
3232            let limit = request
3233                .params
3234                .get("limit")
3235                .and_then(Value::as_u64)
3236                .unwrap_or(50)
3237                .clamp(1, SERVER_HISTORY_CAPACITY as u64) as usize;
3238            match runtime.attach(limit).await {
3239                Ok(attachment) => rpc_ok(
3240                    id,
3241                    serde_json::to_value(FrontendAttachSnapshot {
3242                        descriptor: attachment.descriptor,
3243                        history: attachment.history,
3244                        history_cursor: attachment.history_cursor,
3245                        replay: attachment.replay,
3246                    })
3247                    .unwrap_or_default(),
3248                ),
3249                Err(error) => sdk_runtime_rpc_error(id, -32010, &error),
3250            }
3251        }
3252        crate::FrontendFacadeMethod::SendInput => {
3253            let Some(prompt) = request.params.get("prompt").and_then(Value::as_str) else {
3254                return rpc_error(
3255                    id,
3256                    -32602,
3257                    "frontend.send_input requires a string `params.prompt`",
3258                );
3259            };
3260            let image_urls = match parse_image_urls(&request.params, "frontend.send_input") {
3261                Ok(image_urls) => image_urls,
3262                Err(message) => return rpc_error(id, -32602, message),
3263            };
3264            match runtime
3265                .clone()
3266                .send_input_with_images(prompt.to_string(), image_urls)
3267                .await
3268            {
3269                Ok(()) => rpc_ok(id, json!({"accepted": true})),
3270                Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
3271                    sdk_runtime_rpc_error(id, -32000, &error)
3272                }
3273                Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
3274            }
3275        }
3276        crate::FrontendFacadeMethod::Invoke => {
3277            let operation = request
3278                .params
3279                .get("operation")
3280                .cloned()
3281                .ok_or("frontend.invoke requires `params.operation`")
3282                .and_then(|value| {
3283                    serde_json::from_value(value).map_err(|_| "invalid frontend operation")
3284                });
3285            match operation {
3286                Ok(operation) => match runtime.invoke(operation).await {
3287                    Ok(result) => rpc_ok(id, serde_json::to_value(result).unwrap_or_default()),
3288                    Err(error @ FrontendRuntimeError::UnsupportedOperation(_)) => {
3289                        sdk_runtime_rpc_error(id, -32023, &error)
3290                    }
3291                    Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
3292                        sdk_runtime_rpc_error(id, -32000, &error)
3293                    }
3294                    Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Interrupted)) => {
3295                        sdk_runtime_rpc_error(id, -32001, &error)
3296                    }
3297                    Err(error) => sdk_runtime_rpc_error(id, -32022, &error),
3298                },
3299                Err(message) => rpc_error(id, -32602, message),
3300            }
3301        }
3302        crate::FrontendFacadeMethod::Submit => {
3303            let Some(prompt) = request.params.get("prompt").and_then(Value::as_str) else {
3304                return rpc_error(id, -32602, "submit requires a string `params.prompt`");
3305            };
3306            let image_urls = match request.params.get("image_urls") {
3307                None => Vec::new(),
3308                Some(Value::Array(values)) => {
3309                    let Some(urls) = values.iter().map(Value::as_str).collect::<Option<Vec<_>>>()
3310                    else {
3311                        return rpc_error(
3312                            id,
3313                            -32602,
3314                            "submit requires string entries in `params.image_urls`",
3315                        );
3316                    };
3317                    urls.into_iter().map(str::to_owned).collect()
3318                }
3319                Some(_) => {
3320                    return rpc_error(id, -32602, "submit requires array `params.image_urls`")
3321                }
3322            };
3323            match runtime
3324                .submit_with_images(prompt.to_string(), image_urls)
3325                .await
3326            {
3327                Ok(reply) => rpc_ok(id, json!({"reply":reply})),
3328                Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)) => {
3329                    sdk_runtime_rpc_error(id, -32000, &error)
3330                }
3331                Err(error @ FrontendRuntimeError::Submit(RuntimeSubmitError::Interrupted)) => {
3332                    sdk_runtime_rpc_error(id, -32001, &error)
3333                }
3334                Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
3335            }
3336        }
3337        crate::FrontendFacadeMethod::Interrupt => match runtime.interrupt().await {
3338            Ok(interrupted) => rpc_ok(id, json!({"interrupted":interrupted})),
3339            Err(error) => sdk_runtime_rpc_error(id, -32002, &error),
3340        },
3341        crate::FrontendFacadeMethod::Steer => {
3342            let Some(prompt) = request.params.get("prompt").and_then(Value::as_str) else {
3343                return rpc_error(id, -32602, "steer requires a string `params.prompt`");
3344            };
3345            match runtime.steer(prompt.to_string()).await {
3346                Ok(()) => rpc_ok(id, json!({"queued":true})),
3347                Err(error @ FrontendRuntimeError::UnsupportedAction(_)) => {
3348                    sdk_runtime_rpc_error(id, -32020, &error)
3349                }
3350                Err(error) => sdk_runtime_rpc_error(id, -32022, &error),
3351            }
3352        }
3353        crate::FrontendFacadeMethod::Respond => {
3354            let response = request
3355                .params
3356                .get("response")
3357                .cloned()
3358                .ok_or("respond requires `params.response`")
3359                .and_then(|value| serde_json::from_value(value).map_err(|_| "invalid response"));
3360            match response {
3361                Ok(response) => match runtime.respond(response).await {
3362                    Ok(()) => rpc_ok(id, json!({"accepted":true})),
3363                    Err(error @ FrontendRuntimeError::UnsupportedAction(_)) => {
3364                        sdk_runtime_rpc_error(id, -32020, &error)
3365                    }
3366                    Err(error) => sdk_runtime_rpc_error(id, -32022, &error),
3367                },
3368                Err(message) => rpc_error(id, -32602, message),
3369            }
3370        }
3371        crate::FrontendFacadeMethod::Lease
3372        | crate::FrontendFacadeMethod::AcquireControl
3373        | crate::FrontendFacadeMethod::TakeControl
3374        | crate::FrontendFacadeMethod::Heartbeat
3375        | crate::FrontendFacadeMethod::Detach
3376        | crate::FrontendFacadeMethod::Close => rpc_error(
3377            id,
3378            -32020,
3379            format!(
3380                "frontend action `{}` requires a coordinated runtime",
3381                method.id()
3382            ),
3383        ),
3384    }
3385}
3386
3387fn parse_image_urls(params: &Value, operation: &str) -> std::result::Result<Vec<String>, String> {
3388    let urls = match params.get("image_urls") {
3389        None => Ok(Vec::new()),
3390        Some(Value::Array(values)) => values
3391            .iter()
3392            .map(|value| {
3393                value.as_str().map(str::to_owned).ok_or_else(|| {
3394                    format!("{operation} requires string entries in `params.image_urls`")
3395                })
3396            })
3397            .collect(),
3398        Some(_) => Err(format!("{operation} requires array `params.image_urls`")),
3399    }?;
3400    validate_frontend_image_urls(urls, operation)
3401}
3402
3403fn validate_frontend_image_urls(
3404    urls: Vec<String>,
3405    operation: &str,
3406) -> std::result::Result<Vec<String>, String> {
3407    if urls.len() > 4 {
3408        return Err(format!("{operation} accepts at most 4 images"));
3409    }
3410    let mut total = 0usize;
3411    for url in &urls {
3412        if !(url.starts_with("data:image/")
3413            || url.starts_with("https://")
3414            || url.starts_with("http://"))
3415        {
3416            return Err(format!(
3417                "{operation} images must be image data URLs or HTTP(S) URLs"
3418            ));
3419        }
3420        if url.len() > 12 * 1024 * 1024 {
3421            return Err(format!("{operation} image exceeds the encoded size limit"));
3422        }
3423        total = total.saturating_add(url.len());
3424    }
3425    if total > 32 * 1024 * 1024 {
3426        return Err(format!(
3427            "{operation} images exceed the encoded total size limit"
3428        ));
3429    }
3430    Ok(urls)
3431}
3432
3433/// Lifetime handle for a frontend-only authenticated HTTP listener.
3434#[cfg(feature = "adapter-api")]
3435pub(crate) struct FrontendHttpServer {
3436    address: SocketAddr,
3437    task: tokio::task::JoinHandle<()>,
3438}
3439
3440#[cfg(feature = "adapter-api")]
3441impl FrontendHttpServer {
3442    /// Actually-bound loopback address.
3443    pub(crate) fn address(&self) -> SocketAddr {
3444        self.address
3445    }
3446}
3447
3448#[cfg(feature = "adapter-api")]
3449impl Drop for FrontendHttpServer {
3450    fn drop(&mut self) {
3451        self.task.abort();
3452    }
3453}
3454
3455/// Publish only the versioned frontend contract for a non-Agent runtime.
3456/// The caller owns the runtime and this returned listener lease.
3457#[cfg(feature = "adapter-api")]
3458pub(crate) async fn run_frontend_http(
3459    runtime: Arc<dyn FrontendRuntime>,
3460    events: broadcast::Sender<FrontendEvent>,
3461    bind: &str,
3462    token: Arc<str>,
3463    runtime_id: impl Into<String>,
3464) -> std::io::Result<FrontendHttpServer> {
3465    let listener = TcpListener::bind(bind).await?;
3466    let address = listener.local_addr()?;
3467    let coordinator = CoordinatedRuntime::new(runtime);
3468    // A REGISTRY, not a frozen list: this listener is the one the live-runtime
3469    // receipt points at, so it mints and revokes scoped frontend credentials
3470    // against its own bootstrap bearer.
3471    let credentials =
3472        RuntimeHttpCredentialRegistry::new(runtime_id, vec![RuntimeHttpCredential::owner(token)])?;
3473    let task = tokio::spawn(async move {
3474        while let Ok((stream, _)) = listener.accept().await {
3475            let coordinator = coordinator.clone();
3476            let events = events.clone();
3477            let credentials = credentials.clone();
3478            tokio::spawn(async move {
3479                let _ = handle_frontend_http_conn(stream, coordinator, events, credentials).await;
3480            });
3481        }
3482    });
3483    Ok(FrontendHttpServer { address, task })
3484}
3485
3486/// Lifetime handle for an authenticated `frontend.v2` WebSocket listener.
3487/// Dropping the handle detaches the listener without closing its SDK runtime.
3488#[cfg(feature = "adapter-api")]
3489pub struct FrontendWebSocketServer {
3490    address: SocketAddr,
3491    task: tokio::task::JoinHandle<()>,
3492}
3493
3494#[cfg(feature = "adapter-api")]
3495impl FrontendWebSocketServer {
3496    /// Actually-bound listener address.
3497    pub fn address(&self) -> SocketAddr {
3498        self.address
3499    }
3500}
3501
3502#[cfg(feature = "adapter-api")]
3503impl Drop for FrontendWebSocketServer {
3504    fn drop(&mut self) {
3505        self.task.abort();
3506    }
3507}
3508
3509/// Publish the language-neutral facade over authenticated WebSocket RPC.
3510/// The endpoint accepts only `/frontend/v2`, reuses the SDK coordinator, and
3511/// emits canonical events as `frontend.v2.event` notifications.
3512#[cfg(feature = "adapter-api")]
3513pub async fn run_frontend_websocket(
3514    engine: Arc<RpcEngine>,
3515    bind: &str,
3516    credentials: Vec<RuntimeHttpCredential>,
3517) -> std::io::Result<FrontendWebSocketServer> {
3518    let runtime: Arc<dyn FrontendRuntime> = engine.clone();
3519    let events = engine.frontend_events.clone();
3520    run_frontend_websocket_runtime_inner(runtime, events, Some(engine), bind, credentials).await
3521}
3522
3523#[cfg(feature = "adapter-api")]
3524async fn run_frontend_websocket_runtime_inner(
3525    runtime: Arc<dyn FrontendRuntime>,
3526    events: broadcast::Sender<FrontendEvent>,
3527    shutdown_engine: Option<Arc<RpcEngine>>,
3528    bind: &str,
3529    credentials: Vec<RuntimeHttpCredential>,
3530) -> std::io::Result<FrontendWebSocketServer> {
3531    if credentials.is_empty()
3532        || credentials
3533            .iter()
3534            .any(|credential| credential.token.is_empty())
3535    {
3536        return Err(std::io::Error::new(
3537            std::io::ErrorKind::InvalidInput,
3538            "at least one non-empty runtime WebSocket credential is required",
3539        ));
3540    }
3541    let listener = TcpListener::bind(bind).await?;
3542    let address = listener.local_addr()?;
3543    let coordinator = CoordinatedRuntime::new(runtime);
3544    let credentials: Arc<[RuntimeHttpCredential]> = credentials.into();
3545    let task = tokio::spawn(async move {
3546        loop {
3547            tokio::select! {
3548                biased;
3549                _ = wait_for_optional_runtime_shutdown(shutdown_engine.as_ref()) => break,
3550                accepted = listener.accept() => {
3551                    let Ok((stream, _)) = accepted else { continue };
3552                    let coordinator = coordinator.clone();
3553                    let credentials = credentials.clone();
3554                    let events = events.clone();
3555                    let shutdown_engine = shutdown_engine.clone();
3556                    tokio::spawn(async move {
3557                        let _ = handle_frontend_websocket(stream, events, shutdown_engine, coordinator, credentials).await;
3558                    });
3559                }
3560            }
3561        }
3562    });
3563    Ok(FrontendWebSocketServer { address, task })
3564}
3565
3566#[cfg(feature = "adapter-api")]
3567async fn wait_for_optional_runtime_shutdown(engine: Option<&Arc<RpcEngine>>) {
3568    match engine {
3569        Some(engine) => engine.wait_for_shutdown().await,
3570        None => std::future::pending().await,
3571    }
3572}
3573
3574#[cfg(feature = "adapter-api")]
3575#[allow(clippy::result_large_err)] // tungstenite's handshake callback fixes this error type.
3576async fn handle_frontend_websocket(
3577    stream: tokio::net::TcpStream,
3578    events: broadcast::Sender<FrontendEvent>,
3579    shutdown_engine: Option<Arc<RpcEngine>>,
3580    coordinator: Arc<CoordinatedRuntime>,
3581    credentials: Arc<[RuntimeHttpCredential]>,
3582) -> Result<(), tokio_tungstenite::tungstenite::Error> {
3583    use std::sync::Mutex as SyncMutex;
3584    use tokio_tungstenite::tungstenite::handshake::server::{ErrorResponse, Request, Response};
3585
3586    let selected = Arc::new(SyncMutex::new(None::<Arc<CoordinatedRuntimeClient>>));
3587    let selected_by_callback = selected.clone();
3588    let socket = tokio_tungstenite::accept_hdr_async(
3589        stream,
3590        move |request: &Request, response: Response| -> Result<Response, ErrorResponse> {
3591            let reject = |status, message: &str| {
3592                tokio_tungstenite::tungstenite::http::Response::builder()
3593                    .status(status)
3594                    .body(Some(message.to_string()))
3595                    .expect("static WebSocket rejection is valid")
3596            };
3597            if request.uri().path() != "/frontend/v2" {
3598                return Err(reject(404, "frontend WebSocket route not found"));
3599            }
3600            let token = request
3601                .headers()
3602                .get("authorization")
3603                .and_then(|value| value.to_str().ok())
3604                .and_then(|value| value.strip_prefix("Bearer "));
3605            let Some(credential) = token.and_then(|token| {
3606                credentials.iter().find(|credential| {
3607                    constant_time_eq(token.as_bytes(), credential.token.as_bytes())
3608                })
3609            }) else {
3610                return Err(reject(401, "missing or invalid bearer token"));
3611            };
3612            let client_id = request
3613                .headers()
3614                .get("x-supercode-client-id")
3615                .and_then(|value| value.to_str().ok())
3616                .unwrap_or("legacy-websocket-owner");
3617            let Ok(client_id) = RuntimeClientId::parse(client_id) else {
3618                return Err(reject(400, "invalid runtime client id"));
3619            };
3620            let mut authorization = credential.authorization.clone();
3621            if let Some(requested) = request
3622                .headers()
3623                .get("x-supercode-permissions")
3624                .and_then(|value| value.to_str().ok())
3625            {
3626                let Ok(requested) = RuntimeAuthorization::parse_header(requested) else {
3627                    return Err(reject(400, "invalid runtime authorization grant"));
3628                };
3629                authorization = authorization.restrict_to(&requested);
3630            }
3631            *selected_by_callback
3632                .lock()
3633                .unwrap_or_else(std::sync::PoisonError::into_inner) =
3634                Some(coordinator.client(client_id, authorization));
3635            Ok(response)
3636        },
3637    )
3638    .await?;
3639    let client = selected
3640        .lock()
3641        .unwrap_or_else(std::sync::PoisonError::into_inner)
3642        .take()
3643        .expect("successful WebSocket handshake selects a runtime client");
3644    if let Err(error) = client.observe() {
3645        let mut socket = socket;
3646        let value = sdk_runtime_rpc_error(Value::Null, -32002, &error).to_string();
3647        socket
3648            .send(tokio_tungstenite::tungstenite::Message::Text(value.into()))
3649            .await?;
3650        socket.close(None).await?;
3651        return Ok(());
3652    }
3653
3654    let mut events = events.subscribe();
3655    let (mut writer, mut reader) = socket.split();
3656    loop {
3657        tokio::select! {
3658            biased;
3659            incoming = reader.next() => match incoming {
3660                Some(Ok(tokio_tungstenite::tungstenite::Message::Text(text))) => {
3661                    let response = match serde_json::from_str::<RpcRequest>(&text) {
3662                        Ok(request) => coordinated_runtime_rpc(client.clone(), request).await,
3663                        Err(error) => rpc_error(Value::Null, -32700, format!("parse error: {error}")),
3664                    };
3665                    writer.send(tokio_tungstenite::tungstenite::Message::Text(response.to_string().into())).await?;
3666                }
3667                Some(Ok(tokio_tungstenite::tungstenite::Message::Ping(payload))) => {
3668                    writer.send(tokio_tungstenite::tungstenite::Message::Pong(payload)).await?;
3669                }
3670                Some(Ok(tokio_tungstenite::tungstenite::Message::Close(_))) | None => break,
3671                Some(Ok(_)) => {}
3672                Some(Err(error)) => {
3673                    client.detach();
3674                    return Err(error);
3675                }
3676            },
3677            event = events.recv() => match event {
3678                Ok(event) => {
3679                    let notification = json!({
3680                        "jsonrpc":"2.0",
3681                        "method":"frontend.v2.event",
3682                        "params":{"event":event},
3683                    });
3684                    writer.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await?;
3685                }
3686                Err(broadcast::error::RecvError::Lagged(count)) => {
3687                    let notification = json!({
3688                        "jsonrpc":"2.0",
3689                        "method":"frontend.v2.event",
3690                        "params":{"error":{"name":"transport","message":format!("event replay gap: {count}")}},
3691                    });
3692                    writer.send(tokio_tungstenite::tungstenite::Message::Text(notification.to_string().into())).await?;
3693                    break;
3694                }
3695                Err(broadcast::error::RecvError::Closed) => break,
3696            },
3697            _ = wait_for_optional_runtime_shutdown(shutdown_engine.as_ref()) => break,
3698        }
3699    }
3700    client.detach();
3701    Ok(())
3702}
3703
3704/// Bind `bind` (`host:port`; `:0` for an OS-assigned ephemeral port) and
3705/// serve the HTTP transport (D8 "remote attach") in a background task until
3706/// `engine` signals shutdown. Returns the actually-bound address (so a
3707/// caller that asked for port `0` can learn the real port). Every
3708/// connection is authenticated per-request via `token` — see
3709/// `check_auth`. The LOOPBACK-BY-DEFAULT policy decision is the caller's
3710/// (see the module doc) — this fn binds whatever address it's given.
3711#[cfg(feature = "adapter-api")]
3712pub async fn run_http(
3713    engine: Arc<RpcEngine>,
3714    bind: &str,
3715    token: Arc<str>,
3716) -> std::io::Result<SocketAddr> {
3717    run_http_authorized(engine, bind, vec![RuntimeHttpCredential::owner(token)]).await
3718}
3719
3720/// Bind an SDK HTTP runtime with multiple independently scoped bearer
3721/// credentials. The token bytes remain server-private; each successful
3722/// authentication produces the exact authorization grant projected by the
3723/// shared runtime coordinator.
3724#[cfg(feature = "adapter-api")]
3725pub async fn run_http_authorized(
3726    engine: Arc<RpcEngine>,
3727    bind: &str,
3728    credentials: Vec<RuntimeHttpCredential>,
3729) -> std::io::Result<SocketAddr> {
3730    run_http_authorized_with_lease_ttl(
3731        engine,
3732        bind,
3733        credentials,
3734        crate::DEFAULT_RUNTIME_LEASE_TTL_MS,
3735    )
3736    .await
3737}
3738
3739/// Test/embedder variant of [`run_http_authorized`] with an explicit
3740/// controller lease duration.
3741#[cfg(feature = "adapter-api")]
3742pub async fn run_http_authorized_with_lease_ttl(
3743    engine: Arc<RpcEngine>,
3744    bind: &str,
3745    credentials: Vec<RuntimeHttpCredential>,
3746    lease_ttl_ms: u64,
3747) -> std::io::Result<SocketAddr> {
3748    if credentials.is_empty()
3749        || credentials
3750            .iter()
3751            .any(|credential| credential.token.is_empty())
3752    {
3753        return Err(std::io::Error::new(
3754            std::io::ErrorKind::InvalidInput,
3755            "at least one non-empty runtime HTTP credential is required",
3756        ));
3757    }
3758    if lease_ttl_ms == 0 {
3759        return Err(std::io::Error::new(
3760            std::io::ErrorKind::InvalidInput,
3761            "runtime lease TTL must be non-zero",
3762        ));
3763    }
3764    let listener = TcpListener::bind(bind).await?;
3765    let local_addr = listener.local_addr()?;
3766    let credentials = RuntimeHttpCredentialRegistry::new(engine.session_id(), credentials)?;
3767    let eng = engine;
3768    let runtime: Arc<dyn FrontendRuntime> = eng.clone();
3769    let coordinator = CoordinatedRuntime::with_lease_ttl(runtime, lease_ttl_ms);
3770    tokio::spawn(async move {
3771        loop {
3772            tokio::select! {
3773                biased;
3774                _ = eng.wait_for_shutdown() => break,
3775                accepted = listener.accept() => {
3776                    let Ok((stream, _addr)) = accepted else { continue };
3777                    let eng = eng.clone();
3778                    let coordinator = coordinator.clone();
3779                    let credentials = credentials.clone();
3780                    tokio::spawn(async move {
3781                        let _ = handle_http_conn(stream, eng, coordinator, credentials).await;
3782                    });
3783                }
3784            }
3785        }
3786    });
3787    Ok(local_addr)
3788}