Skip to main content

supercode_harness/runtime/
hosted.rs

1//! Multi-frontend host for one harness-native runtime connection.
2//!
3//! The native ACP/RPC/app-server process remains the single model loop. This
4//! host serializes control, broadcasts the lossless native events to SDK
5//! clients, and projects the same events onto Supercode's frontend contract
6//! so a terminal can join an editor-created session without resuming it.
7
8use std::collections::VecDeque;
9use std::sync::atomic::{AtomicBool, Ordering};
10use std::sync::{Arc, Mutex as StdMutex};
11
12use async_trait::async_trait;
13use serde_json::{json, Value};
14use tokio::sync::{broadcast, mpsc, oneshot};
15
16use super::{HarnessEvent, RuntimeCapabilities, RuntimeConnection, RuntimeHandle, RuntimeInput};
17use crate::frontend::{
18    FrontendActions, FrontendAttachment, FrontendConnectionState, FrontendDisplayCapabilities,
19    FrontendEvent, FrontendRuntime, FrontendRuntimeDescriptor, FrontendRuntimeError,
20    FrontendTurnState, FRONTEND_REPLAY_CAPACITY, FRONTEND_RUNTIME_SCHEMA_VERSION,
21};
22use crate::server::RuntimeSubmitError;
23use crate::{Error, FrontendResponse, Result};
24
25enum HostCommand {
26    Submit {
27        input: RuntimeInput,
28        reply: oneshot::Sender<Result<Option<String>>>,
29    },
30    Interrupt {
31        reply: oneshot::Sender<Result<()>>,
32    },
33    Steer {
34        text: String,
35        reply: oneshot::Sender<Result<()>>,
36    },
37    Respond {
38        request_id: Value,
39        response: Value,
40        reply: oneshot::Sender<Result<()>>,
41    },
42    Shutdown {
43        reply: oneshot::Sender<Result<()>>,
44    },
45}
46
47struct ProjectionState {
48    next_sequence: u64,
49    replay: VecDeque<FrontendEvent>,
50}
51
52/// Shared owner of a single harness-native runtime.
53pub struct HostedHarnessRuntime {
54    handle: RuntimeHandle,
55    capabilities: RuntimeCapabilities,
56    commands: mpsc::Sender<HostCommand>,
57    raw_events: broadcast::Sender<HarnessEvent>,
58    frontend_events: broadcast::Sender<FrontendEvent>,
59    projection: StdMutex<ProjectionState>,
60    busy: AtomicBool,
61    closed: AtomicBool,
62}
63
64impl HostedHarnessRuntime {
65    /// Promote one exclusive native connection into a multi-frontend host and
66    /// return the first SDK connection to it.
67    pub fn spawn(
68        runtime: Box<dyn RuntimeConnection>,
69        capabilities: RuntimeCapabilities,
70    ) -> (Arc<Self>, HostedHarnessConnection) {
71        let handle = runtime.handle().clone();
72        let (commands, command_rx) = mpsc::channel(32);
73        let (raw_events, raw_rx) = broadcast::channel(1024);
74        let (frontend_events, _) = broadcast::channel(1024);
75        let host = Arc::new(Self {
76            handle: handle.clone(),
77            capabilities,
78            commands,
79            raw_events,
80            frontend_events,
81            projection: StdMutex::new(ProjectionState {
82                next_sequence: 1,
83                replay: VecDeque::new(),
84            }),
85            busy: AtomicBool::new(false),
86            closed: AtomicBool::new(false),
87        });
88        tokio::spawn(run_native_runtime(
89            runtime,
90            Arc::downgrade(&host),
91            command_rx,
92        ));
93        let connection = HostedHarnessConnection {
94            host: host.clone(),
95            handle,
96            events: raw_rx,
97            closed: false,
98        };
99        (host, connection)
100    }
101
102    /// Subscribe to the canonical sequenced frontend event stream.
103    pub fn frontend_sender(&self) -> broadcast::Sender<FrontendEvent> {
104        self.frontend_events.clone()
105    }
106
107    /// Shut down the one native process owned by this host.
108    pub async fn shutdown(&self) -> Result<()> {
109        if self.closed.load(Ordering::SeqCst) {
110            return Ok(());
111        }
112        let (reply, response) = oneshot::channel();
113        self.commands
114            .send(HostCommand::Shutdown { reply })
115            .await
116            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
117        response
118            .await
119            .map_err(|_| Error::Other("hosted harness runtime stopped before shutdown".into()))?
120    }
121
122    fn claim_submit(&self) -> Result<()> {
123        if self.closed.load(Ordering::SeqCst) {
124            return Err(Error::Other("hosted harness runtime is closed".into()));
125        }
126        if self
127            .busy
128            .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
129            .is_err()
130        {
131            return Err(Error::Other("a harness turn is already in progress".into()));
132        }
133        Ok(())
134    }
135
136    async fn submit_native_claimed(&self, input: RuntimeInput) -> Result<Option<String>> {
137        self.publish(json!({"type":"user_message", "text":input.text}));
138        self.publish(json!({"type":"turn_started"}));
139        let (reply, response) = oneshot::channel();
140        if self
141            .commands
142            .send(HostCommand::Submit { input, reply })
143            .await
144            .is_err()
145        {
146            self.busy.store(false, Ordering::SeqCst);
147            self.publish(
148                json!({"type":"turn_failed", "message":"Hosted harness runtime is closed."}),
149            );
150            self.publish(json!({"type":"turn_completed"}));
151            self.mark_closed("Harness runtime command channel closed.");
152            return Err(Error::Other("hosted harness runtime is closed".into()));
153        }
154        match response.await {
155            Ok(Ok(turn)) => Ok(turn),
156            Ok(Err(error)) => {
157                self.busy.store(false, Ordering::SeqCst);
158                self.publish(json!({"type":"turn_failed", "message":error.to_string()}));
159                self.publish(json!({"type":"turn_completed"}));
160                Err(error)
161            }
162            Err(_) => {
163                self.busy.store(false, Ordering::SeqCst);
164                self.publish(json!({"type":"turn_failed", "message":"Hosted harness runtime stopped before accepting input."}));
165                self.publish(json!({"type":"turn_completed"}));
166                self.mark_closed("Harness runtime stopped before accepting input.");
167                Err(Error::Other(
168                    "hosted harness runtime stopped before accepting input".into(),
169                ))
170            }
171        }
172    }
173
174    async fn submit_native(&self, text: String) -> Result<Option<String>> {
175        self.claim_submit()?;
176        self.submit_native_claimed(RuntimeInput {
177            text,
178            image_urls: Vec::new(),
179        })
180        .await
181    }
182
183    async fn interrupt_native(&self) -> Result<()> {
184        let (reply, response) = oneshot::channel();
185        self.commands
186            .send(HostCommand::Interrupt { reply })
187            .await
188            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
189        response
190            .await
191            .map_err(|_| Error::Other("hosted harness runtime stopped before interrupt".into()))?
192    }
193
194    async fn steer_native(&self, text: String) -> Result<()> {
195        let (reply, response) = oneshot::channel();
196        self.commands
197            .send(HostCommand::Steer { text, reply })
198            .await
199            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
200        response
201            .await
202            .map_err(|_| Error::Other("hosted harness runtime stopped before steering".into()))?
203    }
204
205    async fn respond_native(&self, request_id: Value, response: Value) -> Result<()> {
206        let (reply, completed) = oneshot::channel();
207        self.commands
208            .send(HostCommand::Respond {
209                request_id,
210                response,
211                reply,
212            })
213            .await
214            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
215        completed
216            .await
217            .map_err(|_| Error::Other("hosted harness runtime stopped before response".into()))?
218    }
219
220    fn publish(&self, payload: Value) {
221        let event = {
222            let mut projection = self
223                .projection
224                .lock()
225                .unwrap_or_else(std::sync::PoisonError::into_inner);
226            let event = FrontendEvent::new(projection.next_sequence, payload);
227            projection.next_sequence = projection.next_sequence.saturating_add(1);
228            projection.replay.push_back(event.clone());
229            while projection.replay.len() > FRONTEND_REPLAY_CAPACITY {
230                projection.replay.pop_front();
231            }
232            event
233        };
234        let _ = self.frontend_events.send(event);
235    }
236
237    fn accept_native_event(&self, event: HarnessEvent) {
238        let _ = self.raw_events.send(event.clone());
239        for payload in project_native_event(self.handle.harness.as_str(), &event) {
240            let terminal = matches!(
241                payload.get("type").and_then(Value::as_str),
242                Some("turn_succeeded" | "turn_interrupted" | "turn_failed")
243            );
244            if terminal {
245                self.busy.store(false, Ordering::SeqCst);
246            }
247            self.publish(payload);
248        }
249    }
250
251    fn mark_closed(&self, message: impl Into<String>) {
252        if self.closed.swap(true, Ordering::SeqCst) {
253            return;
254        }
255        let message = message.into();
256        self.busy.store(false, Ordering::SeqCst);
257        // The raw SDK connection and the portable frontend must observe the
258        // same terminal edge. Keeping the sender alive inside `self` otherwise
259        // leaves the SDK receiver waiting forever after native EOF.
260        let _ = self.raw_events.send(HarnessEvent {
261            sequence: None,
262            kind: "transport_closed".into(),
263            payload: json!({"message":message, "terminal":true}),
264        });
265        self.publish(json!({"type":"runtime_disconnected", "message":message}));
266    }
267
268    fn descriptor(&self) -> FrontendRuntimeDescriptor {
269        FrontendRuntimeDescriptor {
270            schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
271            session_id: self.handle.runtime_id.clone(),
272            source_harness: Some(self.handle.harness.as_str().to_string()),
273            emulation_profile: None,
274            active_modules: Vec::new(),
275            commands: Vec::new(),
276            operations: Vec::new(),
277            actions: FrontendActions {
278                submit: self.capabilities.send_input,
279                interrupt: self.capabilities.interrupt,
280                steer: self.capabilities.steer,
281                // The owner/editor retains the native request-response UI.
282                // A terminal may observe requests but cannot guess native
283                // response envelopes through the portable frontend contract.
284                respond: false,
285                detach: true,
286                close: false,
287            },
288            display: FrontendDisplayCapabilities {
289                event_kinds: vec![
290                    "user_message".into(),
291                    "turn_started".into(),
292                    "turn_succeeded".into(),
293                    "turn_interrupted".into(),
294                    "turn_failed".into(),
295                    "text_delta".into(),
296                    "reasoning".into(),
297                    "tool_call_started".into(),
298                    "tool_call_completed".into(),
299                    "native_event".into(),
300                    "runtime_disconnected".into(),
301                ],
302                opaque_fallback: true,
303            },
304            model: self.handle.harness.as_str().to_string(),
305            turn_state: if self.busy.load(Ordering::SeqCst) {
306                FrontendTurnState::Busy
307            } else {
308                FrontendTurnState::Idle
309            },
310            connection_state: if self.closed.load(Ordering::SeqCst) {
311                FrontendConnectionState::ShuttingDown
312            } else {
313                FrontendConnectionState::Connected
314            },
315            extensions: Default::default(),
316        }
317    }
318}
319
320#[async_trait]
321impl FrontendRuntime for HostedHarnessRuntime {
322    async fn describe(
323        &self,
324    ) -> std::result::Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
325        Ok(self.descriptor())
326    }
327
328    async fn attach(
329        &self,
330        _history_limit: usize,
331    ) -> std::result::Result<FrontendAttachment, FrontendRuntimeError> {
332        let live = self.frontend_events.subscribe();
333        let projection = self
334            .projection
335            .lock()
336            .unwrap_or_else(std::sync::PoisonError::into_inner);
337        let replay = projection.replay.clone();
338        if let Some(first) = replay.front() {
339            if first.sequence > 1 {
340                return Err(FrontendRuntimeError::ReplayGap(first.sequence - 1));
341            }
342        }
343        Ok(FrontendAttachment::new(
344            self.descriptor(),
345            Vec::new(),
346            0,
347            replay,
348            live,
349            None,
350        ))
351    }
352
353    async fn send_input(
354        self: Arc<Self>,
355        prompt: String,
356    ) -> std::result::Result<(), FrontendRuntimeError> {
357        self.claim_submit().map_err(hosted_submit_error)?;
358        tokio::spawn(async move {
359            let _ = self
360                .submit_native_claimed(RuntimeInput {
361                    text: prompt,
362                    image_urls: Vec::new(),
363                })
364                .await;
365        });
366        Ok(())
367    }
368
369    async fn send_input_with_images(
370        self: Arc<Self>,
371        prompt: String,
372        image_urls: Vec<String>,
373    ) -> std::result::Result<(), FrontendRuntimeError> {
374        self.claim_submit().map_err(hosted_submit_error)?;
375        tokio::spawn(async move {
376            let _ = self
377                .submit_native_claimed(RuntimeInput {
378                    text: prompt,
379                    image_urls,
380                })
381                .await;
382        });
383        Ok(())
384    }
385
386    async fn submit(&self, prompt: String) -> std::result::Result<String, FrontendRuntimeError> {
387        self.submit_native(prompt)
388            .await
389            .map(|turn| turn.unwrap_or_default())
390            .map_err(hosted_submit_error)
391    }
392
393    async fn interrupt(&self) -> std::result::Result<bool, FrontendRuntimeError> {
394        if !self.busy.load(Ordering::SeqCst) {
395            return Ok(false);
396        }
397        self.interrupt_native()
398            .await
399            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
400        // Acceptance of the interrupt request does not mean the native turn
401        // has stopped. Keep the shared runtime busy until its terminal event,
402        // so another frontend cannot race a prompt into turn teardown.
403        Ok(true)
404    }
405
406    async fn steer(&self, prompt: String) -> std::result::Result<(), FrontendRuntimeError> {
407        if !self.busy.load(Ordering::SeqCst) || !self.capabilities.steer {
408            return Err(FrontendRuntimeError::UnsupportedAction("steer"));
409        }
410        self.steer_native(prompt)
411            .await
412            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
413    }
414
415    async fn respond(
416        &self,
417        _response: FrontendResponse,
418    ) -> std::result::Result<(), FrontendRuntimeError> {
419        Err(FrontendRuntimeError::UnsupportedAction("respond"))
420    }
421}
422
423fn hosted_submit_error(error: Error) -> FrontendRuntimeError {
424    if error.to_string().contains("already in progress") {
425        FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)
426    } else {
427        FrontendRuntimeError::Transport(error.to_string())
428    }
429}
430
431/// SDK connection to a hosted native runtime. Closing it shuts down the
432/// owner runtime; terminal HTTP attachments are non-owning frontend leases.
433pub struct HostedHarnessConnection {
434    host: Arc<HostedHarnessRuntime>,
435    handle: RuntimeHandle,
436    events: broadcast::Receiver<HarnessEvent>,
437    closed: bool,
438}
439
440#[async_trait]
441impl RuntimeConnection for HostedHarnessConnection {
442    fn handle(&self) -> &RuntimeHandle {
443        &self.handle
444    }
445
446    async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
447        self.host.claim_submit()?;
448        self.host.submit_native_claimed(input).await
449    }
450
451    async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
452        match self.events.recv().await {
453            Ok(event) => Ok(Some(event)),
454            Err(broadcast::error::RecvError::Lagged(count)) => Err(Error::Other(format!(
455                "hosted harness event stream lost {count} event(s)"
456            ))),
457            Err(broadcast::error::RecvError::Closed) => Ok(None),
458        }
459    }
460
461    async fn interrupt(&mut self) -> Result<()> {
462        self.host.interrupt_native().await
463    }
464
465    async fn steer(&mut self, text: String) -> Result<()> {
466        self.host.steer_native(text).await
467    }
468
469    async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
470        self.host.respond_native(request_id, response).await
471    }
472
473    async fn close(&mut self) -> Result<()> {
474        if self.closed {
475            return Ok(());
476        }
477        self.closed = true;
478        self.host.shutdown().await
479    }
480}
481
482async fn run_native_runtime(
483    mut runtime: Box<dyn RuntimeConnection>,
484    host: std::sync::Weak<HostedHarnessRuntime>,
485    mut commands: mpsc::Receiver<HostCommand>,
486) {
487    loop {
488        tokio::select! {
489            command = commands.recv() => {
490                let Some(command) = command else {
491                    let _ = runtime.close().await;
492                    return;
493                };
494                match command {
495                    HostCommand::Submit { input, reply } => {
496                        let _ = reply.send(runtime.send_input(input).await);
497                    }
498                    HostCommand::Interrupt { reply } => {
499                        let _ = reply.send(runtime.interrupt().await);
500                    }
501                    HostCommand::Steer { text, reply } => {
502                        let _ = reply.send(runtime.steer(text).await);
503                    }
504                    HostCommand::Respond { request_id, response, reply } => {
505                        let _ = reply.send(runtime.respond(request_id, response).await);
506                    }
507                    HostCommand::Shutdown { reply } => {
508                        let result = runtime.close().await;
509                        let _ = reply.send(result);
510                        if let Some(host) = host.upgrade() {
511                            host.mark_closed("Harness runtime closed.");
512                        }
513                        return;
514                    }
515                }
516            }
517            event = runtime.next_event() => {
518                let Some(host) = host.upgrade() else {
519                    let _ = runtime.close().await;
520                    return;
521                };
522                match event {
523                    Ok(Some(event)) => host.accept_native_event(event),
524                    Ok(None) => {
525                        host.mark_closed("Harness runtime transport closed.");
526                        return;
527                    }
528                    Err(error) => {
529                        host.mark_closed(error.to_string());
530                        return;
531                    }
532                }
533            }
534        }
535    }
536}
537
538fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
539    let key = event.kind.to_ascii_lowercase().replace('-', "_");
540    let payload = &event.payload;
541    if key == "transport_closed" {
542        return vec![
543            json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
544        ];
545    }
546    if key == "transport_error" || key == "error" {
547        return vec![
548            json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
549        ];
550    }
551    if key == "session/update" {
552        let update = payload
553            .pointer("/params/update")
554            .or_else(|| payload.get("update"))
555            .unwrap_or(payload);
556        let update_kind = update
557            .get("sessionUpdate")
558            .or_else(|| update.get("type"))
559            .and_then(Value::as_str)
560            .unwrap_or_default()
561            .to_ascii_lowercase();
562        if update_kind == "agent_message_chunk" {
563            return vec![
564                json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
565            ];
566        }
567        if update_kind == "agent_thought_chunk" {
568            return vec![
569                json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
570            ];
571        }
572        if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
573            return vec![project_tool(update, payload)];
574        }
575    }
576    if key == "supercode/acp_request_completed" {
577        let failure = payload
578            .pointer("/params/error")
579            .or_else(|| payload.get("error"));
580        return completion(failure.and_then(extract_text));
581    }
582    match harness {
583        "codex" => project_codex(&key, payload),
584        "claude-code" => project_claude(&key, payload),
585        "pi" => project_pi(&key, payload),
586        "opencode" => project_opencode(&key, payload),
587        _ => project_generic(&key, payload),
588    }
589}
590
591fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
592    if key == "turn/started" {
593        return vec![native_payload(key, payload)];
594    }
595    if key == "turn/completed" {
596        let status = payload
597            .pointer("/params/turn/status")
598            .or_else(|| payload.pointer("/turn/status"))
599            .and_then(Value::as_str)
600            .unwrap_or("completed")
601            .to_ascii_lowercase();
602        return completion(
603            (status.contains("fail") || status.contains("error") || status.contains("cancel"))
604                .then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
605        );
606    }
607    if key.ends_with("/delta") {
608        let text = payload
609            .pointer("/params/delta")
610            .or_else(|| payload.get("delta"))
611            .and_then(extract_text)
612            .or_else(|| extract_text(payload))
613            .unwrap_or_default();
614        return vec![
615            json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
616        ];
617    }
618    if key.contains("commandexecution")
619        || key.contains("mcptool")
620        || key.contains("filechange")
621        || key.contains("tool")
622    {
623        let source = payload
624            .pointer("/params/item")
625            .or_else(|| payload.get("item"))
626            .unwrap_or(payload);
627        return vec![project_tool(source, payload)];
628    }
629    vec![native_payload(key, payload)]
630}
631
632fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
633    if key == "result" {
634        let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
635            || payload.get("subtype").and_then(Value::as_str) == Some("error");
636        return completion(
637            failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
638        );
639    }
640    if key == "stream_event" {
641        let stream = payload
642            .get("event")
643            .or_else(|| payload.get("stream_event"))
644            .unwrap_or(payload);
645        if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
646            let delta = stream.get("delta").unwrap_or(stream);
647            let reasoning = delta
648                .get("type")
649                .and_then(Value::as_str)
650                .is_some_and(|kind| kind.contains("thinking"));
651            return vec![
652                json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":extract_text(delta).unwrap_or_default(), "raw":payload}),
653            ];
654        }
655    }
656    if key == "assistant" {
657        let content = payload
658            .pointer("/message/content")
659            .or_else(|| payload.get("content"))
660            .unwrap_or(payload);
661        return vec![
662            json!({"type":"text_delta", "text":extract_text(content).unwrap_or_default(), "raw":payload}),
663        ];
664    }
665    project_generic(key, payload)
666}
667
668fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
669    match key {
670        "agent_start" => vec![native_payload(key, payload)],
671        "agent_end" => completion(payload.get("error").and_then(extract_text)),
672        "message_update" => {
673            let update = payload
674                .get("assistantMessageEvent")
675                .or_else(|| payload.get("event"))
676                .unwrap_or(payload);
677            let reasoning = update
678                .get("type")
679                .and_then(Value::as_str)
680                .is_some_and(|kind| kind.contains("thinking"));
681            vec![
682                json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":extract_text(update).unwrap_or_default(), "raw":payload}),
683            ]
684        }
685        _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
686        _ => vec![native_payload(key, payload)],
687    }
688}
689
690fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
691    if key == "session.idle" {
692        return completion(None);
693    }
694    if key == "session.status" {
695        let status = payload
696            .pointer("/properties/status/type")
697            .or_else(|| payload.pointer("/status/type"))
698            .and_then(Value::as_str)
699            .unwrap_or_default();
700        if status == "busy" {
701            return vec![native_payload(key, payload)];
702        }
703        if status == "idle" {
704            return completion(None);
705        }
706    }
707    if key == "message.part.updated" {
708        let part = payload
709            .pointer("/properties/part")
710            .or_else(|| payload.get("part"))
711            .unwrap_or(payload);
712        if let Some(delta) = payload
713            .pointer("/properties/delta")
714            .or_else(|| payload.get("delta"))
715            .and_then(Value::as_str)
716        {
717            let reasoning = part
718                .get("type")
719                .and_then(Value::as_str)
720                .is_some_and(|kind| kind.contains("reasoning"));
721            return vec![
722                json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
723            ];
724        }
725        if part
726            .get("type")
727            .and_then(Value::as_str)
728            .is_some_and(|kind| kind.contains("tool"))
729        {
730            return vec![project_tool(part, payload)];
731        }
732    }
733    if key == "session.error" {
734        return completion(Some(
735            extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
736        ));
737    }
738    vec![native_payload(key, payload)]
739}
740
741fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
742    match key {
743        "turn_started" | "turn/started" | "agent_start" => {
744            vec![native_payload(key, payload)]
745        }
746        "turn_completed" | "turn/completed" | "agent_end" => {
747            completion(payload.get("error").and_then(extract_text))
748        }
749        "output_delta" | "content_delta" => vec![
750            json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
751        ],
752        "reasoning_delta" => vec![
753            json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
754        ],
755        "tool" => vec![project_tool(payload, payload)],
756        _ => vec![native_payload(key, payload)],
757    }
758}
759
760fn completion(error: Option<String>) -> Vec<Value> {
761    match error {
762        Some(message) => vec![
763            json!({"type":"turn_failed", "message":message}),
764            json!({"type":"turn_completed"}),
765        ],
766        None => vec![
767            json!({"type":"turn_succeeded"}),
768            json!({"type":"turn_completed"}),
769        ],
770    }
771}
772
773fn project_tool(source: &Value, raw: &Value) -> Value {
774    let status = source
775        .get("status")
776        .or_else(|| source.get("state"))
777        .or_else(|| source.get("sessionUpdate"))
778        .and_then(Value::as_str)
779        .unwrap_or_default()
780        .to_ascii_lowercase();
781    let completed = status.contains("complete")
782        || status.contains("result")
783        || status.contains("success")
784        || status.contains("error")
785        || status.contains("fail");
786    let arguments = source
787        .get("arguments")
788        .or_else(|| source.get("input"))
789        .or_else(|| source.get("rawInput"))
790        .cloned()
791        .unwrap_or(Value::Null);
792    json!({
793        "type": if completed { "tool_call_completed" } else { "tool_call_started" },
794        "id": source.get("toolCallId").or_else(|| source.get("tool_call_id")).or_else(|| source.get("callId")).or_else(|| source.get("id")).cloned().unwrap_or(Value::Null),
795        "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
796        "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
797        "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
798        "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
799        "raw": raw,
800    })
801}
802
803fn native_payload(kind: &str, payload: &Value) -> Value {
804    json!({"type":"native_event", "kind":kind, "raw":payload})
805}
806
807fn extract_text(value: &Value) -> Option<String> {
808    match value {
809        Value::String(text) => Some(text.clone()),
810        Value::Array(values) => {
811            let text = values
812                .iter()
813                .filter_map(extract_text)
814                .collect::<Vec<_>>()
815                .join("\n");
816            (!text.is_empty()).then_some(text)
817        }
818        Value::Object(object) => {
819            for key in ["text", "delta", "content", "message", "result", "error"] {
820                if let Some(text) = object.get(key).and_then(Value::as_str) {
821                    return Some(text.to_string());
822                }
823            }
824            for key in [
825                "delta",
826                "content",
827                "message",
828                "error",
829                "data",
830                "part",
831                "params",
832                "properties",
833                "update",
834                "event",
835            ] {
836                if let Some(text) = object.get(key).and_then(extract_text) {
837                    return Some(text);
838                }
839            }
840            None
841        }
842        _ => None,
843    }
844}
845
846#[cfg(test)]
847mod tests {
848    use super::*;
849    use crate::{HarnessId, RuntimeEndpoint};
850
851    struct ControlledRuntime {
852        handle: RuntimeHandle,
853        events: mpsc::UnboundedReceiver<HarnessEvent>,
854    }
855
856    #[async_trait]
857    impl RuntimeConnection for ControlledRuntime {
858        fn handle(&self) -> &RuntimeHandle {
859            &self.handle
860        }
861
862        async fn send_input(&mut self, _input: RuntimeInput) -> Result<Option<String>> {
863            Ok(Some("turn-1".into()))
864        }
865
866        async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
867            Ok(self.events.recv().await)
868        }
869
870        async fn interrupt(&mut self) -> Result<()> {
871            Ok(())
872        }
873
874        async fn respond(&mut self, _request_id: Value, _response: Value) -> Result<()> {
875            Ok(())
876        }
877
878        async fn close(&mut self) -> Result<()> {
879            Ok(())
880        }
881    }
882
883    fn controlled_runtime() -> (
884        Box<dyn RuntimeConnection>,
885        mpsc::UnboundedSender<HarnessEvent>,
886    ) {
887        let (events, event_rx) = mpsc::unbounded_channel();
888        (
889            Box::new(ControlledRuntime {
890                handle: RuntimeHandle {
891                    harness: HarnessId::from(HarnessId::PI),
892                    runtime_id: "shared-runtime".into(),
893                    endpoint: RuntimeEndpoint::LocalProcess {
894                        pid: None,
895                        command: vec!["controlled-runtime".into()],
896                        protocol: "test".into(),
897                    },
898                },
899                events: event_rx,
900            }),
901            events,
902        )
903    }
904
905    fn capabilities() -> RuntimeCapabilities {
906        RuntimeCapabilities {
907            start_session: true,
908            resume_session: true,
909            attach_existing_process: false,
910            send_input: true,
911            stream_events: true,
912            interrupt: true,
913            steer: false,
914            respond_to_requests: false,
915        }
916    }
917
918    #[tokio::test]
919    async fn native_eof_closes_the_raw_owner_connection() {
920        let (runtime, events) = controlled_runtime();
921        let (_host, mut connection) = HostedHarnessRuntime::spawn(runtime, capabilities());
922        drop(events);
923
924        let event =
925            tokio::time::timeout(std::time::Duration::from_secs(1), connection.next_event())
926                .await
927                .expect("raw owner should not hang after native EOF")
928                .unwrap()
929                .expect("EOF is projected as an explicit terminal event");
930        assert_eq!(event.kind, "transport_closed");
931        assert_eq!(event.payload["terminal"], true);
932    }
933
934    #[tokio::test]
935    async fn adapters_without_an_operation_route_preserve_the_requested_id() {
936        let (runtime, _events) = controlled_runtime();
937        let (host, _owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
938        let operation_id = "prompt:not-advertised".to_string();
939        let error = FrontendRuntime::invoke(
940            host.as_ref(),
941            crate::FrontendOperationInvocation::Prompt {
942                operation_id: operation_id.clone(),
943                arguments: String::new(),
944            },
945        )
946        .await
947        .unwrap_err();
948
949        assert!(
950            matches!(error, FrontendRuntimeError::UnsupportedOperation(id) if id == operation_id)
951        );
952    }
953
954    #[tokio::test]
955    async fn interrupt_stays_busy_until_the_native_terminal_event() {
956        let (runtime, events) = controlled_runtime();
957        let (host, mut owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
958        let mut terminal = FrontendRuntime::attach(host.as_ref(), 100).await.unwrap();
959
960        assert_eq!(
961            FrontendRuntime::submit(host.as_ref(), "hello".into())
962                .await
963                .unwrap(),
964            "turn-1"
965        );
966        assert_eq!(terminal.next_event().await.unwrap().kind, "user_message");
967        assert_eq!(terminal.next_event().await.unwrap().kind, "turn_started");
968        assert!(FrontendRuntime::interrupt(host.as_ref()).await.unwrap());
969        assert_eq!(
970            FrontendRuntime::describe(host.as_ref())
971                .await
972                .unwrap()
973                .turn_state,
974            FrontendTurnState::Busy
975        );
976        assert!(
977            tokio::time::timeout(std::time::Duration::from_millis(20), terminal.next_event())
978                .await
979                .is_err(),
980            "interrupt acceptance must not manufacture turn completion"
981        );
982
983        events
984            .send(HarnessEvent {
985                sequence: None,
986                kind: "agent_end".into(),
987                payload: json!({}),
988            })
989            .unwrap();
990        assert_eq!(terminal.next_event().await.unwrap().kind, "turn_succeeded");
991        assert_eq!(terminal.next_event().await.unwrap().kind, "turn_completed");
992        assert_eq!(
993            FrontendRuntime::describe(host.as_ref())
994                .await
995                .unwrap()
996                .turn_state,
997            FrontendTurnState::Idle
998        );
999        owner.close().await.unwrap();
1000    }
1001
1002    #[test]
1003    fn native_start_events_do_not_duplicate_the_hosted_turn_boundary() {
1004        let event = HarnessEvent {
1005            sequence: None,
1006            kind: "turn/started".into(),
1007            payload: json!({"method":"turn/started"}),
1008        };
1009        let projected = project_native_event(HarnessId::CODEX, &event);
1010        assert_eq!(projected.len(), 1);
1011        assert_eq!(projected[0]["type"], "native_event");
1012    }
1013}