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