1use std::collections::{HashMap, 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
54pub 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 history: Vec<ChatMessage>,
65 pending_requests: StdMutex<HashMap<u64, Value>>,
72 busy: AtomicBool,
73 closed: AtomicBool,
74}
75
76impl HostedHarnessRuntime {
77 pub fn spawn(
80 runtime: Box<dyn RuntimeConnection>,
81 capabilities: RuntimeCapabilities,
82 ) -> (Arc<Self>, HostedHarnessConnection) {
83 let handle = runtime.handle().clone();
84 let (commands, command_rx) = mpsc::channel(32);
85 let (raw_events, raw_rx) = broadcast::channel(1024);
86 let (frontend_events, _) = broadcast::channel(1024);
87 let host = Arc::new(Self {
88 handle: handle.clone(),
89 capabilities,
90 commands,
91 raw_events,
92 frontend_events,
93 projection: StdMutex::new(ProjectionState {
94 next_sequence: 1,
95 replay: VecDeque::new(),
96 }),
97 history: recorded_history(&handle),
98 pending_requests: StdMutex::new(HashMap::new()),
99 busy: AtomicBool::new(false),
100 closed: AtomicBool::new(false),
101 });
102 tokio::spawn(run_native_runtime(
103 runtime,
104 Arc::downgrade(&host),
105 command_rx,
106 ));
107 let connection = HostedHarnessConnection {
108 host: host.clone(),
109 handle,
110 events: raw_rx,
111 closed: false,
112 };
113 (host, connection)
114 }
115
116 pub fn frontend_sender(&self) -> broadcast::Sender<FrontendEvent> {
118 self.frontend_events.clone()
119 }
120
121 pub async fn shutdown(&self) -> Result<()> {
123 if self.closed.load(Ordering::SeqCst) {
124 return Ok(());
125 }
126 let (reply, response) = oneshot::channel();
127 self.commands
128 .send(HostCommand::Shutdown { reply })
129 .await
130 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
131 response
132 .await
133 .map_err(|_| Error::Other("hosted harness runtime stopped before shutdown".into()))?
134 }
135
136 fn claim_submit(&self) -> Result<()> {
137 if self.closed.load(Ordering::SeqCst) {
138 return Err(Error::Other("hosted harness runtime is closed".into()));
139 }
140 if self
141 .busy
142 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
143 .is_err()
144 {
145 return Err(Error::Other("a harness turn is already in progress".into()));
146 }
147 Ok(())
148 }
149
150 async fn submit_native_claimed(&self, input: RuntimeInput) -> Result<Option<String>> {
151 self.publish(json!({"type":"user_message", "text":input.text}));
152 self.publish(json!({"type":"turn_started"}));
153 let (reply, response) = oneshot::channel();
154 if self
155 .commands
156 .send(HostCommand::Submit { input, reply })
157 .await
158 .is_err()
159 {
160 self.busy.store(false, Ordering::SeqCst);
161 self.publish(
162 json!({"type":"turn_failed", "message":"Hosted harness runtime is closed."}),
163 );
164 self.publish(json!({"type":"turn_completed"}));
165 self.mark_closed("Harness runtime command channel closed.");
166 return Err(Error::Other("hosted harness runtime is closed".into()));
167 }
168 match response.await {
169 Ok(Ok(turn)) => Ok(turn),
170 Ok(Err(error)) => {
171 self.busy.store(false, Ordering::SeqCst);
172 self.publish(json!({"type":"turn_failed", "message":error.to_string()}));
173 self.publish(json!({"type":"turn_completed"}));
174 Err(error)
175 }
176 Err(_) => {
177 self.busy.store(false, Ordering::SeqCst);
178 self.publish(json!({"type":"turn_failed", "message":"Hosted harness runtime stopped before accepting input."}));
179 self.publish(json!({"type":"turn_completed"}));
180 self.mark_closed("Harness runtime stopped before accepting input.");
181 Err(Error::Other(
182 "hosted harness runtime stopped before accepting input".into(),
183 ))
184 }
185 }
186 }
187
188 async fn submit_native(&self, text: String) -> Result<Option<String>> {
189 self.claim_submit()?;
190 self.submit_native_claimed(RuntimeInput {
191 text,
192 image_urls: Vec::new(),
193 })
194 .await
195 }
196
197 async fn interrupt_native(&self) -> Result<()> {
198 let (reply, response) = oneshot::channel();
199 self.commands
200 .send(HostCommand::Interrupt { reply })
201 .await
202 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
203 response
204 .await
205 .map_err(|_| Error::Other("hosted harness runtime stopped before interrupt".into()))?
206 }
207
208 async fn steer_native(&self, text: String) -> Result<()> {
209 let (reply, response) = oneshot::channel();
210 self.commands
211 .send(HostCommand::Steer { text, reply })
212 .await
213 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
214 response
215 .await
216 .map_err(|_| Error::Other("hosted harness runtime stopped before steering".into()))?
217 }
218
219 async fn respond_native(&self, request_id: Value, response: Value) -> Result<()> {
220 self.respond_native_as(request_id, response, None).await
221 }
222
223 async fn respond_native_as(
227 &self,
228 request_id: Value,
229 response: Value,
230 canonical: Option<Value>,
231 ) -> Result<()> {
232 let (reply, completed) = oneshot::channel();
233 self.commands
234 .send(HostCommand::Respond {
235 request_id: request_id.clone(),
236 response: response.clone(),
237 reply,
238 })
239 .await
240 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
241 completed
242 .await
243 .map_err(|_| Error::Other("hosted harness runtime stopped before response".into()))??;
244 let resolved = {
248 let mut pending = self
249 .pending_requests
250 .lock()
251 .unwrap_or_else(std::sync::PoisonError::into_inner);
252 let found = pending
253 .iter()
254 .find(|(_, native)| *native == &request_id)
255 .map(|(id, _)| *id);
256 found.and_then(|id| pending.remove(&id).map(|_| id))
257 };
258 if let Some(id) = resolved {
259 self.publish(json!({
260 "type": "request_resolved",
261 "request_id": id,
262 "response": canonical.unwrap_or(response),
263 }));
264 }
265 Ok(())
266 }
267
268 fn native_request_id(&self, request_id: u64) -> Option<Value> {
270 self.pending_requests
271 .lock()
272 .unwrap_or_else(std::sync::PoisonError::into_inner)
273 .get(&request_id)
274 .cloned()
275 }
276
277 fn publish(&self, payload: Value) {
278 let event = {
279 let mut projection = self
280 .projection
281 .lock()
282 .unwrap_or_else(std::sync::PoisonError::into_inner);
283 let event = FrontendEvent::new(projection.next_sequence, payload);
284 projection.next_sequence = projection.next_sequence.saturating_add(1);
285 projection.replay.push_back(event.clone());
286 while projection.replay.len() > FRONTEND_REPLAY_CAPACITY {
287 projection.replay.pop_front();
288 }
289 event
290 };
291 let _ = self.frontend_events.send(event);
292 }
293
294 fn accept_native_event(&self, event: HarnessEvent) {
295 let _ = self.raw_events.send(event.clone());
296 for payload in project_native_event(self.handle.harness.as_str(), &event) {
297 let terminal = matches!(
298 payload.get("type").and_then(Value::as_str),
299 Some("turn_succeeded" | "turn_interrupted" | "turn_failed")
300 );
301 if terminal {
302 self.busy.store(false, Ordering::SeqCst);
303 }
304 if payload.get("type").and_then(Value::as_str) == Some("request") {
305 if let (Some(id), Some(native)) = (
306 payload.pointer("/request/id").and_then(Value::as_u64),
307 payload
308 .pointer("/request/payload/native_request_id")
309 .cloned(),
310 ) {
311 self.pending_requests
312 .lock()
313 .unwrap_or_else(std::sync::PoisonError::into_inner)
314 .insert(id, native);
315 }
316 }
317 self.publish(payload);
318 }
319 }
320
321 fn mark_closed(&self, message: impl Into<String>) {
322 if self.closed.swap(true, Ordering::SeqCst) {
323 return;
324 }
325 let message = message.into();
326 self.busy.store(false, Ordering::SeqCst);
327 let _ = self.raw_events.send(HarnessEvent {
331 sequence: None,
332 kind: "transport_closed".into(),
333 payload: json!({"message":message, "terminal":true}),
334 });
335 self.publish(json!({"type":"runtime_disconnected", "message":message}));
336 }
337
338 fn descriptor(&self) -> FrontendRuntimeDescriptor {
339 FrontendRuntimeDescriptor {
340 schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
341 session_id: self.handle.runtime_id.clone(),
342 source_harness: Some(self.handle.harness.as_str().to_string()),
343 emulation_profile: None,
344 active_modules: Vec::new(),
345 commands: Vec::new(),
346 operations: Vec::new(),
347 actions: FrontendActions {
348 submit: self.capabilities.send_input,
349 interrupt: self.capabilities.interrupt,
350 steer: self.capabilities.steer,
351 respond: !hosted_answerable_decisions(self.handle.harness.as_str()).is_empty(),
371 detach: true,
372 close: false,
373 },
374 display: FrontendDisplayCapabilities {
375 event_kinds: vec![
376 "user_message".into(),
377 "turn_started".into(),
378 "turn_succeeded".into(),
379 "turn_interrupted".into(),
380 "turn_failed".into(),
381 "text_delta".into(),
382 "reasoning".into(),
383 "tool_call_started".into(),
384 "tool_call_completed".into(),
385 "request".into(),
388 "native_event".into(),
389 "runtime_disconnected".into(),
390 ],
391 opaque_fallback: true,
392 },
393 model: self.handle.harness.as_str().to_string(),
394 turn_state: if self.busy.load(Ordering::SeqCst) {
395 FrontendTurnState::Busy
396 } else {
397 FrontendTurnState::Idle
398 },
399 connection_state: if self.closed.load(Ordering::SeqCst) {
400 FrontendConnectionState::ShuttingDown
401 } else {
402 FrontendConnectionState::Connected
403 },
404 extensions: Default::default(),
405 }
406 }
407}
408
409#[async_trait]
410impl FrontendRuntime for HostedHarnessRuntime {
411 async fn describe(
412 &self,
413 ) -> std::result::Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
414 Ok(self.descriptor())
415 }
416
417 async fn attach(
418 &self,
419 _history_limit: usize,
420 ) -> std::result::Result<FrontendAttachment, FrontendRuntimeError> {
421 let live = self.frontend_events.subscribe();
422 let projection = self
423 .projection
424 .lock()
425 .unwrap_or_else(std::sync::PoisonError::into_inner);
426 let replay = projection.replay.clone();
427 if let Some(first) = replay.front() {
428 if first.sequence > 1 {
429 return Err(FrontendRuntimeError::ReplayGap(first.sequence - 1));
430 }
431 }
432 Ok(FrontendAttachment::new(
433 self.descriptor(),
434 self.history.clone(),
435 0,
436 replay,
437 live,
438 None,
439 ))
440 }
441
442 async fn send_input(
443 self: Arc<Self>,
444 prompt: String,
445 ) -> std::result::Result<(), FrontendRuntimeError> {
446 self.claim_submit().map_err(hosted_submit_error)?;
447 tokio::spawn(async move {
448 let _ = self
449 .submit_native_claimed(RuntimeInput {
450 text: prompt,
451 image_urls: Vec::new(),
452 })
453 .await;
454 });
455 Ok(())
456 }
457
458 async fn send_input_with_images(
459 self: Arc<Self>,
460 prompt: String,
461 image_urls: Vec<String>,
462 ) -> std::result::Result<(), FrontendRuntimeError> {
463 self.claim_submit().map_err(hosted_submit_error)?;
464 tokio::spawn(async move {
465 let _ = self
466 .submit_native_claimed(RuntimeInput {
467 text: prompt,
468 image_urls,
469 })
470 .await;
471 });
472 Ok(())
473 }
474
475 async fn submit(&self, prompt: String) -> std::result::Result<String, FrontendRuntimeError> {
476 self.submit_native(prompt)
477 .await
478 .map(|turn| turn.unwrap_or_default())
479 .map_err(hosted_submit_error)
480 }
481
482 async fn interrupt(&self) -> std::result::Result<bool, FrontendRuntimeError> {
483 if !self.busy.load(Ordering::SeqCst) {
484 return Ok(false);
485 }
486 self.interrupt_native()
487 .await
488 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
489 Ok(true)
493 }
494
495 async fn steer(&self, prompt: String) -> std::result::Result<(), FrontendRuntimeError> {
496 if !self.busy.load(Ordering::SeqCst) || !self.capabilities.steer {
497 return Err(FrontendRuntimeError::UnsupportedAction("steer"));
498 }
499 self.steer_native(prompt)
500 .await
501 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
502 }
503
504 async fn respond(
505 &self,
506 response: FrontendResponse,
507 ) -> std::result::Result<(), FrontendRuntimeError> {
508 let harness = self.handle.harness.as_str();
509 let FrontendResponse::Approval {
510 request_id,
511 decision,
512 } = &response
513 else {
514 return Err(FrontendRuntimeError::UnsupportedAction(
515 "respond: a hosted harness answers approval requests only",
516 ));
517 };
518 let Some(native_id) = self.native_request_id(*request_id) else {
519 return Err(FrontendRuntimeError::UnknownRequest(*request_id));
520 };
521 let body = hosted_permission_reply(harness, *decision)
522 .map_err(FrontendRuntimeError::InvalidResponse)?;
523 let canonical = serde_json::to_value(&response).unwrap_or(Value::Null);
524 self.respond_native_as(native_id, body, Some(canonical))
525 .await
526 .map_err(|error| FrontendRuntimeError::Execution {
527 operation: crate::SdkOperation::Respond,
528 message: error.to_string(),
529 })
530 }
531}
532
533fn hosted_submit_error(error: Error) -> FrontendRuntimeError {
534 if error.to_string().contains("already in progress") {
535 FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)
536 } else {
537 FrontendRuntimeError::Transport(error.to_string())
538 }
539}
540
541pub struct HostedHarnessConnection {
544 host: Arc<HostedHarnessRuntime>,
545 handle: RuntimeHandle,
546 events: broadcast::Receiver<HarnessEvent>,
547 closed: bool,
548}
549
550#[async_trait]
551impl RuntimeConnection for HostedHarnessConnection {
552 fn handle(&self) -> &RuntimeHandle {
553 &self.handle
554 }
555
556 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
557 self.host.claim_submit()?;
558 self.host.submit_native_claimed(input).await
559 }
560
561 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
562 match self.events.recv().await {
563 Ok(event) => Ok(Some(event)),
564 Err(broadcast::error::RecvError::Lagged(count)) => Err(Error::Other(format!(
565 "hosted harness event stream lost {count} event(s)"
566 ))),
567 Err(broadcast::error::RecvError::Closed) => Ok(None),
568 }
569 }
570
571 async fn interrupt(&mut self) -> Result<()> {
572 self.host.interrupt_native().await
573 }
574
575 async fn steer(&mut self, text: String) -> Result<()> {
576 self.host.steer_native(text).await
577 }
578
579 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
580 self.host.respond_native(request_id, response).await
581 }
582
583 async fn close(&mut self) -> Result<()> {
584 if self.closed {
585 return Ok(());
586 }
587 self.host.shutdown().await?;
588 self.closed = true;
589 Ok(())
590 }
591}
592
593async fn run_native_runtime(
594 mut runtime: Box<dyn RuntimeConnection>,
595 host: std::sync::Weak<HostedHarnessRuntime>,
596 mut commands: mpsc::Receiver<HostCommand>,
597) {
598 loop {
599 tokio::select! {
600 command = commands.recv() => {
601 let Some(command) = command else {
602 let _ = runtime.close().await;
603 return;
604 };
605 match command {
606 HostCommand::Submit { input, reply } => {
607 let _ = reply.send(runtime.send_input(input).await);
608 }
609 HostCommand::Interrupt { reply } => {
610 let _ = reply.send(runtime.interrupt().await);
611 }
612 HostCommand::Steer { text, reply } => {
613 let _ = reply.send(runtime.steer(text).await);
614 }
615 HostCommand::Respond { request_id, response, reply } => {
616 let _ = reply.send(runtime.respond(request_id, response).await);
617 }
618 HostCommand::Shutdown { reply } => {
619 let result = runtime.close().await;
620 let closed = result.is_ok();
621 let _ = reply.send(result);
622 if closed {
623 if let Some(host) = host.upgrade() {
624 host.mark_closed("Harness runtime closed.");
625 }
626 return;
627 }
628 }
629 }
630 }
631 event = runtime.next_event() => {
632 let Some(host) = host.upgrade() else {
633 let _ = runtime.close().await;
634 return;
635 };
636 match event {
637 Ok(Some(event)) => host.accept_native_event(event),
638 Ok(None) => {
639 host.mark_closed("Harness runtime transport closed.");
640 return;
641 }
642 Err(error) => {
643 host.mark_closed(error.to_string());
644 return;
645 }
646 }
647 }
648 }
649 }
650}
651
652pub(crate) fn hosted_answerable_decisions(harness: &str) -> &'static [&'static str] {
661 match harness {
662 "claude-code" => &["allow", "deny"],
672 _ => &[],
673 }
674}
675
676fn hosted_permission_reply(
678 harness: &str,
679 decision: crate::FrontendApprovalDecision,
680) -> std::result::Result<Value, String> {
681 use crate::FrontendApprovalDecision as Decision;
682 let accepted = hosted_answerable_decisions(harness);
683 let wire = match decision {
684 Decision::Allow => "allow",
685 Decision::AllowForSession => "allow_for_session",
686 Decision::Deny => "deny",
687 };
688 if !accepted.contains(&wire) {
689 return Err(format!(
690 "`{wire}` is not a decision a hosted {harness} runtime can carry; this harness accepts \
691 {}",
692 if accepted.is_empty() {
693 "no portable decision — its requests are observe-only".to_string()
694 } else {
695 accepted
696 .iter()
697 .map(|decision| format!("`{decision}`"))
698 .collect::<Vec<_>>()
699 .join(" or ")
700 },
701 ));
702 }
703 match harness {
704 "claude-code" => Ok(match decision {
705 Decision::Allow => json!({"behavior": "allow"}),
706 Decision::Deny => json!({
708 "behavior": "deny",
709 "message": "denied through the Volter Harness frontend",
710 }),
711 Decision::AllowForSession => unreachable!("refused above"),
712 }),
713 _ => Err(format!(
714 "no hosted permission reply is defined for {harness}"
715 )),
716 }
717}
718
719fn recorded_history(handle: &RuntimeHandle) -> Vec<ChatMessage> {
724 let catalog = HarnessCatalog::new();
725 let query = DiscoveryQuery {
726 harnesses: vec![handle.harness.clone()],
727 query: Some(handle.runtime_id.clone()),
728 ..DiscoveryQuery::default()
729 };
730 let Some(descriptor) = catalog.discover(&query).ok().and_then(|found| {
731 found
732 .into_iter()
733 .find(|descriptor| descriptor.locator.session_id == handle.runtime_id)
734 }) else {
735 return Vec::new();
736 };
737 catalog
738 .load_display_view(&descriptor.locator, Fidelity::Semantic, 500)
739 .map(|session| session.messages)
740 .unwrap_or_default()
741}
742
743fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
744 let key = event.kind.to_ascii_lowercase().replace('-', "_");
745 let payload = &event.payload;
746 if key == "transport_closed" {
747 return vec![
748 json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
749 ];
750 }
751 if key == "transport_error" || key == "error" {
752 return vec![
753 json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
754 ];
755 }
756 if key == "session/update" {
757 let update = payload
758 .pointer("/params/update")
759 .or_else(|| payload.get("update"))
760 .unwrap_or(payload);
761 let update_kind = update
762 .get("sessionUpdate")
763 .or_else(|| update.get("type"))
764 .and_then(Value::as_str)
765 .unwrap_or_default()
766 .to_ascii_lowercase();
767 if update_kind == "agent_message_chunk" {
768 return vec![
769 json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
770 ];
771 }
772 if update_kind == "agent_thought_chunk" {
773 return vec![
774 json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
775 ];
776 }
777 if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
778 return vec![project_tool(update, payload)];
779 }
780 }
781 if key == "supercode/acp_request_completed" {
782 let failure = payload
783 .pointer("/params/error")
784 .or_else(|| payload.get("error"));
785 return completion(failure.and_then(extract_text));
786 }
787 match harness {
788 "codex" => project_codex(&key, payload),
789 "claude-code" => project_claude(&key, payload),
790 "pi" => project_pi(&key, payload),
791 "opencode" => project_opencode(&key, payload),
792 _ => project_generic(&key, payload),
793 }
794}
795
796fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
797 if key == "turn/started" {
798 return vec![native_payload(key, payload)];
799 }
800 if key == "turn/completed" {
801 let status = payload
802 .pointer("/params/turn/status")
803 .or_else(|| payload.pointer("/turn/status"))
804 .and_then(Value::as_str)
805 .unwrap_or("completed")
806 .to_ascii_lowercase();
807 return completion(
808 (status.contains("fail") || status.contains("error") || status.contains("cancel"))
809 .then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
810 );
811 }
812 if key.ends_with("/delta") {
813 let text = payload
814 .pointer("/params/delta")
815 .or_else(|| payload.get("delta"))
816 .and_then(extract_text)
817 .or_else(|| extract_text(payload))
818 .unwrap_or_default();
819 return vec![
820 json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
821 ];
822 }
823 if key.contains("commandexecution")
824 || key.contains("mcptool")
825 || key.contains("filechange")
826 || key.contains("tool")
827 {
828 let source = payload
829 .pointer("/params/item")
830 .or_else(|| payload.get("item"))
831 .unwrap_or(payload);
832 return vec![project_tool(source, payload)];
833 }
834 vec![native_payload(key, payload)]
835}
836
837fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
849 if key == "result" {
850 let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
851 || payload
852 .get("subtype")
853 .and_then(Value::as_str)
854 .is_some_and(|subtype| subtype != "success");
855 return completion(
856 failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
857 );
858 }
859 if key == "stream_event" {
860 let stream = payload
865 .get("event")
866 .or_else(|| payload.get("stream_event"))
867 .unwrap_or(payload);
868 if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
869 return vec![native_payload(key, payload)];
870 }
871 }
872 if key == "assistant" {
873 return project_claude_blocks(payload, true);
874 }
875 if key == "user" {
876 return project_claude_blocks(payload, false);
877 }
878 if key == "control_request" {
879 return project_claude_control_request(payload);
880 }
881 if key == "system" && claude_system_is_telemetry(payload) {
882 return Vec::new();
883 }
884 if key == "tool_progress" {
885 return Vec::new();
889 }
890 if key == "rate_limit_event" {
891 let status = payload
896 .pointer("/rate_limit_info/status")
897 .and_then(Value::as_str)
898 .unwrap_or_default();
899 if status.starts_with("allowed") {
900 return Vec::new();
901 }
902 }
903 project_generic(key, payload)
904}
905
906fn claude_system_is_telemetry(payload: &Value) -> bool {
926 match payload
927 .get("subtype")
928 .and_then(Value::as_str)
929 .unwrap_or_default()
930 {
931 "init" | "thinking_tokens" | "task_started" | "task_progress" | "task_updated"
932 | "hook_started" | "hook_progress" | "vcs_state_changed" => true,
933 "task_notification" => payload
934 .get("output_file")
935 .and_then(Value::as_str)
936 .unwrap_or_default()
937 .is_empty(),
938 "hook_response" => payload.get("exit_code").and_then(Value::as_i64) == Some(0),
939 _ => false,
940 }
941}
942
943fn claude_blocks(payload: &Value) -> Option<&Vec<Value>> {
945 payload
946 .pointer("/message/content")
947 .or_else(|| payload.get("content"))
948 .and_then(Value::as_array)
949}
950
951fn project_claude_blocks(payload: &Value, assistant: bool) -> Vec<Value> {
952 let Some(blocks) = claude_blocks(payload) else {
953 let text = extract_text(payload).unwrap_or_default();
956 if text.is_empty() {
957 return vec![native_payload(
958 if assistant { "assistant" } else { "user" },
959 payload,
960 )];
961 }
962 return vec![
963 json!({"type":if assistant { "text_delta" } else { "user_message" }, "text":text, "raw":payload}),
964 ];
965 };
966 let mut projected = Vec::new();
967 for block in blocks {
968 match block
969 .get("type")
970 .and_then(Value::as_str)
971 .unwrap_or_default()
972 {
973 "text" => {
974 let text = block
975 .get("text")
976 .and_then(Value::as_str)
977 .unwrap_or_default();
978 if !text.is_empty() {
979 projected.push(json!({
980 "type": if assistant { "text_delta" } else { "user_message" },
981 "text": text,
982 "raw": block,
983 }));
984 }
985 }
986 "thinking" | "redacted_thinking" => {
987 let text = extract_text(block).unwrap_or_default();
988 if !text.is_empty() {
989 projected.push(json!({"type":"reasoning", "text":text, "raw":block}));
990 }
991 }
992 "tool_use" => projected.push(json!({
993 "type": "tool_call_started",
994 "id": block.get("id").cloned().unwrap_or(Value::Null),
995 "name": block.get("name").cloned().unwrap_or(Value::Null),
996 "arguments": block.get("input").map(Value::to_string).unwrap_or_default(),
997 "raw": block,
998 })),
999 "tool_result" => projected.push(json!({
1000 "type": "tool_call_completed",
1004 "id": block.get("tool_use_id").cloned().unwrap_or(Value::Null),
1005 "name": Value::Null,
1006 "output": extract_text(block.get("content").unwrap_or(block)).unwrap_or_default(),
1007 "is_error": block.get("is_error").and_then(Value::as_bool).unwrap_or(false),
1008 "raw": block,
1009 })),
1010 _ => projected.push(native_payload(
1011 if assistant { "assistant" } else { "user" },
1012 block,
1013 )),
1014 }
1015 }
1016 projected
1017}
1018
1019fn project_claude_control_request(payload: &Value) -> Vec<Value> {
1028 let request = payload.get("request").unwrap_or(payload);
1029 if request.get("subtype").and_then(Value::as_str) != Some("can_use_tool") {
1030 return vec![native_payload("control_request", payload)];
1031 }
1032 let native_id = payload
1033 .get("request_id")
1034 .and_then(Value::as_str)
1035 .unwrap_or_default();
1036 vec![json!({
1037 "type": "request",
1038 "request": {
1039 "id": claude_request_id(native_id),
1040 "kind": "approval",
1041 "payload": {
1042 "tool": request.get("tool_name").or_else(|| request.get("toolName")).cloned().unwrap_or(Value::Null),
1043 "arguments": request.get("input").cloned().unwrap_or(Value::Null),
1044 "native_request_id": native_id,
1045 "decisions": hosted_answerable_decisions("claude-code"),
1049 },
1050 },
1051 })]
1052}
1053
1054fn claude_request_id(native_id: &str) -> u64 {
1061 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1062 for byte in native_id.as_bytes() {
1063 hash ^= u64::from(*byte);
1064 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1065 }
1066 hash & ((1_u64 << 53) - 1)
1068}
1069
1070fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
1071 match key {
1072 "agent_start" => vec![native_payload(key, payload)],
1073 "agent_end" => completion(payload.get("error").and_then(extract_text)),
1074 "message_update" => {
1075 let update = payload
1080 .get("assistantMessageEvent")
1081 .or_else(|| payload.get("event"))
1082 .unwrap_or(payload);
1083 let kind = update
1084 .get("type")
1085 .and_then(Value::as_str)
1086 .unwrap_or_default();
1087 match kind {
1088 "text_delta" => vec![
1089 json!({"type":"text_delta", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1090 ],
1091 "thinking_delta" => vec![
1092 json!({"type":"reasoning", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1093 ],
1094 _ => vec![native_payload(key, payload)],
1095 }
1096 }
1097 _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
1098 _ => vec![native_payload(key, payload)],
1099 }
1100}
1101
1102fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
1103 if key == "session.idle" {
1104 return completion(None);
1105 }
1106 if key == "session.status" {
1107 let status = payload
1108 .pointer("/properties/status/type")
1109 .or_else(|| payload.pointer("/status/type"))
1110 .and_then(Value::as_str)
1111 .unwrap_or_default();
1112 if status == "busy" {
1113 return vec![native_payload(key, payload)];
1114 }
1115 if status == "idle" {
1116 return completion(None);
1117 }
1118 }
1119 if key == "message.part.updated" {
1120 let part = payload
1121 .pointer("/properties/part")
1122 .or_else(|| payload.get("part"))
1123 .unwrap_or(payload);
1124 if let Some(delta) = payload
1125 .pointer("/properties/delta")
1126 .or_else(|| payload.get("delta"))
1127 .and_then(Value::as_str)
1128 {
1129 let reasoning = part
1130 .get("type")
1131 .and_then(Value::as_str)
1132 .is_some_and(|kind| kind.contains("reasoning"));
1133 return vec![
1134 json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1135 ];
1136 }
1137 if part
1138 .get("type")
1139 .and_then(Value::as_str)
1140 .is_some_and(|kind| kind.contains("tool"))
1141 {
1142 return vec![project_tool(part, payload)];
1143 }
1144 }
1145 if key == "session.error" {
1146 return completion(Some(
1147 extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
1148 ));
1149 }
1150 vec![native_payload(key, payload)]
1151}
1152
1153fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
1154 match key {
1155 "turn_started" | "turn/started" | "agent_start" => {
1156 vec![native_payload(key, payload)]
1157 }
1158 "turn_completed" | "turn/completed" | "agent_end" => {
1159 completion(payload.get("error").and_then(extract_text))
1160 }
1161 "output_delta" | "content_delta" => vec![
1162 json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1163 ],
1164 "reasoning_delta" => vec![
1165 json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1166 ],
1167 "tool" => vec![project_tool(payload, payload)],
1168 _ => vec![native_payload(key, payload)],
1169 }
1170}
1171
1172fn completion(error: Option<String>) -> Vec<Value> {
1173 match error {
1174 Some(message) => vec![
1175 json!({"type":"turn_failed", "message":message}),
1176 json!({"type":"turn_completed"}),
1177 ],
1178 None => vec![
1179 json!({"type":"turn_succeeded"}),
1180 json!({"type":"turn_completed"}),
1181 ],
1182 }
1183}
1184
1185fn project_tool(source: &Value, raw: &Value) -> Value {
1186 let status = source
1187 .get("status")
1188 .or_else(|| source.get("state"))
1189 .or_else(|| source.get("sessionUpdate"))
1190 .and_then(Value::as_str)
1191 .unwrap_or_default()
1192 .to_ascii_lowercase();
1193 let completed = status.contains("complete")
1194 || status.contains("result")
1195 || status.contains("success")
1196 || status.contains("error")
1197 || status.contains("fail");
1198 let arguments = source
1199 .get("arguments")
1200 .or_else(|| source.get("input"))
1201 .or_else(|| source.get("rawInput"))
1202 .cloned()
1203 .unwrap_or(Value::Null);
1204 json!({
1205 "type": if completed { "tool_call_completed" } else { "tool_call_started" },
1206 "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),
1207 "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
1208 "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
1209 "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
1210 "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
1211 "raw": raw,
1212 })
1213}
1214
1215fn native_payload(kind: &str, payload: &Value) -> Value {
1216 json!({"type":"native_event", "kind":kind, "raw":payload})
1217}
1218
1219fn extract_text(value: &Value) -> Option<String> {
1220 match value {
1221 Value::String(text) => Some(text.clone()),
1222 Value::Array(values) => {
1223 let text = values
1224 .iter()
1225 .filter_map(extract_text)
1226 .collect::<Vec<_>>()
1227 .join("\n");
1228 (!text.is_empty()).then_some(text)
1229 }
1230 Value::Object(object) => {
1231 for key in ["text", "delta", "content", "message", "result", "error"] {
1232 if let Some(text) = object.get(key).and_then(Value::as_str) {
1233 return Some(text.to_string());
1234 }
1235 }
1236 for key in [
1237 "delta",
1238 "content",
1239 "message",
1240 "error",
1241 "data",
1242 "part",
1243 "params",
1244 "properties",
1245 "update",
1246 "event",
1247 ] {
1248 if let Some(text) = object.get(key).and_then(extract_text) {
1249 return Some(text);
1250 }
1251 }
1252 None
1253 }
1254 _ => None,
1255 }
1256}