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            // pi streams one assistant message as `text_start` → `text_delta`*
674            // → `text_end` (the same for `thinking_*`). Only the `*_delta`
675            // events carry NEW text; `text_end` repeats the whole block as
676            // `content`, and projecting it too doubled every reply.
677            let update = payload
678                .get("assistantMessageEvent")
679                .or_else(|| payload.get("event"))
680                .unwrap_or(payload);
681            let kind = update
682                .get("type")
683                .and_then(Value::as_str)
684                .unwrap_or_default();
685            match kind {
686                "text_delta" => vec![
687                    json!({"type":"text_delta", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
688                ],
689                "thinking_delta" => vec![
690                    json!({"type":"reasoning", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
691                ],
692                _ => vec![native_payload(key, payload)],
693            }
694        }
695        _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
696        _ => vec![native_payload(key, payload)],
697    }
698}
699
700fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
701    if key == "session.idle" {
702        return completion(None);
703    }
704    if key == "session.status" {
705        let status = payload
706            .pointer("/properties/status/type")
707            .or_else(|| payload.pointer("/status/type"))
708            .and_then(Value::as_str)
709            .unwrap_or_default();
710        if status == "busy" {
711            return vec![native_payload(key, payload)];
712        }
713        if status == "idle" {
714            return completion(None);
715        }
716    }
717    if key == "message.part.updated" {
718        let part = payload
719            .pointer("/properties/part")
720            .or_else(|| payload.get("part"))
721            .unwrap_or(payload);
722        if let Some(delta) = payload
723            .pointer("/properties/delta")
724            .or_else(|| payload.get("delta"))
725            .and_then(Value::as_str)
726        {
727            let reasoning = part
728                .get("type")
729                .and_then(Value::as_str)
730                .is_some_and(|kind| kind.contains("reasoning"));
731            return vec![
732                json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
733            ];
734        }
735        if part
736            .get("type")
737            .and_then(Value::as_str)
738            .is_some_and(|kind| kind.contains("tool"))
739        {
740            return vec![project_tool(part, payload)];
741        }
742    }
743    if key == "session.error" {
744        return completion(Some(
745            extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
746        ));
747    }
748    vec![native_payload(key, payload)]
749}
750
751fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
752    match key {
753        "turn_started" | "turn/started" | "agent_start" => {
754            vec![native_payload(key, payload)]
755        }
756        "turn_completed" | "turn/completed" | "agent_end" => {
757            completion(payload.get("error").and_then(extract_text))
758        }
759        "output_delta" | "content_delta" => vec![
760            json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
761        ],
762        "reasoning_delta" => vec![
763            json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
764        ],
765        "tool" => vec![project_tool(payload, payload)],
766        _ => vec![native_payload(key, payload)],
767    }
768}
769
770fn completion(error: Option<String>) -> Vec<Value> {
771    match error {
772        Some(message) => vec![
773            json!({"type":"turn_failed", "message":message}),
774            json!({"type":"turn_completed"}),
775        ],
776        None => vec![
777            json!({"type":"turn_succeeded"}),
778            json!({"type":"turn_completed"}),
779        ],
780    }
781}
782
783fn project_tool(source: &Value, raw: &Value) -> Value {
784    let status = source
785        .get("status")
786        .or_else(|| source.get("state"))
787        .or_else(|| source.get("sessionUpdate"))
788        .and_then(Value::as_str)
789        .unwrap_or_default()
790        .to_ascii_lowercase();
791    let completed = status.contains("complete")
792        || status.contains("result")
793        || status.contains("success")
794        || status.contains("error")
795        || status.contains("fail");
796    let arguments = source
797        .get("arguments")
798        .or_else(|| source.get("input"))
799        .or_else(|| source.get("rawInput"))
800        .cloned()
801        .unwrap_or(Value::Null);
802    json!({
803        "type": if completed { "tool_call_completed" } else { "tool_call_started" },
804        "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),
805        "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
806        "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
807        "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
808        "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
809        "raw": raw,
810    })
811}
812
813fn native_payload(kind: &str, payload: &Value) -> Value {
814    json!({"type":"native_event", "kind":kind, "raw":payload})
815}
816
817fn extract_text(value: &Value) -> Option<String> {
818    match value {
819        Value::String(text) => Some(text.clone()),
820        Value::Array(values) => {
821            let text = values
822                .iter()
823                .filter_map(extract_text)
824                .collect::<Vec<_>>()
825                .join("\n");
826            (!text.is_empty()).then_some(text)
827        }
828        Value::Object(object) => {
829            for key in ["text", "delta", "content", "message", "result", "error"] {
830                if let Some(text) = object.get(key).and_then(Value::as_str) {
831                    return Some(text.to_string());
832                }
833            }
834            for key in [
835                "delta",
836                "content",
837                "message",
838                "error",
839                "data",
840                "part",
841                "params",
842                "properties",
843                "update",
844                "event",
845            ] {
846                if let Some(text) = object.get(key).and_then(extract_text) {
847                    return Some(text);
848                }
849            }
850            None
851        }
852        _ => None,
853    }
854}
855
856#[cfg(test)]
857mod tests {
858    #[test]
859    fn pi_message_update_projects_only_the_deltas() {
860        use serde_json::json;
861        let delta = json!({"type":"message_update","assistantMessageEvent":{"type":"text_delta","contentIndex":0,"delta":"ack: hi"}});
862        let end = json!({"type":"message_update","assistantMessageEvent":{"type":"text_end","contentIndex":0,"content":"ack: hi"}});
863        let start = json!({"type":"message_update","assistantMessageEvent":{"type":"text_start","contentIndex":0}});
864        let projected = super::project_pi("message_update", &delta);
865        assert_eq!(projected[0]["type"], "text_delta");
866        assert_eq!(projected[0]["text"], "ack: hi");
867        for repeat in [&end, &start] {
868            let projected = super::project_pi("message_update", repeat);
869            assert_eq!(
870                projected[0]["type"], "native_event",
871                "{repeat} carries no new text"
872            );
873        }
874        let thinking = json!({"type":"message_update","assistantMessageEvent":{"type":"thinking_delta","contentIndex":0,"delta":"hm"}});
875        let projected = super::project_pi("message_update", &thinking);
876        assert_eq!(projected[0]["type"], "reasoning");
877        assert_eq!(projected[0]["text"], "hm");
878    }
879
880    use super::*;
881    use crate::{HarnessId, RuntimeEndpoint};
882
883    struct ControlledRuntime {
884        handle: RuntimeHandle,
885        events: mpsc::UnboundedReceiver<HarnessEvent>,
886    }
887
888    #[async_trait]
889    impl RuntimeConnection for ControlledRuntime {
890        fn handle(&self) -> &RuntimeHandle {
891            &self.handle
892        }
893
894        async fn send_input(&mut self, _input: RuntimeInput) -> Result<Option<String>> {
895            Ok(Some("turn-1".into()))
896        }
897
898        async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
899            Ok(self.events.recv().await)
900        }
901
902        async fn interrupt(&mut self) -> Result<()> {
903            Ok(())
904        }
905
906        async fn respond(&mut self, _request_id: Value, _response: Value) -> Result<()> {
907            Ok(())
908        }
909
910        async fn close(&mut self) -> Result<()> {
911            Ok(())
912        }
913    }
914
915    fn controlled_runtime() -> (
916        Box<dyn RuntimeConnection>,
917        mpsc::UnboundedSender<HarnessEvent>,
918    ) {
919        let (events, event_rx) = mpsc::unbounded_channel();
920        (
921            Box::new(ControlledRuntime {
922                handle: RuntimeHandle {
923                    harness: HarnessId::from(HarnessId::PI),
924                    runtime_id: "shared-runtime".into(),
925                    endpoint: RuntimeEndpoint::LocalProcess {
926                        pid: None,
927                        command: vec!["controlled-runtime".into()],
928                        protocol: "test".into(),
929                    },
930                },
931                events: event_rx,
932            }),
933            events,
934        )
935    }
936
937    fn capabilities() -> RuntimeCapabilities {
938        RuntimeCapabilities {
939            start_session: true,
940            resume_session: true,
941            attach_existing_process: false,
942            send_input: true,
943            stream_events: true,
944            interrupt: true,
945            steer: false,
946            respond_to_requests: false,
947        }
948    }
949
950    #[tokio::test]
951    async fn native_eof_closes_the_raw_owner_connection() {
952        let (runtime, events) = controlled_runtime();
953        let (_host, mut connection) = HostedHarnessRuntime::spawn(runtime, capabilities());
954        drop(events);
955
956        let event =
957            tokio::time::timeout(std::time::Duration::from_secs(1), connection.next_event())
958                .await
959                .expect("raw owner should not hang after native EOF")
960                .unwrap()
961                .expect("EOF is projected as an explicit terminal event");
962        assert_eq!(event.kind, "transport_closed");
963        assert_eq!(event.payload["terminal"], true);
964    }
965
966    #[tokio::test]
967    async fn adapters_without_an_operation_route_preserve_the_requested_id() {
968        let (runtime, _events) = controlled_runtime();
969        let (host, _owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
970        let operation_id = "prompt:not-advertised".to_string();
971        let error = FrontendRuntime::invoke(
972            host.as_ref(),
973            crate::FrontendOperationInvocation::Prompt {
974                operation_id: operation_id.clone(),
975                arguments: String::new(),
976            },
977        )
978        .await
979        .unwrap_err();
980
981        assert!(
982            matches!(error, FrontendRuntimeError::UnsupportedOperation(id) if id == operation_id)
983        );
984    }
985
986    #[tokio::test]
987    async fn interrupt_stays_busy_until_the_native_terminal_event() {
988        let (runtime, events) = controlled_runtime();
989        let (host, mut owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
990        let mut terminal = FrontendRuntime::attach(host.as_ref(), 100).await.unwrap();
991
992        assert_eq!(
993            FrontendRuntime::submit(host.as_ref(), "hello".into())
994                .await
995                .unwrap(),
996            "turn-1"
997        );
998        assert_eq!(terminal.next_event().await.unwrap().kind, "user_message");
999        assert_eq!(terminal.next_event().await.unwrap().kind, "turn_started");
1000        assert!(FrontendRuntime::interrupt(host.as_ref()).await.unwrap());
1001        assert_eq!(
1002            FrontendRuntime::describe(host.as_ref())
1003                .await
1004                .unwrap()
1005                .turn_state,
1006            FrontendTurnState::Busy
1007        );
1008        assert!(
1009            tokio::time::timeout(std::time::Duration::from_millis(20), terminal.next_event())
1010                .await
1011                .is_err(),
1012            "interrupt acceptance must not manufacture turn completion"
1013        );
1014
1015        events
1016            .send(HarnessEvent {
1017                sequence: None,
1018                kind: "agent_end".into(),
1019                payload: json!({}),
1020            })
1021            .unwrap();
1022        assert_eq!(terminal.next_event().await.unwrap().kind, "turn_succeeded");
1023        assert_eq!(terminal.next_event().await.unwrap().kind, "turn_completed");
1024        assert_eq!(
1025            FrontendRuntime::describe(host.as_ref())
1026                .await
1027                .unwrap()
1028                .turn_state,
1029            FrontendTurnState::Idle
1030        );
1031        owner.close().await.unwrap();
1032    }
1033
1034    #[test]
1035    fn native_start_events_do_not_duplicate_the_hosted_turn_boundary() {
1036        let event = HarnessEvent {
1037            sequence: None,
1038            kind: "turn/started".into(),
1039            payload: json!({"method":"turn/started"}),
1040        };
1041        let projected = project_native_event(HarnessId::CODEX, &event);
1042        assert_eq!(projected.len(), 1);
1043        assert_eq!(projected[0]["type"], "native_event");
1044    }
1045}