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. And a Pi run its session will retry
807/// (`agent_end` with `willRetry`) leaves the turn open: a retry cancelled in its backoff sends no `agent_end` after,
808/// only `auto_retry_end` with `success: false`, which then ends it.
809#[derive(Debug, Default)]
810pub struct NativeProjection {
811    opencode_reasoning_parts: HashSet<String>,
812    opencode_failed: bool,
813    pi_retrying: bool,
814}
815
816impl NativeProjection {
817    /// `event` of `harness` as `project_native_event` projects it, with the session's memory applied.
818    pub fn project(&mut self, harness: &str, event: &HarnessEvent) -> Vec<Value> {
819        if harness == "pi" {
820            let payload = &event.payload;
821            match event.kind.as_str() {
822                "agent_end" => {
823                    self.pi_retrying =
824                        payload.get("willRetry").and_then(Value::as_bool) == Some(true);
825                }
826                // after a retried run's own `agent_end` the turn is already closed; only a retry that never ran
827                // (cancelled in its backoff) is ended here
828                "auto_retry_end"
829                    if self.pi_retrying
830                        && payload.get("success").and_then(Value::as_bool) == Some(false) =>
831                {
832                    self.pi_retrying = false;
833                    return completion(Some(
834                        payload
835                            .get("finalError")
836                            .and_then(Value::as_str)
837                            .unwrap_or("Pi's retry did not run.")
838                            .to_string(),
839                    ));
840                }
841                _ => {}
842            }
843        }
844        if harness == "opencode" {
845            let key = event.kind.to_ascii_lowercase();
846            let payload = &event.payload;
847            match key.as_str() {
848                "message.part.updated" => {
849                    let part = payload
850                        .pointer("/properties/part")
851                        .or_else(|| payload.get("part"))
852                        .unwrap_or(payload);
853                    if let (Some(id), Some(kind)) = (
854                        part.get("id").and_then(Value::as_str),
855                        part.get("type").and_then(Value::as_str),
856                    ) {
857                        if kind.contains("reasoning") {
858                            self.opencode_reasoning_parts.insert(id.to_string());
859                        }
860                    }
861                }
862                "message.part.delta" => {
863                    let reasoning = payload
864                        .pointer("/properties/partID")
865                        .and_then(Value::as_str)
866                        .is_some_and(|id| self.opencode_reasoning_parts.contains(id));
867                    if reasoning {
868                        let text = payload
869                            .pointer("/properties/delta")
870                            .and_then(Value::as_str)
871                            .unwrap_or_default();
872                        return vec![json!({"type":"reasoning", "text":text, "raw":payload})];
873                    }
874                }
875                "session.error" => {
876                    self.opencode_failed = true;
877                    return vec![
878                        json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into())}),
879                    ];
880                }
881                "session.idle" if self.opencode_failed => {
882                    self.opencode_failed = false;
883                    self.opencode_reasoning_parts.clear();
884                    return vec![json!({"type":"turn_completed"})];
885                }
886                "session.idle" => self.opencode_reasoning_parts.clear(),
887                // a new turn starts clean: a failure no idle followed does not carry into it
888                "session.status"
889                    if payload
890                        .pointer("/properties/status/type")
891                        .or_else(|| payload.pointer("/status/type"))
892                        .and_then(Value::as_str)
893                        == Some("busy") =>
894                {
895                    self.opencode_failed = false;
896                }
897                _ => {}
898            }
899        }
900        project_native_event(harness, event)
901    }
902}
903
904/// One native event of `harness` as the protocol-neutral events every reader renders (`text_delta`, `reasoning`,
905/// `turn_succeeded`/`turn_failed`, `turn_completed`, `runtime_disconnected`, tools, approvals): ACP's and each
906/// harness's own (Codex, Claude Code, Pi, OpenCode, generic) alike.
907fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
908    let key = event.kind.to_ascii_lowercase().replace('-', "_");
909    let payload = &event.payload;
910    if key == "transport_closed" {
911        return vec![
912            json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
913        ];
914    }
915    let retrying = payload
916        .pointer("/params/willRetry")
917        .or_else(|| payload.get("willRetry"))
918        .and_then(Value::as_bool)
919        == Some(true);
920    if key == "error" && retrying {
921        // the harness says it is retrying (Codex app-server's `willRetry`): the turn goes on, and ends on its own
922        return vec![native_payload(&key, payload)];
923    }
924    if key == "transport_error" || key == "error" {
925        return vec![
926            json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
927        ];
928    }
929    if key == "session/update" {
930        let update = payload
931            .pointer("/params/update")
932            .or_else(|| payload.get("update"))
933            .unwrap_or(payload);
934        let update_kind = update
935            .get("sessionUpdate")
936            .or_else(|| update.get("type"))
937            .and_then(Value::as_str)
938            .unwrap_or_default()
939            .to_ascii_lowercase();
940        if update_kind == "agent_message_chunk" {
941            return vec![
942                json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
943            ];
944        }
945        if update_kind == "agent_thought_chunk" {
946            return vec![
947                json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
948            ];
949        }
950        if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
951            return vec![project_tool(update, payload)];
952        }
953    }
954    if key == "supercode/acp_request_completed" {
955        let failure = payload
956            .pointer("/params/error")
957            .or_else(|| payload.get("error"));
958        return completion(failure.and_then(extract_text));
959    }
960    match harness {
961        "codex" => project_codex(&key, payload),
962        "claude-code" => project_claude(&key, payload),
963        "pi" => project_pi(&key, payload),
964        "opencode" => project_opencode(&key, payload),
965        _ => project_generic(&key, payload),
966    }
967}
968
969fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
970    if key == "turn/started" {
971        return vec![native_payload(key, payload)];
972    }
973    if key == "turn/completed" {
974        let status = payload
975            .pointer("/params/turn/status")
976            .or_else(|| payload.pointer("/turn/status"))
977            .and_then(Value::as_str)
978            .unwrap_or("completed")
979            .to_ascii_lowercase();
980        return completion(
981            (status.contains("fail") || status.contains("error") || status.contains("cancel"))
982                .then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
983        );
984    }
985    if key.ends_with("/delta") && !key.contains("agentmessage") && !key.contains("reasoning") {
986        // a plan's or another item's stream is not the reply (`item/plan/delta`)
987        return vec![native_payload(key, payload)];
988    }
989    if key.ends_with("/delta") {
990        let text = payload
991            .pointer("/params/delta")
992            .or_else(|| payload.get("delta"))
993            .and_then(extract_text)
994            .or_else(|| extract_text(payload))
995            .unwrap_or_default();
996        return vec![
997            json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
998        ];
999    }
1000    if key.contains("commandexecution")
1001        || key.contains("mcptool")
1002        || key.contains("filechange")
1003        || key.contains("tool")
1004    {
1005        let source = payload
1006            .pointer("/params/item")
1007            .or_else(|| payload.get("item"))
1008            .unwrap_or(payload);
1009        return vec![project_tool(source, payload)];
1010    }
1011    vec![native_payload(key, payload)]
1012}
1013
1014/// Claude Code's stream-json, as `ClaudeCodeBackend` launches it: `--print
1015/// --input-format stream-json --output-format stream-json --verbose
1016/// --permission-prompt-tool stdio` (`crates/harness/src/runtime/adapters.rs`).
1017///
1018/// That argv carries no `--include-partial-messages`, so nothing arrives as a
1019/// `stream_event`: a turn is whole `assistant` and `user` events whose
1020/// `message.content` is the Anthropic content-block array, plus `system`,
1021/// `control_request` and a final `result`. The blocks are where the turn
1022/// actually lives — text, thinking, `tool_use`, and `tool_result` — so a
1023/// projector that reads only the envelope sees a turn with no tools in it and
1024/// leaves every one of them to the `native_event` fallback. This walks them.
1025fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
1026    if key == "result" {
1027        let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
1028            || payload
1029                .get("subtype")
1030                .and_then(Value::as_str)
1031                .is_some_and(|subtype| subtype != "success");
1032        return completion(
1033            failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
1034        );
1035    }
1036    if key == "stream_event" {
1037        // Only reachable when a caller adds `--include-partial-messages`. The
1038        // whole-message `assistant` event still follows, so text would arrive
1039        // twice; the partial path deliberately projects ONLY what the whole
1040        // message cannot carry on its own — nothing today.
1041        let stream = payload
1042            .get("event")
1043            .or_else(|| payload.get("stream_event"))
1044            .unwrap_or(payload);
1045        if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
1046            return vec![native_payload(key, payload)];
1047        }
1048    }
1049    if key == "assistant" {
1050        return project_claude_blocks(payload, true);
1051    }
1052    if key == "user" {
1053        return project_claude_blocks(payload, false);
1054    }
1055    if key == "control_request" {
1056        return project_claude_control_request(payload);
1057    }
1058    if key == "system" && claude_system_is_telemetry(payload) {
1059        return Vec::new();
1060    }
1061    if key == "tool_progress" {
1062        // The CLI's heartbeat for a tool still running (`elapsed_time_seconds`
1063        // every 30 s). The tool's own `tool_call_started` row already shows it
1064        // running; Claude Code's pane shows this only as a spinner's clock.
1065        return Vec::new();
1066    }
1067    if key == "rate_limit_event" {
1068        // Quota telemetry the CLI emits mid-turn. `allowed`/`allowed_warning`
1069        // is nothing a reader can act on and lands in the middle of the
1070        // assistant's own sentences; a `rejected` status is a real refusal, so
1071        // it falls through and stays visible.
1072        let status = payload
1073            .pointer("/rate_limit_info/status")
1074            .and_then(Value::as_str)
1075            .unwrap_or_default();
1076        if status.starts_with("allowed") {
1077            return Vec::new();
1078        }
1079    }
1080    project_generic(key, payload)
1081}
1082
1083/// The `system` events Claude Code's own pane never writes into its
1084/// transcript, so a frontend showing them shows more than the pane does.
1085///
1086/// Measured on a hosted orchestrator turn (1077 events, 218 tool calls): 233
1087/// `thinking_tokens`, 115 `task_started` and 115 `task_notification` (one pair
1088/// per Bash call, every one foreground with an empty `output_file`), 18
1089/// `task_progress`, one `task_updated`, one `vcs_state_changed`.
1090///
1091/// - `init`: the CLI announcing its tools, model and cwd; session setup.
1092/// - `thinking_tokens`: a running estimate the pane shows only on its spinner.
1093/// - `task_started` / `task_progress` / `task_updated`: the CLI's bookkeeping
1094///   for a tool call the `tool_call_*` rows already show.
1095/// - `task_notification` for a foreground task (no `output_file`): the tail
1096///   of that same call. A background task's notification names the file its
1097///   output went to and is the only word that it finished, so it stays.
1098/// - `hook_started` / `hook_progress`, and `hook_response` that exited 0: a
1099///   hook that ran and passed is invisible in the pane; a failing one is not.
1100/// - `vcs_state_changed`: the CLI noticing a commit or push its own tool
1101///   call made.
1102fn claude_system_is_telemetry(payload: &Value) -> bool {
1103    match payload
1104        .get("subtype")
1105        .and_then(Value::as_str)
1106        .unwrap_or_default()
1107    {
1108        "init" | "thinking_tokens" | "task_started" | "task_progress" | "task_updated"
1109        | "hook_started" | "hook_progress" | "vcs_state_changed" => true,
1110        "task_notification" => payload
1111            .get("output_file")
1112            .and_then(Value::as_str)
1113            .unwrap_or_default()
1114            .is_empty(),
1115        "hook_response" => payload.get("exit_code").and_then(Value::as_i64) == Some(0),
1116        _ => false,
1117    }
1118}
1119
1120/// The Anthropic content-block array an `assistant` or `user` event carries.
1121fn claude_blocks(payload: &Value) -> Option<&Vec<Value>> {
1122    payload
1123        .pointer("/message/content")
1124        .or_else(|| payload.get("content"))
1125        .and_then(Value::as_array)
1126}
1127
1128fn project_claude_blocks(payload: &Value, assistant: bool) -> Vec<Value> {
1129    let Some(blocks) = claude_blocks(payload) else {
1130        // A content-less envelope (or a plain string content) keeps the old
1131        // whole-payload reading rather than vanishing.
1132        let text = extract_text(payload).unwrap_or_default();
1133        if text.is_empty() {
1134            return vec![native_payload(
1135                if assistant { "assistant" } else { "user" },
1136                payload,
1137            )];
1138        }
1139        return vec![
1140            json!({"type":if assistant { "text_delta" } else { "user_message" }, "text":text, "raw":payload}),
1141        ];
1142    };
1143    let mut projected = Vec::new();
1144    for block in blocks {
1145        match block
1146            .get("type")
1147            .and_then(Value::as_str)
1148            .unwrap_or_default()
1149        {
1150            "text" => {
1151                let text = block
1152                    .get("text")
1153                    .and_then(Value::as_str)
1154                    .unwrap_or_default();
1155                if !text.is_empty() {
1156                    projected.push(json!({
1157                        "type": if assistant { "text_delta" } else { "user_message" },
1158                        "text": text,
1159                        "raw": block,
1160                    }));
1161                }
1162            }
1163            "thinking" | "redacted_thinking" => {
1164                let text = extract_text(block).unwrap_or_default();
1165                if !text.is_empty() {
1166                    projected.push(json!({"type":"reasoning", "text":text, "raw":block}));
1167                }
1168            }
1169            "tool_use" => projected.push(json!({
1170                "type": "tool_call_started",
1171                "id": block.get("id").cloned().unwrap_or(Value::Null),
1172                "name": block.get("name").cloned().unwrap_or(Value::Null),
1173                "arguments": block.get("input").map(Value::to_string).unwrap_or_default(),
1174                "raw": block,
1175            })),
1176            "tool_result" => projected.push(json!({
1177                // The result names only the call it answers — the tool's own
1178                // name was stated when it started, and a frontend correlates
1179                // the pair by that id.
1180                "type": "tool_call_completed",
1181                "id": block.get("tool_use_id").cloned().unwrap_or(Value::Null),
1182                "name": Value::Null,
1183                "output": extract_text(block.get("content").unwrap_or(block)).unwrap_or_default(),
1184                "is_error": block.get("is_error").and_then(Value::as_bool).unwrap_or(false),
1185                "raw": block,
1186            })),
1187            _ => projected.push(native_payload(
1188                if assistant { "assistant" } else { "user" },
1189                block,
1190            )),
1191        }
1192    }
1193    projected
1194}
1195
1196/// Claude Code's `can_use_tool` control request as the canonical `request`.
1197///
1198/// OBSERVATION ONLY, and deliberately: the hosted descriptor advertises
1199/// `respond: false` because the portable contract cannot carry this harness's
1200/// native response envelope — the owner/editor answers it through
1201/// `harness.v1.approvals.resolve`. Emitting it is what that same decision
1202/// allows ("a terminal may observe requests"), and a frontend renders controls
1203/// only for actions the descriptor advertises, so no answer button appears.
1204fn project_claude_control_request(payload: &Value) -> Vec<Value> {
1205    let request = payload.get("request").unwrap_or(payload);
1206    if request.get("subtype").and_then(Value::as_str) != Some("can_use_tool") {
1207        return vec![native_payload("control_request", payload)];
1208    }
1209    let native_id = payload
1210        .get("request_id")
1211        .and_then(Value::as_str)
1212        .unwrap_or_default();
1213    vec![json!({
1214        "type": "request",
1215        "request": {
1216            "id": claude_request_id(native_id),
1217            "kind": "approval",
1218            "payload": {
1219                "tool": request.get("tool_name").or_else(|| request.get("toolName")).cloned().unwrap_or(Value::Null),
1220                "arguments": request.get("input").cloned().unwrap_or(Value::Null),
1221                "native_request_id": native_id,
1222                // What this harness's own protocol can carry, so a frontend
1223                // offers exactly those answers and no button that would be
1224                // refused.
1225                "decisions": hosted_answerable_decisions("claude-code"),
1226            },
1227        },
1228    })]
1229}
1230
1231/// A stable JSON-safe number for a native request id that is a string.
1232///
1233/// The contract's `FrontendRequest.id` is numeric and Claude Code's
1234/// `request_id` is a uuid, so the id a frontend displays and correlates on is
1235/// this digest. The native id travels beside it, because that is the one the
1236/// runtime's own resolve door needs.
1237fn claude_request_id(native_id: &str) -> u64 {
1238    let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1239    for byte in native_id.as_bytes() {
1240        hash ^= u64::from(*byte);
1241        hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1242    }
1243    // Stay inside the range JSON numbers carry exactly.
1244    hash & ((1_u64 << 53) - 1)
1245}
1246
1247fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
1248    match key {
1249        "agent_start" => vec![native_payload(key, payload)],
1250        // a run Pi will retry (its session's `willRetry`, an overloaded or 5xx model call) is not the turn's end:
1251        // the retry is a run of its own, and the turn ends with it
1252        "agent_end" if payload.get("willRetry").and_then(Value::as_bool) == Some(true) => {
1253            vec![native_payload(key, payload)]
1254        }
1255        "agent_end" => completion(
1256            payload
1257                .get("error")
1258                .and_then(extract_text)
1259                .or_else(|| pi_run_error(payload)),
1260        ),
1261        // a command Pi refuses (`{"type":"response","success":false,"error":…}`): a refused prompt has no turn to end
1262        "response"
1263            if payload.get("success").and_then(Value::as_bool) == Some(false)
1264                && payload.get("command").and_then(Value::as_str) == Some("prompt") =>
1265        {
1266            completion(Some(
1267                payload
1268                    .get("error")
1269                    .and_then(extract_text)
1270                    .unwrap_or_else(|| "Pi refused the prompt.".into()),
1271            ))
1272        }
1273        "message_update" => {
1274            // pi streams one assistant message as `text_start` → `text_delta`*
1275            // → `text_end` (the same for `thinking_*`). Only the `*_delta`
1276            // events carry NEW text; `text_end` repeats the whole block as
1277            // `content`, and projecting it too doubled every reply.
1278            let update = payload
1279                .get("assistantMessageEvent")
1280                .or_else(|| payload.get("event"))
1281                .unwrap_or(payload);
1282            let kind = update
1283                .get("type")
1284                .and_then(Value::as_str)
1285                .unwrap_or_default();
1286            match kind {
1287                "text_delta" => vec![
1288                    json!({"type":"text_delta", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1289                ],
1290                "thinking_delta" => vec![
1291                    json!({"type":"reasoning", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1292                ],
1293                _ => vec![native_payload(key, payload)],
1294            }
1295        }
1296        _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
1297        _ => vec![native_payload(key, payload)],
1298    }
1299}
1300
1301fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
1302    if key == "session.idle" {
1303        return completion(None);
1304    }
1305    if key == "session.status" {
1306        let status = payload
1307            .pointer("/properties/status/type")
1308            .or_else(|| payload.pointer("/status/type"))
1309            .and_then(Value::as_str)
1310            .unwrap_or_default();
1311        if status == "busy" {
1312            return vec![native_payload(key, payload)];
1313        }
1314        // `session.status{idle}` precedes `session.idle` in every turn: the turn ends once, on `session.idle`
1315        return vec![native_payload(key, payload)];
1316    }
1317    if key == "message.part.delta" {
1318        // OpenCode 1.2's text stream: `{field, delta}` per part; the part's final `message.part.updated` repeats it whole
1319        let delta = payload
1320            .pointer("/properties/delta")
1321            .or_else(|| payload.get("delta"))
1322            .and_then(Value::as_str)
1323            .unwrap_or_default();
1324        let field = payload
1325            .pointer("/properties/field")
1326            .and_then(Value::as_str)
1327            .unwrap_or("text");
1328        return vec![
1329            json!({"type":if field.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1330        ];
1331    }
1332    if key == "message.part.updated" {
1333        let part = payload
1334            .pointer("/properties/part")
1335            .or_else(|| payload.get("part"))
1336            .unwrap_or(payload);
1337        if let Some(delta) = payload
1338            .pointer("/properties/delta")
1339            .or_else(|| payload.get("delta"))
1340            .and_then(Value::as_str)
1341        {
1342            let reasoning = part
1343                .get("type")
1344                .and_then(Value::as_str)
1345                .is_some_and(|kind| kind.contains("reasoning"));
1346            return vec![
1347                json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1348            ];
1349        }
1350        if part
1351            .get("type")
1352            .and_then(Value::as_str)
1353            .is_some_and(|kind| kind.contains("tool"))
1354        {
1355            return vec![project_tool(part, payload)];
1356        }
1357    }
1358    if key == "session.error" {
1359        return completion(Some(
1360            extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
1361        ));
1362    }
1363    vec![native_payload(key, payload)]
1364}
1365
1366fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
1367    match key {
1368        "turn_started" | "turn/started" | "agent_start" => {
1369            vec![native_payload(key, payload)]
1370        }
1371        "turn_completed" | "turn/completed" | "agent_end" => {
1372            completion(payload.get("error").and_then(extract_text))
1373        }
1374        "output_delta" | "content_delta" => vec![
1375            json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1376        ],
1377        "reasoning_delta" => vec![
1378            json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1379        ],
1380        "tool" => vec![project_tool(payload, payload)],
1381        _ => vec![native_payload(key, payload)],
1382    }
1383}
1384
1385/// Why a Pi run ended in failure: Pi's `agent_end` is `{messages}` and carries
1386/// no error of its own; a failed model call is its last assistant message,
1387/// `stopReason: "error"` with the provider's words in `errorMessage` (an
1388/// interrupt is `"aborted"`, not a failure).
1389fn pi_run_error(payload: &Value) -> Option<String> {
1390    let last = payload
1391        .get("messages")?
1392        .as_array()?
1393        .iter()
1394        .rev()
1395        .find(|message| message.get("role").and_then(Value::as_str) == Some("assistant"))?;
1396    if last.get("stopReason").and_then(Value::as_str) != Some("error") {
1397        return None;
1398    }
1399    Some(
1400        last.get("errorMessage")
1401            .and_then(Value::as_str)
1402            .filter(|text| !text.is_empty())
1403            .unwrap_or("Pi's model call failed.")
1404            .to_string(),
1405    )
1406}
1407
1408fn completion(error: Option<String>) -> Vec<Value> {
1409    match error {
1410        Some(message) => vec![
1411            json!({"type":"turn_failed", "message":message}),
1412            json!({"type":"turn_completed"}),
1413        ],
1414        None => vec![
1415            json!({"type":"turn_succeeded"}),
1416            json!({"type":"turn_completed"}),
1417        ],
1418    }
1419}
1420
1421fn project_tool(source: &Value, raw: &Value) -> Value {
1422    let status = source
1423        .get("status")
1424        .or_else(|| source.get("state"))
1425        .or_else(|| source.get("sessionUpdate"))
1426        .and_then(Value::as_str)
1427        .unwrap_or_default()
1428        .to_ascii_lowercase();
1429    let completed = status.contains("complete")
1430        || status.contains("result")
1431        || status.contains("success")
1432        || status.contains("error")
1433        || status.contains("fail");
1434    let arguments = source
1435        .get("arguments")
1436        .or_else(|| source.get("input"))
1437        .or_else(|| source.get("rawInput"))
1438        .cloned()
1439        .unwrap_or(Value::Null);
1440    json!({
1441        "type": if completed { "tool_call_completed" } else { "tool_call_started" },
1442        "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),
1443        "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
1444        "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
1445        "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
1446        "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
1447        "raw": raw,
1448    })
1449}
1450
1451fn native_payload(kind: &str, payload: &Value) -> Value {
1452    json!({"type":"native_event", "kind":kind, "raw":payload})
1453}
1454
1455fn extract_text(value: &Value) -> Option<String> {
1456    match value {
1457        Value::String(text) => Some(text.clone()),
1458        Value::Array(values) => {
1459            let text = values
1460                .iter()
1461                .filter_map(extract_text)
1462                .collect::<Vec<_>>()
1463                .join("\n");
1464            (!text.is_empty()).then_some(text)
1465        }
1466        Value::Object(object) => {
1467            for key in ["text", "delta", "content", "message", "result", "error"] {
1468                if let Some(text) = object.get(key).and_then(Value::as_str) {
1469                    return Some(text.to_string());
1470                }
1471            }
1472            for key in [
1473                "delta",
1474                "content",
1475                "message",
1476                "error",
1477                "data",
1478                "part",
1479                "params",
1480                "properties",
1481                "update",
1482                "event",
1483            ] {
1484                if let Some(text) = object.get(key).and_then(extract_text) {
1485                    return Some(text);
1486                }
1487            }
1488            None
1489        }
1490        _ => None,
1491    }
1492}