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