1use std::cell::RefCell;
4use std::collections::{BTreeMap, HashMap, HashSet, VecDeque};
5use std::io::{Read, Write};
6use std::sync::atomic::{AtomicU64, Ordering};
7use std::sync::{Arc, Mutex};
8use std::time::{Duration, Instant};
9
10use crate::ipc_binary::{self, BinaryFrame};
11use crate::runtime_protocol::{BridgeResponse, RuntimeEvent};
12use crate::session::RuntimeEventEnvelope;
13use agentos_runtime::accounting::{Reservation, ResourceClass};
14use agentos_runtime::RuntimeContext;
15
16static SYNCRPC_LAT: std::sync::OnceLock<std::sync::Mutex<(u64, u64, u64)>> =
22 std::sync::OnceLock::new();
23
24fn syncrpc_lat_enabled() -> bool {
25 std::env::var("AGENTOS_SYNCRPC_LAT").as_deref() == Ok("1")
26}
27
28fn record_syncrpc_lat(ns: u64) {
29 let m = SYNCRPC_LAT.get_or_init(|| std::sync::Mutex::new((0, 0, 0)));
30 let Ok(mut a) = m.lock() else {
31 return;
32 };
33 a.0 += 1;
34 a.1 = a.1.wrapping_add(ns / 1000);
35 a.2 = a.2.max(ns / 1000);
36 if a.0 % 25 == 0 {
37 if let Ok(path) = std::env::var("AGENTOS_SYNCRPC_LAT_FILE") {
38 let _ = std::fs::write(
39 &path,
40 format!(
41 "calls={} total_us={} avg_us={} max_us={}\n",
42 a.0,
43 a.1,
44 a.1 / a.0,
45 a.2
46 ),
47 );
48 }
49 }
50}
51
52#[derive(Debug, Default, Clone)]
53struct SyncBridgeHostPhaseStats {
54 calls: u64,
55 total_us: u64,
56 max_us: u64,
57}
58
59static SYNC_BRIDGE_HOST_PHASES: std::sync::OnceLock<
60 std::sync::Mutex<BTreeMap<String, SyncBridgeHostPhaseStats>>,
61> = std::sync::OnceLock::new();
62static SYNC_BRIDGE_CALL_METHODS: std::sync::OnceLock<std::sync::Mutex<HashMap<u64, String>>> =
63 std::sync::OnceLock::new();
64
65fn sync_bridge_host_phases_enabled() -> bool {
66 std::env::var("AGENTOS_SYNC_BRIDGE_HOST_PHASES").as_deref() == Ok("1")
67}
68
69pub(crate) fn record_sync_bridge_host_phase(method: &str, stage: &str, elapsed: Duration) {
70 if !sync_bridge_host_phases_enabled() {
71 return;
72 }
73 let stats = SYNC_BRIDGE_HOST_PHASES.get_or_init(|| std::sync::Mutex::new(BTreeMap::new()));
74 let Ok(mut stats) = stats.lock() else {
75 return;
76 };
77 let elapsed_us = elapsed.as_micros() as u64;
78 let key = format!("{method}:{stage}");
79 let entry = stats.entry(key).or_default();
80 entry.calls += 1;
81 entry.total_us = entry.total_us.wrapping_add(elapsed_us);
82 entry.max_us = entry.max_us.max(elapsed_us);
83
84 if let Ok(path) = std::env::var("AGENTOS_SYNC_BRIDGE_HOST_PHASES_FILE") {
85 let mut lines = String::new();
86 for (key, value) in stats.iter() {
87 let Some((method, stage)) = key.split_once(':') else {
88 continue;
89 };
90 let avg_us = value.total_us.checked_div(value.calls).unwrap_or(0);
91 lines.push_str(&format!(
92 "method={method} stage={stage} calls={} total_us={} avg_us={} max_us={}\n",
93 value.calls, value.total_us, avg_us, value.max_us
94 ));
95 }
96 let _ = std::fs::write(path, lines);
97 }
98}
99
100fn track_sync_bridge_call_method(call_id: u64, method: &str) {
101 if !sync_bridge_host_phases_enabled() {
102 return;
103 }
104 let methods = SYNC_BRIDGE_CALL_METHODS.get_or_init(|| std::sync::Mutex::new(HashMap::new()));
105 let Ok(mut methods) = methods.lock() else {
106 return;
107 };
108 if methods.len() > 4096 {
109 methods.clear();
110 }
111 methods.insert(call_id, method.to_owned());
112}
113
114fn cleanup_sync_bridge_call_tracking(call_id: u64) {
115 if let Some(methods) = SYNC_BRIDGE_CALL_METHODS.get() {
116 if let Ok(mut methods) = methods.lock() {
117 methods.remove(&call_id);
118 }
119 }
120}
121
122pub trait RuntimeEventSender: Send {
125 fn send_event(&self, event: RuntimeEvent) -> Result<(), String>;
126}
127
128pub struct ChannelRuntimeEventSender {
132 pub tx: crate::session::RuntimeEventSender,
133 output_generation: Option<u64>,
134 #[allow(dead_code)]
137 frame_buf: RefCell<Vec<u8>>,
138}
139
140impl ChannelRuntimeEventSender {
141 pub fn new(
142 tx: impl Into<crate::session::RuntimeEventSender>,
143 output_generation: Option<u64>,
144 ) -> Self {
145 ChannelRuntimeEventSender {
146 tx: tx.into(),
147 output_generation,
148 frame_buf: RefCell::new(Vec::with_capacity(256)),
149 }
150 }
151}
152
153impl RuntimeEventSender for ChannelRuntimeEventSender {
154 fn send_event(&self, event: RuntimeEvent) -> Result<(), String> {
155 self.tx
156 .send(RuntimeEventEnvelope {
157 output_generation: self.output_generation,
158 event,
159 })
160 .map_err(|error| format!("runtime event send failed: {error}"))
161 }
162}
163
164#[allow(dead_code)]
166pub struct WriterRuntimeEventSender {
167 writer: Mutex<Box<dyn Write + Send>>,
168}
169
170impl RuntimeEventSender for WriterRuntimeEventSender {
171 fn send_event(&self, event: RuntimeEvent) -> Result<(), String> {
172 let mut w = self.writer.lock().unwrap();
173 let frame: BinaryFrame = event.into();
174 ipc_binary::write_frame(&mut *w, &frame).map_err(|e| format!("write error: {}", e))
175 }
176}
177
178pub trait BridgeResponseReceiver: Send {
181 fn recv_response(&self, expected_call_id: u64) -> Result<BridgeResponse, String>;
182}
183
184#[derive(Debug, Clone, PartialEq, Eq)]
185pub struct SyncBridgeCallResponse {
186 pub status: u8,
187 pub payload: Vec<u8>,
188 pub reservation: Option<agentos_runtime::accounting::SharedReservation>,
189}
190
191#[allow(dead_code)]
194pub struct ReaderBridgeResponseReceiver {
195 reader: Mutex<Box<dyn Read + Send>>,
196}
197
198impl ReaderBridgeResponseReceiver {
199 #[allow(dead_code)]
200 pub fn new(reader: Box<dyn Read + Send>) -> Self {
201 ReaderBridgeResponseReceiver {
202 reader: Mutex::new(reader),
203 }
204 }
205}
206
207impl BridgeResponseReceiver for ReaderBridgeResponseReceiver {
208 fn recv_response(&self, expected_call_id: u64) -> Result<BridgeResponse, String> {
209 let mut reader = self.reader.lock().unwrap();
210 let frame = ipc_binary::read_frame(&mut *reader)
211 .map_err(|e| format!("failed to read BridgeResponse: {}", e))?;
212 match frame {
213 BinaryFrame::BridgeResponse {
214 call_id,
215 status,
216 payload,
217 ..
218 } => {
219 if call_id != expected_call_id {
220 return Err(format!(
221 "call_id mismatch: expected {}, got {}",
222 expected_call_id, call_id
223 ));
224 }
225 Ok(BridgeResponse {
226 call_id,
227 status,
228 payload,
229 reservation: None,
230 })
231 }
232 _ => Err("expected BridgeResponse, got different message type".into()),
233 }
234 }
235}
236
237const MAX_PENDING_BRIDGE_CALLS: usize = 16_384;
238pub(crate) const DEFAULT_BRIDGE_CALL_TIMEOUT: Duration = Duration::from_secs(30);
239
240#[derive(Clone, Copy, Debug, Eq, PartialEq)]
241enum BridgeCallTargetKind {
242 Sync,
243 Async,
244}
245
246struct BridgeCallTarget {
247 session_id: String,
248 session_generation: Option<u64>,
249 kind: BridgeCallTargetKind,
250 sender: crossbeam_channel::Sender<BridgeResponse>,
251 _call_reservation: Reservation,
252 _request_reservation: Reservation,
253 response_resources: Arc<agentos_runtime::accounting::ResourceLedger>,
254 response_reservation: Reservation,
255 max_response_bytes: usize,
256 deadline: Instant,
257 timeout: Duration,
258 _deadline_cancellation: tokio::sync::oneshot::Sender<()>,
259 host_visible: bool,
260}
261
262const BRIDGE_TERMINAL_RESPONSE_RESERVATION_BYTES: usize = 4 * 1024;
263
264fn grow_response_reservation(
265 target: &mut BridgeCallTarget,
266 payload_bytes: usize,
267) -> Result<(), String> {
268 let additional = payload_bytes.saturating_sub(target.response_reservation.amount());
269 if additional == 0 {
270 return Ok(());
271 }
272 let extra = target
273 .response_resources
274 .reserve(ResourceClass::BridgeResponseBytes, additional)
275 .map_err(|error| error.to_string())?;
276 target.response_reservation.merge(extra).map_err(|_| {
277 String::from(
278 "ERR_AGENTOS_BRIDGE_RESPONSE_ACCOUNTING: failed to merge response-byte reservations",
279 )
280 })
281}
282
283fn transfer_response_reservation(
284 mut reservation: Reservation,
285 payload_bytes: usize,
286) -> agentos_runtime::accounting::SharedReservation {
287 let unused = reservation
288 .amount()
289 .checked_sub(payload_bytes)
290 .expect("response payload must fit its admission reservation");
291 if unused != 0 {
292 drop(
293 reservation
294 .split(unused)
295 .expect("unused response capacity must remain transferable"),
296 );
297 }
298 agentos_runtime::accounting::SharedReservation::new(reservation)
299}
300
301fn bounded_terminal_error_payload(error: &str, maximum_bytes: usize) -> Vec<u8> {
302 if error.len() <= maximum_bytes {
303 return error.as_bytes().to_vec();
304 }
305
306 let code = error.split_once(':').map_or(error, |(code, _)| code);
310 if code.len() <= maximum_bytes {
311 return code.as_bytes().to_vec();
312 }
313 code.as_bytes()[..maximum_bytes.min(code.len())].to_vec()
314}
315
316fn bridge_call_timeout_error(call_id: u64, timeout: Duration) -> String {
317 format!(
318 "ERR_AGENTOS_BRIDGE_CALL_TIMEOUT: bridge call_id {call_id} exceeded its {} ms deadline; raise limits.reactor.operationDeadlineMs",
319 timeout.as_millis()
320 )
321}
322
323fn deliver_bridge_call_timeout(call_id: u64, target: BridgeCallTarget) -> Result<String, String> {
324 let error = bridge_call_timeout_error(call_id, target.timeout);
325 let payload = bounded_terminal_error_payload(
326 &error,
327 target
328 .max_response_bytes
329 .min(target.response_reservation.amount()),
330 );
331 let reservation = transfer_response_reservation(target.response_reservation, payload.len());
332 target
333 .sender
334 .try_send(BridgeResponse {
335 call_id,
336 status: 1,
337 payload,
338 reservation: Some(reservation),
339 })
340 .map_err(|delivery_error| {
341 format!(
342 "ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY: failed to deliver timeout for call_id {call_id}: {delivery_error}; original error: {error}"
343 )
344 })?;
345 Ok(error)
346}
347
348#[derive(Clone, Debug, Eq, PartialEq)]
349struct RetiredBridgeCall {
350 session_id: String,
351 session_generation: Option<u64>,
352 retirement_epoch: u64,
353}
354
355#[derive(Default)]
356struct RetiredBridgeCalls {
357 by_call_id: HashMap<u64, RetiredBridgeCall>,
358 order: VecDeque<(u64, u64)>,
359 next_epoch: u64,
360}
361
362impl RetiredBridgeCalls {
363 fn insert(&mut self, call_id: u64, target: &BridgeCallTarget, limit: usize) {
364 while self.order.len() >= limit {
365 let Some((oldest_call_id, oldest_epoch)) = self.order.pop_front() else {
366 break;
367 };
368 if self
369 .by_call_id
370 .get(&oldest_call_id)
371 .is_some_and(|retired| retired.retirement_epoch == oldest_epoch)
372 {
373 self.by_call_id.remove(&oldest_call_id);
374 }
375 }
376
377 self.next_epoch = self.next_epoch.wrapping_add(1);
378 let retirement_epoch = self.next_epoch;
379 self.by_call_id.insert(
380 call_id,
381 RetiredBridgeCall {
382 session_id: target.session_id.clone(),
383 session_generation: target.session_generation,
384 retirement_epoch,
385 },
386 );
387 self.order.push_back((call_id, retirement_epoch));
388 }
389
390 fn remove(&mut self, call_id: u64) {
391 self.by_call_id.remove(&call_id);
392 }
393}
394
395pub struct BridgeCallRegistry {
400 pending: Mutex<HashMap<u64, BridgeCallTarget>>,
401 retired: Mutex<RetiredBridgeCalls>,
405 max_pending: usize,
406}
407
408impl BridgeCallRegistry {
409 pub fn new(max_pending: usize) -> Self {
410 Self {
411 pending: Mutex::new(HashMap::new()),
412 retired: Mutex::new(RetiredBridgeCalls::default()),
413 max_pending: max_pending.max(1),
414 }
415 }
416
417 pub fn with_default_limit() -> Self {
418 Self::new(MAX_PENDING_BRIDGE_CALLS)
419 }
420
421 #[allow(clippy::too_many_arguments)] fn register(
423 &self,
424 runtime: &RuntimeContext,
425 request_bytes: usize,
426 max_response_bytes: usize,
427 call_id: u64,
428 session_id: &str,
429 session_generation: Option<u64>,
430 kind: BridgeCallTargetKind,
431 sender: crossbeam_channel::Sender<BridgeResponse>,
432 timeout: Duration,
433 ) -> Result<tokio::sync::oneshot::Receiver<()>, String> {
434 if timeout.is_zero() {
435 return Err(String::from(
436 "ERR_AGENTOS_BRIDGE_CALL_TIMEOUT_INVALID: limits.reactor.operationDeadlineMs must be greater than zero",
437 ));
438 }
439 let call_reservation = runtime
440 .resources()
441 .reserve(ResourceClass::BridgeCalls, 1)
442 .map_err(|error| error.to_string())?;
443 let request_reservation = runtime
444 .resources()
445 .reserve(ResourceClass::BridgeRequestBytes, request_bytes)
446 .map_err(|error| error.to_string())?;
447 let configured_response_limit = runtime
448 .resources()
449 .usage(ResourceClass::BridgeResponseBytes)
450 .limit
451 .ok_or_else(|| {
452 String::from(
453 "ERR_AGENTOS_RESOURCE_UNBOUNDED: bridge response bytes require limits.reactor.maxBridgeResponseBytes",
454 )
455 })?;
456 if max_response_bytes > configured_response_limit {
457 return Err(format!(
458 "ERR_AGENTOS_BRIDGE_RESPONSE_LIMIT: declared response maximum of {max_response_bytes} bytes exceeds configured limit of {configured_response_limit}; raise limits.reactor.maxBridgeResponseBytes"
459 ));
460 }
461 let admission_response_bytes =
468 max_response_bytes.min(BRIDGE_TERMINAL_RESPONSE_RESERVATION_BYTES);
469 let response_reservation = runtime
470 .resources()
471 .reserve(ResourceClass::BridgeResponseBytes, admission_response_bytes)
472 .map_err(|error| error.to_string())?;
473 let mut pending = self.pending.lock().map_err(|_| {
474 String::from("ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: bridge registry lock poisoned")
475 })?;
476 if pending.len() >= self.max_pending {
477 return Err(format!(
478 "ERR_AGENTOS_BRIDGE_CALL_LIMIT: bridge call registry exceeded limit of {} pending calls; raise runtime.resources.maxBridgeCalls",
479 self.max_pending
480 ));
481 }
482 if kind == BridgeCallTargetKind::Async {
483 let response_capacity = sender.capacity().ok_or_else(|| {
484 String::from(
485 "ERR_AGENTOS_BRIDGE_RESPONSE_LANE_UNBOUNDED: async bridge response lanes must be bounded",
486 )
487 })?;
488 let queued_responses = sender.len();
489 let registered_targets = pending
490 .values()
491 .filter(|target| {
492 target.kind == BridgeCallTargetKind::Async
493 && target.sender.same_channel(&sender)
494 })
495 .count();
496 let occupied = queued_responses.saturating_add(registered_targets);
497 if occupied >= response_capacity {
498 return Err(format!(
499 "ERR_AGENTOS_BRIDGE_RESPONSE_LANE_LIMIT: async response lane has {queued_responses} queued responses and {registered_targets} registered calls at capacity {response_capacity}; admission must reserve a response slot before host visibility"
500 ));
501 }
502 }
503 if pending.contains_key(&call_id) {
504 return Err(format!(
505 "ERR_AGENTOS_BRIDGE_DUPLICATE_CALL_ID: duplicate bridge call_id {call_id}"
506 ));
507 }
508 self.retired
509 .lock()
510 .map_err(|_| {
511 String::from(
512 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: retired bridge registry lock poisoned",
513 )
514 })?
515 .remove(call_id);
516 let (deadline_cancellation, deadline_cancelled) = tokio::sync::oneshot::channel();
517 pending.insert(
518 call_id,
519 BridgeCallTarget {
520 session_id: session_id.to_owned(),
521 session_generation,
522 kind,
523 sender,
524 _call_reservation: call_reservation,
525 _request_reservation: request_reservation,
526 response_resources: Arc::clone(runtime.resources()),
527 response_reservation,
528 max_response_bytes,
529 deadline: Instant::now()
530 .checked_add(timeout)
531 .unwrap_or_else(Instant::now),
532 timeout,
533 _deadline_cancellation: deadline_cancellation,
534 host_visible: false,
535 },
536 );
537 Ok(deadline_cancelled)
538 }
539
540 fn mark_host_visible(&self, call_id: u64) -> Result<(), String> {
543 let mut pending = self.pending.lock().map_err(|_| {
544 String::from("ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: bridge registry lock poisoned")
545 })?;
546 let target = pending.get_mut(&call_id).ok_or_else(|| {
547 format!(
548 "ERR_AGENTOS_BRIDGE_ROUTE_RETIRED: bridge call_id {call_id} was canceled before host publication"
549 )
550 })?;
551 target.host_visible = true;
552 Ok(())
553 }
554
555 pub fn register_sync(
556 &self,
557 runtime: &RuntimeContext,
558 request_bytes: usize,
559 max_response_bytes: usize,
560 call_id: u64,
561 session_id: &str,
562 session_generation: Option<u64>,
563 ) -> Result<crossbeam_channel::Receiver<BridgeResponse>, String> {
564 self.register_sync_with_timeout(
565 runtime,
566 request_bytes,
567 max_response_bytes,
568 call_id,
569 session_id,
570 session_generation,
571 DEFAULT_BRIDGE_CALL_TIMEOUT,
572 )
573 }
574
575 #[allow(clippy::too_many_arguments)]
576 fn register_sync_with_timeout(
577 &self,
578 runtime: &RuntimeContext,
579 request_bytes: usize,
580 max_response_bytes: usize,
581 call_id: u64,
582 session_id: &str,
583 session_generation: Option<u64>,
584 timeout: Duration,
585 ) -> Result<crossbeam_channel::Receiver<BridgeResponse>, String> {
586 let (sender, receiver) = crossbeam_channel::bounded(1);
587 let _deadline_cancelled = self.register(
588 runtime,
589 request_bytes,
590 max_response_bytes,
591 call_id,
592 session_id,
593 session_generation,
594 BridgeCallTargetKind::Sync,
595 sender,
596 timeout,
597 )?;
598 Ok(receiver)
599 }
600
601 #[allow(clippy::too_many_arguments)] pub fn register_async(
603 &self,
604 runtime: &RuntimeContext,
605 request_bytes: usize,
606 max_response_bytes: usize,
607 call_id: u64,
608 session_id: &str,
609 session_generation: Option<u64>,
610 sender: crossbeam_channel::Sender<BridgeResponse>,
611 ) -> Result<(), String> {
612 self.register(
613 runtime,
614 request_bytes,
615 max_response_bytes,
616 call_id,
617 session_id,
618 session_generation,
619 BridgeCallTargetKind::Async,
620 sender,
621 DEFAULT_BRIDGE_CALL_TIMEOUT,
622 )
623 .map(drop)
624 }
625
626 #[allow(clippy::too_many_arguments)]
627 fn register_async_with_timeout(
628 &self,
629 runtime: &RuntimeContext,
630 request_bytes: usize,
631 max_response_bytes: usize,
632 call_id: u64,
633 session_id: &str,
634 session_generation: Option<u64>,
635 sender: crossbeam_channel::Sender<BridgeResponse>,
636 timeout: Duration,
637 ) -> Result<tokio::sync::oneshot::Receiver<()>, String> {
638 self.register(
639 runtime,
640 request_bytes,
641 max_response_bytes,
642 call_id,
643 session_id,
644 session_generation,
645 BridgeCallTargetKind::Async,
646 sender,
647 timeout,
648 )
649 }
650
651 fn timeout(&self, call_id: u64) -> Result<bool, String> {
652 let mut pending = self.pending.lock().map_err(|_| {
653 String::from("ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: bridge registry lock poisoned")
654 })?;
655 if pending
656 .get(&call_id)
657 .is_none_or(|target| Instant::now() < target.deadline)
658 {
659 return Ok(false);
660 }
661 let target = pending
662 .remove(&call_id)
663 .expect("expired bridge target must remain registered while locked");
664 if target.host_visible {
665 self.retired
666 .lock()
667 .map_err(|_| {
668 String::from(
669 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: retired bridge registry lock poisoned",
670 )
671 })?
672 .insert(call_id, &target, self.max_pending);
673 }
674 deliver_bridge_call_timeout(call_id, target)?;
675 Ok(true)
676 }
677
678 pub fn settle(
679 &self,
680 supplied_session_id: &str,
681 supplied_generation: Option<u64>,
682 mut response: BridgeResponse,
683 ) -> Result<(), String> {
684 let call_id = response.call_id;
685 let mut pending = self.pending.lock().map_err(|_| {
686 String::from("ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: bridge registry lock poisoned")
687 })?;
688 let target = match pending.get(&response.call_id) {
689 Some(target) => target,
690 None => {
691 let retired = self.retired.lock().map_err(|_| {
692 String::from(
693 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: retired bridge registry lock poisoned",
694 )
695 })?;
696 if let Some(retired) = retired.by_call_id.get(&response.call_id) {
697 if retired.session_id == supplied_session_id
698 && retired.session_generation == supplied_generation
699 {
700 return Err(format!(
701 "ERR_AGENTOS_BRIDGE_STALE_COMPLETION: response for canceled host-visible bridge call_id {} in session {} generation {:?}",
702 response.call_id, supplied_session_id, supplied_generation
703 ));
704 }
705 return Err(format!(
706 "ERR_AGENTOS_BRIDGE_STALE_GENERATION: response call_id {} named session {} generation {:?}, expected {} generation {:?}",
707 response.call_id,
708 supplied_session_id,
709 supplied_generation,
710 retired.session_id,
711 retired.session_generation
712 ));
713 }
714 return Err(format!(
715 "ERR_AGENTOS_BRIDGE_UNKNOWN_CALL_ID: response for unknown bridge call_id {}",
716 response.call_id
717 ));
718 }
719 };
720 if target.session_id != supplied_session_id {
721 return Err(format!(
722 "ERR_AGENTOS_BRIDGE_STALE_GENERATION: response call_id {} named session {}, expected {}",
723 response.call_id, supplied_session_id, target.session_id
724 ));
725 }
726 if target.session_generation != supplied_generation {
727 return Err(format!(
728 "ERR_AGENTOS_BRIDGE_STALE_GENERATION: response call_id {} generation {:?}, expected {:?}",
729 response.call_id, supplied_generation, target.session_generation
730 ));
731 }
732 let deadline_expired = Instant::now() >= target.deadline;
740 let mut target = pending
741 .remove(&call_id)
742 .expect("validated bridge target must remain registered while locked");
743
744 if deadline_expired {
745 if target.host_visible {
746 self.retired
747 .lock()
748 .map_err(|_| {
749 String::from(
750 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: retired bridge registry lock poisoned",
751 )
752 })?
753 .insert(call_id, &target, self.max_pending);
754 }
755 let error = deliver_bridge_call_timeout(call_id, target)?;
756 return Err(error);
757 }
758
759 if response.payload.len() > target.max_response_bytes {
760 let error = format!(
761 "ERR_AGENTOS_BRIDGE_RESPONSE_LIMIT: response call_id {} contains {} bytes, exceeding its declared maximum of {}; raise limits.reactor.maxBridgeResponseBytes",
762 response.call_id,
763 response.payload.len(),
764 target.max_response_bytes
765 );
766 let payload = bounded_terminal_error_payload(
767 &error,
768 target
769 .max_response_bytes
770 .min(target.response_reservation.amount()),
771 );
772 let reservation =
773 transfer_response_reservation(target.response_reservation, payload.len());
774 target
775 .sender
776 .try_send(BridgeResponse {
777 call_id: response.call_id,
778 status: 1,
779 payload,
780 reservation: Some(reservation),
781 })
782 .map_err(|delivery_error| {
783 format!(
784 "ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY: failed to deliver oversize response error for call_id {}: {delivery_error}; original error: {error}",
785 response.call_id
786 )
787 })?;
788 return Err(error);
789 }
790
791 if let Some(reservation) = response.reservation.take() {
792 let error = format!(
793 "ERR_AGENTOS_BRIDGE_RESPONSE_ACCOUNTING: response call_id {} carries a producer-side {:?} reservation of {} bytes; bridge response ownership must come from the call's admission reservation",
794 response.call_id,
795 reservation.resource(),
796 reservation.amount()
797 );
798 let payload = bounded_terminal_error_payload(
799 &error,
800 target
801 .max_response_bytes
802 .min(target.response_reservation.amount()),
803 );
804 let response_reservation =
805 transfer_response_reservation(target.response_reservation, payload.len());
806 target
807 .sender
808 .try_send(BridgeResponse {
809 call_id: response.call_id,
810 status: 1,
811 payload,
812 reservation: Some(response_reservation),
813 })
814 .map_err(|delivery_error| {
815 format!(
816 "ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY: failed to deliver accounting error for call_id {}: {delivery_error}; original error: {error}",
817 response.call_id
818 )
819 })?;
820 return Err(error);
821 }
822
823 if let Err(limit_error) = grow_response_reservation(&mut target, response.payload.len()) {
824 let error = format!(
825 "ERR_AGENTOS_BRIDGE_RESPONSE_LIMIT: response call_id {} could not reserve {} concrete bytes: {}; raise limits.reactor.maxBridgeResponseBytes",
826 response.call_id,
827 response.payload.len(),
828 limit_error
829 );
830 let payload = bounded_terminal_error_payload(
831 &error,
832 target
833 .max_response_bytes
834 .min(target.response_reservation.amount()),
835 );
836 let response_reservation =
837 transfer_response_reservation(target.response_reservation, payload.len());
838 target
839 .sender
840 .try_send(BridgeResponse {
841 call_id: response.call_id,
842 status: 1,
843 payload,
844 reservation: Some(response_reservation),
845 })
846 .map_err(|delivery_error| {
847 format!(
848 "ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY: failed to deliver response-capacity error for call_id {}: {delivery_error}; original error: {error}",
849 response.call_id
850 )
851 })?;
852 return Err(error);
853 }
854
855 response.reservation = Some(transfer_response_reservation(
856 target.response_reservation,
857 response.payload.len(),
858 ));
859
860 target.sender.try_send(response).map_err(|error| {
861 format!(
862 "ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY: {:?} response target for session {} generation {:?} rejected settlement: {error}",
863 target.kind, target.session_id, target.session_generation
864 )
865 })?;
866 Ok(())
867 }
868
869 fn cancel_with_visibility(&self, call_id: u64, force_unpublished: bool) {
870 match self.pending.lock() {
871 Ok(mut pending) => {
872 let target = pending.remove(&call_id);
873 match self.retired.lock() {
874 Ok(mut retired) => {
875 if force_unpublished {
876 retired.remove(call_id);
881 } else if let Some(target) =
882 target.as_ref().filter(|target| target.host_visible)
883 {
884 retired.insert(call_id, target, self.max_pending);
885 }
886 }
887 Err(_) => eprintln!(
888 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: could not update retirement for bridge call_id {call_id}"
889 ),
890 }
891 }
892 Err(_) => eprintln!(
893 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: could not cancel bridge call_id {call_id}"
894 ),
895 }
896 }
897
898 pub fn cancel(&self, call_id: u64) {
899 self.cancel_with_visibility(call_id, false);
900 }
901
902 fn cancel_unpublished(&self, call_id: u64) {
903 self.cancel_with_visibility(call_id, true);
904 }
905
906 pub fn cancel_session(&self, session_id: &str, session_generation: Option<u64>) {
907 match self.pending.lock() {
908 Ok(mut pending) => {
909 let canceled_call_ids = pending
910 .iter()
911 .filter_map(|(call_id, target)| {
912 (target.session_id == session_id
913 && (session_generation.is_none()
914 || target.session_generation == session_generation))
915 .then_some(*call_id)
916 })
917 .collect::<Vec<_>>();
918 let mut retired = match self.retired.lock() {
919 Ok(retired) => retired,
920 Err(_) => {
921 eprintln!(
922 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: could not retire bridge calls for session {session_id} generation {session_generation:?}"
923 );
924 pending.retain(|_, target| {
925 target.session_id != session_id
926 || (session_generation.is_some()
927 && target.session_generation != session_generation)
928 });
929 return;
930 }
931 };
932 for call_id in canceled_call_ids {
933 if let Some(target) = pending.remove(&call_id) {
934 if target.host_visible {
935 retired.insert(call_id, &target, self.max_pending);
936 }
937 }
938 }
939 }
940 Err(_) => eprintln!(
941 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: could not cancel bridge calls for session {session_id} generation {session_generation:?}"
942 ),
943 }
944 }
945
946 pub fn clear(&self) {
947 match self.pending.lock() {
948 Ok(mut pending) => match self.retired.lock() {
949 Ok(mut retired) => {
950 for (call_id, target) in pending.drain() {
951 if target.host_visible {
952 retired.insert(call_id, &target, self.max_pending);
953 }
954 }
955 }
956 Err(_) => {
957 eprintln!(
958 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: could not retire cleared bridge calls"
959 );
960 pending.clear();
961 }
962 },
963 Err(_) => eprintln!(
964 "ERR_AGENTOS_BRIDGE_REGISTRY_POISONED: could not clear bridge call registry"
965 ),
966 }
967 }
968
969 #[cfg(test)]
970 pub fn pending_len(&self) -> usize {
971 self.pending
972 .lock()
973 .map(|pending| pending.len())
974 .unwrap_or(0)
975 }
976
977 #[cfg(test)]
978 pub fn retired_len(&self) -> usize {
979 self.retired
980 .lock()
981 .map(|retired| retired.by_call_id.len())
982 .unwrap_or(0)
983 }
984}
985
986pub type CallIdRouter = Arc<BridgeCallRegistry>;
989
990pub type SharedCallIdCounter = Arc<AtomicU64>;
994
995struct PendingBridgeRoute {
999 registry: CallIdRouter,
1000 call_id: Option<u64>,
1001}
1002
1003impl PendingBridgeRoute {
1004 fn new(registry: CallIdRouter, call_id: u64) -> Self {
1005 Self {
1006 registry,
1007 call_id: Some(call_id),
1008 }
1009 }
1010
1011 fn disarm(mut self) {
1012 self.call_id = None;
1013 }
1014
1015 fn mark_host_visible(&self) -> Result<(), String> {
1016 if let Some(call_id) = self.call_id {
1017 self.registry.mark_host_visible(call_id)
1018 } else {
1019 Ok(())
1020 }
1021 }
1022
1023 fn cancel_unpublished(mut self) {
1024 if let Some(call_id) = self.call_id.take() {
1025 self.registry.cancel_unpublished(call_id);
1026 }
1027 }
1028}
1029
1030impl Drop for PendingBridgeRoute {
1031 fn drop(&mut self) {
1032 if let Some(call_id) = self.call_id.take() {
1033 self.registry.cancel(call_id);
1034 }
1035 }
1036}
1037
1038pub struct BridgeCallContext {
1044 sender: Box<dyn RuntimeEventSender>,
1046 response_rx: Option<Mutex<Box<dyn BridgeResponseReceiver>>>,
1048 pub session_id: String,
1050 next_call_id: Arc<AtomicU64>,
1053 pending_calls: Mutex<HashSet<u64>>,
1055 track_pending_calls: bool,
1059 call_id_router: Option<CallIdRouter>,
1061 session_generation: Option<u64>,
1062 async_response_tx: Option<crossbeam_channel::Sender<BridgeResponse>>,
1063 abort_rx: Option<crossbeam_channel::Receiver<()>>,
1064 pause_control: Option<Arc<crate::session::SessionPauseControl>>,
1068 runtime: Option<agentos_runtime::RuntimeContext>,
1072 bridge_call_timeout: Duration,
1073}
1074
1075pub struct PreparedAsyncBridgeCall {
1076 pub(crate) call_id: u64,
1077 event: RuntimeEvent,
1078 pending_route: Option<PendingBridgeRoute>,
1079}
1080
1081#[allow(dead_code)]
1084struct StubRuntimeEventSender;
1085
1086impl RuntimeEventSender for StubRuntimeEventSender {
1087 fn send_event(&self, _event: RuntimeEvent) -> Result<(), String> {
1088 panic!(
1089 "stub bridge function called during snapshot creation — bridge IIFE must not call bridge functions at setup time"
1090 )
1091 }
1092}
1093
1094#[allow(dead_code)]
1097struct StubBridgeResponseReceiver;
1098
1099impl BridgeResponseReceiver for StubBridgeResponseReceiver {
1100 fn recv_response(&self, _expected_call_id: u64) -> Result<BridgeResponse, String> {
1101 panic!(
1102 "stub bridge function called during snapshot creation — bridge IIFE must not call bridge functions at setup time"
1103 )
1104 }
1105}
1106
1107#[allow(dead_code)]
1108impl BridgeCallContext {
1109 pub fn stub() -> Self {
1113 BridgeCallContext {
1114 sender: Box::new(StubRuntimeEventSender),
1115 response_rx: Some(Mutex::new(Box::new(StubBridgeResponseReceiver))),
1116 session_id: "stub".into(),
1117 next_call_id: Arc::new(AtomicU64::new(1)),
1118 pending_calls: Mutex::new(HashSet::new()),
1119 track_pending_calls: should_track_pending_sync_calls(),
1120 call_id_router: None,
1121 session_generation: None,
1122 async_response_tx: None,
1123 abort_rx: None,
1124 pause_control: None,
1125 runtime: None,
1126 bridge_call_timeout: DEFAULT_BRIDGE_CALL_TIMEOUT,
1127 }
1128 }
1129
1130 pub fn new(
1133 writer: Box<dyn Write + Send>,
1134 reader: Box<dyn Read + Send>,
1135 session_id: String,
1136 ) -> Self {
1137 BridgeCallContext {
1138 sender: Box::new(WriterRuntimeEventSender {
1139 writer: Mutex::new(writer),
1140 }),
1141 response_rx: Some(Mutex::new(Box::new(ReaderBridgeResponseReceiver::new(
1142 reader,
1143 )))),
1144 session_id,
1145 next_call_id: Arc::new(AtomicU64::new(1)),
1146 pending_calls: Mutex::new(HashSet::new()),
1147 track_pending_calls: should_track_pending_sync_calls(),
1148 call_id_router: None,
1149 session_generation: None,
1150 async_response_tx: None,
1151 abort_rx: None,
1152 pause_control: None,
1153 runtime: None,
1154 bridge_call_timeout: DEFAULT_BRIDGE_CALL_TIMEOUT,
1155 }
1156 }
1157
1158 pub fn with_receiver(
1162 sender: Box<dyn RuntimeEventSender>,
1163 response_rx: Box<dyn BridgeResponseReceiver>,
1164 session_id: String,
1165 _router: CallIdRouter,
1166 shared_call_id: SharedCallIdCounter,
1167 ) -> Self {
1168 BridgeCallContext {
1169 sender,
1170 response_rx: Some(Mutex::new(response_rx)),
1171 session_id,
1172 next_call_id: shared_call_id,
1173 pending_calls: Mutex::new(HashSet::new()),
1174 track_pending_calls: should_track_pending_sync_calls(),
1175 call_id_router: None,
1176 session_generation: None,
1177 async_response_tx: None,
1178 abort_rx: None,
1179 pause_control: None,
1180 runtime: None,
1181 bridge_call_timeout: DEFAULT_BRIDGE_CALL_TIMEOUT,
1182 }
1183 }
1184
1185 #[allow(clippy::too_many_arguments)]
1186 pub(crate) fn with_registry(
1187 sender: Box<dyn RuntimeEventSender>,
1188 session_id: String,
1189 session_generation: Option<u64>,
1190 registry: CallIdRouter,
1191 shared_call_id: SharedCallIdCounter,
1192 async_response_tx: crossbeam_channel::Sender<BridgeResponse>,
1193 abort_rx: crossbeam_channel::Receiver<()>,
1194 runtime: agentos_runtime::RuntimeContext,
1195 pause_control: Arc<crate::session::SessionPauseControl>,
1196 bridge_call_timeout: Duration,
1197 ) -> Self {
1198 BridgeCallContext {
1199 sender,
1200 response_rx: None,
1201 session_id,
1202 next_call_id: shared_call_id,
1203 pending_calls: Mutex::new(HashSet::new()),
1204 track_pending_calls: should_track_pending_sync_calls(),
1205 call_id_router: Some(registry),
1206 session_generation,
1207 async_response_tx: Some(async_response_tx),
1208 abort_rx: Some(abort_rx),
1209 pause_control: Some(pause_control),
1210 runtime: Some(runtime),
1211 bridge_call_timeout,
1212 }
1213 }
1214
1215 pub(crate) fn runtime_context(&self) -> Option<&agentos_runtime::RuntimeContext> {
1216 self.runtime.as_ref()
1217 }
1218
1219 pub(crate) fn timer_task_owner(&self) -> Option<agentos_runtime::TaskOwner> {
1220 self.session_generation
1221 .map(|generation| agentos_runtime::TaskOwner::Vm { generation })
1222 }
1223
1224 pub fn sync_call_response(
1230 &self,
1231 method: &str,
1232 args: Vec<u8>,
1233 ) -> Result<Option<SyncBridgeCallResponse>, String> {
1234 let max_response_bytes = self.configured_bridge_response_bytes(method)?;
1235 self.sync_call_response_with_max_response_bytes(method, args, max_response_bytes)
1236 }
1237
1238 pub fn sync_call_response_with_max_response_bytes(
1239 &self,
1240 method: &str,
1241 args: Vec<u8>,
1242 max_response_bytes: usize,
1243 ) -> Result<Option<SyncBridgeCallResponse>, String> {
1244 let call_id = self.next_call_id.fetch_add(1, Ordering::Relaxed);
1245 track_sync_bridge_call_method(call_id, method);
1246
1247 if self.track_pending_calls {
1250 let mut pending = self.pending_calls.lock().unwrap();
1251 if !pending.insert(call_id) {
1252 return Err(format!("duplicate call_id: {}", call_id));
1253 }
1254 }
1255
1256 let (direct_response_rx, mut pending_route) = if let Some(ref registry) =
1257 self.call_id_router
1258 {
1259 let phase_start = Instant::now();
1260 let runtime = self.runtime.as_ref().ok_or_else(|| {
1261 String::from(
1262 "ERR_AGENTOS_RUNTIME_NOT_INJECTED: direct bridge calls require a session RuntimeContext",
1263 )
1264 })?;
1265 let receiver = match registry.register_sync_with_timeout(
1266 runtime,
1267 args.len(),
1268 max_response_bytes,
1269 call_id,
1270 &self.session_id,
1271 self.session_generation,
1272 self.bridge_call_timeout,
1273 ) {
1274 Ok(receiver) => receiver,
1275 Err(error) => {
1276 self.remove_pending_call(call_id);
1277 return Err(error);
1278 }
1279 };
1280 record_sync_bridge_host_phase(method, "host_register_route", phase_start.elapsed());
1281 (
1282 Some(receiver),
1283 Some(PendingBridgeRoute::new(Arc::clone(registry), call_id)),
1284 )
1285 } else {
1286 (None, None)
1287 };
1288
1289 let bridge_call = RuntimeEvent::BridgeCall {
1291 session_id: self.session_id.clone(),
1292 call_id,
1293 method: method.to_string(),
1294 payload: args,
1295 };
1296
1297 let __lat = syncrpc_lat_enabled().then(Instant::now);
1298 let phase_start = Instant::now();
1299 if let Some(route) = pending_route.as_ref() {
1300 if let Err(error) = route.mark_host_visible() {
1301 self.remove_pending_call(call_id);
1302 return Err(error);
1303 }
1304 }
1305 if let Err(e) = self.sender.send_event(bridge_call) {
1306 if let Some(route) = pending_route.take() {
1307 route.cancel_unpublished();
1308 }
1309 self.remove_pending_call(call_id);
1310 return Err(format!("failed to write BridgeCall: {}", e));
1311 }
1312 record_sync_bridge_host_phase(method, "host_send_event", phase_start.elapsed());
1313
1314 let response = if let Some(receiver) = direct_response_rx {
1316 let phase_start = Instant::now();
1317 let result = if let Some(abort_rx) = self.abort_rx.as_ref() {
1318 crossbeam_channel::select! {
1319 recv(receiver) -> response => response.map_err(|_| {
1320 String::from("bridge response target closed before settlement")
1321 }),
1322 recv(abort_rx) -> _ => Err(String::from("execution aborted")),
1323 default(self.bridge_call_timeout) => {
1324 let registry = self.call_id_router.as_ref().expect("direct response route");
1325 registry.timeout(call_id)?;
1326 receiver.recv().map_err(|_| {
1327 String::from("bridge response target closed during deadline settlement")
1328 })
1329 },
1330 }
1331 } else {
1332 match receiver.recv_timeout(self.bridge_call_timeout) {
1333 Ok(response) => Ok(response),
1334 Err(crossbeam_channel::RecvTimeoutError::Disconnected) => Err(String::from(
1335 "bridge response target closed before settlement",
1336 )),
1337 Err(crossbeam_channel::RecvTimeoutError::Timeout) => {
1338 let registry = self.call_id_router.as_ref().expect("direct response route");
1339 registry.timeout(call_id)?;
1340 receiver.recv().map_err(|_| {
1341 String::from("bridge response target closed during deadline settlement")
1342 })
1343 }
1344 }
1345 };
1346 match result {
1347 Ok(frame) => {
1348 if let Some(control) = &self.pause_control {
1349 control.wait_while_paused();
1350 }
1351 record_sync_bridge_host_phase(
1352 method,
1353 "host_recv_response",
1354 phase_start.elapsed(),
1355 );
1356 frame
1357 }
1358 Err(e) => {
1359 self.remove_pending_call(call_id);
1360 return Err(e);
1361 }
1362 }
1363 } else {
1364 let rx = self
1365 .response_rx
1366 .as_ref()
1367 .expect("legacy bridge context has a response receiver")
1368 .lock()
1369 .unwrap();
1370 let phase_start = Instant::now();
1371 match rx.recv_response(call_id) {
1372 Ok(frame) => {
1373 record_sync_bridge_host_phase(
1374 method,
1375 "host_recv_response",
1376 phase_start.elapsed(),
1377 );
1378 frame
1379 }
1380 Err(e) => {
1381 self.remove_pending_call(call_id);
1382 return Err(e);
1383 }
1384 }
1385 };
1386 if let Some(t) = __lat {
1387 record_syncrpc_lat(t.elapsed().as_nanos() as u64);
1388 }
1389
1390 let phase_start = Instant::now();
1391 self.remove_pending_call(call_id);
1392 record_sync_bridge_host_phase(method, "host_cleanup", phase_start.elapsed());
1393
1394 let phase_start = Instant::now();
1396 if response.status == 1 {
1397 let result = Err(String::from_utf8_lossy(&response.payload).to_string());
1398 record_sync_bridge_host_phase(method, "host_extract_response", phase_start.elapsed());
1399 result
1400 } else if response.payload.is_empty() && response.status != 2 {
1401 record_sync_bridge_host_phase(method, "host_extract_response", phase_start.elapsed());
1402 Ok(None)
1403 } else {
1404 let result = Ok(Some(SyncBridgeCallResponse {
1405 status: response.status,
1406 payload: response.payload,
1407 reservation: response.reservation,
1408 }));
1409 record_sync_bridge_host_phase(method, "host_extract_response", phase_start.elapsed());
1410 result
1411 }
1412 }
1413
1414 pub fn sync_call(&self, method: &str, args: Vec<u8>) -> Result<Option<Vec<u8>>, String> {
1415 self.sync_call_response(method, args)
1416 .map(|response| response.map(|response| response.payload))
1417 }
1418
1419 pub fn prepare_async_call(
1420 &self,
1421 method: &str,
1422 args: Vec<u8>,
1423 ) -> Result<PreparedAsyncBridgeCall, String> {
1424 let max_response_bytes = self.configured_bridge_response_bytes(method)?;
1425 self.prepare_async_call_with_max_response_bytes(method, args, max_response_bytes)
1426 }
1427
1428 pub fn prepare_async_call_with_max_response_bytes(
1429 &self,
1430 method: &str,
1431 args: Vec<u8>,
1432 max_response_bytes: usize,
1433 ) -> Result<PreparedAsyncBridgeCall, String> {
1434 let call_id = self.next_call_id.fetch_add(1, Ordering::Relaxed);
1435
1436 let pending_route = if let Some(ref registry) = self.call_id_router {
1437 let runtime = self.runtime.as_ref().ok_or_else(|| {
1438 String::from(
1439 "ERR_AGENTOS_RUNTIME_NOT_INJECTED: direct bridge calls require a session RuntimeContext",
1440 )
1441 })?;
1442 let sender = self.async_response_tx.as_ref().ok_or_else(|| {
1443 String::from(
1444 "ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY: async response lane is unavailable",
1445 )
1446 })?;
1447 let mut deadline_cancelled = registry.register_async_with_timeout(
1448 runtime,
1449 args.len(),
1450 max_response_bytes,
1451 call_id,
1452 &self.session_id,
1453 self.session_generation,
1454 sender.clone(),
1455 self.bridge_call_timeout,
1456 )?;
1457 let deadline_registry = Arc::clone(registry);
1458 let timeout = self.bridge_call_timeout;
1459 if let Err(error) = runtime.spawn(agentos_runtime::TaskClass::Timer, async move {
1460 if tokio::time::timeout(timeout, &mut deadline_cancelled)
1461 .await
1462 .is_err()
1463 {
1464 if let Err(error) = deadline_registry.timeout(call_id) {
1465 eprintln!("{error}");
1466 }
1467 }
1468 }) {
1469 registry.cancel_unpublished(call_id);
1473 return Err(format!(
1474 "ERR_AGENTOS_BRIDGE_DEADLINE_TASK: failed to arm bridge call_id {call_id} deadline: {error}"
1475 ));
1476 }
1477 Some(PendingBridgeRoute::new(Arc::clone(registry), call_id))
1478 } else {
1479 None
1480 };
1481
1482 Ok(PreparedAsyncBridgeCall {
1483 call_id,
1484 event: RuntimeEvent::BridgeCall {
1485 session_id: self.session_id.clone(),
1486 call_id,
1487 method: method.to_string(),
1488 payload: args,
1489 },
1490 pending_route,
1491 })
1492 }
1493
1494 fn configured_bridge_response_bytes(&self, method: &str) -> Result<usize, String> {
1495 if self.call_id_router.is_none() {
1496 return Ok(0);
1497 }
1498 let runtime = self.runtime.as_ref().ok_or_else(|| {
1499 String::from(
1500 "ERR_AGENTOS_RUNTIME_NOT_INJECTED: direct bridge calls require a session RuntimeContext",
1501 )
1502 })?;
1503 let configured = runtime
1504 .resources()
1505 .usage(ResourceClass::BridgeResponseBytes)
1506 .limit
1507 .ok_or_else(|| {
1508 String::from(
1509 "ERR_AGENTOS_RESOURCE_UNBOUNDED: bridge response bytes require limits.reactor.maxBridgeResponseBytes",
1510 )
1511 })?;
1512 Ok(configured.min(crate::bridge::declared_bridge_response_bytes(method, None)))
1513 }
1514
1515 pub fn dispatch_async_call(&self, prepared: PreparedAsyncBridgeCall) -> Result<u64, String> {
1516 let PreparedAsyncBridgeCall {
1517 call_id,
1518 event,
1519 mut pending_route,
1520 } = prepared;
1521 if let Some(route) = pending_route.as_ref() {
1522 route.mark_host_visible()?;
1523 }
1524 if let Err(e) = self.sender.send_event(event) {
1525 if let Some(route) = pending_route.take() {
1526 route.cancel_unpublished();
1527 }
1528 return Err(format!("failed to write BridgeCall: {}", e));
1529 }
1530 if let Some(pending_route) = pending_route {
1531 pending_route.disarm();
1532 }
1533 Ok(call_id)
1534 }
1535
1536 pub fn async_send(&self, method: &str, args: Vec<u8>) -> Result<u64, String> {
1538 let prepared = self.prepare_async_call(method, args)?;
1539 self.dispatch_async_call(prepared)
1540 }
1541
1542 fn remove_pending_call(&self, call_id: u64) {
1543 cleanup_sync_bridge_call_tracking(call_id);
1544 if self.track_pending_calls {
1545 self.pending_calls.lock().unwrap().remove(&call_id);
1546 }
1547 }
1548
1549 pub fn is_call_pending(&self, call_id: u64) -> bool {
1551 if !self.track_pending_calls {
1552 return false;
1553 }
1554 self.pending_calls.lock().unwrap().contains(&call_id)
1555 }
1556
1557 pub fn pending_count(&self) -> usize {
1559 if !self.track_pending_calls {
1560 return 0;
1561 }
1562 self.pending_calls.lock().unwrap().len()
1563 }
1564}
1565
1566impl Drop for BridgeCallContext {
1567 fn drop(&mut self) {
1568 if let Some(registry) = self.call_id_router.as_ref() {
1569 registry.cancel_session(&self.session_id, self.session_generation);
1570 }
1571 }
1572}
1573
1574fn should_track_pending_sync_calls() -> bool {
1575 std::env::var("AGENTOS_TRACK_PENDING_SYNC_CALLS").as_deref() == Ok("1")
1576}
1577
1578#[cfg(test)]
1579mod tests {
1580 use super::*;
1581 use std::io::Cursor;
1582 use std::sync::Arc;
1583
1584 fn test_runtime_context() -> RuntimeContext {
1585 agentos_runtime::SidecarRuntime::process(&agentos_runtime::RuntimeConfig::default())
1586 .expect("test process runtime")
1587 .context()
1588 }
1589
1590 fn limited_bridge_runtime(
1591 max_calls: usize,
1592 max_request_bytes: usize,
1593 max_response_bytes: usize,
1594 ) -> (
1595 RuntimeContext,
1596 Arc<agentos_runtime::accounting::ResourceLedger>,
1597 ) {
1598 let process = test_runtime_context();
1599 let resources = Arc::new(agentos_runtime::accounting::ResourceLedger::child(
1600 "bridge-test-vm",
1601 [
1602 (
1603 ResourceClass::BridgeCalls,
1604 agentos_runtime::accounting::ResourceLimit::new(
1605 max_calls,
1606 "limits.reactor.maxBridgeCalls",
1607 ),
1608 ),
1609 (
1610 ResourceClass::BridgeRequestBytes,
1611 agentos_runtime::accounting::ResourceLimit::new(
1612 max_request_bytes,
1613 "limits.reactor.maxBridgeRequestBytes",
1614 ),
1615 ),
1616 (
1617 ResourceClass::BridgeResponseBytes,
1618 agentos_runtime::accounting::ResourceLimit::new(
1619 max_response_bytes,
1620 "limits.reactor.maxBridgeResponseBytes",
1621 ),
1622 ),
1623 ],
1624 Arc::clone(process.resources()),
1625 ));
1626 (process.scoped_for_vm(Arc::clone(&resources), 7), resources)
1627 }
1628
1629 fn assert_ledger_settles_to_zero(resources: &agentos_runtime::accounting::ResourceLedger) {
1630 let deadline = Instant::now() + Duration::from_secs(1);
1631 while !resources.is_zero() && Instant::now() < deadline {
1632 std::thread::yield_now();
1633 }
1634 assert!(
1635 resources.is_zero(),
1636 "bridge reservations and supervised deadline task must reconcile"
1637 );
1638 }
1639
1640 struct SharedWriter(Arc<Mutex<Vec<u8>>>);
1642
1643 impl Write for SharedWriter {
1644 fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
1645 self.0.lock().unwrap().write(buf)
1646 }
1647 fn flush(&mut self) -> std::io::Result<()> {
1648 self.0.lock().unwrap().flush()
1649 }
1650 }
1651
1652 struct RejectingRuntimeEventSender;
1653
1654 impl RuntimeEventSender for RejectingRuntimeEventSender {
1655 fn send_event(&self, _event: RuntimeEvent) -> Result<(), String> {
1656 Err(String::from("injected event-lane failure"))
1657 }
1658 }
1659
1660 fn make_response_bytes(
1662 call_id: u64,
1663 result: Option<Vec<u8>>,
1664 error: Option<String>,
1665 ) -> Vec<u8> {
1666 let mut buf = Vec::new();
1667 let (status, payload) = if let Some(err) = error {
1668 (1u8, err.into_bytes())
1669 } else if let Some(res) = result {
1670 (0u8, res)
1671 } else {
1672 (0u8, vec![])
1673 };
1674 ipc_binary::write_frame(
1675 &mut buf,
1676 &BinaryFrame::BridgeResponse {
1677 session_id: String::new(),
1678 call_id,
1679 status,
1680 payload,
1681 },
1682 )
1683 .unwrap();
1684 buf
1685 }
1686
1687 #[test]
1688 fn sync_call_success_with_result() {
1689 let response_bytes = make_response_bytes(1, Some(vec![0x93, 0x01, 0x02, 0x03]), None);
1690 let writer_buf = Arc::new(Mutex::new(Vec::new()));
1691
1692 let ctx = BridgeCallContext::new(
1693 Box::new(SharedWriter(Arc::clone(&writer_buf))),
1694 Box::new(Cursor::new(response_bytes)),
1695 "test-session-abc".into(),
1696 );
1697
1698 let result = ctx.sync_call("_fsReadFile", vec![0x91, 0xa3, 0x66, 0x6f, 0x6f]);
1699 assert!(result.is_ok());
1700 assert_eq!(result.unwrap(), Some(vec![0x93, 0x01, 0x02, 0x03]));
1701
1702 let written = writer_buf.lock().unwrap();
1704 let call = ipc_binary::read_frame(&mut Cursor::new(&*written)).unwrap();
1705 match call {
1706 BinaryFrame::BridgeCall {
1707 call_id,
1708 session_id,
1709 method,
1710 payload,
1711 ..
1712 } => {
1713 assert_eq!(call_id, 1);
1714 assert_eq!(session_id, "test-session-abc");
1715 assert_eq!(method, "_fsReadFile");
1716 assert_eq!(payload, vec![0x91, 0xa3, 0x66, 0x6f, 0x6f]);
1717 }
1718 _ => panic!("expected BridgeCall"),
1719 }
1720 }
1721
1722 #[test]
1723 fn sync_call_success_null_result() {
1724 let response_bytes = make_response_bytes(1, None, None);
1725 let ctx = BridgeCallContext::new(
1726 Box::new(Vec::new()),
1727 Box::new(Cursor::new(response_bytes)),
1728 "session-1".into(),
1729 );
1730
1731 let result = ctx.sync_call("_log", vec![0xc0]).unwrap();
1732 assert_eq!(result, None);
1733 }
1734
1735 #[test]
1736 fn sync_call_error_response() {
1737 let response_bytes = make_response_bytes(1, None, Some("ENOENT: no such file".into()));
1738 let ctx = BridgeCallContext::new(
1739 Box::new(Vec::new()),
1740 Box::new(Cursor::new(response_bytes)),
1741 "session-1".into(),
1742 );
1743
1744 let result = ctx.sync_call("_fsReadFile", vec![0xc0]);
1745 assert!(result.is_err());
1746 assert_eq!(result.unwrap_err(), "ENOENT: no such file");
1747 }
1748
1749 #[test]
1750 fn sync_call_call_id_increments() {
1751 let mut response_bytes = make_response_bytes(1, Some(vec![0xa1, 0x61]), None);
1753 response_bytes.extend_from_slice(&make_response_bytes(2, Some(vec![0xa1, 0x62]), None));
1754
1755 let ctx = BridgeCallContext::new(
1756 Box::new(Vec::new()),
1757 Box::new(Cursor::new(response_bytes)),
1758 "session-1".into(),
1759 );
1760
1761 let r1 = ctx.sync_call("_fn1", vec![]).unwrap();
1762 let r2 = ctx.sync_call("_fn2", vec![]).unwrap();
1763 assert_eq!(r1, Some(vec![0xa1, 0x61]));
1764 assert_eq!(r2, Some(vec![0xa1, 0x62]));
1765 }
1766
1767 #[test]
1768 fn sync_call_pending_cleanup_on_read_error() {
1769 let ctx = BridgeCallContext::new(
1771 Box::new(Vec::new()),
1772 Box::new(Cursor::new(Vec::new())),
1773 "session-1".into(),
1774 );
1775
1776 assert_eq!(ctx.pending_count(), 0);
1777 let _ = ctx.sync_call("_fn", vec![]);
1778 assert_eq!(ctx.pending_count(), 0);
1779 }
1780
1781 #[test]
1782 fn sync_call_id_mismatch_rejected() {
1783 let response_bytes = make_response_bytes(99, Some(vec![0xc0]), None);
1785 let ctx = BridgeCallContext::new(
1786 Box::new(Vec::new()),
1787 Box::new(Cursor::new(response_bytes)),
1788 "session-1".into(),
1789 );
1790
1791 let result = ctx.sync_call("_fn", vec![]);
1792 assert!(result.is_err());
1793 assert!(result.unwrap_err().contains("call_id mismatch"));
1794 }
1795
1796 #[test]
1797 fn sync_call_unexpected_message_type_rejected() {
1798 let mut response_bytes = Vec::new();
1800 ipc_binary::write_frame(
1801 &mut response_bytes,
1802 &BinaryFrame::TerminateExecution {
1803 session_id: "session-1".into(),
1804 },
1805 )
1806 .unwrap();
1807
1808 let ctx = BridgeCallContext::new(
1809 Box::new(Vec::new()),
1810 Box::new(Cursor::new(response_bytes)),
1811 "session-1".into(),
1812 );
1813
1814 let result = ctx.sync_call("_fn", vec![]);
1815 assert!(result.is_err());
1816 assert!(result.unwrap_err().contains("expected BridgeResponse"));
1817 }
1818
1819 #[test]
1820 fn async_send_writes_bridge_call() {
1821 let writer_buf = Arc::new(Mutex::new(Vec::new()));
1822 let ctx = BridgeCallContext::new(
1823 Box::new(SharedWriter(Arc::clone(&writer_buf))),
1824 Box::new(Cursor::new(Vec::new())),
1825 "test-session-abc".into(),
1826 );
1827
1828 let call_id = ctx
1829 .async_send("_asyncFn", vec![0x91, 0xa3, 0x66, 0x6f, 0x6f])
1830 .unwrap();
1831 assert_eq!(call_id, 1);
1832
1833 let written = writer_buf.lock().unwrap();
1835 let call = ipc_binary::read_frame(&mut Cursor::new(&*written)).unwrap();
1836 match call {
1837 BinaryFrame::BridgeCall {
1838 call_id,
1839 session_id,
1840 method,
1841 payload,
1842 ..
1843 } => {
1844 assert_eq!(call_id, 1);
1845 assert_eq!(session_id, "test-session-abc");
1846 assert_eq!(method, "_asyncFn");
1847 assert_eq!(payload, vec![0x91, 0xa3, 0x66, 0x6f, 0x6f]);
1848 }
1849 _ => panic!("expected BridgeCall"),
1850 }
1851 }
1852
1853 #[test]
1854 fn async_send_increments_call_id() {
1855 let ctx = BridgeCallContext::new(
1856 Box::new(Vec::new()),
1857 Box::new(Cursor::new(Vec::new())),
1858 "session-1".into(),
1859 );
1860
1861 let id1 = ctx.async_send("_fn1", vec![]).unwrap();
1862 let id2 = ctx.async_send("_fn2", vec![]).unwrap();
1863 assert_eq!(id1, 1);
1864 assert_eq!(id2, 2);
1865 }
1866
1867 #[test]
1868 fn async_send_shares_counter_with_sync() {
1869 let response_bytes = make_response_bytes(1, Some(vec![0xc0]), None);
1871 let ctx = BridgeCallContext::new(
1872 Box::new(Vec::new()),
1873 Box::new(Cursor::new(response_bytes)),
1874 "session-1".into(),
1875 );
1876
1877 let _ = ctx.sync_call("_sync", vec![]);
1878 let id = ctx.async_send("_async", vec![]).unwrap();
1879 assert_eq!(id, 2);
1880 }
1881
1882 #[test]
1883 fn channel_runtime_event_sender_delivers_frames() {
1884 let (tx, rx) = crossbeam_channel::unbounded();
1885 let sender = super::ChannelRuntimeEventSender::new(tx, None);
1886
1887 let event = RuntimeEvent::BridgeCall {
1888 session_id: "sess-1".into(),
1889 call_id: 42,
1890 method: "_fsReadFile".into(),
1891 payload: vec![0x01, 0x02],
1892 };
1893 sender.send_event(event.clone()).expect("send_event");
1894
1895 let received = rx.recv().expect("recv");
1897 assert_eq!(received.output_generation, None);
1898 assert_eq!(received.event, event);
1899 }
1900
1901 #[test]
1902 fn channel_runtime_event_sender_no_mutex_contention() {
1903 let (tx, rx) = crossbeam_channel::unbounded();
1905 let handles: Vec<_> = (0..4)
1906 .map(|i| {
1907 let sender = super::ChannelRuntimeEventSender::new(tx.clone(), None);
1908 std::thread::spawn(move || {
1909 for j in 0..10 {
1910 let event = RuntimeEvent::BridgeCall {
1911 session_id: format!("sess-{}", i),
1912 call_id: (i * 100 + j) as u64,
1913 method: "_fn".into(),
1914 payload: vec![],
1915 };
1916 sender.send_event(event).expect("send_event");
1917 }
1918 })
1919 })
1920 .collect();
1921 drop(tx); for h in handles {
1924 h.join().expect("thread join");
1925 }
1926
1927 let mut count = 0;
1929 while rx.try_recv().is_ok() {
1930 count += 1;
1931 }
1932 assert_eq!(count, 40);
1933 }
1934
1935 #[test]
1936 fn channel_runtime_event_sender_with_bridge_context() {
1937 let (tx, rx) = crossbeam_channel::unbounded();
1939
1940 let response_bytes = make_response_bytes(1, Some(vec![0xAB, 0xCD]), None);
1942 let router: super::CallIdRouter = Arc::new(super::BridgeCallRegistry::with_default_limit());
1943
1944 let ctx = BridgeCallContext::with_receiver(
1945 Box::new(super::ChannelRuntimeEventSender::new(tx, None)),
1946 Box::new(super::ReaderBridgeResponseReceiver::new(Box::new(
1947 Cursor::new(response_bytes),
1948 ))),
1949 "test-session".into(),
1950 router,
1951 Arc::new(std::sync::atomic::AtomicU64::new(1)),
1952 );
1953
1954 let result = ctx.sync_call("_fsReadFile", vec![0x01]).unwrap();
1955 assert_eq!(result, Some(vec![0xAB, 0xCD]));
1956
1957 let event = rx.recv().expect("recv bridge call");
1959 match event.event {
1960 RuntimeEvent::BridgeCall { method, .. } => assert_eq!(method, "_fsReadFile"),
1961 _ => panic!("expected BridgeCall"),
1962 }
1963 }
1964
1965 #[test]
1966 fn sync_call_success_clears_call_id_route() {
1967 let (tx, _rx) = crossbeam_channel::unbounded();
1968 let response_bytes = make_response_bytes(1, Some(vec![0xAB, 0xCD]), None);
1969 let router: super::CallIdRouter = Arc::new(super::BridgeCallRegistry::with_default_limit());
1970
1971 let ctx = BridgeCallContext::with_receiver(
1972 Box::new(super::ChannelRuntimeEventSender::new(tx, None)),
1973 Box::new(super::ReaderBridgeResponseReceiver::new(Box::new(
1974 Cursor::new(response_bytes),
1975 ))),
1976 "test-session".into(),
1977 Arc::clone(&router),
1978 Arc::new(std::sync::atomic::AtomicU64::new(1)),
1979 );
1980
1981 let result = ctx.sync_call("_fsReadFile", vec![0x01]).unwrap();
1982 assert_eq!(result, Some(vec![0xAB, 0xCD]));
1983 assert!(
1984 router.pending_len() == 0,
1985 "sync bridge response completion should clear the call_id route"
1986 );
1987 }
1988
1989 #[test]
1990 fn bridge_registry_settles_a_call_specific_waiter_directly() {
1991 let registry = BridgeCallRegistry::new(2);
1992 let runtime = test_runtime_context();
1993 let waiter = registry
1994 .register_sync(&runtime, 0, 1, 7, "session-a", Some(3))
1995 .expect("register direct waiter");
1996
1997 registry
1998 .settle(
1999 "session-a",
2000 Some(3),
2001 BridgeResponse {
2002 call_id: 7,
2003 status: 0,
2004 payload: vec![0xA7],
2005 reservation: None,
2006 },
2007 )
2008 .expect("settle direct waiter");
2009
2010 assert_eq!(waiter.recv().expect("direct response").payload, vec![0xA7]);
2011 assert_eq!(registry.pending_len(), 0);
2012 }
2013
2014 #[test]
2015 fn bridge_registry_rejects_stale_generation_without_consuming_target() {
2016 let registry = BridgeCallRegistry::new(1);
2017 let runtime = test_runtime_context();
2018 let waiter = registry
2019 .register_sync(&runtime, 0, 1, 8, "session-a", Some(4))
2020 .expect("register direct waiter");
2021 let response = BridgeResponse {
2022 call_id: 8,
2023 status: 0,
2024 payload: vec![0xA8],
2025 reservation: None,
2026 };
2027
2028 let error = registry
2029 .settle("session-a", Some(5), response.clone())
2030 .expect_err("stale response must be rejected");
2031 assert!(error.contains("ERR_AGENTOS_BRIDGE_STALE_GENERATION"));
2032 assert_eq!(registry.pending_len(), 1);
2033
2034 registry
2035 .settle("session-a", Some(4), response)
2036 .expect("correct generation should still settle");
2037 assert_eq!(waiter.recv().expect("direct response").call_id, 8);
2038 }
2039
2040 #[test]
2041 fn bridge_registry_only_classifies_canceled_host_visible_routes_as_stale() {
2042 let registry = BridgeCallRegistry::new(2);
2043 let runtime = test_runtime_context();
2044
2045 let _visible_waiter = registry
2046 .register_sync(&runtime, 0, 1, 20, "session-a", Some(4))
2047 .expect("register host-visible route");
2048 registry
2049 .mark_host_visible(20)
2050 .expect("mark route host-visible");
2051 registry.cancel_session("session-a", Some(4));
2052 let error = registry
2053 .settle(
2054 "session-a",
2055 Some(4),
2056 BridgeResponse {
2057 call_id: 20,
2058 status: 0,
2059 payload: Vec::new(),
2060 reservation: None,
2061 },
2062 )
2063 .expect_err("canceled host-visible route must reject a stale completion");
2064 assert!(error.contains("ERR_AGENTOS_BRIDGE_STALE_COMPLETION"));
2065
2066 let mismatched = registry
2067 .settle(
2068 "session-a",
2069 Some(5),
2070 BridgeResponse {
2071 call_id: 20,
2072 status: 0,
2073 payload: Vec::new(),
2074 reservation: None,
2075 },
2076 )
2077 .expect_err("a different generation must not inherit stale-completion status");
2078 assert!(mismatched.contains("ERR_AGENTOS_BRIDGE_STALE_GENERATION"));
2079
2080 let _unpublished_waiter = registry
2081 .register_sync(&runtime, 0, 1, 21, "session-a", Some(4))
2082 .expect("register unpublished route");
2083 registry.cancel(21);
2084 let unknown = registry
2085 .settle(
2086 "session-a",
2087 Some(4),
2088 BridgeResponse {
2089 call_id: 21,
2090 status: 0,
2091 payload: Vec::new(),
2092 reservation: None,
2093 },
2094 )
2095 .expect_err("unpublished route must not authorize a stale response");
2096 assert!(unknown.contains("ERR_AGENTOS_BRIDGE_UNKNOWN_CALL_ID"));
2097 assert_eq!(registry.retired_len(), 1);
2098 }
2099
2100 #[test]
2101 fn bridge_registry_does_not_hide_duplicate_settlement_as_teardown() {
2102 let registry = BridgeCallRegistry::new(1);
2103 let runtime = test_runtime_context();
2104 let waiter = registry
2105 .register_sync(&runtime, 0, 1, 22, "session-a", Some(4))
2106 .expect("register direct waiter");
2107 registry
2108 .mark_host_visible(22)
2109 .expect("mark route host-visible");
2110 let response = BridgeResponse {
2111 call_id: 22,
2112 status: 0,
2113 payload: vec![0xA2],
2114 reservation: None,
2115 };
2116 registry
2117 .settle("session-a", Some(4), response.clone())
2118 .expect("settle direct waiter");
2119 drop(waiter.recv().expect("receive direct response"));
2120
2121 let duplicate = registry
2122 .settle("session-a", Some(4), response)
2123 .expect_err("duplicate settlement must remain a hard error");
2124 assert!(duplicate.contains("ERR_AGENTOS_BRIDGE_UNKNOWN_CALL_ID"));
2125 assert_eq!(registry.retired_len(), 0);
2126 }
2127
2128 #[test]
2129 fn bridge_registry_bounds_retired_identity_history_and_releases_reservations() {
2130 let registry = BridgeCallRegistry::new(1);
2131 let (runtime, resources) = limited_bridge_runtime(1, 4, 4);
2132
2133 for call_id in [30, 31] {
2134 let _waiter = registry
2135 .register_sync(&runtime, 2, 4, call_id, "session-a", Some(4))
2136 .expect("register route for cancellation");
2137 registry
2138 .mark_host_visible(call_id)
2139 .expect("mark route host-visible");
2140 registry.cancel_session("session-a", Some(4));
2141 assert!(
2142 resources.is_zero(),
2143 "cancellation must release reservations"
2144 );
2145 assert_eq!(registry.retired_len(), 1);
2146 }
2147
2148 let evicted = registry
2149 .settle(
2150 "session-a",
2151 Some(4),
2152 BridgeResponse {
2153 call_id: 30,
2154 status: 0,
2155 payload: Vec::new(),
2156 reservation: None,
2157 },
2158 )
2159 .expect_err("oldest retired identity must be evicted at the configured bound");
2160 assert!(evicted.contains("ERR_AGENTOS_BRIDGE_UNKNOWN_CALL_ID"));
2161
2162 let retained = registry
2163 .settle(
2164 "session-a",
2165 Some(4),
2166 BridgeResponse {
2167 call_id: 31,
2168 status: 0,
2169 payload: Vec::new(),
2170 reservation: None,
2171 },
2172 )
2173 .expect_err("newest retired identity must remain classified");
2174 assert!(retained.contains("ERR_AGENTOS_BRIDGE_STALE_COMPLETION"));
2175 }
2176
2177 #[test]
2178 fn bridge_registry_clear_releases_all_pre_reserved_response_capacity() {
2179 let registry = BridgeCallRegistry::new(2);
2180 let (runtime, resources) = limited_bridge_runtime(2, 4, 8);
2181
2182 for (call_id, session_id) in [(32, "session-a"), (33, "session-b")] {
2183 let _waiter = registry
2184 .register_sync(&runtime, 1, 4, call_id, session_id, Some(4))
2185 .expect("pre-reserve response capacity before teardown");
2186 registry
2187 .mark_host_visible(call_id)
2188 .expect("mark route host-visible");
2189 }
2190 assert_eq!(resources.usage(ResourceClass::BridgeCalls).used, 2);
2191 assert_eq!(resources.usage(ResourceClass::BridgeRequestBytes).used, 2);
2192 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 8);
2193
2194 registry.clear();
2195
2196 assert_eq!(registry.pending_len(), 0);
2197 assert_eq!(registry.retired_len(), 2);
2198 assert!(resources.is_zero());
2199 }
2200
2201 #[test]
2202 fn bridge_registry_stale_completion_cannot_consume_reused_session_generation() {
2203 let registry = BridgeCallRegistry::new(2);
2204 let runtime = test_runtime_context();
2205 let _old_waiter = registry
2206 .register_sync(&runtime, 0, 1, 40, "reused-session", Some(4))
2207 .expect("register old generation route");
2208 registry
2209 .mark_host_visible(40)
2210 .expect("mark old route host-visible");
2211 registry.cancel_session("reused-session", Some(4));
2212
2213 let current_waiter = registry
2214 .register_sync(&runtime, 0, 1, 41, "reused-session", Some(5))
2215 .expect("register current generation route");
2216 registry
2217 .mark_host_visible(41)
2218 .expect("mark current route host-visible");
2219
2220 let stale = registry
2221 .settle(
2222 "reused-session",
2223 Some(4),
2224 BridgeResponse {
2225 call_id: 40,
2226 status: 0,
2227 payload: Vec::new(),
2228 reservation: None,
2229 },
2230 )
2231 .expect_err("retired generation must reject its late response");
2232 assert!(stale.contains("ERR_AGENTOS_BRIDGE_STALE_COMPLETION"));
2233 assert_eq!(registry.pending_len(), 1);
2234
2235 registry
2236 .settle(
2237 "reused-session",
2238 Some(5),
2239 BridgeResponse {
2240 call_id: 41,
2241 status: 0,
2242 payload: vec![0xA5],
2243 reservation: None,
2244 },
2245 )
2246 .expect("settle current generation");
2247 assert_eq!(
2248 current_waiter
2249 .recv()
2250 .expect("current generation response")
2251 .payload,
2252 vec![0xA5]
2253 );
2254 }
2255
2256 #[test]
2257 fn sync_call_teardown_after_publication_records_stale_completion_proof() {
2258 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(2));
2259 let runtime = test_runtime_context();
2260 let (event_tx, event_rx) = crossbeam_channel::unbounded();
2261 let (async_tx, _async_rx) = crossbeam_channel::bounded(1);
2262 let (abort_tx, abort_rx) = crossbeam_channel::bounded(1);
2263 let ctx = BridgeCallContext::with_registry(
2264 Box::new(ChannelRuntimeEventSender::new(event_tx, Some(4))),
2265 String::from("session-a"),
2266 Some(4),
2267 Arc::clone(®istry),
2268 Arc::new(AtomicU64::new(50)),
2269 async_tx,
2270 abort_rx,
2271 runtime,
2272 Arc::new(crate::session::SessionPauseControl::default()),
2273 DEFAULT_BRIDGE_CALL_TIMEOUT,
2274 );
2275
2276 let guest = std::thread::spawn(move || {
2277 ctx.sync_call_response_with_max_response_bytes("_fsReadFile", Vec::new(), 1)
2278 });
2279 let event = event_rx
2280 .recv_timeout(Duration::from_secs(1))
2281 .expect("host must observe bridge call before teardown");
2282 assert!(matches!(
2283 event.event,
2284 RuntimeEvent::BridgeCall { call_id: 50, .. }
2285 ));
2286
2287 registry.cancel_session("session-a", Some(4));
2288 let _ = abort_tx.try_send(());
2289 let guest_error = guest
2290 .join()
2291 .expect("join guest bridge call")
2292 .expect_err("teardown must abort sync call");
2293 assert!(
2294 guest_error.contains("bridge response target closed")
2295 || guest_error.contains("execution aborted")
2296 );
2297
2298 let stale = registry
2299 .settle(
2300 "session-a",
2301 Some(4),
2302 BridgeResponse {
2303 call_id: 50,
2304 status: 0,
2305 payload: Vec::new(),
2306 reservation: None,
2307 },
2308 )
2309 .expect_err("late host response must use retirement proof");
2310 assert!(stale.contains("ERR_AGENTOS_BRIDGE_STALE_COMPLETION"));
2311 }
2312
2313 #[test]
2314 fn sync_bridge_call_deadline_releases_route_and_rejects_late_response() {
2315 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(2));
2316 let runtime = test_runtime_context();
2317 let (event_tx, event_rx) = crossbeam_channel::unbounded();
2318 let (async_tx, _async_rx) = crossbeam_channel::bounded(1);
2319 let (_abort_tx, abort_rx) = crossbeam_channel::bounded(1);
2320 let ctx = BridgeCallContext::with_registry(
2321 Box::new(ChannelRuntimeEventSender::new(event_tx, Some(4))),
2322 String::from("session-deadline"),
2323 Some(4),
2324 Arc::clone(®istry),
2325 Arc::new(AtomicU64::new(60)),
2326 async_tx,
2327 abort_rx,
2328 runtime,
2329 Arc::new(crate::session::SessionPauseControl::default()),
2330 Duration::from_millis(10),
2331 );
2332
2333 let error = ctx
2334 .sync_call_response_with_max_response_bytes("_fsReadFile", Vec::new(), 128)
2335 .expect_err("missing host response must hit the per-call deadline");
2336 assert!(error.contains("ERR_AGENTOS_BRIDGE_CALL_TIMEOUT"));
2337 assert_eq!(registry.pending_len(), 0);
2338 assert_eq!(registry.retired_len(), 1);
2339 assert!(matches!(
2340 event_rx.recv().expect("host-visible call").event,
2341 RuntimeEvent::BridgeCall { call_id: 60, .. }
2342 ));
2343 let late = registry
2344 .settle(
2345 "session-deadline",
2346 Some(4),
2347 BridgeResponse {
2348 call_id: 60,
2349 status: 0,
2350 payload: Vec::new(),
2351 reservation: None,
2352 },
2353 )
2354 .expect_err("late host response must not resurrect the timed-out route");
2355 assert!(late.contains("ERR_AGENTOS_BRIDGE_STALE_COMPLETION"));
2356 }
2357
2358 #[test]
2359 fn async_bridge_call_deadline_settles_dedicated_response_lane() {
2360 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(2));
2361 let runtime = test_runtime_context();
2362 let (event_tx, event_rx) = crossbeam_channel::unbounded();
2363 let (async_tx, async_rx) = crossbeam_channel::bounded(1);
2364 let (_abort_tx, abort_rx) = crossbeam_channel::bounded(1);
2365 let ctx = BridgeCallContext::with_registry(
2366 Box::new(ChannelRuntimeEventSender::new(event_tx, Some(5))),
2367 String::from("async-deadline"),
2368 Some(5),
2369 Arc::clone(®istry),
2370 Arc::new(AtomicU64::new(70)),
2371 async_tx,
2372 abort_rx,
2373 runtime,
2374 Arc::new(crate::session::SessionPauseControl::default()),
2375 Duration::from_millis(10),
2376 );
2377
2378 assert_eq!(ctx.async_send("_toolCall", Vec::new()).unwrap(), 70);
2379 assert!(matches!(
2380 event_rx.recv().expect("host-visible async call").event,
2381 RuntimeEvent::BridgeCall { call_id: 70, .. }
2382 ));
2383 let timeout = async_rx
2384 .recv_timeout(Duration::from_secs(1))
2385 .expect("deadline must settle the async lane");
2386 assert_eq!(timeout.status, 1);
2387 assert!(
2388 String::from_utf8_lossy(&timeout.payload).contains("ERR_AGENTOS_BRIDGE_CALL_TIMEOUT")
2389 );
2390 assert_eq!(registry.pending_len(), 0);
2391 }
2392
2393 #[test]
2394 fn direct_sync_response_waits_at_paused_execution_boundary() {
2395 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(2));
2396 let runtime = test_runtime_context();
2397 let (event_tx, event_rx) = crossbeam_channel::unbounded();
2398 let (async_tx, _async_rx) = crossbeam_channel::bounded(1);
2399 let (_abort_tx, abort_rx) = crossbeam_channel::bounded(1);
2400 let pause_control = Arc::new(crate::session::SessionPauseControl::default());
2401 let ctx = BridgeCallContext::with_registry(
2402 Box::new(ChannelRuntimeEventSender::new(event_tx, Some(4))),
2403 String::from("session-a"),
2404 Some(4),
2405 Arc::clone(®istry),
2406 Arc::new(AtomicU64::new(50)),
2407 async_tx,
2408 abort_rx,
2409 runtime,
2410 Arc::clone(&pause_control),
2411 DEFAULT_BRIDGE_CALL_TIMEOUT,
2412 );
2413
2414 pause_control.pause();
2415 let (result_tx, result_rx) = std::sync::mpsc::channel();
2416 let guest = std::thread::spawn(move || {
2417 let _ = result_tx.send(ctx.sync_call_response_with_max_response_bytes(
2418 "_fsReadFile",
2419 Vec::new(),
2420 1,
2421 ));
2422 });
2423 let event = event_rx
2424 .recv_timeout(Duration::from_secs(1))
2425 .expect("host must observe the direct bridge call");
2426 assert!(matches!(
2427 event.event,
2428 RuntimeEvent::BridgeCall { call_id: 50, .. }
2429 ));
2430
2431 registry
2432 .settle(
2433 "session-a",
2434 Some(4),
2435 BridgeResponse {
2436 call_id: 50,
2437 status: 0,
2438 payload: Vec::new(),
2439 reservation: None,
2440 },
2441 )
2442 .expect("settle direct response while paused");
2443 assert!(
2444 matches!(
2445 result_rx.recv_timeout(Duration::from_millis(50)),
2446 Err(std::sync::mpsc::RecvTimeoutError::Timeout)
2447 ),
2448 "direct response must not let synchronous JavaScript cross a paused boundary"
2449 );
2450
2451 pause_control.resume();
2452 result_rx
2453 .recv_timeout(Duration::from_secs(1))
2454 .expect("resume must release the direct response")
2455 .expect("direct sync bridge call must complete successfully");
2456 guest.join().expect("join direct bridge caller");
2457 }
2458
2459 #[test]
2460 fn bridge_registry_enforces_its_configured_bound() {
2461 let registry = BridgeCallRegistry::new(1);
2462 let runtime = test_runtime_context();
2463 let _waiter = registry
2464 .register_sync(&runtime, 0, 1, 9, "session-a", None)
2465 .expect("register first waiter");
2466
2467 let error = registry
2468 .register_sync(&runtime, 0, 1, 10, "session-b", None)
2469 .expect_err("second waiter must exceed configured bound");
2470 assert!(error.contains("ERR_AGENTOS_BRIDGE_CALL_LIMIT"));
2471 assert!(error.contains("runtime.resources.maxBridgeCalls"));
2472 }
2473
2474 #[test]
2475 fn bridge_registry_accounts_and_releases_settled_call_reservations() {
2476 let registry = BridgeCallRegistry::new(2);
2477 let (runtime, resources) = limited_bridge_runtime(2, 8, 8);
2478 let waiter = registry
2479 .register_sync(&runtime, 2, 4, 70, "session-a", Some(7))
2480 .expect("reserve bridge call admission");
2481
2482 assert_eq!(resources.usage(ResourceClass::BridgeCalls).used, 1);
2483 assert_eq!(resources.usage(ResourceClass::BridgeRequestBytes).used, 2);
2484 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 4);
2485
2486 registry
2487 .settle(
2488 "session-a",
2489 Some(7),
2490 BridgeResponse {
2491 call_id: 70,
2492 status: 0,
2493 payload: vec![0xA7; 4],
2494 reservation: None,
2495 },
2496 )
2497 .expect("settle admitted response");
2498 let response = waiter.recv().expect("receive response");
2499 assert_eq!(response.payload.len(), 4);
2500 assert_eq!(resources.usage(ResourceClass::BridgeCalls).used, 0);
2501 assert_eq!(resources.usage(ResourceClass::BridgeRequestBytes).used, 0);
2502 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 4);
2503 drop(response);
2504 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 0);
2505
2506 let metrics = runtime.metrics().snapshot();
2507 assert!(
2508 metrics.resources[agentos_runtime::metrics::ResourceMetricClass::BridgeCalls.index()]
2509 .high_water
2510 >= 1
2511 );
2512 assert!(
2513 metrics.buffers[agentos_runtime::metrics::BufferMetricClass::Bridge.index()].high_water
2514 >= 4
2515 );
2516 }
2517
2518 #[test]
2519 fn bridge_registry_delivery_failure_consumes_route_and_releases_accounting() {
2520 let registry = BridgeCallRegistry::new(2);
2521 let (runtime, resources) = limited_bridge_runtime(2, 8, 8);
2522 let waiter = registry
2523 .register_sync(&runtime, 2, 4, 75, "session-a", Some(7))
2524 .expect("register route whose receiver will be cancelled");
2525 drop(waiter);
2526
2527 let error = registry
2528 .settle(
2529 "session-a",
2530 Some(7),
2531 BridgeResponse {
2532 call_id: 75,
2533 status: 0,
2534 payload: vec![0xA7; 4],
2535 reservation: None,
2536 },
2537 )
2538 .expect_err("disconnected response target must reject delivery");
2539 assert!(error.contains("ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY"));
2540 assert_eq!(registry.pending_len(), 0);
2541 assert!(resources.is_zero());
2542
2543 let waiter = registry
2544 .register_sync(&runtime, 2, 4, 76, "session-a", Some(7))
2545 .expect("register route whose terminal error cannot be delivered");
2546 drop(waiter);
2547 let error = registry
2548 .settle(
2549 "session-a",
2550 Some(7),
2551 BridgeResponse {
2552 call_id: 76,
2553 status: 0,
2554 payload: vec![0; 5],
2555 reservation: None,
2556 },
2557 )
2558 .expect_err("terminal oversize error delivery must fail closed");
2559 assert!(error.contains("ERR_AGENTOS_BRIDGE_RESPONSE_DELIVERY"));
2560 assert_eq!(registry.pending_len(), 0);
2561 assert!(resources.is_zero());
2562 }
2563
2564 #[test]
2565 fn bridge_registry_rejects_producer_side_response_reservations() {
2566 let registry = BridgeCallRegistry::new(2);
2567 let (runtime, resources) = limited_bridge_runtime(2, 8, 1024);
2568 let waiter = registry
2569 .register_sync(&runtime, 1, 512, 77, "session-a", Some(7))
2570 .expect("pre-reserve declared response maximum");
2571 let producer_reservation = resources
2572 .reserve(ResourceClass::BridgeResponseBytes, 2)
2573 .expect("temporary producer encoding charge");
2574 assert_eq!(
2575 resources.usage(ResourceClass::BridgeResponseBytes).used,
2576 514
2577 );
2578
2579 let error = registry
2580 .settle(
2581 "session-a",
2582 Some(7),
2583 BridgeResponse {
2584 call_id: 77,
2585 status: 0,
2586 payload: vec![0xA7; 2],
2587 reservation: Some(agentos_runtime::accounting::SharedReservation::new(
2588 producer_reservation,
2589 )),
2590 },
2591 )
2592 .expect_err("producer-side response charging must fail closed");
2593 assert!(error.contains("ERR_AGENTOS_BRIDGE_RESPONSE_ACCOUNTING"));
2594
2595 let terminal = waiter.recv().expect("receive accounting error");
2596 assert_eq!(terminal.status, 1);
2597 assert!(String::from_utf8_lossy(&terminal.payload)
2598 .contains("ERR_AGENTOS_BRIDGE_RESPONSE_ACCOUNTING"));
2599 drop(terminal);
2600 assert!(resources.is_zero());
2601 }
2602
2603 #[test]
2604 fn bridge_registry_bounds_and_accounts_terminal_error_payloads() {
2605 let registry = BridgeCallRegistry::new(1);
2606 let (runtime, resources) = limited_bridge_runtime(1, 1, 4);
2607 let waiter = registry
2608 .register_sync(&runtime, 0, 4, 78, "session-a", Some(7))
2609 .expect("reserve a deliberately tiny response budget");
2610
2611 registry
2612 .settle(
2613 "session-a",
2614 Some(7),
2615 BridgeResponse {
2616 call_id: 78,
2617 status: 0,
2618 payload: vec![0; 5],
2619 reservation: None,
2620 },
2621 )
2622 .expect_err("oversize response must settle a bounded terminal error");
2623
2624 let terminal = waiter.recv().expect("receive bounded terminal error");
2625 assert_eq!(terminal.status, 1);
2626 assert_eq!(terminal.payload.len(), 4);
2627 assert_eq!(terminal.reservation.as_ref().unwrap().amount(), 4);
2628 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 4);
2629 drop(terminal);
2630 assert!(resources.is_zero());
2631 }
2632
2633 #[test]
2634 fn bridge_registry_limit_cancel_and_oversize_paths_leave_zero_accounting() {
2635 let registry = BridgeCallRegistry::new(4);
2636 let (runtime, resources) = limited_bridge_runtime(1, 4, 512);
2637 let _waiter = registry
2638 .register_sync(&runtime, 2, 4, 80, "session-a", Some(7))
2639 .expect("reserve first bridge call");
2640
2641 let error = registry
2642 .register_sync(&runtime, 1, 1, 81, "session-a", Some(7))
2643 .expect_err("second call must exceed VM call limit");
2644 assert!(error.contains("limits.reactor.maxBridgeCalls"));
2645 assert_eq!(resources.usage(ResourceClass::BridgeCalls).used, 1);
2646 assert_eq!(resources.usage(ResourceClass::BridgeRequestBytes).used, 2);
2647 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 4);
2648
2649 registry.cancel(80);
2650 assert_eq!(resources.usage(ResourceClass::BridgeCalls).used, 0);
2651 assert_eq!(resources.usage(ResourceClass::BridgeRequestBytes).used, 0);
2652 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 0);
2653
2654 let error = match registry.register_sync(&runtime, 5, 1, 81, "session-a", Some(7)) {
2655 Ok(_) => panic!("oversize request must fail admission"),
2656 Err(error) => error,
2657 };
2658 assert!(error.contains("limits.reactor.maxBridgeRequestBytes"));
2659 assert!(resources.is_zero());
2660
2661 let error = match registry.register_sync(&runtime, 2, 513, 81, "session-a", Some(7)) {
2662 Ok(_) => panic!("oversize declared response must fail admission"),
2663 Err(error) => error,
2664 };
2665 assert!(error.contains("limits.reactor.maxBridgeResponseBytes"));
2666 assert!(resources.is_zero());
2667
2668 let oversize_waiter = registry
2669 .register_sync(&runtime, 2, 512, 82, "session-a", Some(7))
2670 .expect("reserve bridge call for oversize response");
2671 let error = registry
2672 .settle(
2673 "session-a",
2674 Some(7),
2675 BridgeResponse {
2676 call_id: 82,
2677 status: 0,
2678 payload: vec![0; 513],
2679 reservation: None,
2680 },
2681 )
2682 .expect_err("oversize response must fail closed");
2683 assert!(error.contains("ERR_AGENTOS_BRIDGE_RESPONSE_LIMIT"));
2684 let terminal = oversize_waiter
2685 .recv()
2686 .expect("oversize response must settle its waiter with an error");
2687 assert_eq!(terminal.status, 1);
2688 assert!(String::from_utf8_lossy(&terminal.payload)
2689 .contains("ERR_AGENTOS_BRIDGE_RESPONSE_LIMIT"));
2690 assert_eq!(
2691 terminal.reservation.as_ref().unwrap().amount(),
2692 terminal.payload.len()
2693 );
2694 assert_eq!(resources.usage(ResourceClass::BridgeCalls).used, 0);
2695 assert_eq!(resources.usage(ResourceClass::BridgeRequestBytes).used, 0);
2696 assert_eq!(
2697 resources.usage(ResourceClass::BridgeResponseBytes).used,
2698 terminal.payload.len()
2699 );
2700 drop(terminal);
2701 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 0);
2702
2703 let _waiter = registry
2704 .register_sync(&runtime, 1, 2, 83, "session-a", Some(7))
2705 .expect("reserve bridge call for teardown");
2706 registry.cancel_session("session-a", Some(7));
2707 assert!(resources.is_zero());
2708 }
2709
2710 #[test]
2711 fn bridge_registry_pre_reserves_overlapping_declared_maxima_until_consumed() {
2712 let registry = BridgeCallRegistry::new(4);
2713 let (runtime, resources) = limited_bridge_runtime(4, 16, 8);
2714 let first = registry
2715 .register_sync(&runtime, 1, 4, 90, "session-a", Some(7))
2716 .expect("pre-reserve first declared maximum");
2717 let second = registry
2718 .register_sync(&runtime, 1, 4, 91, "session-a", Some(7))
2719 .expect("pre-reserve second declared maximum");
2720 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 8);
2721
2722 let error = registry
2723 .register_sync(&runtime, 1, 1, 92, "session-a", Some(7))
2724 .expect_err("a call without guaranteed response capacity must fail admission");
2725 assert!(error.contains("limits.reactor.maxBridgeResponseBytes"));
2726 assert_eq!(registry.pending_len(), 2);
2727 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 8);
2728
2729 registry
2730 .settle(
2731 "session-a",
2732 Some(7),
2733 BridgeResponse {
2734 call_id: 90,
2735 status: 0,
2736 payload: vec![0xA5],
2737 reservation: None,
2738 },
2739 )
2740 .expect("transfer first pre-reservation to its completion");
2741 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 5);
2742
2743 registry
2744 .settle(
2745 "session-a",
2746 Some(7),
2747 BridgeResponse {
2748 call_id: 91,
2749 status: 0,
2750 payload: vec![0xB6; 2],
2751 reservation: None,
2752 },
2753 )
2754 .expect("every admitted response must complete without reacquiring capacity");
2755 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 3);
2756
2757 let first_response = first.recv().expect("first completion");
2758 let second_response = second.recv().expect("second completion");
2759 assert_eq!(first_response.reservation.as_ref().unwrap().amount(), 1);
2760 assert_eq!(second_response.reservation.as_ref().unwrap().amount(), 2);
2761 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 3);
2762 drop(first_response);
2763 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 2);
2764 drop(second_response);
2765 assert!(resources.is_zero());
2766 }
2767
2768 #[test]
2769 fn bridge_registry_large_declared_responses_do_not_serialize_small_completions() {
2770 let registry = BridgeCallRegistry::new(2);
2771 let (runtime, resources) = limited_bridge_runtime(2, 2, 16 * 1024 * 1024);
2772 let first = registry
2773 .register_sync(&runtime, 1, 16 * 1024 * 1024, 93, "session-a", Some(7))
2774 .expect("admit first call with a large declared response maximum");
2775 let second = registry
2776 .register_sync(&runtime, 1, 16 * 1024 * 1024, 94, "session-a", Some(7))
2777 .expect("a large declared maximum must not monopolize concrete response capacity");
2778 assert_eq!(
2779 resources.usage(ResourceClass::BridgeResponseBytes).used,
2780 8 * 1024,
2781 "each admitted call keeps only its bounded terminal-response reservation"
2782 );
2783
2784 registry
2785 .settle(
2786 "session-a",
2787 Some(7),
2788 BridgeResponse {
2789 call_id: 93,
2790 status: 0,
2791 payload: vec![0xA5],
2792 reservation: None,
2793 },
2794 )
2795 .expect("settle first small concrete response");
2796 registry
2797 .settle(
2798 "session-a",
2799 Some(7),
2800 BridgeResponse {
2801 call_id: 94,
2802 status: 0,
2803 payload: vec![0xB6; 2],
2804 reservation: None,
2805 },
2806 )
2807 .expect("settle second small concrete response");
2808
2809 let first_response = first.recv().expect("first completion");
2810 let second_response = second.recv().expect("second completion");
2811 assert_eq!(first_response.payload, vec![0xA5]);
2812 assert_eq!(second_response.payload, vec![0xB6; 2]);
2813 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 3);
2814 drop(first_response);
2815 drop(second_response);
2816 assert!(resources.is_zero());
2817 }
2818
2819 #[test]
2820 fn bridge_registry_grows_floor_to_concrete_response_size_before_delivery() {
2821 let registry = BridgeCallRegistry::new(1);
2822 let (runtime, resources) = limited_bridge_runtime(1, 1, 8 * 1024);
2823 let waiter = registry
2824 .register_sync(&runtime, 1, 8 * 1024, 95, "session-a", Some(7))
2825 .expect("admit response with the bounded terminal floor");
2826 assert_eq!(
2827 resources.usage(ResourceClass::BridgeResponseBytes).used,
2828 BRIDGE_TERMINAL_RESPONSE_RESERVATION_BYTES
2829 );
2830
2831 registry
2832 .settle(
2833 "session-a",
2834 Some(7),
2835 BridgeResponse {
2836 call_id: 95,
2837 status: 0,
2838 payload: vec![0xA5; 6 * 1024],
2839 reservation: None,
2840 },
2841 )
2842 .expect("grow concrete response reservation before delivery");
2843 let response = waiter.recv().expect("receive grown response");
2844 assert_eq!(response.payload.len(), 6 * 1024);
2845 assert_eq!(response.reservation.as_ref().unwrap().amount(), 6 * 1024);
2846 assert_eq!(
2847 resources.usage(ResourceClass::BridgeResponseBytes).used,
2848 6 * 1024
2849 );
2850 drop(response);
2851 assert!(resources.is_zero());
2852 }
2853
2854 #[test]
2855 fn bridge_registry_growth_failure_delivers_bounded_error_and_releases_floor() {
2856 let registry = BridgeCallRegistry::new(2);
2857 let (runtime, resources) = limited_bridge_runtime(2, 2, 8 * 1024);
2858 let first = registry
2859 .register_sync(&runtime, 1, 8 * 1024, 96, "session-a", Some(7))
2860 .expect("admit first response floor");
2861 let _second = registry
2862 .register_sync(&runtime, 1, 8 * 1024, 97, "session-a", Some(7))
2863 .expect("admit second response floor");
2864 assert_eq!(
2865 resources.usage(ResourceClass::BridgeResponseBytes).used,
2866 8 * 1024
2867 );
2868
2869 let error = registry
2870 .settle(
2871 "session-a",
2872 Some(7),
2873 BridgeResponse {
2874 call_id: 96,
2875 status: 0,
2876 payload: vec![0xA5; 5 * 1024],
2877 reservation: None,
2878 },
2879 )
2880 .expect_err("overlapping floors must prevent unbounded concrete growth");
2881 assert!(error.contains("ERR_AGENTOS_BRIDGE_RESPONSE_LIMIT"));
2882 let terminal = first.recv().expect("receive bounded growth error");
2883 assert_eq!(terminal.status, 1);
2884 assert!(String::from_utf8_lossy(&terminal.payload)
2885 .contains("ERR_AGENTOS_BRIDGE_RESPONSE_LIMIT"));
2886 assert_eq!(
2887 terminal.reservation.as_ref().unwrap().amount(),
2888 terminal.payload.len()
2889 );
2890 assert_eq!(registry.pending_len(), 1);
2891 drop(terminal);
2892 assert_eq!(
2893 resources.usage(ResourceClass::BridgeResponseBytes).used,
2894 BRIDGE_TERMINAL_RESPONSE_RESERVATION_BYTES
2895 );
2896 registry.cancel(97);
2897 assert!(resources.is_zero());
2898 }
2899
2900 #[test]
2901 fn async_bridge_admission_happens_before_host_visibility() {
2902 let (runtime, resources) = limited_bridge_runtime(1, 4, 4);
2903 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(4));
2904 let (event_tx, event_rx) = crossbeam_channel::unbounded();
2905 let (async_tx, _async_rx) = crossbeam_channel::bounded(4);
2906 let (_abort_tx, abort_rx) = crossbeam_channel::bounded(1);
2907 let ctx = BridgeCallContext::with_registry(
2908 Box::new(ChannelRuntimeEventSender::new(event_tx, Some(7))),
2909 String::from("session-a"),
2910 Some(7),
2911 Arc::clone(®istry),
2912 Arc::new(AtomicU64::new(1)),
2913 async_tx,
2914 abort_rx,
2915 runtime,
2916 Arc::new(crate::session::SessionPauseControl::default()),
2917 DEFAULT_BRIDGE_CALL_TIMEOUT,
2918 );
2919
2920 let first = ctx
2921 .prepare_async_call_with_max_response_bytes("_first", Vec::new(), 4)
2922 .expect("admit first async bridge call");
2923 assert!(
2924 event_rx.is_empty(),
2925 "registration must not make the request host-visible"
2926 );
2927 ctx.dispatch_async_call(first)
2928 .expect("dispatch admitted bridge call");
2929 assert_eq!(event_rx.len(), 1);
2930
2931 let error = match ctx.prepare_async_call_with_max_response_bytes("_second", Vec::new(), 1) {
2932 Ok(_) => panic!("VM call admission must fail before dispatch"),
2933 Err(error) => error,
2934 };
2935 assert!(error.contains("limits.reactor.maxBridgeCalls"));
2936 assert_eq!(
2937 event_rx.len(),
2938 1,
2939 "rejected call must never reach the host event lane"
2940 );
2941
2942 drop(ctx);
2943 assert!(registry.pending_len() == 0);
2944 assert_ledger_settles_to_zero(&resources);
2945 }
2946
2947 #[test]
2948 fn dropping_prepared_async_call_cancels_route_and_releases_accounting() {
2949 let (runtime, resources) = limited_bridge_runtime(1, 4, 4);
2950 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(4));
2951 let (event_tx, event_rx) = crossbeam_channel::unbounded();
2952 let (async_tx, _async_rx) = crossbeam_channel::bounded(1);
2953 let (_abort_tx, abort_rx) = crossbeam_channel::bounded(1);
2954 let ctx = BridgeCallContext::with_registry(
2955 Box::new(ChannelRuntimeEventSender::new(event_tx, Some(7))),
2956 String::from("session-a"),
2957 Some(7),
2958 Arc::clone(®istry),
2959 Arc::new(AtomicU64::new(1)),
2960 async_tx,
2961 abort_rx,
2962 runtime,
2963 Arc::new(crate::session::SessionPauseControl::default()),
2964 DEFAULT_BRIDGE_CALL_TIMEOUT,
2965 );
2966
2967 let prepared = ctx
2968 .prepare_async_call_with_max_response_bytes("_cancelled", vec![0xA5; 2], 4)
2969 .expect("admit prepared async bridge call");
2970 assert_eq!(registry.pending_len(), 1);
2971 assert_eq!(resources.usage(ResourceClass::BridgeCalls).used, 1);
2972 assert_eq!(resources.usage(ResourceClass::BridgeRequestBytes).used, 2);
2973 assert_eq!(resources.usage(ResourceClass::BridgeResponseBytes).used, 4);
2974 assert!(event_rx.is_empty());
2975
2976 drop(prepared);
2977 assert_eq!(registry.pending_len(), 0);
2978 assert_ledger_settles_to_zero(&resources);
2979 assert!(event_rx.is_empty());
2980 }
2981
2982 #[test]
2983 fn failed_async_dispatch_cancels_route_and_releases_accounting() {
2984 let (runtime, resources) = limited_bridge_runtime(1, 4, 4);
2985 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(4));
2986 let (async_tx, _async_rx) = crossbeam_channel::bounded(1);
2987 let (_abort_tx, abort_rx) = crossbeam_channel::bounded(1);
2988 let ctx = BridgeCallContext::with_registry(
2989 Box::new(RejectingRuntimeEventSender),
2990 String::from("session-a"),
2991 Some(7),
2992 Arc::clone(®istry),
2993 Arc::new(AtomicU64::new(1)),
2994 async_tx,
2995 abort_rx,
2996 runtime,
2997 Arc::new(crate::session::SessionPauseControl::default()),
2998 DEFAULT_BRIDGE_CALL_TIMEOUT,
2999 );
3000
3001 let prepared = ctx
3002 .prepare_async_call_with_max_response_bytes("_dispatch", vec![0xA5; 2], 4)
3003 .expect("admit prepared async bridge call");
3004 assert_eq!(registry.pending_len(), 1);
3005 let error = ctx
3006 .dispatch_async_call(prepared)
3007 .expect_err("injected event-lane failure must reject dispatch");
3008 assert!(error.contains("injected event-lane failure"));
3009 assert_eq!(registry.pending_len(), 0);
3010 assert_ledger_settles_to_zero(&resources);
3011 }
3012
3013 #[test]
3014 fn rejected_deadline_task_cancels_unpublished_route_and_releases_accounting() {
3015 let (runtime, resources) = limited_bridge_runtime(1, 4, 4);
3016 runtime.close_admission();
3017 let registry: CallIdRouter = Arc::new(BridgeCallRegistry::new(4));
3018 let (event_tx, event_rx) = crossbeam_channel::unbounded();
3019 let (async_tx, _async_rx) = crossbeam_channel::bounded(1);
3020 let (_abort_tx, abort_rx) = crossbeam_channel::bounded(1);
3021 let ctx = BridgeCallContext::with_registry(
3022 Box::new(ChannelRuntimeEventSender::new(event_tx, Some(7))),
3023 String::from("session-a"),
3024 Some(7),
3025 Arc::clone(®istry),
3026 Arc::new(AtomicU64::new(1)),
3027 async_tx,
3028 abort_rx,
3029 runtime,
3030 Arc::new(crate::session::SessionPauseControl::default()),
3031 DEFAULT_BRIDGE_CALL_TIMEOUT,
3032 );
3033
3034 let error = match ctx.prepare_async_call_with_max_response_bytes(
3035 "_timer-rejected",
3036 vec![0xA5; 2],
3037 4,
3038 ) {
3039 Ok(_) => panic!("closed task admission must reject the deadline task"),
3040 Err(error) => error,
3041 };
3042 assert!(error.contains("ERR_AGENTOS_BRIDGE_DEADLINE_TASK"));
3043 assert_eq!(registry.pending_len(), 0);
3044 assert!(event_rx.is_empty());
3045 assert!(resources.is_zero());
3046 }
3047
3048 #[test]
3049 fn async_bridge_admission_reserves_every_response_lane_slot() {
3050 let (runtime, resources) = limited_bridge_runtime(3, 8, 8);
3051 let registry = BridgeCallRegistry::new(4);
3052 let (response_tx, response_rx) = crossbeam_channel::bounded(2);
3053
3054 for call_id in [101, 102] {
3055 registry
3056 .register_async(
3057 &runtime,
3058 1,
3059 1,
3060 call_id,
3061 "session-a",
3062 Some(7),
3063 response_tx.clone(),
3064 )
3065 .expect("each physical response slot admits one async call");
3066 }
3067 let error = registry
3068 .register_async(
3069 &runtime,
3070 1,
3071 1,
3072 103,
3073 "session-a",
3074 Some(7),
3075 response_tx.clone(),
3076 )
3077 .expect_err("an async call without a response slot must fail admission");
3078 assert!(error.contains("ERR_AGENTOS_BRIDGE_RESPONSE_LANE_LIMIT"));
3079
3080 for call_id in [101, 102] {
3081 registry
3082 .settle(
3083 "session-a",
3084 Some(7),
3085 BridgeResponse {
3086 call_id,
3087 status: 0,
3088 payload: vec![call_id as u8],
3089 reservation: None,
3090 },
3091 )
3092 .expect("every admitted response must fit before the lane drains");
3093 }
3094 let error = registry
3095 .register_async(
3096 &runtime,
3097 1,
3098 1,
3099 103,
3100 "session-a",
3101 Some(7),
3102 response_tx.clone(),
3103 )
3104 .expect_err("queued responses must continue to own their lane slots");
3105 assert!(error.contains("ERR_AGENTOS_BRIDGE_RESPONSE_LANE_LIMIT"));
3106
3107 drop(
3108 response_rx
3109 .recv()
3110 .expect("drain one reserved response slot"),
3111 );
3112 registry
3113 .register_async(&runtime, 1, 1, 103, "session-a", Some(7), response_tx)
3114 .expect("draining a response releases exactly one admission slot");
3115 registry.cancel(103);
3116 drop(response_rx.recv().expect("drain remaining response"));
3117 assert!(resources.is_zero());
3118 }
3119
3120 #[test]
3121 fn writer_runtime_event_sender_serializes_events() {
3122 let (tx, rx) = crossbeam_channel::unbounded();
3123 let sender = super::ChannelRuntimeEventSender::new(tx, None);
3124
3125 for i in 0..5 {
3127 let event = RuntimeEvent::BridgeCall {
3128 session_id: "sess-1".into(),
3129 call_id: i,
3130 method: "_fn".into(),
3131 payload: vec![0xAA; 100 * (i as usize + 1)],
3132 };
3133 sender.send_event(event).expect("send_event");
3134 }
3135
3136 for i in 0..5u64 {
3138 let decoded = rx.recv().expect("recv");
3139 match decoded.event {
3140 RuntimeEvent::BridgeCall {
3141 call_id, payload, ..
3142 } => {
3143 assert_eq!(call_id, i);
3144 assert_eq!(payload.len(), 100 * (i as usize + 1));
3145 }
3146 _ => panic!("expected BridgeCall"),
3147 }
3148 }
3149
3150 let small = RuntimeEvent::Log {
3152 session_id: "s".into(),
3153 channel: 0,
3154 message: "x".into(),
3155 };
3156 sender.send_event(small.clone()).expect("send_event");
3157 let decoded = rx.recv().expect("recv");
3158 assert_eq!(decoded.event, small);
3159 }
3160
3161 #[test]
3162 fn stub_context_panics_on_sync_call() {
3163 let ctx = BridgeCallContext::stub();
3164 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3165 let _ = ctx.sync_call("_fsReadFile", vec![]);
3166 }));
3167 assert!(result.is_err(), "stub sync_call should panic");
3168 }
3169
3170 #[test]
3171 fn stub_context_panics_on_async_send() {
3172 let ctx = BridgeCallContext::stub();
3173 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3174 let _ = ctx.async_send("_asyncFn", vec![]);
3175 }));
3176 assert!(result.is_err(), "stub async_send should panic");
3177 }
3178}