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