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