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::{HashMap, HashSet, 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::{
24    ChatMessage, DiscoveryQuery, Error, Fidelity, FrontendResponse, HarnessCatalog, Result,
25};
26
27enum HostCommand {
28    Submit {
29        input: RuntimeInput,
30        reply: oneshot::Sender<Result<Option<String>>>,
31    },
32    Interrupt {
33        reply: oneshot::Sender<Result<()>>,
34    },
35    Steer {
36        text: String,
37        reply: oneshot::Sender<Result<()>>,
38    },
39    Respond {
40        request_id: Value,
41        response: Value,
42        reply: oneshot::Sender<Result<()>>,
43    },
44    Shutdown {
45        reply: oneshot::Sender<Result<()>>,
46    },
47}
48
49struct ProjectionState {
50    next_sequence: u64,
51    replay: VecDeque<FrontendEvent>,
52}
53
54/// Shared owner of a single harness-native runtime.
55pub struct HostedHarnessRuntime {
56    handle: RuntimeHandle,
57    capabilities: RuntimeCapabilities,
58    commands: mpsc::Sender<HostCommand>,
59    raw_events: broadcast::Sender<HarnessEvent>,
60    frontend_events: broadcast::Sender<FrontendEvent>,
61    projection: StdMutex<ProjectionState>,
62    /// The conversation the harness had recorded before this host took it
63    /// over; every frontend event since is in `projection.replay`.
64    history: Vec<ChatMessage>,
65    /// Requests the translator has raised and nobody has answered yet, keyed
66    /// by the canonical id a frontend sees. BOTH doors — `frontend.v2.respond`
67    /// and the owner's `harness.v1.runtimes.respond` — funnel through
68    /// `respond_native`, which is where an entry is dropped and
69    /// `request_resolved` is published, so one request takes exactly one
70    /// answer and the door that did not answer reads the resolution.
71    pending_requests: StdMutex<HashMap<u64, Value>>,
72    /// What the projection must remember across a session's events (`NativeProjection`).
73    native: StdMutex<NativeProjection>,
74    busy: AtomicBool,
75    closed: AtomicBool,
76    registration: StdMutex<Option<crate::LiveRuntimeRegistration>>,
77}
78
79impl HostedHarnessRuntime {
80    /// Promote one exclusive native connection into a multi-frontend host and
81    /// return the first SDK connection to it.
82    pub fn spawn(
83        runtime: Box<dyn RuntimeConnection>,
84        capabilities: RuntimeCapabilities,
85        fresh: bool,
86    ) -> (Arc<Self>, HostedHarnessConnection) {
87        let handle = runtime.handle().clone();
88        let (commands, command_rx) = mpsc::channel(32);
89        let (raw_events, raw_rx) = broadcast::channel(1024);
90        let (frontend_events, _) = broadcast::channel(1024);
91        let host = Arc::new(Self {
92            handle: handle.clone(),
93            capabilities,
94            commands,
95            raw_events,
96            frontend_events,
97            projection: StdMutex::new(ProjectionState {
98                next_sequence: 1,
99                replay: VecDeque::new(),
100            }),
101            history: if fresh {
102                Vec::new()
103            } else {
104                recorded_history(&handle)
105            },
106            pending_requests: StdMutex::new(HashMap::new()),
107            native: StdMutex::new(NativeProjection::default()),
108            busy: AtomicBool::new(false),
109            closed: AtomicBool::new(false),
110            registration: StdMutex::new(None),
111        });
112        tokio::spawn(run_native_runtime(
113            runtime,
114            Arc::downgrade(&host),
115            command_rx,
116        ));
117        let connection = HostedHarnessConnection {
118            host: host.clone(),
119            handle,
120            events: raw_rx,
121            closed: false,
122        };
123        (host, connection)
124    }
125
126    /// Keep the receipt only while the native host remains open.
127    pub(crate) fn retain_registration(&self, registration: crate::LiveRuntimeRegistration) {
128        let mut retained = self
129            .registration
130            .lock()
131            .unwrap_or_else(std::sync::PoisonError::into_inner);
132        if !self.closed.load(Ordering::SeqCst) {
133            *retained = Some(registration);
134        }
135    }
136
137    /// Subscribe to the canonical sequenced frontend event stream.
138    pub fn frontend_sender(&self) -> broadcast::Sender<FrontendEvent> {
139        self.frontend_events.clone()
140    }
141
142    /// Shut down the one native process owned by this host.
143    pub async fn shutdown(&self) -> Result<()> {
144        if self.closed.load(Ordering::SeqCst) {
145            return Ok(());
146        }
147        let (reply, response) = oneshot::channel();
148        self.commands
149            .send(HostCommand::Shutdown { reply })
150            .await
151            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
152        response
153            .await
154            .map_err(|_| Error::Other("hosted harness runtime stopped before shutdown".into()))?
155    }
156
157    fn claim_submit(&self) -> Result<()> {
158        if self.closed.load(Ordering::SeqCst) {
159            return Err(Error::Other("hosted harness runtime is closed".into()));
160        }
161        if self
162            .busy
163            .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
164            .is_err()
165        {
166            return Err(Error::Other("a harness turn is already in progress".into()));
167        }
168        Ok(())
169    }
170
171    async fn submit_native_claimed(&self, input: RuntimeInput) -> Result<Option<String>> {
172        self.publish(json!({"type":"user_message", "text":input.text}));
173        self.publish(json!({"type":"turn_started"}));
174        let (reply, response) = oneshot::channel();
175        if self
176            .commands
177            .send(HostCommand::Submit { input, reply })
178            .await
179            .is_err()
180        {
181            self.busy.store(false, Ordering::SeqCst);
182            self.publish(
183                json!({"type":"turn_failed", "message":"Hosted harness runtime is closed."}),
184            );
185            self.publish(json!({"type":"turn_completed"}));
186            self.mark_closed("Harness runtime command channel closed.");
187            return Err(Error::Other("hosted harness runtime is closed".into()));
188        }
189        match response.await {
190            Ok(Ok(turn)) => Ok(turn),
191            Ok(Err(error)) => {
192                self.busy.store(false, Ordering::SeqCst);
193                self.publish(json!({"type":"turn_failed", "message":error.to_string()}));
194                self.publish(json!({"type":"turn_completed"}));
195                Err(error)
196            }
197            Err(_) => {
198                self.busy.store(false, Ordering::SeqCst);
199                self.publish(json!({"type":"turn_failed", "message":"Hosted harness runtime stopped before accepting input."}));
200                self.publish(json!({"type":"turn_completed"}));
201                self.mark_closed("Harness runtime stopped before accepting input.");
202                Err(Error::Other(
203                    "hosted harness runtime stopped before accepting input".into(),
204                ))
205            }
206        }
207    }
208
209    async fn submit_native(&self, text: String) -> Result<Option<String>> {
210        self.claim_submit()?;
211        self.submit_native_claimed(RuntimeInput {
212            text,
213            image_urls: Vec::new(),
214        })
215        .await
216    }
217
218    async fn interrupt_native(&self) -> Result<()> {
219        let (reply, response) = oneshot::channel();
220        self.commands
221            .send(HostCommand::Interrupt { reply })
222            .await
223            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
224        response
225            .await
226            .map_err(|_| Error::Other("hosted harness runtime stopped before interrupt".into()))?
227    }
228
229    async fn steer_native(&self, text: String) -> Result<()> {
230        let (reply, response) = oneshot::channel();
231        self.commands
232            .send(HostCommand::Steer { text, reply })
233            .await
234            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
235        response
236            .await
237            .map_err(|_| Error::Other("hosted harness runtime stopped before steering".into()))?
238    }
239
240    async fn respond_native(&self, request_id: Value, response: Value) -> Result<()> {
241        self.respond_native_as(request_id, response, None).await
242    }
243
244    /// Answer one native request. `canonical` is the portable response when a
245    /// frontend answered; the owner's door passes `None` and the published
246    /// resolution carries the native body it actually sent.
247    async fn respond_native_as(
248        &self,
249        request_id: Value,
250        response: Value,
251        canonical: Option<Value>,
252    ) -> Result<()> {
253        let (reply, completed) = oneshot::channel();
254        self.commands
255            .send(HostCommand::Respond {
256                request_id: request_id.clone(),
257                response: response.clone(),
258                reply,
259            })
260            .await
261            .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
262        completed
263            .await
264            .map_err(|_| Error::Other("hosted harness runtime stopped before response".into()))??;
265        // Only a request the native runtime ACCEPTED is resolved, and only
266        // once: the first door through here removes it, so a second answer
267        // finds nothing to publish.
268        let resolved = {
269            let mut pending = self
270                .pending_requests
271                .lock()
272                .unwrap_or_else(std::sync::PoisonError::into_inner);
273            let found = pending
274                .iter()
275                .find(|(_, native)| *native == &request_id)
276                .map(|(id, _)| *id);
277            found.and_then(|id| pending.remove(&id).map(|_| id))
278        };
279        if let Some(id) = resolved {
280            self.publish(json!({
281                "type": "request_resolved",
282                "request_id": id,
283                "response": canonical.unwrap_or(response),
284            }));
285        }
286        Ok(())
287    }
288
289    /// The native id a canonical request maps to, while it is still open.
290    fn native_request_id(&self, request_id: u64) -> Option<Value> {
291        self.pending_requests
292            .lock()
293            .unwrap_or_else(std::sync::PoisonError::into_inner)
294            .get(&request_id)
295            .cloned()
296    }
297
298    fn publish(&self, payload: Value) {
299        let event = {
300            let mut projection = self
301                .projection
302                .lock()
303                .unwrap_or_else(std::sync::PoisonError::into_inner);
304            let event = FrontendEvent::new(projection.next_sequence, payload);
305            projection.next_sequence = projection.next_sequence.saturating_add(1);
306            projection.replay.push_back(event.clone());
307            while projection.replay.len() > FRONTEND_REPLAY_CAPACITY {
308                projection.replay.pop_front();
309            }
310            event
311        };
312        let _ = self.frontend_events.send(event);
313    }
314
315    fn accept_native_event(&self, event: HarnessEvent) {
316        let _ = self.raw_events.send(event.clone());
317        let projected = self
318            .native
319            .lock()
320            .unwrap_or_else(std::sync::PoisonError::into_inner)
321            .project(self.handle.harness.as_str(), &event);
322        for payload in projected {
323            let terminal = matches!(
324                payload.get("type").and_then(Value::as_str),
325                Some("turn_succeeded" | "turn_interrupted" | "turn_failed")
326            );
327            if terminal {
328                self.busy.store(false, Ordering::SeqCst);
329            }
330            if payload.get("type").and_then(Value::as_str) == Some("request") {
331                if let (Some(id), Some(native)) = (
332                    payload.pointer("/request/id").and_then(Value::as_u64),
333                    payload
334                        .pointer("/request/payload/native_request_id")
335                        .cloned(),
336                ) {
337                    self.pending_requests
338                        .lock()
339                        .unwrap_or_else(std::sync::PoisonError::into_inner)
340                        .insert(id, native);
341                }
342            }
343            self.publish(payload);
344        }
345    }
346
347    fn mark_closed(&self, message: impl Into<String>) {
348        if self.closed.swap(true, Ordering::SeqCst) {
349            return;
350        }
351        self.registration
352            .lock()
353            .unwrap_or_else(std::sync::PoisonError::into_inner)
354            .take();
355        let message = message.into();
356        self.busy.store(false, Ordering::SeqCst);
357        // The raw SDK connection and the portable frontend must observe the
358        // same terminal edge. Keeping the sender alive inside `self` otherwise
359        // leaves the SDK receiver waiting forever after native EOF.
360        let _ = self.raw_events.send(HarnessEvent {
361            sequence: None,
362            kind: "transport_closed".into(),
363            payload: json!({"message":message, "terminal":true}),
364        });
365        self.publish(json!({"type":"runtime_disconnected", "message":message}));
366    }
367
368    fn descriptor(&self) -> FrontendRuntimeDescriptor {
369        FrontendRuntimeDescriptor {
370            schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
371            session_id: self.handle.runtime_id.clone(),
372            source_harness: Some(self.handle.harness.as_str().to_string()),
373            emulation_profile: None,
374            active_modules: Vec::new(),
375            commands: Vec::new(),
376            operations: Vec::new(),
377            actions: FrontendActions {
378                submit: self.capabilities.send_input,
379                interrupt: self.capabilities.interrupt,
380                steer: self.capabilities.steer,
381                // ANSWERABLE where this harness has a reply translation, and
382                // only there. `project_native_event` turns a harness's native
383                // events into the canonical stream;
384                // `hosted_permission_reply` is its mirror, turning a canonical
385                // decision back into that harness's own reply envelope — so a
386                // portable frontend can answer a request without ever seeing
387                // the native shape. Today that is Claude Code, whose
388                // permission-prompt-tool protocol takes
389                // `{behavior:'allow'}` / `{behavior:'deny', message}`; a
390                // harness with no translation keeps `false` and its requests
391                // stay observe-only, which is what every other hosted harness
392                // still is. The request event states which decisions its own
393                // protocol accepts, so a frontend offers exactly those.
394                //
395                // The owner's `harness.v1.approvals.resolve` is unchanged and
396                // still authoritative: both doors funnel through
397                // `respond_native_as`, which drops the pending entry and
398                // publishes `request_resolved` once, so one request takes one
399                // answer and whichever door answered first wins.
400                respond: !hosted_answerable_decisions(self.handle.harness.as_str()).is_empty(),
401                detach: true,
402                close: false,
403            },
404            display: FrontendDisplayCapabilities {
405                event_kinds: vec![
406                    "user_message".into(),
407                    "turn_started".into(),
408                    "turn_succeeded".into(),
409                    "turn_interrupted".into(),
410                    "turn_failed".into(),
411                    "text_delta".into(),
412                    "reasoning".into(),
413                    "tool_call_started".into(),
414                    "tool_call_completed".into(),
415                    // Observation only — `actions.respond` stays false, so a
416                    // frontend shows the request and no answer control.
417                    "request".into(),
418                    "native_event".into(),
419                    "runtime_disconnected".into(),
420                ],
421                opaque_fallback: true,
422            },
423            model: self.handle.harness.as_str().to_string(),
424            turn_state: if self.busy.load(Ordering::SeqCst) {
425                FrontendTurnState::Busy
426            } else {
427                FrontendTurnState::Idle
428            },
429            connection_state: if self.closed.load(Ordering::SeqCst) {
430                FrontendConnectionState::ShuttingDown
431            } else {
432                FrontendConnectionState::Connected
433            },
434            extensions: Default::default(),
435        }
436    }
437}
438
439#[async_trait]
440impl FrontendRuntime for HostedHarnessRuntime {
441    async fn describe(
442        &self,
443    ) -> std::result::Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
444        Ok(self.descriptor())
445    }
446
447    async fn attach(
448        &self,
449        _history_limit: usize,
450    ) -> std::result::Result<FrontendAttachment, FrontendRuntimeError> {
451        let live = self.frontend_events.subscribe();
452        let projection = self
453            .projection
454            .lock()
455            .unwrap_or_else(std::sync::PoisonError::into_inner);
456        let replay = projection.replay.clone();
457        if let Some(first) = replay.front() {
458            if first.sequence > 1 {
459                return Err(FrontendRuntimeError::ReplayGap(first.sequence - 1));
460            }
461        }
462        Ok(FrontendAttachment::new(
463            self.descriptor(),
464            self.history.clone(),
465            0,
466            replay,
467            live,
468            None,
469        ))
470    }
471
472    async fn send_input(
473        self: Arc<Self>,
474        prompt: String,
475    ) -> std::result::Result<(), FrontendRuntimeError> {
476        self.claim_submit().map_err(hosted_submit_error)?;
477        tokio::spawn(async move {
478            let _ = self
479                .submit_native_claimed(RuntimeInput {
480                    text: prompt,
481                    image_urls: Vec::new(),
482                })
483                .await;
484        });
485        Ok(())
486    }
487
488    async fn send_input_with_images(
489        self: Arc<Self>,
490        prompt: String,
491        image_urls: Vec<String>,
492    ) -> std::result::Result<(), FrontendRuntimeError> {
493        self.claim_submit().map_err(hosted_submit_error)?;
494        tokio::spawn(async move {
495            let _ = self
496                .submit_native_claimed(RuntimeInput {
497                    text: prompt,
498                    image_urls,
499                })
500                .await;
501        });
502        Ok(())
503    }
504
505    async fn submit(&self, prompt: String) -> std::result::Result<String, FrontendRuntimeError> {
506        self.submit_native(prompt)
507            .await
508            .map(|turn| turn.unwrap_or_default())
509            .map_err(hosted_submit_error)
510    }
511
512    async fn interrupt(&self) -> std::result::Result<bool, FrontendRuntimeError> {
513        if !self.busy.load(Ordering::SeqCst) {
514            return Ok(false);
515        }
516        self.interrupt_native()
517            .await
518            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
519        // Acceptance of the interrupt request does not mean the native turn
520        // has stopped. Keep the shared runtime busy until its terminal event,
521        // so another frontend cannot race a prompt into turn teardown.
522        Ok(true)
523    }
524
525    async fn steer(&self, prompt: String) -> std::result::Result<(), FrontendRuntimeError> {
526        if !self.busy.load(Ordering::SeqCst) || !self.capabilities.steer {
527            return Err(FrontendRuntimeError::UnsupportedAction("steer"));
528        }
529        self.steer_native(prompt)
530            .await
531            .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
532    }
533
534    async fn respond(
535        &self,
536        response: FrontendResponse,
537    ) -> std::result::Result<(), FrontendRuntimeError> {
538        let harness = self.handle.harness.as_str();
539        let FrontendResponse::Approval {
540            request_id,
541            decision,
542        } = &response
543        else {
544            return Err(FrontendRuntimeError::UnsupportedAction(
545                "respond: a hosted harness answers approval requests only",
546            ));
547        };
548        let Some(native_id) = self.native_request_id(*request_id) else {
549            return Err(FrontendRuntimeError::UnknownRequest(*request_id));
550        };
551        let body = hosted_permission_reply(harness, *decision)
552            .map_err(FrontendRuntimeError::InvalidResponse)?;
553        let canonical = serde_json::to_value(&response).unwrap_or(Value::Null);
554        self.respond_native_as(native_id, body, Some(canonical))
555            .await
556            .map_err(|error| FrontendRuntimeError::Execution {
557                operation: crate::SdkOperation::Respond,
558                message: error.to_string(),
559            })
560    }
561}
562
563fn hosted_submit_error(error: Error) -> FrontendRuntimeError {
564    if error.to_string().contains("already in progress") {
565        FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)
566    } else {
567        FrontendRuntimeError::Transport(error.to_string())
568    }
569}
570
571/// SDK connection to a hosted native runtime. Closing it shuts down the
572/// owner runtime; terminal HTTP attachments are non-owning frontend leases.
573pub struct HostedHarnessConnection {
574    host: Arc<HostedHarnessRuntime>,
575    handle: RuntimeHandle,
576    events: broadcast::Receiver<HarnessEvent>,
577    closed: bool,
578}
579
580#[async_trait]
581impl RuntimeConnection for HostedHarnessConnection {
582    fn handle(&self) -> &RuntimeHandle {
583        &self.handle
584    }
585
586    async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
587        self.host.claim_submit()?;
588        self.host.submit_native_claimed(input).await
589    }
590
591    async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
592        match self.events.recv().await {
593            Ok(event) => Ok(Some(event)),
594            Err(broadcast::error::RecvError::Lagged(count)) => Err(Error::Other(format!(
595                "hosted harness event stream lost {count} event(s)"
596            ))),
597            Err(broadcast::error::RecvError::Closed) => Ok(None),
598        }
599    }
600
601    async fn interrupt(&mut self) -> Result<()> {
602        self.host.interrupt_native().await
603    }
604
605    async fn steer(&mut self, text: String) -> Result<()> {
606        self.host.steer_native(text).await
607    }
608
609    async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
610        self.host.respond_native(request_id, response).await
611    }
612
613    async fn close(&mut self) -> Result<()> {
614        if self.closed {
615            return Ok(());
616        }
617        self.host.shutdown().await?;
618        self.closed = true;
619        Ok(())
620    }
621}
622
623async fn run_native_runtime(
624    mut runtime: Box<dyn RuntimeConnection>,
625    host: std::sync::Weak<HostedHarnessRuntime>,
626    mut commands: mpsc::Receiver<HostCommand>,
627) {
628    loop {
629        tokio::select! {
630            command = commands.recv() => {
631                let Some(command) = command else {
632                    let _ = runtime.close().await;
633                    return;
634                };
635                match command {
636                    HostCommand::Submit { input, reply } => {
637                        let _ = reply.send(runtime.send_input(input).await);
638                    }
639                    HostCommand::Interrupt { reply } => {
640                        let _ = reply.send(runtime.interrupt().await);
641                    }
642                    HostCommand::Steer { text, reply } => {
643                        let _ = reply.send(runtime.steer(text).await);
644                    }
645                    HostCommand::Respond { request_id, response, reply } => {
646                        let _ = reply.send(runtime.respond(request_id, response).await);
647                    }
648                    HostCommand::Shutdown { reply } => {
649                        let result = runtime.close().await;
650                        let closed = result.is_ok();
651                        let _ = reply.send(result);
652                        if closed {
653                            if let Some(host) = host.upgrade() {
654                                host.mark_closed("Harness runtime closed.");
655                            }
656                            return;
657                        }
658                    }
659                }
660            }
661            event = runtime.next_event() => {
662                let Some(host) = host.upgrade() else {
663                    let _ = runtime.close().await;
664                    return;
665                };
666                match event {
667                    Ok(Some(event)) if event.kind == "transport_closed" => {
668                        host.mark_closed("Harness runtime transport closed.");
669                        return;
670                    }
671                    Ok(Some(event)) => host.accept_native_event(event),
672                    Ok(None) => {
673                        host.mark_closed("Harness runtime transport closed.");
674                        return;
675                    }
676                    Err(error) => {
677                        host.mark_closed(error.to_string());
678                        return;
679                    }
680                }
681            }
682        }
683    }
684}
685
686/// Which harnesses can answer a `request` through the PORTABLE door, and the
687/// decisions each one's own protocol accepts.
688///
689/// This is the mirror of `project_native_event`: that side translates a
690/// harness's native events into the canonical stream, this side translates a
691/// canonical decision back into the harness's own reply envelope. A harness
692/// with no entry here keeps `respond: false` and its request stays
693/// observe-only.
694pub(crate) fn hosted_answerable_decisions(harness: &str) -> &'static [&'static str] {
695    match harness {
696        // Claude Code's permission-prompt-tool protocol, as its own validator
697        // states it: `{behavior:'allow', updatedInput?:object}` or
698        // `{behavior:'deny', message:string}`
699        // (docs/interop/research/orc2-claude-respond-receipt-2026-09-04.json).
700        // There is no always-allow BEHAVIOR — the CLI carries session
701        // permission through a separate `updatedPermissions` field — so
702        // `allow_for_session` is not offered rather than quietly downgraded to
703        // a once-allow the person did not choose. `crates/harness/src/approvals.rs`
704        // refuses it by name on the owner's door for the same reason.
705        "claude-code" => &["allow", "deny"],
706        _ => &[],
707    }
708}
709
710/// Translate one canonical approval decision into the harness's own envelope.
711fn hosted_permission_reply(
712    harness: &str,
713    decision: crate::FrontendApprovalDecision,
714) -> std::result::Result<Value, String> {
715    use crate::FrontendApprovalDecision as Decision;
716    let accepted = hosted_answerable_decisions(harness);
717    let wire = match decision {
718        Decision::Allow => "allow",
719        Decision::AllowForSession => "allow_for_session",
720        Decision::Deny => "deny",
721    };
722    if !accepted.contains(&wire) {
723        return Err(format!(
724            "`{wire}` is not a decision a hosted {harness} runtime can carry; this harness accepts \
725             {}",
726            if accepted.is_empty() {
727                "no portable decision — its requests are observe-only".to_string()
728            } else {
729                accepted
730                    .iter()
731                    .map(|decision| format!("`{decision}`"))
732                    .collect::<Vec<_>>()
733                    .join(" or ")
734            },
735        ));
736    }
737    match harness {
738        "claude-code" => Ok(match decision {
739            Decision::Allow => json!({"behavior": "allow"}),
740            // Measured: the CLI's validator rejects a `deny` with no message.
741            Decision::Deny => json!({
742                "behavior": "deny",
743                "message": "denied through the Volter Harness frontend",
744            }),
745            Decision::AllowForSession => unreachable!("refused above"),
746        }),
747        _ => Err(format!(
748            "no hosted permission reply is defined for {harness}"
749        )),
750    }
751}
752
753/// Read once, when the host takes a resumed session over, so an opener sees
754/// the conversation before this resume as a pane's scrollback shows it. The
755/// replay begins at the host's first event, so the two never overlap. A
756/// fresh session has no transcript yet and is never looked for: the search is
757/// a scan of every session on the machine (7-11 s on a loaded Mac, and past
758/// the SDK's 30 s request bound at worse moments, so `runtimes.start` failed).
759fn recorded_history(handle: &RuntimeHandle) -> Vec<ChatMessage> {
760    let catalog = HarnessCatalog::new();
761    let Some(locator) = claude_transcript(handle).or_else(|| {
762        let query = DiscoveryQuery {
763            harnesses: vec![handle.harness.clone()],
764            query: Some(handle.runtime_id.clone()),
765            ..DiscoveryQuery::default()
766        };
767        catalog.discover(&query).ok().and_then(|found| {
768            found
769                .into_iter()
770                .find(|descriptor| descriptor.locator.session_id == handle.runtime_id)
771                .map(|descriptor| descriptor.locator)
772        })
773    }) else {
774        return Vec::new();
775    };
776    catalog
777        .load_display_view(&locator, Fidelity::Semantic, 500)
778        .map(|session| session.messages)
779        .unwrap_or_default()
780}
781
782/// A Claude Code session's transcript by its file name
783/// (`projects/<folder>/<session id>.jsonl`): one look in each project folder
784/// instead of reading every session's header.
785fn claude_transcript(handle: &RuntimeHandle) -> Option<crate::SessionLocator> {
786    if handle.harness.as_str() != crate::HarnessId::CLAUDE_CODE {
787        return None;
788    }
789    let file = format!("{}.jsonl", handle.runtime_id);
790    let projects = crate::HarnessHomes::default().claude_code;
791    std::fs::read_dir(projects)
792        .ok()?
793        .flatten()
794        .map(|folder| folder.path().join(&file))
795        .find(|path| path.is_file())
796        .map(|path| crate::SessionLocator {
797            harness: handle.harness.clone(),
798            session_id: handle.runtime_id.clone(),
799            storage: crate::StorageLocator::File { path },
800        })
801}
802
803/// The projection of one session's events, with what it must remember between them: OpenCode streams reasoning and
804/// text alike as `message.part.delta` with `field: "text"` (its processor, `reasoning-delta` and `text-delta`), and only
805/// the part an earlier `message.part.updated` declared tells them apart; and a turn OpenCode failed (`session.error`)
806/// still goes idle after, which closes it without saying it succeeded.
807#[derive(Debug, Default)]
808pub struct NativeProjection {
809    opencode_reasoning_parts: HashSet<String>,
810    opencode_failed: bool,
811}
812
813impl NativeProjection {
814    /// `event` of `harness` as `project_native_event` projects it, with the session's memory applied.
815    pub fn project(&mut self, harness: &str, event: &HarnessEvent) -> Vec<Value> {
816        if harness == "opencode" {
817            let key = event.kind.to_ascii_lowercase();
818            let payload = &event.payload;
819            match key.as_str() {
820                "message.part.updated" => {
821                    let part = payload
822                        .pointer("/properties/part")
823                        .or_else(|| payload.get("part"))
824                        .unwrap_or(payload);
825                    if let (Some(id), Some(kind)) = (
826                        part.get("id").and_then(Value::as_str),
827                        part.get("type").and_then(Value::as_str),
828                    ) {
829                        if kind.contains("reasoning") {
830                            self.opencode_reasoning_parts.insert(id.to_string());
831                        }
832                    }
833                }
834                "message.part.delta" => {
835                    let reasoning = payload
836                        .pointer("/properties/partID")
837                        .and_then(Value::as_str)
838                        .is_some_and(|id| self.opencode_reasoning_parts.contains(id));
839                    if reasoning {
840                        let text = payload
841                            .pointer("/properties/delta")
842                            .and_then(Value::as_str)
843                            .unwrap_or_default();
844                        return vec![json!({"type":"reasoning", "text":text, "raw":payload})];
845                    }
846                }
847                "session.error" => {
848                    self.opencode_failed = true;
849                    return vec![
850                        json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into())}),
851                    ];
852                }
853                "session.idle" if self.opencode_failed => {
854                    self.opencode_failed = false;
855                    self.opencode_reasoning_parts.clear();
856                    return vec![json!({"type":"turn_completed"})];
857                }
858                "session.idle" => self.opencode_reasoning_parts.clear(),
859                // a new turn starts clean: a failure no idle followed does not carry into it
860                "session.status"
861                    if payload
862                        .pointer("/properties/status/type")
863                        .or_else(|| payload.pointer("/status/type"))
864                        .and_then(Value::as_str)
865                        == Some("busy") =>
866                {
867                    self.opencode_failed = false;
868                }
869                _ => {}
870            }
871        }
872        project_native_event(harness, event)
873    }
874}
875
876/// One native event of `harness` as the protocol-neutral events every reader renders (`text_delta`, `reasoning`,
877/// `turn_succeeded`/`turn_failed`, `turn_completed`, `runtime_disconnected`, tools, approvals): ACP's and each
878/// harness's own (Codex, Claude Code, Pi, OpenCode, generic) alike.
879fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
880    let key = event.kind.to_ascii_lowercase().replace('-', "_");
881    let payload = &event.payload;
882    if key == "transport_closed" {
883        return vec![
884            json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
885        ];
886    }
887    let retrying = payload
888        .pointer("/params/willRetry")
889        .or_else(|| payload.get("willRetry"))
890        .and_then(Value::as_bool)
891        == Some(true);
892    if key == "error" && retrying {
893        // the harness says it is retrying (Codex app-server's `willRetry`): the turn goes on, and ends on its own
894        return vec![native_payload(&key, payload)];
895    }
896    if key == "transport_error" || key == "error" {
897        return vec![
898            json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
899        ];
900    }
901    if key == "session/update" {
902        let update = payload
903            .pointer("/params/update")
904            .or_else(|| payload.get("update"))
905            .unwrap_or(payload);
906        let update_kind = update
907            .get("sessionUpdate")
908            .or_else(|| update.get("type"))
909            .and_then(Value::as_str)
910            .unwrap_or_default()
911            .to_ascii_lowercase();
912        if update_kind == "agent_message_chunk" {
913            return vec![
914                json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
915            ];
916        }
917        if update_kind == "agent_thought_chunk" {
918            return vec![
919                json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
920            ];
921        }
922        if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
923            return vec![project_tool(update, payload)];
924        }
925    }
926    if key == "supercode/acp_request_completed" {
927        let failure = payload
928            .pointer("/params/error")
929            .or_else(|| payload.get("error"));
930        return completion(failure.and_then(extract_text));
931    }
932    match harness {
933        "codex" => project_codex(&key, payload),
934        "claude-code" => project_claude(&key, payload),
935        "pi" => project_pi(&key, payload),
936        "opencode" => project_opencode(&key, payload),
937        _ => project_generic(&key, payload),
938    }
939}
940
941fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
942    if key == "turn/started" {
943        return vec![native_payload(key, payload)];
944    }
945    if key == "turn/completed" {
946        let status = payload
947            .pointer("/params/turn/status")
948            .or_else(|| payload.pointer("/turn/status"))
949            .and_then(Value::as_str)
950            .unwrap_or("completed")
951            .to_ascii_lowercase();
952        return completion(
953            (status.contains("fail") || status.contains("error") || status.contains("cancel"))
954                .then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
955        );
956    }
957    if key.ends_with("/delta") && !key.contains("agentmessage") && !key.contains("reasoning") {
958        // a plan's or another item's stream is not the reply (`item/plan/delta`)
959        return vec![native_payload(key, payload)];
960    }
961    if key.ends_with("/delta") {
962        let text = payload
963            .pointer("/params/delta")
964            .or_else(|| payload.get("delta"))
965            .and_then(extract_text)
966            .or_else(|| extract_text(payload))
967            .unwrap_or_default();
968        return vec![
969            json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
970        ];
971    }
972    if key.contains("commandexecution")
973        || key.contains("mcptool")
974        || key.contains("filechange")
975        || key.contains("tool")
976    {
977        let source = payload
978            .pointer("/params/item")
979            .or_else(|| payload.get("item"))
980            .unwrap_or(payload);
981        return vec![project_tool(source, payload)];
982    }
983    vec![native_payload(key, payload)]
984}
985
986/// Claude Code's stream-json, as `ClaudeCodeBackend` launches it: `--print
987/// --input-format stream-json --output-format stream-json --verbose
988/// --permission-prompt-tool stdio` (`crates/harness/src/runtime/adapters.rs`).
989///
990/// That argv carries no `--include-partial-messages`, so nothing arrives as a
991/// `stream_event`: a turn is whole `assistant` and `user` events whose
992/// `message.content` is the Anthropic content-block array, plus `system`,
993/// `control_request` and a final `result`. The blocks are where the turn
994/// actually lives — text, thinking, `tool_use`, and `tool_result` — so a
995/// projector that reads only the envelope sees a turn with no tools in it and
996/// leaves every one of them to the `native_event` fallback. This walks them.
997fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
998    if key == "result" {
999        let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
1000            || payload
1001                .get("subtype")
1002                .and_then(Value::as_str)
1003                .is_some_and(|subtype| subtype != "success");
1004        return completion(
1005            failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
1006        );
1007    }
1008    if key == "stream_event" {
1009        // Only reachable when a caller adds `--include-partial-messages`. The
1010        // whole-message `assistant` event still follows, so text would arrive
1011        // twice; the partial path deliberately projects ONLY what the whole
1012        // message cannot carry on its own — nothing today.
1013        let stream = payload
1014            .get("event")
1015            .or_else(|| payload.get("stream_event"))
1016            .unwrap_or(payload);
1017        if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
1018            return vec![native_payload(key, payload)];
1019        }
1020    }
1021    if key == "assistant" {
1022        return project_claude_blocks(payload, true);
1023    }
1024    if key == "user" {
1025        return project_claude_blocks(payload, false);
1026    }
1027    if key == "control_request" {
1028        return project_claude_control_request(payload);
1029    }
1030    if key == "system" && claude_system_is_telemetry(payload) {
1031        return Vec::new();
1032    }
1033    if key == "tool_progress" {
1034        // The CLI's heartbeat for a tool still running (`elapsed_time_seconds`
1035        // every 30 s). The tool's own `tool_call_started` row already shows it
1036        // running; Claude Code's pane shows this only as a spinner's clock.
1037        return Vec::new();
1038    }
1039    if key == "rate_limit_event" {
1040        // Quota telemetry the CLI emits mid-turn. `allowed`/`allowed_warning`
1041        // is nothing a reader can act on and lands in the middle of the
1042        // assistant's own sentences; a `rejected` status is a real refusal, so
1043        // it falls through and stays visible.
1044        let status = payload
1045            .pointer("/rate_limit_info/status")
1046            .and_then(Value::as_str)
1047            .unwrap_or_default();
1048        if status.starts_with("allowed") {
1049            return Vec::new();
1050        }
1051    }
1052    project_generic(key, payload)
1053}
1054
1055/// The `system` events Claude Code's own pane never writes into its
1056/// transcript, so a frontend showing them shows more than the pane does.
1057///
1058/// Measured on a hosted orchestrator turn (1077 events, 218 tool calls): 233
1059/// `thinking_tokens`, 115 `task_started` and 115 `task_notification` (one pair
1060/// per Bash call, every one foreground with an empty `output_file`), 18
1061/// `task_progress`, one `task_updated`, one `vcs_state_changed`.
1062///
1063/// - `init`: the CLI announcing its tools, model and cwd; session setup.
1064/// - `thinking_tokens`: a running estimate the pane shows only on its spinner.
1065/// - `task_started` / `task_progress` / `task_updated`: the CLI's bookkeeping
1066///   for a tool call the `tool_call_*` rows already show.
1067/// - `task_notification` for a foreground task (no `output_file`): the tail
1068///   of that same call. A background task's notification names the file its
1069///   output went to and is the only word that it finished, so it stays.
1070/// - `hook_started` / `hook_progress`, and `hook_response` that exited 0: a
1071///   hook that ran and passed is invisible in the pane; a failing one is not.
1072/// - `vcs_state_changed`: the CLI noticing a commit or push its own tool
1073///   call made.
1074fn claude_system_is_telemetry(payload: &Value) -> bool {
1075    match payload
1076        .get("subtype")
1077        .and_then(Value::as_str)
1078        .unwrap_or_default()
1079    {
1080        "init" | "thinking_tokens" | "task_started" | "task_progress" | "task_updated"
1081        | "hook_started" | "hook_progress" | "vcs_state_changed" => true,
1082        "task_notification" => payload
1083            .get("output_file")
1084            .and_then(Value::as_str)
1085            .unwrap_or_default()
1086            .is_empty(),
1087        "hook_response" => payload.get("exit_code").and_then(Value::as_i64) == Some(0),
1088        _ => false,
1089    }
1090}
1091
1092/// The Anthropic content-block array an `assistant` or `user` event carries.
1093fn claude_blocks(payload: &Value) -> Option<&Vec<Value>> {
1094    payload
1095        .pointer("/message/content")
1096        .or_else(|| payload.get("content"))
1097        .and_then(Value::as_array)
1098}
1099
1100fn project_claude_blocks(payload: &Value, assistant: bool) -> Vec<Value> {
1101    let Some(blocks) = claude_blocks(payload) else {
1102        // A content-less envelope (or a plain string content) keeps the old
1103        // whole-payload reading rather than vanishing.
1104        let text = extract_text(payload).unwrap_or_default();
1105        if text.is_empty() {
1106            return vec![native_payload(
1107                if assistant { "assistant" } else { "user" },
1108                payload,
1109            )];
1110        }
1111        return vec![
1112            json!({"type":if assistant { "text_delta" } else { "user_message" }, "text":text, "raw":payload}),
1113        ];
1114    };
1115    let mut projected = Vec::new();
1116    for block in blocks {
1117        match block
1118            .get("type")
1119            .and_then(Value::as_str)
1120            .unwrap_or_default()
1121        {
1122            "text" => {
1123                let text = block
1124                    .get("text")
1125                    .and_then(Value::as_str)
1126                    .unwrap_or_default();
1127                if !text.is_empty() {
1128                    projected.push(json!({
1129                        "type": if assistant { "text_delta" } else { "user_message" },
1130                        "text": text,
1131                        "raw": block,
1132                    }));
1133                }
1134            }
1135            "thinking" | "redacted_thinking" => {
1136                let text = extract_text(block).unwrap_or_default();
1137                if !text.is_empty() {
1138                    projected.push(json!({"type":"reasoning", "text":text, "raw":block}));
1139                }
1140            }
1141            "tool_use" => projected.push(json!({
1142                "type": "tool_call_started",
1143                "id": block.get("id").cloned().unwrap_or(Value::Null),
1144                "name": block.get("name").cloned().unwrap_or(Value::Null),
1145                "arguments": block.get("input").map(Value::to_string).unwrap_or_default(),
1146                "raw": block,
1147            })),
1148            "tool_result" => projected.push(json!({
1149                // The result names only the call it answers — the tool's own
1150                // name was stated when it started, and a frontend correlates
1151                // the pair by that id.
1152                "type": "tool_call_completed",
1153                "id": block.get("tool_use_id").cloned().unwrap_or(Value::Null),
1154                "name": Value::Null,
1155                "output": extract_text(block.get("content").unwrap_or(block)).unwrap_or_default(),
1156                "is_error": block.get("is_error").and_then(Value::as_bool).unwrap_or(false),
1157                "raw": block,
1158            })),
1159            _ => projected.push(native_payload(
1160                if assistant { "assistant" } else { "user" },
1161                block,
1162            )),
1163        }
1164    }
1165    projected
1166}
1167
1168/// Claude Code's `can_use_tool` control request as the canonical `request`.
1169///
1170/// OBSERVATION ONLY, and deliberately: the hosted descriptor advertises
1171/// `respond: false` because the portable contract cannot carry this harness's
1172/// native response envelope — the owner/editor answers it through
1173/// `harness.v1.approvals.resolve`. Emitting it is what that same decision
1174/// allows ("a terminal may observe requests"), and a frontend renders controls
1175/// only for actions the descriptor advertises, so no answer button appears.
1176fn project_claude_control_request(payload: &Value) -> Vec<Value> {
1177    let request = payload.get("request").unwrap_or(payload);
1178    if request.get("subtype").and_then(Value::as_str) != Some("can_use_tool") {
1179        return vec![native_payload("control_request", payload)];
1180    }
1181    let native_id = payload
1182        .get("request_id")
1183        .and_then(Value::as_str)
1184        .unwrap_or_default();
1185    vec![json!({
1186        "type": "request",
1187        "request": {
1188            "id": claude_request_id(native_id),
1189            "kind": "approval",
1190            "payload": {
1191                "tool": request.get("tool_name").or_else(|| request.get("toolName")).cloned().unwrap_or(Value::Null),
1192                "arguments": request.get("input").cloned().unwrap_or(Value::Null),
1193                "native_request_id": native_id,
1194                // What this harness's own protocol can carry, so a frontend
1195                // offers exactly those answers and no button that would be
1196                // refused.
1197                "decisions": hosted_answerable_decisions("claude-code"),
1198            },
1199        },
1200    })]
1201}
1202
1203/// A stable JSON-safe number for a native request id that is a string.
1204///
1205/// The contract's `FrontendRequest.id` is numeric and Claude Code's
1206/// `request_id` is a uuid, so the id a frontend displays and correlates on is
1207/// this digest. The native id travels beside it, because that is the one the
1208/// runtime's own resolve door needs.
1209fn claude_request_id(native_id: &str) -> u64 {
1210    let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1211    for byte in native_id.as_bytes() {
1212        hash ^= u64::from(*byte);
1213        hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1214    }
1215    // Stay inside the range JSON numbers carry exactly.
1216    hash & ((1_u64 << 53) - 1)
1217}
1218
1219fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
1220    match key {
1221        "agent_start" => vec![native_payload(key, payload)],
1222        "agent_end" => completion(payload.get("error").and_then(extract_text)),
1223        // a command Pi refuses (`{"type":"response","success":false,"error":…}`): a refused prompt has no turn to end
1224        "response"
1225            if payload.get("success").and_then(Value::as_bool) == Some(false)
1226                && payload.get("command").and_then(Value::as_str) == Some("prompt") =>
1227        {
1228            completion(Some(
1229                payload
1230                    .get("error")
1231                    .and_then(extract_text)
1232                    .unwrap_or_else(|| "Pi refused the prompt.".into()),
1233            ))
1234        }
1235        "message_update" => {
1236            // pi streams one assistant message as `text_start` → `text_delta`*
1237            // → `text_end` (the same for `thinking_*`). Only the `*_delta`
1238            // events carry NEW text; `text_end` repeats the whole block as
1239            // `content`, and projecting it too doubled every reply.
1240            let update = payload
1241                .get("assistantMessageEvent")
1242                .or_else(|| payload.get("event"))
1243                .unwrap_or(payload);
1244            let kind = update
1245                .get("type")
1246                .and_then(Value::as_str)
1247                .unwrap_or_default();
1248            match kind {
1249                "text_delta" => vec![
1250                    json!({"type":"text_delta", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1251                ],
1252                "thinking_delta" => vec![
1253                    json!({"type":"reasoning", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1254                ],
1255                _ => vec![native_payload(key, payload)],
1256            }
1257        }
1258        _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
1259        _ => vec![native_payload(key, payload)],
1260    }
1261}
1262
1263fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
1264    if key == "session.idle" {
1265        return completion(None);
1266    }
1267    if key == "session.status" {
1268        let status = payload
1269            .pointer("/properties/status/type")
1270            .or_else(|| payload.pointer("/status/type"))
1271            .and_then(Value::as_str)
1272            .unwrap_or_default();
1273        if status == "busy" {
1274            return vec![native_payload(key, payload)];
1275        }
1276        // `session.status{idle}` precedes `session.idle` in every turn: the turn ends once, on `session.idle`
1277        return vec![native_payload(key, payload)];
1278    }
1279    if key == "message.part.delta" {
1280        // OpenCode 1.2's text stream: `{field, delta}` per part; the part's final `message.part.updated` repeats it whole
1281        let delta = payload
1282            .pointer("/properties/delta")
1283            .or_else(|| payload.get("delta"))
1284            .and_then(Value::as_str)
1285            .unwrap_or_default();
1286        let field = payload
1287            .pointer("/properties/field")
1288            .and_then(Value::as_str)
1289            .unwrap_or("text");
1290        return vec![
1291            json!({"type":if field.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1292        ];
1293    }
1294    if key == "message.part.updated" {
1295        let part = payload
1296            .pointer("/properties/part")
1297            .or_else(|| payload.get("part"))
1298            .unwrap_or(payload);
1299        if let Some(delta) = payload
1300            .pointer("/properties/delta")
1301            .or_else(|| payload.get("delta"))
1302            .and_then(Value::as_str)
1303        {
1304            let reasoning = part
1305                .get("type")
1306                .and_then(Value::as_str)
1307                .is_some_and(|kind| kind.contains("reasoning"));
1308            return vec![
1309                json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1310            ];
1311        }
1312        if part
1313            .get("type")
1314            .and_then(Value::as_str)
1315            .is_some_and(|kind| kind.contains("tool"))
1316        {
1317            return vec![project_tool(part, payload)];
1318        }
1319    }
1320    if key == "session.error" {
1321        return completion(Some(
1322            extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
1323        ));
1324    }
1325    vec![native_payload(key, payload)]
1326}
1327
1328fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
1329    match key {
1330        "turn_started" | "turn/started" | "agent_start" => {
1331            vec![native_payload(key, payload)]
1332        }
1333        "turn_completed" | "turn/completed" | "agent_end" => {
1334            completion(payload.get("error").and_then(extract_text))
1335        }
1336        "output_delta" | "content_delta" => vec![
1337            json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1338        ],
1339        "reasoning_delta" => vec![
1340            json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1341        ],
1342        "tool" => vec![project_tool(payload, payload)],
1343        _ => vec![native_payload(key, payload)],
1344    }
1345}
1346
1347fn completion(error: Option<String>) -> Vec<Value> {
1348    match error {
1349        Some(message) => vec![
1350            json!({"type":"turn_failed", "message":message}),
1351            json!({"type":"turn_completed"}),
1352        ],
1353        None => vec![
1354            json!({"type":"turn_succeeded"}),
1355            json!({"type":"turn_completed"}),
1356        ],
1357    }
1358}
1359
1360fn project_tool(source: &Value, raw: &Value) -> Value {
1361    let status = source
1362        .get("status")
1363        .or_else(|| source.get("state"))
1364        .or_else(|| source.get("sessionUpdate"))
1365        .and_then(Value::as_str)
1366        .unwrap_or_default()
1367        .to_ascii_lowercase();
1368    let completed = status.contains("complete")
1369        || status.contains("result")
1370        || status.contains("success")
1371        || status.contains("error")
1372        || status.contains("fail");
1373    let arguments = source
1374        .get("arguments")
1375        .or_else(|| source.get("input"))
1376        .or_else(|| source.get("rawInput"))
1377        .cloned()
1378        .unwrap_or(Value::Null);
1379    json!({
1380        "type": if completed { "tool_call_completed" } else { "tool_call_started" },
1381        "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),
1382        "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
1383        "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
1384        "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
1385        "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
1386        "raw": raw,
1387    })
1388}
1389
1390fn native_payload(kind: &str, payload: &Value) -> Value {
1391    json!({"type":"native_event", "kind":kind, "raw":payload})
1392}
1393
1394fn extract_text(value: &Value) -> Option<String> {
1395    match value {
1396        Value::String(text) => Some(text.clone()),
1397        Value::Array(values) => {
1398            let text = values
1399                .iter()
1400                .filter_map(extract_text)
1401                .collect::<Vec<_>>()
1402                .join("\n");
1403            (!text.is_empty()).then_some(text)
1404        }
1405        Value::Object(object) => {
1406            for key in ["text", "delta", "content", "message", "result", "error"] {
1407                if let Some(text) = object.get(key).and_then(Value::as_str) {
1408                    return Some(text.to_string());
1409                }
1410            }
1411            for key in [
1412                "delta",
1413                "content",
1414                "message",
1415                "error",
1416                "data",
1417                "part",
1418                "params",
1419                "properties",
1420                "update",
1421                "event",
1422            ] {
1423                if let Some(text) = object.get(key).and_then(extract_text) {
1424                    return Some(text);
1425                }
1426            }
1427            None
1428        }
1429        _ => None,
1430    }
1431}