Skip to main content

supercode_harness/
frontend.rs

1//! Protocol-neutral frontend contract for one SDK-owned Supercode runtime.
2//!
3//! Terminal, HTTP, ACP, and future browser frontends consume this contract;
4//! none of them owns an [`crate::Agent`] or a second model loop.  Events keep
5//! their complete JSON payload and gain a monotonic sequence so a frontend can
6//! cross the history-replay/live-stream boundary without duplicates.
7
8use std::collections::{BTreeMap, VecDeque};
9#[cfg(feature = "adapter-api")]
10use std::sync::atomic::AtomicBool;
11use std::sync::atomic::{AtomicU64, Ordering};
12use std::sync::Arc;
13#[cfg(feature = "adapter-api")]
14use std::sync::Weak;
15
16use async_trait::async_trait;
17#[cfg(feature = "adapter-api")]
18use futures::StreamExt;
19use serde::{Deserialize, Serialize};
20#[cfg(feature = "adapter-api")]
21use serde_json::json;
22use serde_json::Value;
23use tokio::sync::broadcast;
24
25#[cfg(feature = "adapter-api")]
26use crate::sdk::RuntimeSubmitError;
27pub use crate::sdk::SdkError as FrontendRuntimeError;
28pub use crate::sdk::SdkEvent as FrontendEvent;
29pub use crate::sdk::SdkRuntime as FrontendRuntime;
30use crate::server::RpcEngine;
31use crate::ChatMessage;
32
33/// Frontend contract schema version.
34pub const FRONTEND_RUNTIME_SCHEMA_VERSION: u32 = 2;
35
36/// Runtime lifecycle-event schema version.
37///
38/// Operation descriptors evolve the attach contract independently from the
39/// established event payloads consumed by machine frontends.
40pub(crate) const FRONTEND_EVENT_SCHEMA_VERSION: u32 = 1;
41
42/// Maximum sequenced events retained between canonical history snapshots.
43pub const FRONTEND_REPLAY_CAPACITY: usize = 4096;
44
45/// Whether a model/tool turn currently owns the runtime.
46#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
47#[serde(rename_all = "snake_case")]
48pub enum FrontendTurnState {
49    /// The runtime accepts a new turn.
50    Idle,
51    /// A user, scheduler, or tool turn is active.
52    Busy,
53}
54
55/// Frontend-visible runtime connection state.
56#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
57#[serde(rename_all = "snake_case")]
58pub enum FrontendConnectionState {
59    /// The SDK runtime is reachable.
60    Connected,
61    /// Graceful shutdown has been requested.
62    ShuttingDown,
63}
64
65/// Actions the current runtime adapter can actually perform.
66#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
67pub struct FrontendActions {
68    /// Submit a new user turn.
69    pub submit: bool,
70    /// Interrupt an active turn.
71    pub interrupt: bool,
72    /// Queue a steering instruction during a turn.
73    pub steer: bool,
74    /// Answer an approval, elicitation, or other protocol request.
75    pub respond: bool,
76    /// Detach without stopping the runtime.
77    pub detach: bool,
78    /// Close the SDK-owned runtime.
79    pub close: bool,
80}
81
82/// Display semantics emitted by the runtime.
83#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
84pub struct FrontendDisplayCapabilities {
85    /// Known normalized event kinds at this schema version.
86    pub event_kinds: Vec<String>,
87    /// Whether unknown payloads remain available for generic rendering.
88    pub opaque_fallback: bool,
89}
90
91/// One runtime-provided command surfaced by a composer.
92#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
93pub struct FrontendCommandDescriptor {
94    /// Command name without the leading slash.
95    pub name: String,
96    /// Optional short help text.
97    pub description: Option<String>,
98    /// Optional argument usage shown beside the command.
99    #[serde(default, skip_serializing_if = "Option::is_none")]
100    pub argument_hint: Option<String>,
101}
102
103/// Stable family for an explicitly invocable frontend operation.
104///
105/// Families without a production [`FrontendRuntime::invoke`] implementation
106/// are never advertised. Keeping the full vocabulary here lets frontends
107/// render future file/model/session/subagent/image/reduction controls from the
108/// catalog without inferring them from composable modules.
109#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
110#[serde(rename_all = "snake_case")]
111pub enum FrontendOperationKind {
112    /// Invoke a trusted runtime prompt template.
113    Prompt,
114    /// Attach or inspect a file through a typed runtime route.
115    File,
116    /// Inspect or switch the active model through a typed runtime route.
117    Model,
118    /// Perform a session operation through a typed runtime route.
119    Session,
120    /// Perform a subagent operation through a typed runtime route.
121    Subagent,
122    /// Attach an image through a typed runtime route.
123    Image,
124    /// Perform a reversible reduction operation through a typed runtime route.
125    Reduction,
126}
127
128/// One operation the runtime can genuinely invoke.
129#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
130pub struct FrontendOperationDescriptor {
131    /// Stable runtime-scoped identifier supplied back during invocation.
132    pub id: String,
133    /// Typed operation family.
134    pub kind: FrontendOperationKind,
135    /// Optional slash-command trigger rendered by terminal composers.
136    pub command: Option<FrontendCommandDescriptor>,
137}
138
139/// Typed invocation accepted by [`FrontendRuntime::invoke`].
140#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
141#[serde(tag = "kind", rename_all = "snake_case")]
142pub enum FrontendOperationInvocation {
143    /// Expand and submit one advertised trusted prompt template.
144    Prompt {
145        /// Identifier from [`FrontendOperationDescriptor::id`].
146        operation_id: String,
147        /// Free text replacing the prompt template's `{args}` placeholder.
148        arguments: String,
149    },
150}
151
152impl FrontendOperationInvocation {
153    /// Identifier supplied by the runtime catalog.
154    pub fn operation_id(&self) -> &str {
155        match self {
156            Self::Prompt { operation_id, .. } => operation_id,
157        }
158    }
159}
160
161/// Typed result returned by [`FrontendRuntime::invoke`].
162#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
163#[serde(tag = "kind", rename_all = "snake_case")]
164pub enum FrontendOperationResult {
165    /// Reply from a prompt-template turn.
166    Prompt {
167        /// Final assistant reply.
168        reply: String,
169    },
170}
171
172/// Source/emulation identity supplied by the session-loading surface.
173#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
174pub struct FrontendRuntimeMetadata {
175    /// Source harness whose session semantics are being continued.
176    pub source_harness: Option<String>,
177    /// Resolved composable preset/profile name, when one was selected.
178    pub emulation_profile: Option<String>,
179}
180
181/// Complete frontend-facing description of one SDK-owned runtime.
182#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
183pub struct FrontendRuntimeDescriptor {
184    /// Contract schema version.
185    pub schema_version: u32,
186    /// Stable SDK runtime/session identity.
187    pub session_id: String,
188    /// Source harness whose semantics are being emulated.
189    pub source_harness: Option<String>,
190    /// Resolved composable preset/profile name.
191    pub emulation_profile: Option<String>,
192    /// Active composable modules, using their stable config keys.
193    pub active_modules: Vec<String>,
194    /// Runtime-provided composer commands.
195    pub commands: Vec<FrontendCommandDescriptor>,
196    /// Explicit typed operation catalog. Missing on schema-v1 peers.
197    #[serde(default)]
198    pub operations: Vec<FrontendOperationDescriptor>,
199    /// Supported control actions.
200    pub actions: FrontendActions,
201    /// Display/event capabilities.
202    pub display: FrontendDisplayCapabilities,
203    /// Current model label.
204    pub model: String,
205    /// Current turn state.
206    pub turn_state: FrontendTurnState,
207    /// Current connection state.
208    pub connection_state: FrontendConnectionState,
209    /// Compatible client/adapter metadata with no canonical-session meaning.
210    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
211    pub extensions: BTreeMap<String, Value>,
212}
213
214/// Serializable half of an attachment returned by an out-of-process runtime.
215/// The live receiver is transport-owned and joined to this snapshot locally.
216#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
217pub struct FrontendAttachSnapshot {
218    /// Runtime description captured at attachment time.
219    pub descriptor: FrontendRuntimeDescriptor,
220    /// Bounded canonical history through `history_cursor`.
221    pub history: Vec<ChatMessage>,
222    /// Highest event sequence represented by `history`.
223    pub history_cursor: u64,
224    /// Events after the canonical history boundary and before the response.
225    pub replay: VecDeque<FrontendEvent>,
226}
227
228/// Kind of interactive request surfaced by the SDK runtime.
229#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
230#[serde(rename_all = "snake_case")]
231pub enum FrontendRequestKind {
232    /// A tool or sandbox action needs a policy-authorized human decision.
233    Approval,
234    /// An MCP server requested structured user input.
235    Elicitation,
236    /// Another versioned runtime request not known to this frontend build.
237    /// Its complete payload remains available for a generic overlay.
238    #[serde(other)]
239    Other,
240}
241
242/// One pending interactive request, emitted as a sequenced frontend event.
243#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
244pub struct FrontendRequest {
245    /// Runtime-scoped request identifier used exactly once by `respond`.
246    pub id: u64,
247    /// Typed request category.
248    pub kind: FrontendRequestKind,
249    /// Complete request payload, including raw tool arguments or schema.
250    pub payload: Value,
251}
252
253/// Typed approval decision accepted by [`FrontendRuntime::respond`].
254#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
255#[serde(rename_all = "snake_case")]
256pub enum FrontendApprovalDecision {
257    /// Refuse this request.
258    Deny,
259    /// Allow only this request.
260    Allow,
261    /// Allow this request and cache the exact policy key for the session.
262    AllowForSession,
263}
264
265/// MCP elicitation outcome accepted by a frontend response.
266#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
267#[serde(rename_all = "snake_case")]
268pub enum FrontendElicitationAction {
269    /// Submit structured content.
270    Accept,
271    /// Explicitly decline the request.
272    Decline,
273    /// Dismiss the request without a decision.
274    Cancel,
275}
276
277/// Typed response to one SDK-owned interactive request.
278#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
279#[serde(tag = "kind", rename_all = "snake_case")]
280pub enum FrontendResponse {
281    /// Answer an approval request.
282    Approval {
283        /// Identifier from [`FrontendRequest::id`].
284        request_id: u64,
285        /// Human decision.
286        decision: FrontendApprovalDecision,
287    },
288    /// Answer an MCP elicitation request.
289    Elicitation {
290        /// Identifier from [`FrontendRequest::id`].
291        request_id: u64,
292        /// MCP elicitation outcome.
293        action: FrontendElicitationAction,
294        /// Structured content for `accept`.
295        content: Option<Value>,
296    },
297    /// Answer a generic runtime request without discarding its payload.
298    Other {
299        /// Identifier from [`FrontendRequest::id`].
300        request_id: u64,
301        /// Generic accept/decline/cancel outcome.
302        action: FrontendElicitationAction,
303        /// Optional structured response content.
304        content: Option<Value>,
305    },
306}
307
308impl FrontendResponse {
309    pub(crate) fn request_id(&self) -> u64 {
310        match self {
311            Self::Approval { request_id, .. }
312            | Self::Elicitation { request_id, .. }
313            | Self::Other { request_id, .. } => *request_id,
314        }
315    }
316}
317
318/// Atomic history/replay/live attachment to one runtime.
319pub struct FrontendAttachment {
320    /// Runtime description captured at attachment time.
321    pub descriptor: FrontendRuntimeDescriptor,
322    /// Bounded canonical history through `history_cursor`.
323    pub history: Vec<ChatMessage>,
324    /// Highest event sequence already represented by `history`.
325    pub history_cursor: u64,
326    pub(crate) replay: VecDeque<FrontendEvent>,
327    live: broadcast::Receiver<FrontendEvent>,
328    delivered: u64,
329    acknowledged: Option<Arc<AtomicU64>>,
330    _transport_lease: Option<Arc<()>>,
331}
332
333impl FrontendAttachment {
334    /// Build an in-process attachment from a serialized snapshot and a live
335    /// SDK event receiver. Runtime adapters use this constructor in tests and
336    /// protocol bridges without acquiring transport ownership.
337    pub fn from_snapshot(
338        snapshot: FrontendAttachSnapshot,
339        live: broadcast::Receiver<FrontendEvent>,
340    ) -> Self {
341        Self::from_snapshot_after(snapshot, live, 0)
342    }
343
344    /// Build an attachment that resumes after a sequence acknowledged by a
345    /// prior transport connection. Snapshot replay and any overlapping live
346    /// events at or below the cursor are skipped without changing canonical
347    /// history or event payloads.
348    pub fn from_snapshot_after(
349        snapshot: FrontendAttachSnapshot,
350        live: broadcast::Receiver<FrontendEvent>,
351        acknowledged_sequence: u64,
352    ) -> Self {
353        let delivered = snapshot.history_cursor.max(acknowledged_sequence);
354        Self::new_with_delivered(
355            snapshot.descriptor,
356            snapshot.history,
357            snapshot.history_cursor,
358            snapshot.replay,
359            live,
360            None,
361            delivered,
362        )
363    }
364
365    pub(crate) fn new(
366        descriptor: FrontendRuntimeDescriptor,
367        history: Vec<ChatMessage>,
368        history_cursor: u64,
369        replay: VecDeque<FrontendEvent>,
370        live: broadcast::Receiver<FrontendEvent>,
371        transport_lease: Option<Arc<()>>,
372    ) -> Self {
373        let delivered = history_cursor;
374        Self::new_with_delivered(
375            descriptor,
376            history,
377            history_cursor,
378            replay,
379            live,
380            transport_lease,
381            delivered,
382        )
383    }
384
385    fn new_with_delivered(
386        descriptor: FrontendRuntimeDescriptor,
387        history: Vec<ChatMessage>,
388        history_cursor: u64,
389        replay: VecDeque<FrontendEvent>,
390        live: broadcast::Receiver<FrontendEvent>,
391        transport_lease: Option<Arc<()>>,
392        delivered: u64,
393    ) -> Self {
394        Self {
395            descriptor,
396            history,
397            history_cursor,
398            replay,
399            live,
400            delivered,
401            acknowledged: None,
402            _transport_lease: transport_lease,
403        }
404    }
405
406    #[cfg(feature = "adapter-acp")]
407    pub(crate) fn with_acknowledgement(mut self, acknowledged: Arc<AtomicU64>) -> Self {
408        acknowledged.fetch_max(self.history_cursor, Ordering::SeqCst);
409        self.acknowledged = Some(acknowledged);
410        self
411    }
412
413    fn acknowledge(&self, event: &FrontendEvent) {
414        if !event_advances_acknowledgement(event) {
415            return;
416        }
417        if let Some(acknowledged) = &self.acknowledged {
418            acknowledged.fetch_max(event.sequence, Ordering::SeqCst);
419        }
420    }
421
422    /// Receive the next event not already represented by the history or a
423    /// prior replay item. Duplicate events queued during attachment are
424    /// skipped by sequence.
425    pub async fn next_event(&mut self) -> Result<FrontendEvent, FrontendRuntimeError> {
426        loop {
427            let event = match self.next_replay_event() {
428                Some(event) => return Ok(event),
429                None => match self.live.recv().await {
430                    Ok(event) => event,
431                    Err(broadcast::error::RecvError::Lagged(count)) => {
432                        return Err(FrontendRuntimeError::ReplayGap(count));
433                    }
434                    Err(broadcast::error::RecvError::Closed) => {
435                        return Err(FrontendRuntimeError::Closed);
436                    }
437                },
438            };
439            if event.sequence <= self.delivered {
440                continue;
441            }
442            self.delivered = event.sequence;
443            self.acknowledge(&event);
444            return Ok(event);
445        }
446    }
447
448    /// Drain one event from the finite attachment replay without waiting for
449    /// live input. Interactive frontends use this to project the complete
450    /// atomic snapshot before accepting keystrokes, so a historical resolved
451    /// request never appears transiently actionable.
452    pub fn next_replay_event(&mut self) -> Option<FrontendEvent> {
453        while let Some(event) = self.replay.pop_front() {
454            if event.sequence <= self.delivered {
455                continue;
456            }
457            self.delivered = event.sequence;
458            self.acknowledge(&event);
459            return Some(event);
460        }
461        None
462    }
463}
464
465pub(crate) fn event_advances_acknowledgement(event: &FrontendEvent) -> bool {
466    event
467        .payload
468        .pointer("/_meta/supercode/transient")
469        .and_then(Value::as_bool)
470        != Some(true)
471}
472
473/// State protected by `RpcEngine`'s short synchronous projection lock.
474pub(crate) struct FrontendProjectionState {
475    pub(crate) history: Vec<ChatMessage>,
476    pub(crate) history_cursor: u64,
477    pub(crate) next_sequence: u64,
478    pub(crate) replay: VecDeque<FrontendEvent>,
479}
480
481#[async_trait]
482impl FrontendRuntime for RpcEngine {
483    async fn describe(&self) -> Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
484        Ok(self.frontend_descriptor())
485    }
486
487    async fn attach(
488        &self,
489        history_limit: usize,
490    ) -> Result<FrontendAttachment, FrontendRuntimeError> {
491        self.frontend_attach(history_limit)
492    }
493
494    async fn send_input(self: Arc<Self>, prompt: String) -> Result<(), FrontendRuntimeError> {
495        RpcEngine::send_input(&self, prompt)?;
496        Ok(())
497    }
498
499    async fn send_input_with_images(
500        self: Arc<Self>,
501        prompt: String,
502        image_urls: Vec<String>,
503    ) -> Result<(), FrontendRuntimeError> {
504        RpcEngine::send_input_with_images(&self, prompt, image_urls)?;
505        Ok(())
506    }
507
508    async fn submit(&self, prompt: String) -> Result<String, FrontendRuntimeError> {
509        Ok(RpcEngine::submit(self, prompt).await?)
510    }
511
512    async fn submit_with_images(
513        &self,
514        prompt: String,
515        image_urls: Vec<String>,
516    ) -> Result<String, FrontendRuntimeError> {
517        Ok(RpcEngine::submit_with_images(self, prompt, image_urls).await?)
518    }
519
520    async fn interrupt(&self) -> Result<bool, FrontendRuntimeError> {
521        Ok(RpcEngine::interrupt(self).await)
522    }
523
524    async fn steer(&self, prompt: String) -> Result<(), FrontendRuntimeError> {
525        RpcEngine::steer(self, prompt)
526    }
527
528    async fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
529        RpcEngine::respond(self, response)
530    }
531
532    async fn invoke(
533        &self,
534        operation: FrontendOperationInvocation,
535    ) -> Result<FrontendOperationResult, FrontendRuntimeError> {
536        RpcEngine::invoke(self, operation).await
537    }
538
539    async fn close(&self) -> Result<(), FrontendRuntimeError> {
540        RpcEngine::shutdown(self).await;
541        Ok(())
542    }
543}
544
545/// Authenticated HTTP implementation of [`FrontendRuntime`].
546///
547/// It owns only an RPC/SSE connection. The remote [`RpcEngine`] remains the
548/// sole owner of the agent loop, transcript, scheduler, and persistence.
549#[cfg(feature = "adapter-api")]
550pub struct HttpFrontendRuntime {
551    base_url: String,
552    token: String,
553    client_id: crate::RuntimeClientId,
554    authorization: crate::RuntimeAuthorization,
555    client: reqwest::Client,
556    events: broadcast::Sender<FrontendEvent>,
557    next_id: AtomicU64,
558    lifecycle: Arc<()>,
559    disconnected: AtomicBool,
560}
561
562#[cfg(feature = "adapter-api")]
563impl HttpFrontendRuntime {
564    /// Authenticate, verify the frontend descriptor, and establish the
565    /// sequenced SSE stream before returning.
566    pub async fn connect(
567        base_url: impl Into<String>,
568        token: impl Into<String>,
569    ) -> Result<Arc<Self>, FrontendRuntimeError> {
570        let mut random = [0_u8; 16];
571        getrandom::getrandom(&mut random).map_err(|error| {
572            FrontendRuntimeError::Transport(format!(
573                "cannot generate runtime client identity: {error}"
574            ))
575        })?;
576        let suffix = random
577            .iter()
578            .map(|byte| format!("{byte:02x}"))
579            .collect::<String>();
580        let client_id = crate::RuntimeClientId::parse(format!("http-{suffix}"))
581            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
582        Self::connect_with_client_id(base_url, token, client_id).await
583    }
584
585    /// Connect with a caller-owned stable client identity. Reconnect tests
586    /// and external bindings use this to retain deterministic lease state.
587    pub async fn connect_with_client_id(
588        base_url: impl Into<String>,
589        token: impl Into<String>,
590        client_id: crate::RuntimeClientId,
591    ) -> Result<Arc<Self>, FrontendRuntimeError> {
592        Self::connect_with_authorization(
593            base_url,
594            token,
595            client_id,
596            crate::RuntimeAuthorization::owner(),
597        )
598        .await
599    }
600
601    /// Connect while requesting an exact subset of the bearer credential's
602    /// permissions. The server intersects this with the authenticated grant;
603    /// this header can narrow authority but can never elevate it.
604    pub async fn connect_with_authorization(
605        base_url: impl Into<String>,
606        token: impl Into<String>,
607        client_id: crate::RuntimeClientId,
608        authorization: crate::RuntimeAuthorization,
609    ) -> Result<Arc<Self>, FrontendRuntimeError> {
610        Self::connect_inner(base_url, token, client_id, authorization, true)
611            .await
612            .map(|(runtime, _)| runtime)
613    }
614
615    /// Authenticated metadata probe that does not open an event stream or
616    /// register an observer, returning the descriptor the connect handshake
617    /// already fetched. Used by the local runtime registry, which runs this on
618    /// every `harness serve` tick for every followed session: asking the same
619    /// runtime to describe itself twice for one read is pure load on that
620    /// path, and each round trip costs its own loopback connection.
621    pub(crate) async fn probe_described(
622        base_url: impl Into<String>,
623        token: impl Into<String>,
624        client_id: crate::RuntimeClientId,
625    ) -> Result<(Arc<Self>, FrontendRuntimeDescriptor), FrontendRuntimeError> {
626        Self::connect_inner(
627            base_url,
628            token,
629            client_id,
630            crate::RuntimeAuthorization::observer(),
631            false,
632        )
633        .await
634    }
635
636    async fn connect_inner(
637        base_url: impl Into<String>,
638        token: impl Into<String>,
639        client_id: crate::RuntimeClientId,
640        authorization: crate::RuntimeAuthorization,
641        stream_events: bool,
642    ) -> Result<(Arc<Self>, FrontendRuntimeDescriptor), FrontendRuntimeError> {
643        let runtime = Arc::new(Self {
644            base_url: base_url.into().trim_end_matches('/').to_string(),
645            token: token.into(),
646            client_id,
647            authorization,
648            client: reqwest::Client::new(),
649            events: broadcast::channel(1024).0,
650            next_id: AtomicU64::new(1),
651            lifecycle: Arc::new(()),
652            disconnected: AtomicBool::new(false),
653        });
654        // Validate auth and schema before opening a long-lived connection.
655        let descriptor: FrontendRuntimeDescriptor = runtime
656            .rpc_typed(crate::FrontendFacadeMethod::Describe.wire_name(), json!({}))
657            .await?;
658        if stream_events {
659            Self::start_event_stream(&runtime).await?;
660        }
661        Ok((runtime, descriptor))
662    }
663
664    async fn start_event_stream(runtime: &Arc<Self>) -> Result<(), FrontendRuntimeError> {
665        let (ready_tx, ready_rx) = tokio::sync::oneshot::channel();
666        let weak = Arc::downgrade(runtime);
667        let lifecycle = Arc::downgrade(&runtime.lifecycle);
668        tokio::spawn(async move {
669            Self::run_event_stream(weak, lifecycle, ready_tx).await;
670        });
671        ready_rx.await.map_err(|_| {
672            FrontendRuntimeError::Transport("frontend event stream exited before startup".into())
673        })?
674    }
675
676    async fn run_event_stream(
677        weak: Weak<Self>,
678        lifecycle: Weak<()>,
679        ready: tokio::sync::oneshot::Sender<Result<(), FrontendRuntimeError>>,
680    ) {
681        let Some(runtime) = weak.upgrade() else {
682            let _ = ready.send(Err(FrontendRuntimeError::Closed));
683            return;
684        };
685        let request = runtime
686            .client
687            .get(format!("{}/frontend/events", runtime.base_url))
688            .bearer_auth(&runtime.token)
689            .header("x-supercode-client-id", runtime.client_id.as_str())
690            .header(
691                "x-supercode-permissions",
692                runtime.authorization.header_value(),
693            );
694        let events = runtime.events.clone();
695        drop(runtime);
696        let response = request.send().await;
697        let response = match response {
698            Ok(response) if response.status().is_success() => response,
699            Ok(response) => {
700                let _ = ready.send(Err(FrontendRuntimeError::Transport(format!(
701                    "frontend event stream returned {}",
702                    response.status()
703                ))));
704                return;
705            }
706            Err(error) => {
707                let _ = ready.send(Err(FrontendRuntimeError::Transport(error.to_string())));
708                return;
709            }
710        };
711        let _ = ready.send(Ok(()));
712        let mut stream = response.bytes_stream();
713        let mut pending = Vec::<u8>::new();
714        let mut liveness = tokio::time::interval(std::time::Duration::from_millis(100));
715        loop {
716            let chunk = tokio::select! {
717                _ = liveness.tick() => {
718                    if lifecycle.strong_count() == 0 {
719                        break;
720                    }
721                    if weak
722                        .upgrade()
723                        .is_some_and(|runtime| runtime.disconnected.load(Ordering::SeqCst))
724                    {
725                        break;
726                    }
727                    continue;
728                }
729                chunk = stream.next() => chunk,
730            };
731            let Some(chunk) = chunk else {
732                break;
733            };
734            let Ok(chunk) = chunk else {
735                break;
736            };
737            pending.extend_from_slice(&chunk);
738            while let Some(position) = pending.iter().position(|byte| *byte == b'\n') {
739                let line = pending.drain(..=position).collect::<Vec<_>>();
740                let line = String::from_utf8_lossy(&line);
741                let Some(data) = line.trim_end().strip_prefix("data: ") else {
742                    continue;
743                };
744                if let Ok(event) = serde_json::from_str::<FrontendEvent>(data) {
745                    let _ = events.send(event);
746                }
747            }
748        }
749        if let Some(runtime) = weak.upgrade() {
750            runtime.disconnected.store(true, Ordering::SeqCst);
751            let _ = runtime.events.send(FrontendEvent::new(
752                u64::MAX,
753                json!({
754                    "type": "runtime_disconnected",
755                    "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
756                }),
757            ));
758        }
759    }
760
761    async fn rpc(&self, method: &str, params: Value) -> Result<Value, FrontendRuntimeError> {
762        let id = self.next_id.fetch_add(1, Ordering::SeqCst);
763        let requested_operation = params
764            .pointer("/operation/operation_id")
765            .and_then(Value::as_str)
766            .map(str::to_owned);
767        let response = self
768            .client
769            .post(format!("{}/rpc", self.base_url))
770            .bearer_auth(&self.token)
771            .header("x-supercode-client-id", self.client_id.as_str())
772            .header("x-supercode-permissions", self.authorization.header_value())
773            .json(&json!({"id": id, "method": method, "params": params}))
774            .send()
775            .await
776            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
777        if !response.status().is_success() {
778            return Err(FrontendRuntimeError::Transport(format!(
779                "SDK HTTP RPC returned {}",
780                response.status()
781            )));
782        }
783        let value: Value = response
784            .json()
785            .await
786            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
787        if let Some(error) = value.get("error") {
788            let code = error.get("code").and_then(Value::as_i64);
789            let name = error.get("name").and_then(Value::as_str);
790            let operation = error
791                .get("operation")
792                .and_then(Value::as_str)
793                .and_then(crate::SdkOperation::from_action_name);
794            let message = error
795                .get("message")
796                .and_then(Value::as_str)
797                .unwrap_or("SDK runtime request failed")
798                .to_string();
799            return Err(match (name, code) {
800                (Some("unauthenticated"), _) | (_, Some(-32030)) => {
801                    FrontendRuntimeError::Unauthenticated
802                }
803                (Some("unauthorized"), _) | (_, Some(-32031)) => {
804                    FrontendRuntimeError::Unauthorized {
805                        permission: error
806                            .get("permission")
807                            .and_then(Value::as_str)
808                            .unwrap_or("unknown")
809                            .to_string(),
810                    }
811                }
812                (Some("controller_required"), _) | (_, Some(-32032)) => {
813                    FrontendRuntimeError::ControllerRequired {
814                        holder: error
815                            .get("holder")
816                            .and_then(Value::as_str)
817                            .map(str::to_owned),
818                        expires_at_ms: error.get("expiresAtMs").and_then(Value::as_u64),
819                    }
820                }
821                (Some("lease_expired"), _) | (_, Some(-32033)) => {
822                    FrontendRuntimeError::LeaseExpired
823                }
824                (_, Some(-32023)) => FrontendRuntimeError::UnsupportedOperation(
825                    requested_operation.unwrap_or(message),
826                ),
827                (Some("unsupported_action"), _) => FrontendRuntimeError::UnsupportedAction(
828                    operation
829                        .unwrap_or_else(|| {
830                            crate::SdkOperation::from_action_name(method)
831                                .unwrap_or(crate::SdkOperation::Respond)
832                        })
833                        .action_name(),
834                ),
835                (Some("not_found"), Some(-32021)) => {
836                    let request_id = params
837                        .pointer("/response/request_id")
838                        .and_then(Value::as_u64)
839                        .unwrap_or_default();
840                    FrontendRuntimeError::UnknownRequest(request_id)
841                }
842                (Some("invalid_argument"), _) => FrontendRuntimeError::InvalidResponse(message),
843                (_, Some(-32000)) => RuntimeSubmitError::Busy.into(),
844                (_, Some(-32001)) => RuntimeSubmitError::Interrupted.into(),
845                (_, Some(-32002)) => RuntimeSubmitError::Agent(message).into(),
846                (_, Some(-32020)) => FrontendRuntimeError::UnsupportedAction(
847                    crate::SdkOperation::from_action_name(method)
848                        .unwrap_or(crate::SdkOperation::Respond)
849                        .action_name(),
850                ),
851                (_, Some(-32021)) => {
852                    let request_id = params
853                        .pointer("/response/request_id")
854                        .and_then(Value::as_u64)
855                        .unwrap_or_default();
856                    FrontendRuntimeError::UnknownRequest(request_id)
857                }
858                (_, Some(-32022)) => FrontendRuntimeError::InvalidResponse(message),
859                _ => FrontendRuntimeError::Transport(message),
860            });
861        }
862        Ok(value.get("result").cloned().unwrap_or(Value::Null))
863    }
864
865    async fn rpc_typed<T: serde::de::DeserializeOwned>(
866        &self,
867        method: &str,
868        params: Value,
869    ) -> Result<T, FrontendRuntimeError> {
870        serde_json::from_value(self.rpc(method, params).await?)
871            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
872    }
873
874    /// Current authenticated client identity.
875    pub fn client_id(&self) -> &crate::RuntimeClientId {
876        &self.client_id
877    }
878
879    /// Whether the remote event stream has already ended or this client
880    /// explicitly detached/closed.
881    pub fn is_disconnected(&self) -> bool {
882        self.disconnected.load(Ordering::SeqCst)
883    }
884
885    /// Explicitly acquire the controller lease, displacing another
886    /// interactive client only through this named operation.
887    pub async fn take_control(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
888        self.rpc_typed(
889            crate::FrontendFacadeMethod::TakeControl.wire_name(),
890            json!({}),
891        )
892        .await
893    }
894
895    /// Renew observer activity and any controller lease owned by this client.
896    pub async fn heartbeat(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
897        self.rpc_typed(
898            crate::FrontendFacadeMethod::Heartbeat.wire_name(),
899            json!({}),
900        )
901        .await
902    }
903
904    /// Read the coordinated ownership state.
905    pub async fn lease_snapshot(
906        &self,
907    ) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
908        self.rpc_typed(crate::FrontendFacadeMethod::Lease.wire_name(), json!({}))
909            .await
910    }
911
912    /// Release observer and controller state without stopping the runtime.
913    pub async fn detach(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
914        let snapshot = self
915            .rpc_typed(crate::FrontendFacadeMethod::Detach.wire_name(), json!({}))
916            .await?;
917        self.disconnected.store(true, Ordering::SeqCst);
918        Ok(snapshot)
919    }
920}
921
922#[async_trait]
923#[cfg(feature = "adapter-api")]
924impl FrontendRuntime for HttpFrontendRuntime {
925    async fn describe(&self) -> Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
926        self.rpc_typed(crate::FrontendFacadeMethod::Describe.wire_name(), json!({}))
927            .await
928    }
929
930    async fn attach(
931        &self,
932        history_limit: usize,
933    ) -> Result<FrontendAttachment, FrontendRuntimeError> {
934        if self.disconnected.load(Ordering::SeqCst) {
935            return Err(FrontendRuntimeError::Closed);
936        }
937        // Subscribe locally before asking the server for its atomic snapshot.
938        // Anything concurrently received over SSE is either in snapshot.replay
939        // or queued here; sequence filtering removes the overlap.
940        let live = self.events.subscribe();
941        let snapshot: FrontendAttachSnapshot = self
942            .rpc_typed(
943                crate::FrontendFacadeMethod::Attach.wire_name(),
944                json!({"limit": history_limit}),
945            )
946            .await?;
947        Ok(FrontendAttachment::new(
948            snapshot.descriptor,
949            snapshot.history,
950            snapshot.history_cursor,
951            snapshot.replay,
952            live,
953            Some(self.lifecycle.clone()),
954        ))
955    }
956
957    async fn send_input(self: Arc<Self>, prompt: String) -> Result<(), FrontendRuntimeError> {
958        self.send_input_with_images(prompt, Vec::new()).await
959    }
960
961    async fn send_input_with_images(
962        self: Arc<Self>,
963        prompt: String,
964        image_urls: Vec<String>,
965    ) -> Result<(), FrontendRuntimeError> {
966        self.rpc(
967            crate::FrontendFacadeMethod::SendInput.wire_name(),
968            json!({"prompt": prompt, "image_urls": image_urls}),
969        )
970        .await?;
971        Ok(())
972    }
973
974    async fn submit(&self, prompt: String) -> Result<String, FrontendRuntimeError> {
975        let result = self
976            .rpc(
977                crate::FrontendFacadeMethod::Submit.wire_name(),
978                json!({"prompt": prompt}),
979            )
980            .await?;
981        Ok(result
982            .get("reply")
983            .and_then(Value::as_str)
984            .unwrap_or_default()
985            .to_string())
986    }
987
988    async fn submit_with_images(
989        &self,
990        prompt: String,
991        image_urls: Vec<String>,
992    ) -> Result<String, FrontendRuntimeError> {
993        let result = self
994            .rpc(
995                crate::FrontendFacadeMethod::Submit.wire_name(),
996                json!({"prompt": prompt, "image_urls": image_urls}),
997            )
998            .await?;
999        Ok(result
1000            .get("reply")
1001            .and_then(Value::as_str)
1002            .unwrap_or_default()
1003            .to_string())
1004    }
1005
1006    async fn interrupt(&self) -> Result<bool, FrontendRuntimeError> {
1007        let result = self
1008            .rpc(
1009                crate::FrontendFacadeMethod::Interrupt.wire_name(),
1010                json!({}),
1011            )
1012            .await?;
1013        Ok(result
1014            .get("interrupted")
1015            .and_then(Value::as_bool)
1016            .unwrap_or(false))
1017    }
1018
1019    async fn steer(&self, prompt: String) -> Result<(), FrontendRuntimeError> {
1020        self.rpc(
1021            crate::FrontendFacadeMethod::Steer.wire_name(),
1022            json!({"prompt": prompt}),
1023        )
1024        .await?;
1025        Ok(())
1026    }
1027
1028    async fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
1029        self.rpc(
1030            crate::FrontendFacadeMethod::Respond.wire_name(),
1031            json!({"response": response}),
1032        )
1033        .await?;
1034        Ok(())
1035    }
1036
1037    async fn invoke(
1038        &self,
1039        operation: FrontendOperationInvocation,
1040    ) -> Result<FrontendOperationResult, FrontendRuntimeError> {
1041        self.rpc_typed(
1042            crate::FrontendFacadeMethod::Invoke.wire_name(),
1043            json!({"operation": operation}),
1044        )
1045        .await
1046    }
1047
1048    async fn lease_snapshot(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1049        HttpFrontendRuntime::lease_snapshot(self).await
1050    }
1051
1052    async fn take_control(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1053        HttpFrontendRuntime::take_control(self).await
1054    }
1055
1056    async fn heartbeat(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1057        HttpFrontendRuntime::heartbeat(self).await
1058    }
1059
1060    async fn detach(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1061        HttpFrontendRuntime::detach(self).await
1062    }
1063
1064    async fn close(&self) -> Result<(), FrontendRuntimeError> {
1065        self.rpc(crate::FrontendFacadeMethod::Close.wire_name(), json!({}))
1066            .await?;
1067        self.disconnected.store(true, Ordering::SeqCst);
1068        Ok(())
1069    }
1070}