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