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