1use std::collections::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::{Error, FrontendResponse, Result};
24
25enum HostCommand {
26 Submit {
27 input: RuntimeInput,
28 reply: oneshot::Sender<Result<Option<String>>>,
29 },
30 Interrupt {
31 reply: oneshot::Sender<Result<()>>,
32 },
33 Respond {
34 request_id: Value,
35 response: Value,
36 reply: oneshot::Sender<Result<()>>,
37 },
38 Shutdown {
39 reply: oneshot::Sender<Result<()>>,
40 },
41}
42
43struct ProjectionState {
44 next_sequence: u64,
45 replay: VecDeque<FrontendEvent>,
46}
47
48pub struct HostedHarnessRuntime {
50 handle: RuntimeHandle,
51 capabilities: RuntimeCapabilities,
52 commands: mpsc::Sender<HostCommand>,
53 raw_events: broadcast::Sender<HarnessEvent>,
54 frontend_events: broadcast::Sender<FrontendEvent>,
55 projection: StdMutex<ProjectionState>,
56 busy: AtomicBool,
57 closed: AtomicBool,
58}
59
60impl HostedHarnessRuntime {
61 pub fn spawn(
64 runtime: Box<dyn RuntimeConnection>,
65 capabilities: RuntimeCapabilities,
66 ) -> (Arc<Self>, HostedHarnessConnection) {
67 let handle = runtime.handle().clone();
68 let (commands, command_rx) = mpsc::channel(32);
69 let (raw_events, raw_rx) = broadcast::channel(1024);
70 let (frontend_events, _) = broadcast::channel(1024);
71 let host = Arc::new(Self {
72 handle: handle.clone(),
73 capabilities,
74 commands,
75 raw_events,
76 frontend_events,
77 projection: StdMutex::new(ProjectionState {
78 next_sequence: 1,
79 replay: VecDeque::new(),
80 }),
81 busy: AtomicBool::new(false),
82 closed: AtomicBool::new(false),
83 });
84 tokio::spawn(run_native_runtime(
85 runtime,
86 Arc::downgrade(&host),
87 command_rx,
88 ));
89 let connection = HostedHarnessConnection {
90 host: host.clone(),
91 handle,
92 events: raw_rx,
93 closed: false,
94 };
95 (host, connection)
96 }
97
98 pub fn frontend_sender(&self) -> broadcast::Sender<FrontendEvent> {
100 self.frontend_events.clone()
101 }
102
103 pub async fn shutdown(&self) -> Result<()> {
105 if self.closed.load(Ordering::SeqCst) {
106 return Ok(());
107 }
108 let (reply, response) = oneshot::channel();
109 self.commands
110 .send(HostCommand::Shutdown { reply })
111 .await
112 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
113 response
114 .await
115 .map_err(|_| Error::Other("hosted harness runtime stopped before shutdown".into()))?
116 }
117
118 fn claim_submit(&self) -> Result<()> {
119 if self.closed.load(Ordering::SeqCst) {
120 return Err(Error::Other("hosted harness runtime is closed".into()));
121 }
122 if self
123 .busy
124 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
125 .is_err()
126 {
127 return Err(Error::Other("a harness turn is already in progress".into()));
128 }
129 Ok(())
130 }
131
132 async fn submit_native_claimed(&self, input: RuntimeInput) -> Result<Option<String>> {
133 self.publish(json!({"type":"user_message", "text":input.text}));
134 self.publish(json!({"type":"turn_started"}));
135 let (reply, response) = oneshot::channel();
136 if self
137 .commands
138 .send(HostCommand::Submit { input, reply })
139 .await
140 .is_err()
141 {
142 self.busy.store(false, Ordering::SeqCst);
143 self.publish(
144 json!({"type":"turn_failed", "message":"Hosted harness runtime is closed."}),
145 );
146 self.publish(json!({"type":"turn_completed"}));
147 self.mark_closed("Harness runtime command channel closed.");
148 return Err(Error::Other("hosted harness runtime is closed".into()));
149 }
150 match response.await {
151 Ok(Ok(turn)) => Ok(turn),
152 Ok(Err(error)) => {
153 self.busy.store(false, Ordering::SeqCst);
154 self.publish(json!({"type":"turn_failed", "message":error.to_string()}));
155 self.publish(json!({"type":"turn_completed"}));
156 Err(error)
157 }
158 Err(_) => {
159 self.busy.store(false, Ordering::SeqCst);
160 self.publish(json!({"type":"turn_failed", "message":"Hosted harness runtime stopped before accepting input."}));
161 self.publish(json!({"type":"turn_completed"}));
162 self.mark_closed("Harness runtime stopped before accepting input.");
163 Err(Error::Other(
164 "hosted harness runtime stopped before accepting input".into(),
165 ))
166 }
167 }
168 }
169
170 async fn submit_native(&self, text: String) -> Result<Option<String>> {
171 self.claim_submit()?;
172 self.submit_native_claimed(RuntimeInput {
173 text,
174 image_urls: Vec::new(),
175 })
176 .await
177 }
178
179 async fn interrupt_native(&self) -> Result<()> {
180 let (reply, response) = oneshot::channel();
181 self.commands
182 .send(HostCommand::Interrupt { reply })
183 .await
184 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
185 response
186 .await
187 .map_err(|_| Error::Other("hosted harness runtime stopped before interrupt".into()))?
188 }
189
190 async fn respond_native(&self, request_id: Value, response: Value) -> Result<()> {
191 let (reply, completed) = oneshot::channel();
192 self.commands
193 .send(HostCommand::Respond {
194 request_id,
195 response,
196 reply,
197 })
198 .await
199 .map_err(|_| Error::Other("hosted harness runtime is closed".into()))?;
200 completed
201 .await
202 .map_err(|_| Error::Other("hosted harness runtime stopped before response".into()))?
203 }
204
205 fn publish(&self, payload: Value) {
206 let event = {
207 let mut projection = self
208 .projection
209 .lock()
210 .unwrap_or_else(std::sync::PoisonError::into_inner);
211 let event = FrontendEvent::new(projection.next_sequence, payload);
212 projection.next_sequence = projection.next_sequence.saturating_add(1);
213 projection.replay.push_back(event.clone());
214 while projection.replay.len() > FRONTEND_REPLAY_CAPACITY {
215 projection.replay.pop_front();
216 }
217 event
218 };
219 let _ = self.frontend_events.send(event);
220 }
221
222 fn accept_native_event(&self, event: HarnessEvent) {
223 let _ = self.raw_events.send(event.clone());
224 for payload in project_native_event(self.handle.harness.as_str(), &event) {
225 let terminal = matches!(
226 payload.get("type").and_then(Value::as_str),
227 Some("turn_succeeded" | "turn_interrupted" | "turn_failed")
228 );
229 if terminal {
230 self.busy.store(false, Ordering::SeqCst);
231 }
232 self.publish(payload);
233 }
234 }
235
236 fn mark_closed(&self, message: impl Into<String>) {
237 if self.closed.swap(true, Ordering::SeqCst) {
238 return;
239 }
240 let message = message.into();
241 self.busy.store(false, Ordering::SeqCst);
242 let _ = self.raw_events.send(HarnessEvent {
246 sequence: None,
247 kind: "transport_closed".into(),
248 payload: json!({"message":message, "terminal":true}),
249 });
250 self.publish(json!({"type":"runtime_disconnected", "message":message}));
251 }
252
253 fn descriptor(&self) -> FrontendRuntimeDescriptor {
254 FrontendRuntimeDescriptor {
255 schema_version: FRONTEND_RUNTIME_SCHEMA_VERSION,
256 session_id: self.handle.runtime_id.clone(),
257 source_harness: Some(self.handle.harness.as_str().to_string()),
258 emulation_profile: None,
259 active_modules: Vec::new(),
260 commands: Vec::new(),
261 operations: Vec::new(),
262 actions: FrontendActions {
263 submit: self.capabilities.send_input,
264 interrupt: self.capabilities.interrupt,
265 steer: false,
266 respond: false,
270 detach: true,
271 close: false,
272 },
273 display: FrontendDisplayCapabilities {
274 event_kinds: vec![
275 "user_message".into(),
276 "turn_started".into(),
277 "turn_succeeded".into(),
278 "turn_interrupted".into(),
279 "turn_failed".into(),
280 "text_delta".into(),
281 "reasoning".into(),
282 "tool_call_started".into(),
283 "tool_call_completed".into(),
284 "native_event".into(),
285 "runtime_disconnected".into(),
286 ],
287 opaque_fallback: true,
288 },
289 model: self.handle.harness.as_str().to_string(),
290 turn_state: if self.busy.load(Ordering::SeqCst) {
291 FrontendTurnState::Busy
292 } else {
293 FrontendTurnState::Idle
294 },
295 connection_state: if self.closed.load(Ordering::SeqCst) {
296 FrontendConnectionState::ShuttingDown
297 } else {
298 FrontendConnectionState::Connected
299 },
300 extensions: Default::default(),
301 }
302 }
303}
304
305#[async_trait]
306impl FrontendRuntime for HostedHarnessRuntime {
307 async fn describe(
308 &self,
309 ) -> std::result::Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
310 Ok(self.descriptor())
311 }
312
313 async fn attach(
314 &self,
315 _history_limit: usize,
316 ) -> std::result::Result<FrontendAttachment, FrontendRuntimeError> {
317 let live = self.frontend_events.subscribe();
318 let projection = self
319 .projection
320 .lock()
321 .unwrap_or_else(std::sync::PoisonError::into_inner);
322 let replay = projection.replay.clone();
323 if let Some(first) = replay.front() {
324 if first.sequence > 1 {
325 return Err(FrontendRuntimeError::ReplayGap(first.sequence - 1));
326 }
327 }
328 Ok(FrontendAttachment::new(
329 self.descriptor(),
330 Vec::new(),
331 0,
332 replay,
333 live,
334 None,
335 ))
336 }
337
338 async fn send_input(
339 self: Arc<Self>,
340 prompt: String,
341 ) -> std::result::Result<(), FrontendRuntimeError> {
342 self.claim_submit().map_err(hosted_submit_error)?;
343 tokio::spawn(async move {
344 let _ = self
345 .submit_native_claimed(RuntimeInput {
346 text: prompt,
347 image_urls: Vec::new(),
348 })
349 .await;
350 });
351 Ok(())
352 }
353
354 async fn send_input_with_images(
355 self: Arc<Self>,
356 prompt: String,
357 image_urls: Vec<String>,
358 ) -> std::result::Result<(), FrontendRuntimeError> {
359 self.claim_submit().map_err(hosted_submit_error)?;
360 tokio::spawn(async move {
361 let _ = self
362 .submit_native_claimed(RuntimeInput {
363 text: prompt,
364 image_urls,
365 })
366 .await;
367 });
368 Ok(())
369 }
370
371 async fn submit(&self, prompt: String) -> std::result::Result<String, FrontendRuntimeError> {
372 self.submit_native(prompt)
373 .await
374 .map(|turn| turn.unwrap_or_default())
375 .map_err(hosted_submit_error)
376 }
377
378 async fn interrupt(&self) -> std::result::Result<bool, FrontendRuntimeError> {
379 if !self.busy.load(Ordering::SeqCst) {
380 return Ok(false);
381 }
382 self.interrupt_native()
383 .await
384 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
385 Ok(true)
389 }
390
391 async fn steer(&self, _prompt: String) -> std::result::Result<(), FrontendRuntimeError> {
392 Err(FrontendRuntimeError::UnsupportedAction("steer"))
393 }
394
395 async fn respond(
396 &self,
397 _response: FrontendResponse,
398 ) -> std::result::Result<(), FrontendRuntimeError> {
399 Err(FrontendRuntimeError::UnsupportedAction("respond"))
400 }
401}
402
403fn hosted_submit_error(error: Error) -> FrontendRuntimeError {
404 if error.to_string().contains("already in progress") {
405 FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)
406 } else {
407 FrontendRuntimeError::Transport(error.to_string())
408 }
409}
410
411pub struct HostedHarnessConnection {
414 host: Arc<HostedHarnessRuntime>,
415 handle: RuntimeHandle,
416 events: broadcast::Receiver<HarnessEvent>,
417 closed: bool,
418}
419
420#[async_trait]
421impl RuntimeConnection for HostedHarnessConnection {
422 fn handle(&self) -> &RuntimeHandle {
423 &self.handle
424 }
425
426 async fn send_input(&mut self, input: RuntimeInput) -> Result<Option<String>> {
427 self.host.claim_submit()?;
428 self.host.submit_native_claimed(input).await
429 }
430
431 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
432 match self.events.recv().await {
433 Ok(event) => Ok(Some(event)),
434 Err(broadcast::error::RecvError::Lagged(count)) => Err(Error::Other(format!(
435 "hosted harness event stream lost {count} event(s)"
436 ))),
437 Err(broadcast::error::RecvError::Closed) => Ok(None),
438 }
439 }
440
441 async fn interrupt(&mut self) -> Result<()> {
442 self.host.interrupt_native().await
443 }
444
445 async fn respond(&mut self, request_id: Value, response: Value) -> Result<()> {
446 self.host.respond_native(request_id, response).await
447 }
448
449 async fn close(&mut self) -> Result<()> {
450 if self.closed {
451 return Ok(());
452 }
453 self.closed = true;
454 self.host.shutdown().await
455 }
456}
457
458async fn run_native_runtime(
459 mut runtime: Box<dyn RuntimeConnection>,
460 host: std::sync::Weak<HostedHarnessRuntime>,
461 mut commands: mpsc::Receiver<HostCommand>,
462) {
463 loop {
464 tokio::select! {
465 command = commands.recv() => {
466 let Some(command) = command else {
467 let _ = runtime.close().await;
468 return;
469 };
470 match command {
471 HostCommand::Submit { input, reply } => {
472 let _ = reply.send(runtime.send_input(input).await);
473 }
474 HostCommand::Interrupt { reply } => {
475 let _ = reply.send(runtime.interrupt().await);
476 }
477 HostCommand::Respond { request_id, response, reply } => {
478 let _ = reply.send(runtime.respond(request_id, response).await);
479 }
480 HostCommand::Shutdown { reply } => {
481 let result = runtime.close().await;
482 let _ = reply.send(result);
483 if let Some(host) = host.upgrade() {
484 host.mark_closed("Harness runtime closed.");
485 }
486 return;
487 }
488 }
489 }
490 event = runtime.next_event() => {
491 let Some(host) = host.upgrade() else {
492 let _ = runtime.close().await;
493 return;
494 };
495 match event {
496 Ok(Some(event)) => host.accept_native_event(event),
497 Ok(None) => {
498 host.mark_closed("Harness runtime transport closed.");
499 return;
500 }
501 Err(error) => {
502 host.mark_closed(error.to_string());
503 return;
504 }
505 }
506 }
507 }
508 }
509}
510
511fn project_native_event(harness: &str, event: &HarnessEvent) -> Vec<Value> {
512 let key = event.kind.to_ascii_lowercase().replace('-', "_");
513 let payload = &event.payload;
514 if key == "transport_closed" {
515 return vec![
516 json!({"type":"runtime_disconnected", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime transport closed.".into())}),
517 ];
518 }
519 if key == "transport_error" || key == "error" {
520 return vec![
521 json!({"type":"turn_failed", "message":extract_text(payload).unwrap_or_else(|| "Harness runtime failed.".into()), "raw":payload}),
522 ];
523 }
524 if key == "session/update" {
525 let update = payload
526 .pointer("/params/update")
527 .or_else(|| payload.get("update"))
528 .unwrap_or(payload);
529 let update_kind = update
530 .get("sessionUpdate")
531 .or_else(|| update.get("type"))
532 .and_then(Value::as_str)
533 .unwrap_or_default()
534 .to_ascii_lowercase();
535 if update_kind == "agent_message_chunk" {
536 return vec![
537 json!({"type":"text_delta", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
538 ];
539 }
540 if update_kind == "agent_thought_chunk" {
541 return vec![
542 json!({"type":"reasoning", "text":extract_text(update.get("content").unwrap_or(update)).unwrap_or_default(), "raw":payload}),
543 ];
544 }
545 if matches!(update_kind.as_str(), "tool_call" | "tool_call_update") {
546 return vec![project_tool(update, payload)];
547 }
548 }
549 if key == "supercode/acp_request_completed" {
550 let failure = payload
551 .pointer("/params/error")
552 .or_else(|| payload.get("error"));
553 return completion(failure.and_then(extract_text));
554 }
555 match harness {
556 "codex" => project_codex(&key, payload),
557 "claude-code" => project_claude(&key, payload),
558 "pi" => project_pi(&key, payload),
559 "opencode" => project_opencode(&key, payload),
560 _ => project_generic(&key, payload),
561 }
562}
563
564fn project_codex(key: &str, payload: &Value) -> Vec<Value> {
565 if key == "turn/started" {
566 return vec![native_payload(key, payload)];
567 }
568 if key == "turn/completed" {
569 let status = payload
570 .pointer("/params/turn/status")
571 .or_else(|| payload.pointer("/turn/status"))
572 .and_then(Value::as_str)
573 .unwrap_or("completed")
574 .to_ascii_lowercase();
575 return completion(
576 (status.contains("fail") || status.contains("error") || status.contains("cancel"))
577 .then(|| extract_text(payload).unwrap_or_else(|| status.clone())),
578 );
579 }
580 if key.ends_with("/delta") {
581 let text = payload
582 .pointer("/params/delta")
583 .or_else(|| payload.get("delta"))
584 .and_then(extract_text)
585 .or_else(|| extract_text(payload))
586 .unwrap_or_default();
587 return vec![
588 json!({"type":if key.contains("reasoning") { "reasoning" } else { "text_delta" }, "text":text, "raw":payload}),
589 ];
590 }
591 if key.contains("commandexecution")
592 || key.contains("mcptool")
593 || key.contains("filechange")
594 || key.contains("tool")
595 {
596 let source = payload
597 .pointer("/params/item")
598 .or_else(|| payload.get("item"))
599 .unwrap_or(payload);
600 return vec![project_tool(source, payload)];
601 }
602 vec![native_payload(key, payload)]
603}
604
605fn project_claude(key: &str, payload: &Value) -> Vec<Value> {
606 if key == "result" {
607 let failed = payload.get("is_error").and_then(Value::as_bool) == Some(true)
608 || payload.get("subtype").and_then(Value::as_str) == Some("error");
609 return completion(
610 failed.then(|| extract_text(payload).unwrap_or_else(|| "Claude turn failed.".into())),
611 );
612 }
613 if key == "stream_event" {
614 let stream = payload
615 .get("event")
616 .or_else(|| payload.get("stream_event"))
617 .unwrap_or(payload);
618 if stream.get("type").and_then(Value::as_str) == Some("content_block_delta") {
619 let delta = stream.get("delta").unwrap_or(stream);
620 let reasoning = delta
621 .get("type")
622 .and_then(Value::as_str)
623 .is_some_and(|kind| kind.contains("thinking"));
624 return vec![
625 json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":extract_text(delta).unwrap_or_default(), "raw":payload}),
626 ];
627 }
628 }
629 if key == "assistant" {
630 let content = payload
631 .pointer("/message/content")
632 .or_else(|| payload.get("content"))
633 .unwrap_or(payload);
634 return vec![
635 json!({"type":"text_delta", "text":extract_text(content).unwrap_or_default(), "raw":payload}),
636 ];
637 }
638 project_generic(key, payload)
639}
640
641fn project_pi(key: &str, payload: &Value) -> Vec<Value> {
642 match key {
643 "agent_start" => vec![native_payload(key, payload)],
644 "agent_end" => completion(payload.get("error").and_then(extract_text)),
645 "message_update" => {
646 let update = payload
647 .get("assistantMessageEvent")
648 .or_else(|| payload.get("event"))
649 .unwrap_or(payload);
650 let reasoning = update
651 .get("type")
652 .and_then(Value::as_str)
653 .is_some_and(|kind| kind.contains("thinking"));
654 vec![
655 json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":extract_text(update).unwrap_or_default(), "raw":payload}),
656 ]
657 }
658 _ if key.starts_with("tool_execution_") => vec![project_tool(payload, payload)],
659 _ => vec![native_payload(key, payload)],
660 }
661}
662
663fn project_opencode(key: &str, payload: &Value) -> Vec<Value> {
664 if key == "session.idle" {
665 return completion(None);
666 }
667 if key == "session.status" {
668 let status = payload
669 .pointer("/properties/status/type")
670 .or_else(|| payload.pointer("/status/type"))
671 .and_then(Value::as_str)
672 .unwrap_or_default();
673 if status == "busy" {
674 return vec![native_payload(key, payload)];
675 }
676 if status == "idle" {
677 return completion(None);
678 }
679 }
680 if key == "message.part.updated" {
681 let part = payload
682 .pointer("/properties/part")
683 .or_else(|| payload.get("part"))
684 .unwrap_or(payload);
685 if let Some(delta) = payload
686 .pointer("/properties/delta")
687 .or_else(|| payload.get("delta"))
688 .and_then(Value::as_str)
689 {
690 let reasoning = part
691 .get("type")
692 .and_then(Value::as_str)
693 .is_some_and(|kind| kind.contains("reasoning"));
694 return vec![
695 json!({"type":if reasoning { "reasoning" } else { "text_delta" }, "text":delta, "raw":payload}),
696 ];
697 }
698 if part
699 .get("type")
700 .and_then(Value::as_str)
701 .is_some_and(|kind| kind.contains("tool"))
702 {
703 return vec![project_tool(part, payload)];
704 }
705 }
706 if key == "session.error" {
707 return completion(Some(
708 extract_text(payload).unwrap_or_else(|| "OpenCode session failed.".into()),
709 ));
710 }
711 vec![native_payload(key, payload)]
712}
713
714fn project_generic(key: &str, payload: &Value) -> Vec<Value> {
715 match key {
716 "turn_started" | "turn/started" | "agent_start" => {
717 vec![native_payload(key, payload)]
718 }
719 "turn_completed" | "turn/completed" | "agent_end" => {
720 completion(payload.get("error").and_then(extract_text))
721 }
722 "output_delta" | "content_delta" => vec![
723 json!({"type":"text_delta", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
724 ],
725 "reasoning_delta" => vec![
726 json!({"type":"reasoning", "text":extract_text(payload).unwrap_or_default(), "raw":payload}),
727 ],
728 "tool" => vec![project_tool(payload, payload)],
729 _ => vec![native_payload(key, payload)],
730 }
731}
732
733fn completion(error: Option<String>) -> Vec<Value> {
734 match error {
735 Some(message) => vec![
736 json!({"type":"turn_failed", "message":message}),
737 json!({"type":"turn_completed"}),
738 ],
739 None => vec![
740 json!({"type":"turn_succeeded"}),
741 json!({"type":"turn_completed"}),
742 ],
743 }
744}
745
746fn project_tool(source: &Value, raw: &Value) -> Value {
747 let status = source
748 .get("status")
749 .or_else(|| source.get("state"))
750 .or_else(|| source.get("sessionUpdate"))
751 .and_then(Value::as_str)
752 .unwrap_or_default()
753 .to_ascii_lowercase();
754 let completed = status.contains("complete")
755 || status.contains("result")
756 || status.contains("success")
757 || status.contains("error")
758 || status.contains("fail");
759 let arguments = source
760 .get("arguments")
761 .or_else(|| source.get("input"))
762 .or_else(|| source.get("rawInput"))
763 .cloned()
764 .unwrap_or(Value::Null);
765 json!({
766 "type": if completed { "tool_call_completed" } else { "tool_call_started" },
767 "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),
768 "name": source.get("title").or_else(|| source.get("name")).or_else(|| source.get("tool")).or_else(|| source.get("toolName")).cloned().unwrap_or(Value::Null),
769 "arguments": if arguments.is_string() { arguments } else { Value::String(arguments.to_string()) },
770 "output": extract_text(source.get("result").or_else(|| source.get("output")).or_else(|| source.get("content")).unwrap_or(&Value::Null)).unwrap_or_default(),
771 "is_error": status.contains("error") || status.contains("fail") || status.contains("denied"),
772 "raw": raw,
773 })
774}
775
776fn native_payload(kind: &str, payload: &Value) -> Value {
777 json!({"type":"native_event", "kind":kind, "raw":payload})
778}
779
780fn extract_text(value: &Value) -> Option<String> {
781 match value {
782 Value::String(text) => Some(text.clone()),
783 Value::Array(values) => {
784 let text = values
785 .iter()
786 .filter_map(extract_text)
787 .collect::<Vec<_>>()
788 .join("\n");
789 (!text.is_empty()).then_some(text)
790 }
791 Value::Object(object) => {
792 for key in ["text", "delta", "content", "message", "result", "error"] {
793 if let Some(text) = object.get(key).and_then(Value::as_str) {
794 return Some(text.to_string());
795 }
796 }
797 for key in [
798 "delta",
799 "content",
800 "message",
801 "error",
802 "data",
803 "part",
804 "params",
805 "properties",
806 "update",
807 "event",
808 ] {
809 if let Some(text) = object.get(key).and_then(extract_text) {
810 return Some(text);
811 }
812 }
813 None
814 }
815 _ => None,
816 }
817}
818
819#[cfg(test)]
820mod tests {
821 use super::*;
822 use crate::{HarnessId, RuntimeEndpoint};
823
824 struct ControlledRuntime {
825 handle: RuntimeHandle,
826 events: mpsc::UnboundedReceiver<HarnessEvent>,
827 }
828
829 #[async_trait]
830 impl RuntimeConnection for ControlledRuntime {
831 fn handle(&self) -> &RuntimeHandle {
832 &self.handle
833 }
834
835 async fn send_input(&mut self, _input: RuntimeInput) -> Result<Option<String>> {
836 Ok(Some("turn-1".into()))
837 }
838
839 async fn next_event(&mut self) -> Result<Option<HarnessEvent>> {
840 Ok(self.events.recv().await)
841 }
842
843 async fn interrupt(&mut self) -> Result<()> {
844 Ok(())
845 }
846
847 async fn respond(&mut self, _request_id: Value, _response: Value) -> Result<()> {
848 Ok(())
849 }
850
851 async fn close(&mut self) -> Result<()> {
852 Ok(())
853 }
854 }
855
856 fn controlled_runtime() -> (
857 Box<dyn RuntimeConnection>,
858 mpsc::UnboundedSender<HarnessEvent>,
859 ) {
860 let (events, event_rx) = mpsc::unbounded_channel();
861 (
862 Box::new(ControlledRuntime {
863 handle: RuntimeHandle {
864 harness: HarnessId::from(HarnessId::PI),
865 runtime_id: "shared-runtime".into(),
866 endpoint: RuntimeEndpoint::LocalProcess {
867 pid: None,
868 command: vec!["controlled-runtime".into()],
869 protocol: "test".into(),
870 },
871 },
872 events: event_rx,
873 }),
874 events,
875 )
876 }
877
878 fn capabilities() -> RuntimeCapabilities {
879 RuntimeCapabilities {
880 start_session: true,
881 resume_session: true,
882 attach_existing_process: false,
883 send_input: true,
884 stream_events: true,
885 interrupt: true,
886 respond_to_requests: false,
887 }
888 }
889
890 #[tokio::test]
891 async fn native_eof_closes_the_raw_owner_connection() {
892 let (runtime, events) = controlled_runtime();
893 let (_host, mut connection) = HostedHarnessRuntime::spawn(runtime, capabilities());
894 drop(events);
895
896 let event =
897 tokio::time::timeout(std::time::Duration::from_secs(1), connection.next_event())
898 .await
899 .expect("raw owner should not hang after native EOF")
900 .unwrap()
901 .expect("EOF is projected as an explicit terminal event");
902 assert_eq!(event.kind, "transport_closed");
903 assert_eq!(event.payload["terminal"], true);
904 }
905
906 #[tokio::test]
907 async fn adapters_without_an_operation_route_preserve_the_requested_id() {
908 let (runtime, _events) = controlled_runtime();
909 let (host, _owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
910 let operation_id = "prompt:not-advertised".to_string();
911 let error = FrontendRuntime::invoke(
912 host.as_ref(),
913 crate::FrontendOperationInvocation::Prompt {
914 operation_id: operation_id.clone(),
915 arguments: String::new(),
916 },
917 )
918 .await
919 .unwrap_err();
920
921 assert!(
922 matches!(error, FrontendRuntimeError::UnsupportedOperation(id) if id == operation_id)
923 );
924 }
925
926 #[tokio::test]
927 async fn interrupt_stays_busy_until_the_native_terminal_event() {
928 let (runtime, events) = controlled_runtime();
929 let (host, mut owner) = HostedHarnessRuntime::spawn(runtime, capabilities());
930 let mut terminal = FrontendRuntime::attach(host.as_ref(), 100).await.unwrap();
931
932 assert_eq!(
933 FrontendRuntime::submit(host.as_ref(), "hello".into())
934 .await
935 .unwrap(),
936 "turn-1"
937 );
938 assert_eq!(terminal.next_event().await.unwrap().kind, "user_message");
939 assert_eq!(terminal.next_event().await.unwrap().kind, "turn_started");
940 assert!(FrontendRuntime::interrupt(host.as_ref()).await.unwrap());
941 assert_eq!(
942 FrontendRuntime::describe(host.as_ref())
943 .await
944 .unwrap()
945 .turn_state,
946 FrontendTurnState::Busy
947 );
948 assert!(
949 tokio::time::timeout(std::time::Duration::from_millis(20), terminal.next_event())
950 .await
951 .is_err(),
952 "interrupt acceptance must not manufacture turn completion"
953 );
954
955 events
956 .send(HarnessEvent {
957 sequence: None,
958 kind: "agent_end".into(),
959 payload: json!({}),
960 })
961 .unwrap();
962 assert_eq!(terminal.next_event().await.unwrap().kind, "turn_succeeded");
963 assert_eq!(terminal.next_event().await.unwrap().kind, "turn_completed");
964 assert_eq!(
965 FrontendRuntime::describe(host.as_ref())
966 .await
967 .unwrap()
968 .turn_state,
969 FrontendTurnState::Idle
970 );
971 owner.close().await.unwrap();
972 }
973
974 #[test]
975 fn native_start_events_do_not_duplicate_the_hosted_turn_boundary() {
976 let event = HarnessEvent {
977 sequence: None,
978 kind: "turn/started".into(),
979 payload: json!({"method":"turn/started"}),
980 };
981 let projected = project_native_event(HarnessId::CODEX, &event);
982 assert_eq!(projected.len(), 1);
983 assert_eq!(projected[0]["type"], "native_event");
984 }
985}