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