1use 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
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 native: StdMutex<NativeProjection>,
74 busy: AtomicBool,
75 closed: AtomicBool,
76 registration: StdMutex<Option<crate::LiveRuntimeRegistration>>,
77}
78
79impl HostedHarnessRuntime {
80 pub fn spawn(
83 runtime: Box<dyn RuntimeConnection>,
84 capabilities: RuntimeCapabilities,
85 fresh: bool,
86 ) -> (Arc<Self>, HostedHarnessConnection) {
87 let handle = runtime.handle().clone();
88 let (commands, command_rx) = mpsc::channel(32);
89 let (raw_events, raw_rx) = broadcast::channel(1024);
90 let (frontend_events, _) = broadcast::channel(1024);
91 let host = Arc::new(Self {
92 handle: handle.clone(),
93 capabilities,
94 commands,
95 raw_events,
96 frontend_events,
97 projection: StdMutex::new(ProjectionState {
98 next_sequence: 1,
99 replay: VecDeque::new(),
100 }),
101 history: if fresh {
102 Vec::new()
103 } else {
104 recorded_history(&handle)
105 },
106 pending_requests: StdMutex::new(HashMap::new()),
107 native: StdMutex::new(NativeProjection::default()),
108 busy: AtomicBool::new(false),
109 closed: AtomicBool::new(false),
110 registration: StdMutex::new(None),
111 });
112 tokio::spawn(run_native_runtime(
113 runtime,
114 Arc::downgrade(&host),
115 command_rx,
116 ));
117 let connection = HostedHarnessConnection {
118 host: host.clone(),
119 handle,
120 events: raw_rx,
121 closed: false,
122 };
123 (host, connection)
124 }
125
126 pub(crate) fn retain_registration(&self, registration: crate::LiveRuntimeRegistration) {
128 let mut retained = self
129 .registration
130 .lock()
131 .unwrap_or_else(std::sync::PoisonError::into_inner);
132 if !self.closed.load(Ordering::SeqCst) {
133 *retained = Some(registration);
134 }
135 }
136
137 pub fn frontend_sender(&self) -> broadcast::Sender<FrontendEvent> {
139 self.frontend_events.clone()
140 }
141
142 pub async fn shutdown(&self) -> Result<()> {
144 if self.closed.load(Ordering::SeqCst) {
145 return Ok(());
146 }
147 let (reply, response) = oneshot::channel();
148 self.commands
149 .send(HostCommand::Shutdown { reply })
150 .await
151 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
152 response
153 .await
154 .map_err(|_| Error::Other("hosted harness runtime stopped before shutdown".into()))?
155 }
156
157 fn claim_submit(&self) -> Result<()> {
158 if self.closed.load(Ordering::SeqCst) {
159 return Err(Error::Other("hosted harness runtime is closed".into()));
160 }
161 if self
162 .busy
163 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
164 .is_err()
165 {
166 return Err(Error::Other("a harness turn is already in progress".into()));
167 }
168 Ok(())
169 }
170
171 async fn submit_native_claimed(&self, input: RuntimeInput) -> Result<Option<String>> {
172 self.publish(json!({"type":"user_message", "text":input.text}));
173 self.publish(json!({"type":"turn_started"}));
174 let (reply, response) = oneshot::channel();
175 if self
176 .commands
177 .send(HostCommand::Submit { input, reply })
178 .await
179 .is_err()
180 {
181 self.busy.store(false, Ordering::SeqCst);
182 self.publish(
183 json!({"type":"turn_failed", "message":"Hosted harness runtime is closed."}),
184 );
185 self.publish(json!({"type":"turn_completed"}));
186 self.mark_closed("Harness runtime command channel closed.");
187 return Err(Error::Other("hosted harness runtime is closed".into()));
188 }
189 match response.await {
190 Ok(Ok(turn)) => Ok(turn),
191 Ok(Err(error)) => {
192 self.busy.store(false, Ordering::SeqCst);
193 self.publish(json!({"type":"turn_failed", "message":error.to_string()}));
194 self.publish(json!({"type":"turn_completed"}));
195 Err(error)
196 }
197 Err(_) => {
198 self.busy.store(false, Ordering::SeqCst);
199 self.publish(json!({"type":"turn_failed", "message":"Hosted harness runtime stopped before accepting input."}));
200 self.publish(json!({"type":"turn_completed"}));
201 self.mark_closed("Harness runtime stopped before accepting input.");
202 Err(Error::Other(
203 "hosted harness runtime stopped before accepting input".into(),
204 ))
205 }
206 }
207 }
208
209 async fn submit_native(&self, text: String) -> Result<Option<String>> {
210 self.claim_submit()?;
211 self.submit_native_claimed(RuntimeInput {
212 text,
213 image_urls: Vec::new(),
214 })
215 .await
216 }
217
218 async fn interrupt_native(&self) -> Result<()> {
219 let (reply, response) = oneshot::channel();
220 self.commands
221 .send(HostCommand::Interrupt { reply })
222 .await
223 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
224 response
225 .await
226 .map_err(|_| Error::Other("hosted harness runtime stopped before interrupt".into()))?
227 }
228
229 async fn steer_native(&self, text: String) -> Result<()> {
230 let (reply, response) = oneshot::channel();
231 self.commands
232 .send(HostCommand::Steer { text, reply })
233 .await
234 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
235 response
236 .await
237 .map_err(|_| Error::Other("hosted harness runtime stopped before steering".into()))?
238 }
239
240 async fn respond_native(&self, request_id: Value, response: Value) -> Result<()> {
241 self.respond_native_as(request_id, response, None).await
242 }
243
244 async fn respond_native_as(
248 &self,
249 request_id: Value,
250 response: Value,
251 canonical: Option<Value>,
252 ) -> Result<()> {
253 let (reply, completed) = oneshot::channel();
254 self.commands
255 .send(HostCommand::Respond {
256 request_id: request_id.clone(),
257 response: response.clone(),
258 reply,
259 })
260 .await
261 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
262 completed
263 .await
264 .map_err(|_| Error::Other("hosted harness runtime stopped before response".into()))??;
265 let resolved = {
269 let mut pending = self
270 .pending_requests
271 .lock()
272 .unwrap_or_else(std::sync::PoisonError::into_inner);
273 let found = pending
274 .iter()
275 .find(|(_, native)| *native == &request_id)
276 .map(|(id, _)| *id);
277 found.and_then(|id| pending.remove(&id).map(|_| id))
278 };
279 if let Some(id) = resolved {
280 self.publish(json!({
281 "type": "request_resolved",
282 "request_id": id,
283 "response": canonical.unwrap_or(response),
284 }));
285 }
286 Ok(())
287 }
288
289 fn native_request_id(&self, request_id: u64) -> Option<Value> {
291 self.pending_requests
292 .lock()
293 .unwrap_or_else(std::sync::PoisonError::into_inner)
294 .get(&request_id)
295 .cloned()
296 }
297
298 fn publish(&self, payload: Value) {
299 let event = {
300 let mut projection = self
301 .projection
302 .lock()
303 .unwrap_or_else(std::sync::PoisonError::into_inner);
304 let event = FrontendEvent::new(projection.next_sequence, payload);
305 projection.next_sequence = projection.next_sequence.saturating_add(1);
306 projection.replay.push_back(event.clone());
307 while projection.replay.len() > FRONTEND_REPLAY_CAPACITY {
308 projection.replay.pop_front();
309 }
310 event
311 };
312 let _ = self.frontend_events.send(event);
313 }
314
315 fn accept_native_event(&self, event: HarnessEvent) {
316 let _ = self.raw_events.send(event.clone());
317 let projected = self
318 .native
319 .lock()
320 .unwrap_or_else(std::sync::PoisonError::into_inner)
321 .project(self.handle.harness.as_str(), &event);
322 for payload in projected {
323 let terminal = matches!(
324 payload.get("type").and_then(Value::as_str),
325 Some("turn_succeeded" | "turn_interrupted" | "turn_failed")
326 );
327 if terminal {
328 self.busy.store(false, Ordering::SeqCst);
329 }
330 if payload.get("type").and_then(Value::as_str) == Some("request") {
331 if let (Some(id), Some(native)) = (
332 payload.pointer("/request/id").and_then(Value::as_u64),
333 payload
334 .pointer("/request/payload/native_request_id")
335 .cloned(),
336 ) {
337 self.pending_requests
338 .lock()
339 .unwrap_or_else(std::sync::PoisonError::into_inner)
340 .insert(id, native);
341 }
342 }
343 self.publish(payload);
344 }
345 }
346
347 fn mark_closed(&self, message: impl Into<String>) {
348 if self.closed.swap(true, Ordering::SeqCst) {
349 return;
350 }
351 self.registration
352 .lock()
353 .unwrap_or_else(std::sync::PoisonError::into_inner)
354 .take();
355 let message = message.into();
356 self.busy.store(false, Ordering::SeqCst);
357 let _ = self.raw_events.send(HarnessEvent {
361 sequence: None,
362 kind: "transport_closed".into(),
363 payload: json!({"message":message, "terminal":true}),
364 });
365 self.publish(json!({"type":"runtime_disconnected", "message":message}));
366 }
367
368 fn descriptor(&self) -> FrontendRuntimeDescriptor {
369 FrontendRuntimeDescriptor {
370 schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
371 session_id: self.handle.runtime_id.clone(),
372 source_harness: Some(self.handle.harness.as_str().to_string()),
373 emulation_profile: None,
374 active_modules: Vec::new(),
375 commands: Vec::new(),
376 operations: Vec::new(),
377 actions: FrontendActions {
378 submit: self.capabilities.send_input,
379 interrupt: self.capabilities.interrupt,
380 steer: self.capabilities.steer,
381 respond: !hosted_answerable_decisions(self.handle.harness.as_str()).is_empty(),
401 detach: true,
402 close: false,
403 },
404 display: FrontendDisplayCapabilities {
405 event_kinds: vec![
406 "user_message".into(),
407 "turn_started".into(),
408 "turn_succeeded".into(),
409 "turn_interrupted".into(),
410 "turn_failed".into(),
411 "text_delta".into(),
412 "reasoning".into(),
413 "tool_call_started".into(),
414 "tool_call_completed".into(),
415 "request".into(),
418 "native_event".into(),
419 "runtime_disconnected".into(),
420 ],
421 opaque_fallback: true,
422 },
423 model: self.handle.harness.as_str().to_string(),
424 turn_state: if self.busy.load(Ordering::SeqCst) {
425 FrontendTurnState::Busy
426 } else {
427 FrontendTurnState::Idle
428 },
429 connection_state: if self.closed.load(Ordering::SeqCst) {
430 FrontendConnectionState::ShuttingDown
431 } else {
432 FrontendConnectionState::Connected
433 },
434 extensions: Default::default(),
435 }
436 }
437}
438
439#[async_trait]
440impl FrontendRuntime for HostedHarnessRuntime {
441 async fn describe(
442 &self,
443 ) -> std::result::Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
444 Ok(self.descriptor())
445 }
446
447 async fn attach(
448 &self,
449 _history_limit: usize,
450 ) -> std::result::Result<FrontendAttachment, FrontendRuntimeError> {
451 let live = self.frontend_events.subscribe();
452 let projection = self
453 .projection
454 .lock()
455 .unwrap_or_else(std::sync::PoisonError::into_inner);
456 let replay = projection.replay.clone();
457 if let Some(first) = replay.front() {
458 if first.sequence > 1 {
459 return Err(FrontendRuntimeError::ReplayGap(first.sequence - 1));
460 }
461 }
462 Ok(FrontendAttachment::new(
463 self.descriptor(),
464 self.history.clone(),
465 0,
466 replay,
467 live,
468 None,
469 ))
470 }
471
472 async fn send_input(
473 self: Arc<Self>,
474 prompt: String,
475 ) -> std::result::Result<(), FrontendRuntimeError> {
476 self.claim_submit().map_err(hosted_submit_error)?;
477 tokio::spawn(async move {
478 let _ = self
479 .submit_native_claimed(RuntimeInput {
480 text: prompt,
481 image_urls: Vec::new(),
482 })
483 .await;
484 });
485 Ok(())
486 }
487
488 async fn send_input_with_images(
489 self: Arc<Self>,
490 prompt: String,
491 image_urls: Vec<String>,
492 ) -> std::result::Result<(), FrontendRuntimeError> {
493 self.claim_submit().map_err(hosted_submit_error)?;
494 tokio::spawn(async move {
495 let _ = self
496 .submit_native_claimed(RuntimeInput {
497 text: prompt,
498 image_urls,
499 })
500 .await;
501 });
502 Ok(())
503 }
504
505 async fn submit(&self, prompt: String) -> std::result::Result<String, FrontendRuntimeError> {
506 self.submit_native(prompt)
507 .await
508 .map(|turn| turn.unwrap_or_default())
509 .map_err(hosted_submit_error)
510 }
511
512 async fn interrupt(&self) -> std::result::Result<bool, FrontendRuntimeError> {
513 if !self.busy.load(Ordering::SeqCst) {
514 return Ok(false);
515 }
516 self.interrupt_native()
517 .await
518 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
519 Ok(true)
523 }
524
525 async fn steer(&self, prompt: String) -> std::result::Result<(), FrontendRuntimeError> {
526 if !self.busy.load(Ordering::SeqCst) || !self.capabilities.steer {
527 return Err(FrontendRuntimeError::UnsupportedAction("steer"));
528 }
529 self.steer_native(prompt)
530 .await
531 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
532 }
533
534 async fn respond(
535 &self,
536 response: FrontendResponse,
537 ) -> std::result::Result<(), FrontendRuntimeError> {
538 let harness = self.handle.harness.as_str();
539 let FrontendResponse::Approval {
540 request_id,
541 decision,
542 } = &response
543 else {
544 return Err(FrontendRuntimeError::UnsupportedAction(
545 "respond: a hosted harness answers approval requests only",
546 ));
547 };
548 let Some(native_id) = self.native_request_id(*request_id) else {
549 return Err(FrontendRuntimeError::UnknownRequest(*request_id));
550 };
551 let body = hosted_permission_reply(harness, *decision)
552 .map_err(FrontendRuntimeError::InvalidResponse)?;
553 let canonical = serde_json::to_value(&response).unwrap_or(Value::Null);
554 self.respond_native_as(native_id, body, Some(canonical))
555 .await
556 .map_err(|error| FrontendRuntimeError::Execution {
557 operation: crate::SdkOperation::Respond,
558 message: error.to_string(),
559 })
560 }
561}
562
563fn hosted_submit_error(error: Error) -> FrontendRuntimeError {
564 if error.to_string().contains("already in progress") {
565 FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)
566 } else {
567 FrontendRuntimeError::Transport(error.to_string())
568 }
569}
570
571pub struct HostedHarnessConnection {
574 host: Arc<HostedHarnessRuntime>,
575 handle: RuntimeHandle,
576 events: broadcast::Receiver<HarnessEvent>,
577 closed: bool,
578}
579
580#[async_trait]
581impl RuntimeConnection for HostedHarnessConnection {
582 fn handle(&self) -> &RuntimeHandle {
583 &self.handle
584 }
585
586 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
587 self.host.claim_submit()?;
588 self.host.submit_native_claimed(input).await
589 }
590
591 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
592 match self.events.recv().await {
593 Ok(event) => Ok(Some(event)),
594 Err(broadcast::error::RecvError::Lagged(count)) => Err(Error::Other(format!(
595 "hosted harness event stream lost {count} event(s)"
596 ))),
597 Err(broadcast::error::RecvError::Closed) => Ok(None),
598 }
599 }
600
601 async fn interrupt(&mut self) -> Result<()> {
602 self.host.interrupt_native().await
603 }
604
605 async fn steer(&mut self, text: String) -> Result<()> {
606 self.host.steer_native(text).await
607 }
608
609 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
610 self.host.respond_native(request_id, response).await
611 }
612
613 async fn close(&mut self) -> Result<()> {
614 if self.closed {
615 return Ok(());
616 }
617 self.host.shutdown().await?;
618 self.closed = true;
619 Ok(())
620 }
621}
622
623async fn run_native_runtime(
624 mut runtime: Box<dyn RuntimeConnection>,
625 host: std::sync::Weak<HostedHarnessRuntime>,
626 mut commands: mpsc::Receiver<HostCommand>,
627) {
628 loop {
629 tokio::select! {
630 command = commands.recv() => {
631 let Some(command) = command else {
632 let _ = runtime.close().await;
633 return;
634 };
635 match command {
636 HostCommand::Submit { input, reply } => {
637 let _ = reply.send(runtime.send_input(input).await);
638 }
639 HostCommand::Interrupt { reply } => {
640 let _ = reply.send(runtime.interrupt().await);
641 }
642 HostCommand::Steer { text, reply } => {
643 let _ = reply.send(runtime.steer(text).await);
644 }
645 HostCommand::Respond { request_id, response, reply } => {
646 let _ = reply.send(runtime.respond(request_id, response).await);
647 }
648 HostCommand::Shutdown { reply } => {
649 let result = runtime.close().await;
650 let closed = result.is_ok();
651 let _ = reply.send(result);
652 if closed {
653 if let Some(host) = host.upgrade() {
654 host.mark_closed("Harness runtime closed.");
655 }
656 return;
657 }
658 }
659 }
660 }
661 event = runtime.next_event() => {
662 let Some(host) = host.upgrade() else {
663 let _ = runtime.close().await;
664 return;
665 };
666 match event {
667 Ok(Some(event)) if event.kind == "transport_closed" => {
668 host.mark_closed("Harness runtime transport closed.");
669 return;
670 }
671 Ok(Some(event)) => host.accept_native_event(event),
672 Ok(None) => {
673 host.mark_closed("Harness runtime transport closed.");
674 return;
675 }
676 Err(error) => {
677 host.mark_closed(error.to_string());
678 return;
679 }
680 }
681 }
682 }
683 }
684}
685
686pub(crate) fn hosted_answerable_decisions(harness: &str) -> &'static [&'static str] {
695 match harness {
696 "claude-code" => &["allow", "deny"],
706 _ => &[],
707 }
708}
709
710fn hosted_permission_reply(
712 harness: &str,
713 decision: crate::FrontendApprovalDecision,
714) -> std::result::Result<Value, String> {
715 use crate::FrontendApprovalDecision as Decision;
716 let accepted = hosted_answerable_decisions(harness);
717 let wire = match decision {
718 Decision::Allow => "allow",
719 Decision::AllowForSession => "allow_for_session",
720 Decision::Deny => "deny",
721 };
722 if !accepted.contains(&wire) {
723 return Err(format!(
724 "`{wire}` is not a decision a hosted {harness} runtime can carry; this harness accepts \
725 {}",
726 if accepted.is_empty() {
727 "no portable decision — its requests are observe-only".to_string()
728 } else {
729 accepted
730 .iter()
731 .map(|decision| format!("`{decision}`"))
732 .collect::<Vec<_>>()
733 .join(" or ")
734 },
735 ));
736 }
737 match harness {
738 "claude-code" => Ok(match decision {
739 Decision::Allow => json!({"behavior": "allow"}),
740 Decision::Deny => json!({
742 "behavior": "deny",
743 "message": "denied through the Volter Harness frontend",
744 }),
745 Decision::AllowForSession => unreachable!("refused above"),
746 }),
747 _ => Err(format!(
748 "no hosted permission reply is defined for {harness}"
749 )),
750 }
751}
752
753fn recorded_history(handle: &RuntimeHandle) -> Vec<ChatMessage> {
760 let catalog = HarnessCatalog::new();
761 let Some(locator) = claude_transcript(handle).or_else(|| {
762 let query = DiscoveryQuery {
763 harnesses: vec![handle.harness.clone()],
764 query: Some(handle.runtime_id.clone()),
765 ..DiscoveryQuery::default()
766 };
767 catalog.discover(&query).ok().and_then(|found| {
768 found
769 .into_iter()
770 .find(|descriptor| descriptor.locator.session_id == handle.runtime_id)
771 .map(|descriptor| descriptor.locator)
772 })
773 }) else {
774 return Vec::new();
775 };
776 catalog
777 .load_display_view(&locator, Fidelity::Semantic, 500)
778 .map(|session| session.messages)
779 .unwrap_or_default()
780}
781
782fn claude_transcript(handle: &RuntimeHandle) -> Option<crate::SessionLocator> {
786 if handle.harness.as_str() != crate::HarnessId::CLAUDE_CODE {
787 return None;
788 }
789 let file = format!("{}.jsonl", handle.runtime_id);
790 let projects = crate::HarnessHomes::default().claude_code;
791 std::fs::read_dir(projects)
792 .ok()?
793 .flatten()
794 .map(|folder| folder.path().join(&file))
795 .find(|path| path.is_file())
796 .map(|path| crate::SessionLocator {
797 harness: handle.harness.clone(),
798 session_id: handle.runtime_id.clone(),
799 storage: crate::StorageLocator::File { path },
800 })
801}
802
803#[derive(Debug, Default)]
808pub struct NativeProjection {
809 opencode_reasoning_parts: HashSet<String>,
810 opencode_failed: bool,
811}
812
813impl NativeProjection {
814 pub fn project(&mut self, harness: &str, event: &HarnessEvent) -> Vec<Value> {
816 if harness == "opencode" {
817 let key = event.kind.to_ascii_lowercase();
818 let payload = &event.payload;
819 match key.as_str() {
820 "message.part.updated" => {
821 let part = payload
822 .pointer("/properties/part")
823 .or_else(|| payload.get("part"))
824 .unwrap_or(payload);
825 if let (Some(id), Some(kind)) = (
826 part.get("id").and_then(Value::as_str),
827 part.get("type").and_then(Value::as_str),
828 ) {
829 if kind.contains("reasoning") {
830 self.opencode_reasoning_parts.insert(id.to_string());
831 }
832 }
833 }
834 "message.part.delta" => {
835 let reasoning = payload
836 .pointer("/properties/partID")
837 .and_then(Value::as_str)
838 .is_some_and(|id| self.opencode_reasoning_parts.contains(id));
839 if reasoning {
840 let text = payload
841 .pointer("/properties/delta")
842 .and_then(Value::as_str)
843 .unwrap_or_default();
844 return vec![json!({"type":"reasoning", "text":text, "raw":payload})];
845 }
846 }
847 "session.error" => {
848 self.opencode_failed = true;
849 return vec![
850 json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into())}),
851 ];
852 }
853 "session.idle" if self.opencode_failed => {
854 self.opencode_failed = false;
855 self.opencode_reasoning_parts.clear();
856 return vec![json!({"type":"turn_completed"})];
857 }
858 "session.idle" => self.opencode_reasoning_parts.clear(),
859 "session.status"
861 if payload
862 .pointer("/properties/status/type")
863 .or_else(|| payload.pointer("/status/type"))
864 .and_then(Value::as_str)
865 == Some("busy") =>
866 {
867 self.opencode_failed = false;
868 }
869 _ => {}
870 }
871 }
872 project_native_event(harness, event)
873 }
874}
875
876fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
880 let key = event.kind.to_ascii_lowercase().replace('-', "_");
881 let payload = &event.payload;
882 if key == "transport_closed" {
883 return vec![
884 json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
885 ];
886 }
887 let retrying = payload
888 .pointer("/params/willRetry")
889 .or_else(|| payload.get("willRetry"))
890 .and_then(Value::as_bool)
891 == Some(true);
892 if key == "error" && retrying {
893 return vec![native_payload(&key, payload)];
895 }
896 if key == "transport_error" || key == "error" {
897 return vec![
898 json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
899 ];
900 }
901 if key == "session/update" {
902 let update = payload
903 .pointer("/params/update")
904 .or_else(|| payload.get("update"))
905 .unwrap_or(payload);
906 let update_kind = update
907 .get("sessionUpdate")
908 .or_else(|| update.get("type"))
909 .and_then(Value::as_str)
910 .unwrap_or_default()
911 .to_ascii_lowercase();
912 if update_kind == "agent_message_chunk" {
913 return vec![
914 json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
915 ];
916 }
917 if update_kind == "agent_thought_chunk" {
918 return vec![
919 json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
920 ];
921 }
922 if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
923 return vec![project_tool(update, payload)];
924 }
925 }
926 if key == "supercode/acp_request_completed" {
927 let failure = payload
928 .pointer("/params/error")
929 .or_else(|| payload.get("error"));
930 return completion(failure.and_then(extract_text));
931 }
932 match harness {
933 "codex" => project_codex(&key, payload),
934 "claude-code" => project_claude(&key, payload),
935 "pi" => project_pi(&key, payload),
936 "opencode" => project_opencode(&key, payload),
937 _ => project_generic(&key, payload),
938 }
939}
940
941fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
942 if key == "turn/started" {
943 return vec![native_payload(key, payload)];
944 }
945 if key == "turn/completed" {
946 let status = payload
947 .pointer("/params/turn/status")
948 .or_else(|| payload.pointer("/turn/status"))
949 .and_then(Value::as_str)
950 .unwrap_or("completed")
951 .to_ascii_lowercase();
952 return completion(
953 (status.contains("fail") || status.contains("error") || status.contains("cancel"))
954 .then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
955 );
956 }
957 if key.ends_with("/delta") && !key.contains("agentmessage") && !key.contains("reasoning") {
958 return vec![native_payload(key, payload)];
960 }
961 if key.ends_with("/delta") {
962 let text = payload
963 .pointer("/params/delta")
964 .or_else(|| payload.get("delta"))
965 .and_then(extract_text)
966 .or_else(|| extract_text(payload))
967 .unwrap_or_default();
968 return vec![
969 json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
970 ];
971 }
972 if key.contains("commandexecution")
973 || key.contains("mcptool")
974 || key.contains("filechange")
975 || key.contains("tool")
976 {
977 let source = payload
978 .pointer("/params/item")
979 .or_else(|| payload.get("item"))
980 .unwrap_or(payload);
981 return vec![project_tool(source, payload)];
982 }
983 vec![native_payload(key, payload)]
984}
985
986fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
998 if key == "result" {
999 let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
1000 || payload
1001 .get("subtype")
1002 .and_then(Value::as_str)
1003 .is_some_and(|subtype| subtype != "success");
1004 return completion(
1005 failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
1006 );
1007 }
1008 if key == "stream_event" {
1009 let stream = payload
1014 .get("event")
1015 .or_else(|| payload.get("stream_event"))
1016 .unwrap_or(payload);
1017 if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
1018 return vec![native_payload(key, payload)];
1019 }
1020 }
1021 if key == "assistant" {
1022 return project_claude_blocks(payload, true);
1023 }
1024 if key == "user" {
1025 return project_claude_blocks(payload, false);
1026 }
1027 if key == "control_request" {
1028 return project_claude_control_request(payload);
1029 }
1030 if key == "system" && claude_system_is_telemetry(payload) {
1031 return Vec::new();
1032 }
1033 if key == "tool_progress" {
1034 return Vec::new();
1038 }
1039 if key == "rate_limit_event" {
1040 let status = payload
1045 .pointer("/rate_limit_info/status")
1046 .and_then(Value::as_str)
1047 .unwrap_or_default();
1048 if status.starts_with("allowed") {
1049 return Vec::new();
1050 }
1051 }
1052 project_generic(key, payload)
1053}
1054
1055fn claude_system_is_telemetry(payload: &Value) -> bool {
1075 match payload
1076 .get("subtype")
1077 .and_then(Value::as_str)
1078 .unwrap_or_default()
1079 {
1080 "init" | "thinking_tokens" | "task_started" | "task_progress" | "task_updated"
1081 | "hook_started" | "hook_progress" | "vcs_state_changed" => true,
1082 "task_notification" => payload
1083 .get("output_file")
1084 .and_then(Value::as_str)
1085 .unwrap_or_default()
1086 .is_empty(),
1087 "hook_response" => payload.get("exit_code").and_then(Value::as_i64) == Some(0),
1088 _ => false,
1089 }
1090}
1091
1092fn claude_blocks(payload: &Value) -> Option<&Vec<Value>> {
1094 payload
1095 .pointer("/message/content")
1096 .or_else(|| payload.get("content"))
1097 .and_then(Value::as_array)
1098}
1099
1100fn project_claude_blocks(payload: &Value, assistant: bool) -> Vec<Value> {
1101 let Some(blocks) = claude_blocks(payload) else {
1102 let text = extract_text(payload).unwrap_or_default();
1105 if text.is_empty() {
1106 return vec![native_payload(
1107 if assistant { "assistant" } else { "user" },
1108 payload,
1109 )];
1110 }
1111 return vec![
1112 json!({"type":if assistant { "text_delta" } else { "user_message" }, "text":text, "raw":payload}),
1113 ];
1114 };
1115 let mut projected = Vec::new();
1116 for block in blocks {
1117 match block
1118 .get("type")
1119 .and_then(Value::as_str)
1120 .unwrap_or_default()
1121 {
1122 "text" => {
1123 let text = block
1124 .get("text")
1125 .and_then(Value::as_str)
1126 .unwrap_or_default();
1127 if !text.is_empty() {
1128 projected.push(json!({
1129 "type": if assistant { "text_delta" } else { "user_message" },
1130 "text": text,
1131 "raw": block,
1132 }));
1133 }
1134 }
1135 "thinking" | "redacted_thinking" => {
1136 let text = extract_text(block).unwrap_or_default();
1137 if !text.is_empty() {
1138 projected.push(json!({"type":"reasoning", "text":text, "raw":block}));
1139 }
1140 }
1141 "tool_use" => projected.push(json!({
1142 "type": "tool_call_started",
1143 "id": block.get("id").cloned().unwrap_or(Value::Null),
1144 "name": block.get("name").cloned().unwrap_or(Value::Null),
1145 "arguments": block.get("input").map(Value::to_string).unwrap_or_default(),
1146 "raw": block,
1147 })),
1148 "tool_result" => projected.push(json!({
1149 "type": "tool_call_completed",
1153 "id": block.get("tool_use_id").cloned().unwrap_or(Value::Null),
1154 "name": Value::Null,
1155 "output": extract_text(block.get("content").unwrap_or(block)).unwrap_or_default(),
1156 "is_error": block.get("is_error").and_then(Value::as_bool).unwrap_or(false),
1157 "raw": block,
1158 })),
1159 _ => projected.push(native_payload(
1160 if assistant { "assistant" } else { "user" },
1161 block,
1162 )),
1163 }
1164 }
1165 projected
1166}
1167
1168fn project_claude_control_request(payload: &Value) -> Vec<Value> {
1177 let request = payload.get("request").unwrap_or(payload);
1178 if request.get("subtype").and_then(Value::as_str) != Some("can_use_tool") {
1179 return vec![native_payload("control_request", payload)];
1180 }
1181 let native_id = payload
1182 .get("request_id")
1183 .and_then(Value::as_str)
1184 .unwrap_or_default();
1185 vec![json!({
1186 "type": "request",
1187 "request": {
1188 "id": claude_request_id(native_id),
1189 "kind": "approval",
1190 "payload": {
1191 "tool": request.get("tool_name").or_else(|| request.get("toolName")).cloned().unwrap_or(Value::Null),
1192 "arguments": request.get("input").cloned().unwrap_or(Value::Null),
1193 "native_request_id": native_id,
1194 "decisions": hosted_answerable_decisions("claude-code"),
1198 },
1199 },
1200 })]
1201}
1202
1203fn claude_request_id(native_id: &str) -> u64 {
1210 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1211 for byte in native_id.as_bytes() {
1212 hash ^= u64::from(*byte);
1213 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1214 }
1215 hash & ((1_u64 << 53) - 1)
1217}
1218
1219fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
1220 match key {
1221 "agent_start" => vec![native_payload(key, payload)],
1222 "agent_end" => completion(payload.get("error").and_then(extract_text)),
1223 "response"
1225 if payload.get("success").and_then(Value::as_bool) == Some(false)
1226 && payload.get("command").and_then(Value::as_str) == Some("prompt") =>
1227 {
1228 completion(Some(
1229 payload
1230 .get("error")
1231 .and_then(extract_text)
1232 .unwrap_or_else(|| "Pi refused the prompt.".into()),
1233 ))
1234 }
1235 "message_update" => {
1236 let update = payload
1241 .get("assistantMessageEvent")
1242 .or_else(|| payload.get("event"))
1243 .unwrap_or(payload);
1244 let kind = update
1245 .get("type")
1246 .and_then(Value::as_str)
1247 .unwrap_or_default();
1248 match kind {
1249 "text_delta" => vec![
1250 json!({"type":"text_delta", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1251 ],
1252 "thinking_delta" => vec![
1253 json!({"type":"reasoning", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1254 ],
1255 _ => vec![native_payload(key, payload)],
1256 }
1257 }
1258 _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
1259 _ => vec![native_payload(key, payload)],
1260 }
1261}
1262
1263fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
1264 if key == "session.idle" {
1265 return completion(None);
1266 }
1267 if key == "session.status" {
1268 let status = payload
1269 .pointer("/properties/status/type")
1270 .or_else(|| payload.pointer("/status/type"))
1271 .and_then(Value::as_str)
1272 .unwrap_or_default();
1273 if status == "busy" {
1274 return vec![native_payload(key, payload)];
1275 }
1276 return vec![native_payload(key, payload)];
1278 }
1279 if key == "message.part.delta" {
1280 let delta = payload
1282 .pointer("/properties/delta")
1283 .or_else(|| payload.get("delta"))
1284 .and_then(Value::as_str)
1285 .unwrap_or_default();
1286 let field = payload
1287 .pointer("/properties/field")
1288 .and_then(Value::as_str)
1289 .unwrap_or("text");
1290 return vec![
1291 json!({"type":if field.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1292 ];
1293 }
1294 if key == "message.part.updated" {
1295 let part = payload
1296 .pointer("/properties/part")
1297 .or_else(|| payload.get("part"))
1298 .unwrap_or(payload);
1299 if let Some(delta) = payload
1300 .pointer("/properties/delta")
1301 .or_else(|| payload.get("delta"))
1302 .and_then(Value::as_str)
1303 {
1304 let reasoning = part
1305 .get("type")
1306 .and_then(Value::as_str)
1307 .is_some_and(|kind| kind.contains("reasoning"));
1308 return vec![
1309 json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1310 ];
1311 }
1312 if part
1313 .get("type")
1314 .and_then(Value::as_str)
1315 .is_some_and(|kind| kind.contains("tool"))
1316 {
1317 return vec![project_tool(part, payload)];
1318 }
1319 }
1320 if key == "session.error" {
1321 return completion(Some(
1322 extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
1323 ));
1324 }
1325 vec![native_payload(key, payload)]
1326}
1327
1328fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
1329 match key {
1330 "turn_started" | "turn/started" | "agent_start" => {
1331 vec![native_payload(key, payload)]
1332 }
1333 "turn_completed" | "turn/completed" | "agent_end" => {
1334 completion(payload.get("error").and_then(extract_text))
1335 }
1336 "output_delta" | "content_delta" => vec![
1337 json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1338 ],
1339 "reasoning_delta" => vec![
1340 json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1341 ],
1342 "tool" => vec![project_tool(payload, payload)],
1343 _ => vec![native_payload(key, payload)],
1344 }
1345}
1346
1347fn completion(error: Option<String>) -> Vec<Value> {
1348 match error {
1349 Some(message) => vec![
1350 json!({"type":"turn_failed", "message":message}),
1351 json!({"type":"turn_completed"}),
1352 ],
1353 None => vec![
1354 json!({"type":"turn_succeeded"}),
1355 json!({"type":"turn_completed"}),
1356 ],
1357 }
1358}
1359
1360fn project_tool(source: &Value, raw: &Value) -> Value {
1361 let status = source
1362 .get("status")
1363 .or_else(|| source.get("state"))
1364 .or_else(|| source.get("sessionUpdate"))
1365 .and_then(Value::as_str)
1366 .unwrap_or_default()
1367 .to_ascii_lowercase();
1368 let completed = status.contains("complete")
1369 || status.contains("result")
1370 || status.contains("success")
1371 || status.contains("error")
1372 || status.contains("fail");
1373 let arguments = source
1374 .get("arguments")
1375 .or_else(|| source.get("input"))
1376 .or_else(|| source.get("rawInput"))
1377 .cloned()
1378 .unwrap_or(Value::Null);
1379 json!({
1380 "type": if completed { "tool_call_completed" } else { "tool_call_started" },
1381 "id": source.get("toolCallId").or_else(|| source.get("tool_call_id")).or_else(|| source.get("callId")).or_else(|| source.get("id")).cloned().unwrap_or(Value::Null),
1382 "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
1383 "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
1384 "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
1385 "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
1386 "raw": raw,
1387 })
1388}
1389
1390fn native_payload(kind: &str, payload: &Value) -> Value {
1391 json!({"type":"native_event", "kind":kind, "raw":payload})
1392}
1393
1394fn extract_text(value: &Value) -> Option<String> {
1395 match value {
1396 Value::String(text) => Some(text.clone()),
1397 Value::Array(values) => {
1398 let text = values
1399 .iter()
1400 .filter_map(extract_text)
1401 .collect::<Vec<_>>()
1402 .join("\n");
1403 (!text.is_empty()).then_some(text)
1404 }
1405 Value::Object(object) => {
1406 for key in ["text", "delta", "content", "message", "result", "error"] {
1407 if let Some(text) = object.get(key).and_then(Value::as_str) {
1408 return Some(text.to_string());
1409 }
1410 }
1411 for key in [
1412 "delta",
1413 "content",
1414 "message",
1415 "error",
1416 "data",
1417 "part",
1418 "params",
1419 "properties",
1420 "update",
1421 "event",
1422 ] {
1423 if let Some(text) = object.get(key).and_then(extract_text) {
1424 return Some(text);
1425 }
1426 }
1427 None
1428 }
1429 _ => None,
1430 }
1431}