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)]
810pub struct NativeProjection {
811 opencode_reasoning_parts: HashSet<String>,
812 opencode_failed: bool,
813 pi_retrying: bool,
814}
815
816impl NativeProjection {
817 pub fn project(&mut self, harness: &str, event: &HarnessEvent) -> Vec<Value> {
819 if harness == "pi" {
820 let payload = &event.payload;
821 match event.kind.as_str() {
822 "agent_end" => {
823 self.pi_retrying =
824 payload.get("willRetry").and_then(Value::as_bool) == Some(true);
825 }
826 "auto_retry_end"
829 if self.pi_retrying
830 && payload.get("success").and_then(Value::as_bool) == Some(false) =>
831 {
832 self.pi_retrying = false;
833 return completion(Some(
834 payload
835 .get("finalError")
836 .and_then(Value::as_str)
837 .unwrap_or("Pi's retry did not run.")
838 .to_string(),
839 ));
840 }
841 _ => {}
842 }
843 }
844 if harness == "opencode" {
845 let key = event.kind.to_ascii_lowercase();
846 let payload = &event.payload;
847 match key.as_str() {
848 "message.part.updated" => {
849 let part = payload
850 .pointer("/properties/part")
851 .or_else(|| payload.get("part"))
852 .unwrap_or(payload);
853 if let (Some(id), Some(kind)) = (
854 part.get("id").and_then(Value::as_str),
855 part.get("type").and_then(Value::as_str),
856 ) {
857 if kind.contains("reasoning") {
858 self.opencode_reasoning_parts.insert(id.to_string());
859 }
860 }
861 }
862 "message.part.delta" => {
863 let reasoning = payload
864 .pointer("/properties/partID")
865 .and_then(Value::as_str)
866 .is_some_and(|id| self.opencode_reasoning_parts.contains(id));
867 if reasoning {
868 let text = payload
869 .pointer("/properties/delta")
870 .and_then(Value::as_str)
871 .unwrap_or_default();
872 return vec![json!({"type":"reasoning", "text":text, "raw":payload})];
873 }
874 }
875 "session.error" => {
876 self.opencode_failed = true;
877 return vec![
878 json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into())}),
879 ];
880 }
881 "session.idle" if self.opencode_failed => {
882 self.opencode_failed = false;
883 self.opencode_reasoning_parts.clear();
884 return vec![json!({"type":"turn_completed"})];
885 }
886 "session.idle" => self.opencode_reasoning_parts.clear(),
887 "session.status"
889 if payload
890 .pointer("/properties/status/type")
891 .or_else(|| payload.pointer("/status/type"))
892 .and_then(Value::as_str)
893 == Some("busy") =>
894 {
895 self.opencode_failed = false;
896 }
897 _ => {}
898 }
899 }
900 project_native_event(harness, event)
901 }
902}
903
904fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
908 let key = event.kind.to_ascii_lowercase().replace('-', "_");
909 let payload = &event.payload;
910 if key == "transport_closed" {
911 return vec![
912 json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
913 ];
914 }
915 let retrying = payload
916 .pointer("/params/willRetry")
917 .or_else(|| payload.get("willRetry"))
918 .and_then(Value::as_bool)
919 == Some(true);
920 if key == "error" && retrying {
921 return vec![native_payload(&key, payload)];
923 }
924 if key == "transport_error" || key == "error" {
925 return vec![
926 json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
927 ];
928 }
929 if key == "session/update" {
930 let update = payload
931 .pointer("/params/update")
932 .or_else(|| payload.get("update"))
933 .unwrap_or(payload);
934 let update_kind = update
935 .get("sessionUpdate")
936 .or_else(|| update.get("type"))
937 .and_then(Value::as_str)
938 .unwrap_or_default()
939 .to_ascii_lowercase();
940 if update_kind == "agent_message_chunk" {
941 return vec![
942 json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
943 ];
944 }
945 if update_kind == "agent_thought_chunk" {
946 return vec![
947 json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
948 ];
949 }
950 if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
951 return vec![project_tool(update, payload)];
952 }
953 }
954 if key == "supercode/acp_request_completed" {
955 let failure = payload
956 .pointer("/params/error")
957 .or_else(|| payload.get("error"));
958 return completion(failure.and_then(extract_text));
959 }
960 match harness {
961 "codex" => project_codex(&key, payload),
962 "claude-code" => project_claude(&key, payload),
963 "pi" => project_pi(&key, payload),
964 "opencode" => project_opencode(&key, payload),
965 _ => project_generic(&key, payload),
966 }
967}
968
969fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
970 if key == "turn/started" {
971 return vec![native_payload(key, payload)];
972 }
973 if key == "turn/completed" {
974 let status = payload
975 .pointer("/params/turn/status")
976 .or_else(|| payload.pointer("/turn/status"))
977 .and_then(Value::as_str)
978 .unwrap_or("completed")
979 .to_ascii_lowercase();
980 return completion(
981 (status.contains("fail") || status.contains("error") || status.contains("cancel"))
982 .then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
983 );
984 }
985 if key.ends_with("/delta") && !key.contains("agentmessage") && !key.contains("reasoning") {
986 return vec![native_payload(key, payload)];
988 }
989 if key.ends_with("/delta") {
990 let text = payload
991 .pointer("/params/delta")
992 .or_else(|| payload.get("delta"))
993 .and_then(extract_text)
994 .or_else(|| extract_text(payload))
995 .unwrap_or_default();
996 return vec![
997 json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
998 ];
999 }
1000 if key.contains("commandexecution")
1001 || key.contains("mcptool")
1002 || key.contains("filechange")
1003 || key.contains("tool")
1004 {
1005 let source = payload
1006 .pointer("/params/item")
1007 .or_else(|| payload.get("item"))
1008 .unwrap_or(payload);
1009 return vec![project_tool(source, payload)];
1010 }
1011 vec![native_payload(key, payload)]
1012}
1013
1014fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
1026 if key == "result" {
1027 let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
1028 || payload
1029 .get("subtype")
1030 .and_then(Value::as_str)
1031 .is_some_and(|subtype| subtype != "success");
1032 return completion(
1033 failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
1034 );
1035 }
1036 if key == "stream_event" {
1037 let stream = payload
1042 .get("event")
1043 .or_else(|| payload.get("stream_event"))
1044 .unwrap_or(payload);
1045 if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
1046 return vec![native_payload(key, payload)];
1047 }
1048 }
1049 if key == "assistant" {
1050 return project_claude_blocks(payload, true);
1051 }
1052 if key == "user" {
1053 return project_claude_blocks(payload, false);
1054 }
1055 if key == "control_request" {
1056 return project_claude_control_request(payload);
1057 }
1058 if key == "system" && claude_system_is_telemetry(payload) {
1059 return Vec::new();
1060 }
1061 if key == "tool_progress" {
1062 return Vec::new();
1066 }
1067 if key == "rate_limit_event" {
1068 let status = payload
1073 .pointer("/rate_limit_info/status")
1074 .and_then(Value::as_str)
1075 .unwrap_or_default();
1076 if status.starts_with("allowed") {
1077 return Vec::new();
1078 }
1079 }
1080 project_generic(key, payload)
1081}
1082
1083fn claude_system_is_telemetry(payload: &Value) -> bool {
1103 match payload
1104 .get("subtype")
1105 .and_then(Value::as_str)
1106 .unwrap_or_default()
1107 {
1108 "init" | "thinking_tokens" | "task_started" | "task_progress" | "task_updated"
1109 | "hook_started" | "hook_progress" | "vcs_state_changed" => true,
1110 "task_notification" => payload
1111 .get("output_file")
1112 .and_then(Value::as_str)
1113 .unwrap_or_default()
1114 .is_empty(),
1115 "hook_response" => payload.get("exit_code").and_then(Value::as_i64) == Some(0),
1116 _ => false,
1117 }
1118}
1119
1120fn claude_blocks(payload: &Value) -> Option<&Vec<Value>> {
1122 payload
1123 .pointer("/message/content")
1124 .or_else(|| payload.get("content"))
1125 .and_then(Value::as_array)
1126}
1127
1128fn project_claude_blocks(payload: &Value, assistant: bool) -> Vec<Value> {
1129 let Some(blocks) = claude_blocks(payload) else {
1130 let text = extract_text(payload).unwrap_or_default();
1133 if text.is_empty() {
1134 return vec![native_payload(
1135 if assistant { "assistant" } else { "user" },
1136 payload,
1137 )];
1138 }
1139 return vec![
1140 json!({"type":if assistant { "text_delta" } else { "user_message" }, "text":text, "raw":payload}),
1141 ];
1142 };
1143 let mut projected = Vec::new();
1144 for block in blocks {
1145 match block
1146 .get("type")
1147 .and_then(Value::as_str)
1148 .unwrap_or_default()
1149 {
1150 "text" => {
1151 let text = block
1152 .get("text")
1153 .and_then(Value::as_str)
1154 .unwrap_or_default();
1155 if !text.is_empty() {
1156 projected.push(json!({
1157 "type": if assistant { "text_delta" } else { "user_message" },
1158 "text": text,
1159 "raw": block,
1160 }));
1161 }
1162 }
1163 "thinking" | "redacted_thinking" => {
1164 let text = extract_text(block).unwrap_or_default();
1165 if !text.is_empty() {
1166 projected.push(json!({"type":"reasoning", "text":text, "raw":block}));
1167 }
1168 }
1169 "tool_use" => projected.push(json!({
1170 "type": "tool_call_started",
1171 "id": block.get("id").cloned().unwrap_or(Value::Null),
1172 "name": block.get("name").cloned().unwrap_or(Value::Null),
1173 "arguments": block.get("input").map(Value::to_string).unwrap_or_default(),
1174 "raw": block,
1175 })),
1176 "tool_result" => projected.push(json!({
1177 "type": "tool_call_completed",
1181 "id": block.get("tool_use_id").cloned().unwrap_or(Value::Null),
1182 "name": Value::Null,
1183 "output": extract_text(block.get("content").unwrap_or(block)).unwrap_or_default(),
1184 "is_error": block.get("is_error").and_then(Value::as_bool).unwrap_or(false),
1185 "raw": block,
1186 })),
1187 _ => projected.push(native_payload(
1188 if assistant { "assistant" } else { "user" },
1189 block,
1190 )),
1191 }
1192 }
1193 projected
1194}
1195
1196fn project_claude_control_request(payload: &Value) -> Vec<Value> {
1205 let request = payload.get("request").unwrap_or(payload);
1206 if request.get("subtype").and_then(Value::as_str) != Some("can_use_tool") {
1207 return vec![native_payload("control_request", payload)];
1208 }
1209 let native_id = payload
1210 .get("request_id")
1211 .and_then(Value::as_str)
1212 .unwrap_or_default();
1213 vec![json!({
1214 "type": "request",
1215 "request": {
1216 "id": claude_request_id(native_id),
1217 "kind": "approval",
1218 "payload": {
1219 "tool": request.get("tool_name").or_else(|| request.get("toolName")).cloned().unwrap_or(Value::Null),
1220 "arguments": request.get("input").cloned().unwrap_or(Value::Null),
1221 "native_request_id": native_id,
1222 "decisions": hosted_answerable_decisions("claude-code"),
1226 },
1227 },
1228 })]
1229}
1230
1231fn claude_request_id(native_id: &str) -> u64 {
1238 let mut hash: u64 = 0xcbf2_9ce4_8422_2325;
1239 for byte in native_id.as_bytes() {
1240 hash ^= u64::from(*byte);
1241 hash = hash.wrapping_mul(0x0000_0100_0000_01b3);
1242 }
1243 hash & ((1_u64 << 53) - 1)
1245}
1246
1247fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
1248 match key {
1249 "agent_start" => vec![native_payload(key, payload)],
1250 "agent_end" if payload.get("willRetry").and_then(Value::as_bool) == Some(true) => {
1253 vec![native_payload(key, payload)]
1254 }
1255 "agent_end" => completion(
1256 payload
1257 .get("error")
1258 .and_then(extract_text)
1259 .or_else(|| pi_run_error(payload)),
1260 ),
1261 "response"
1263 if payload.get("success").and_then(Value::as_bool) == Some(false)
1264 && payload.get("command").and_then(Value::as_str) == Some("prompt") =>
1265 {
1266 completion(Some(
1267 payload
1268 .get("error")
1269 .and_then(extract_text)
1270 .unwrap_or_else(|| "Pi refused the prompt.".into()),
1271 ))
1272 }
1273 "message_update" => {
1274 let update = payload
1279 .get("assistantMessageEvent")
1280 .or_else(|| payload.get("event"))
1281 .unwrap_or(payload);
1282 let kind = update
1283 .get("type")
1284 .and_then(Value::as_str)
1285 .unwrap_or_default();
1286 match kind {
1287 "text_delta" => vec![
1288 json!({"type":"text_delta", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1289 ],
1290 "thinking_delta" => vec![
1291 json!({"type":"reasoning", "text":update.get("delta").and_then(Value::as_str).unwrap_or_default(), "raw":payload}),
1292 ],
1293 _ => vec![native_payload(key, payload)],
1294 }
1295 }
1296 _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
1297 _ => vec![native_payload(key, payload)],
1298 }
1299}
1300
1301fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
1302 if key == "session.idle" {
1303 return completion(None);
1304 }
1305 if key == "session.status" {
1306 let status = payload
1307 .pointer("/properties/status/type")
1308 .or_else(|| payload.pointer("/status/type"))
1309 .and_then(Value::as_str)
1310 .unwrap_or_default();
1311 if status == "busy" {
1312 return vec![native_payload(key, payload)];
1313 }
1314 return vec![native_payload(key, payload)];
1316 }
1317 if key == "message.part.delta" {
1318 let delta = payload
1320 .pointer("/properties/delta")
1321 .or_else(|| payload.get("delta"))
1322 .and_then(Value::as_str)
1323 .unwrap_or_default();
1324 let field = payload
1325 .pointer("/properties/field")
1326 .and_then(Value::as_str)
1327 .unwrap_or("text");
1328 return vec![
1329 json!({"type":if field.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1330 ];
1331 }
1332 if key == "message.part.updated" {
1333 let part = payload
1334 .pointer("/properties/part")
1335 .or_else(|| payload.get("part"))
1336 .unwrap_or(payload);
1337 if let Some(delta) = payload
1338 .pointer("/properties/delta")
1339 .or_else(|| payload.get("delta"))
1340 .and_then(Value::as_str)
1341 {
1342 let reasoning = part
1343 .get("type")
1344 .and_then(Value::as_str)
1345 .is_some_and(|kind| kind.contains("reasoning"));
1346 return vec![
1347 json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
1348 ];
1349 }
1350 if part
1351 .get("type")
1352 .and_then(Value::as_str)
1353 .is_some_and(|kind| kind.contains("tool"))
1354 {
1355 return vec![project_tool(part, payload)];
1356 }
1357 }
1358 if key == "session.error" {
1359 return completion(Some(
1360 extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
1361 ));
1362 }
1363 vec![native_payload(key, payload)]
1364}
1365
1366fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
1367 match key {
1368 "turn_started" | "turn/started" | "agent_start" => {
1369 vec![native_payload(key, payload)]
1370 }
1371 "turn_completed" | "turn/completed" | "agent_end" => {
1372 completion(payload.get("error").and_then(extract_text))
1373 }
1374 "output_delta" | "content_delta" => vec![
1375 json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1376 ],
1377 "reasoning_delta" => vec![
1378 json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
1379 ],
1380 "tool" => vec![project_tool(payload, payload)],
1381 _ => vec![native_payload(key, payload)],
1382 }
1383}
1384
1385fn pi_run_error(payload: &Value) -> Option<String> {
1390 let last = payload
1391 .get("messages")?
1392 .as_array()?
1393 .iter()
1394 .rev()
1395 .find(|message| message.get("role").and_then(Value::as_str) == Some("assistant"))?;
1396 if last.get("stopReason").and_then(Value::as_str) != Some("error") {
1397 return None;
1398 }
1399 Some(
1400 last.get("errorMessage")
1401 .and_then(Value::as_str)
1402 .filter(|text| !text.is_empty())
1403 .unwrap_or("Pi's model call failed.")
1404 .to_string(),
1405 )
1406}
1407
1408fn completion(error: Option<String>) -> Vec<Value> {
1409 match error {
1410 Some(message) => vec![
1411 json!({"type":"turn_failed", "message":message}),
1412 json!({"type":"turn_completed"}),
1413 ],
1414 None => vec![
1415 json!({"type":"turn_succeeded"}),
1416 json!({"type":"turn_completed"}),
1417 ],
1418 }
1419}
1420
1421fn project_tool(source: &Value, raw: &Value) -> Value {
1422 let status = source
1423 .get("status")
1424 .or_else(|| source.get("state"))
1425 .or_else(|| source.get("sessionUpdate"))
1426 .and_then(Value::as_str)
1427 .unwrap_or_default()
1428 .to_ascii_lowercase();
1429 let completed = status.contains("complete")
1430 || status.contains("result")
1431 || status.contains("success")
1432 || status.contains("error")
1433 || status.contains("fail");
1434 let arguments = source
1435 .get("arguments")
1436 .or_else(|| source.get("input"))
1437 .or_else(|| source.get("rawInput"))
1438 .cloned()
1439 .unwrap_or(Value::Null);
1440 json!({
1441 "type": if completed { "tool_call_completed" } else { "tool_call_started" },
1442 "id": source.get("toolCallId").or_else(|| source.get("tool_call_id")).or_else(|| source.get("callId")).or_else(|| source.get("id")).cloned().unwrap_or(Value::Null),
1443 "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
1444 "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
1445 "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
1446 "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
1447 "raw": raw,
1448 })
1449}
1450
1451fn native_payload(kind: &str, payload: &Value) -> Value {
1452 json!({"type":"native_event", "kind":kind, "raw":payload})
1453}
1454
1455fn extract_text(value: &Value) -> Option<String> {
1456 match value {
1457 Value::String(text) => Some(text.clone()),
1458 Value::Array(values) => {
1459 let text = values
1460 .iter()
1461 .filter_map(extract_text)
1462 .collect::<Vec<_>>()
1463 .join("\n");
1464 (!text.is_empty()).then_some(text)
1465 }
1466 Value::Object(object) => {
1467 for key in ["text", "delta", "content", "message", "result", "error"] {
1468 if let Some(text) = object.get(key).and_then(Value::as_str) {
1469 return Some(text.to_string());
1470 }
1471 }
1472 for key in [
1473 "delta",
1474 "content",
1475 "message",
1476 "error",
1477 "data",
1478 "part",
1479 "params",
1480 "properties",
1481 "update",
1482 "event",
1483 ] {
1484 if let Some(text) = object.get(key).and_then(extract_text) {
1485 return Some(text);
1486 }
1487 }
1488 None
1489 }
1490 _ => None,
1491 }
1492}