Skip to main content

agentos_v8_runtime/
host_call.rs

1// Sync-blocking bridge call: serialize, write to socket, block on read, deserialize
2
3use 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
16// ── Sync bridge-call round-trip latency (opt-in via AGENTOS_SYNCRPC_LAT=1) ──
17// Measures the guest-observed cost of one host_call round trip (send + block on
18// response). If this is ~ms per call, the embedded-V8 IPC floor is a big part of
19// the evaluate-phase VM tax (~50 fs sync RPCs during SDK init). Writes running
20// (count, total_us, max_us) to AGENTOS_SYNCRPC_LAT_FILE.
21static 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
122/// Trait for sending serialized frames to the host without holding a shared mutex.
123/// Production code uses ChannelRuntimeEventSender (lock-free MPSC); tests use WriterRuntimeEventSender.
124pub trait RuntimeEventSender: Send {
125    fn send_event(&self, event: RuntimeEvent) -> Result<(), String>;
126}
127
128/// Sends frames via a crossbeam channel to a dedicated writer thread.
129/// Maintains a reusable frame buffer that grows to high-water mark,
130/// avoiding per-call allocation for frame construction.
131pub struct ChannelRuntimeEventSender {
132    pub tx: crate::session::RuntimeEventSender,
133    output_generation: Option<u64>,
134    /// Pre-allocated frame buffer reused across send_frame calls.
135    /// Grows to high-water mark; cleared (not deallocated) between calls.
136    #[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/// Sends frames directly to a Write impl (used by tests).
165#[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
178/// Trait for receiving a BridgeResponse directly without re-serialization.
179/// Production code uses a channel-based implementation; tests use a buffer-based one.
180pub 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/// ResponseReceiver that reads frames from a byte buffer via ipc_binary::read_frame.
192/// Used by tests and any code that has a pre-serialized byte stream.
193#[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    // Preserve the stable typed code whenever the declared response capacity
307    // can hold it. Production bridge declarations reserve at least 4 KiB; the
308    // shorter fallback only applies to synthetic or misconfigured callers.
309    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
395/// Bounded call-specific response registry shared by all sessions.
396///
397/// A response is settled directly into the target registered for its globally
398/// unique call ID. It never enters the ordinary session command/event channel.
399pub struct BridgeCallRegistry {
400    pending: Mutex<HashMap<u64, BridgeCallTarget>>,
401    /// Bounded proof that a host-visible call was canceled before settlement.
402    /// This distinguishes expected teardown completions from arbitrary unknown
403    /// or duplicate responses without keeping retired calls forever.
404    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)] // one immutable identity/admission tuple per call route
422    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        // Reserve enough capacity to guarantee that every admitted call can
462        // settle with a typed terminal error. Reserving the full per-method
463        // maximum here would make one 16 MiB-capable call consume the entire
464        // default VM budget and reject unrelated concurrent calls even when all
465        // concrete responses are tiny. Settlement grows this reservation to the
466        // actual payload size before it can enter the response lane.
467        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    /// Marks the point after which cancellation can race a legitimate host
541    /// response. Callers invoke this immediately before publishing the event.
542    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)] // immutable identity/admission tuple for one direct route
602    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        // Identity validation must not consume a legitimate route when a stale
733        // response arrives. Once identity is validated, however, settlement is
734        // terminal: take the target before attempting any delivery so every
735        // success or failure path drops or transfers its call, request, and
736        // response reservations exactly once. Keep the registry lock through
737        // delivery so an async response lane's registered slot cannot be
738        // re-admitted in the gap between the take and try_send.
739        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                            // A session cancellation can race between marking
877                            // the route and the failed event send. Retract any
878                            // conservative tombstone once publication is known
879                            // to have failed.
880                            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
986/// Compatibility name retained while callers migrate from the old
987/// call_id-to-session map to direct call-specific settlement.
988pub type CallIdRouter = Arc<BridgeCallRegistry>;
989
990/// Shared call_id counter type. Sessions sharing a CallIdRouter must use the same
991/// counter to prevent call_id collisions that cause BridgeResponses to be delivered
992/// to the wrong session.
993pub type SharedCallIdCounter = Arc<AtomicU64>;
994
995/// Cancels a registered route unless ownership has crossed the host-visibility
996/// boundary. This closes the prepare-without-dispatch cancellation window for
997/// async bridge calls without requiring callers to remember manual cleanup.
998struct 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
1038/// Context for sync-blocking bridge calls from a V8 session.
1039///
1040/// Holds the frame sender and response receiver, session ID, call_id counter,
1041/// and pending-call tracking. Used by V8 FunctionTemplate callbacks to
1042/// implement the sync-blocking bridge pattern.
1043pub struct BridgeCallContext {
1044    /// Sender for serialized frames to the host (channel-based in production)
1045    sender: Box<dyn RuntimeEventSender>,
1046    /// Receiver for BridgeResponse frames (no re-serialization needed)
1047    response_rx: Option<Mutex<Box<dyn BridgeResponseReceiver>>>,
1048    /// Session ID included in every BridgeCall
1049    pub session_id: String,
1050    /// Monotonically increasing call_id counter. Sessions sharing a CallIdRouter
1051    /// must share the same counter (via Arc) to prevent call_id collisions.
1052    next_call_id: Arc<AtomicU64>,
1053    /// Set of in-flight call_ids (for duplicate rejection)
1054    pending_calls: Mutex<HashSet<u64>>,
1055    /// Opt-in diagnostic tracking for sync call_ids. The atomic call_id counter
1056    /// plus recv_response(call_id) validation are the correctness path; this set
1057    /// is only needed when inspecting in-flight calls.
1058    track_pending_calls: bool,
1059    /// Optional direct call-specific response registry.
1060    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    /// Execution gate shared with the owning session. Direct response routing
1065    /// bypasses the ordinary command lane, but a sync response must still stop
1066    /// at this boundary while the VM is paused.
1067    pause_control: Option<Arc<crate::session::SessionPauseControl>>,
1068    /// Session-injected process scheduler used by local bridge operations such
1069    /// as `node:vm` timeouts. Snapshot/test-only contexts may omit it when they
1070    /// never arm runtime work.
1071    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/// No-op FrameSender for snapshot stub functions.
1082/// Panics if called — stubs must never be invoked during snapshot creation.
1083#[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/// No-op ResponseReceiver for snapshot stub functions.
1095/// Panics if called — stubs must never be invoked during snapshot creation.
1096#[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    /// Create a no-op BridgeCallContext for snapshot stub functions.
1110    /// Panics if sync_call or async_send is called — stubs exist only for
1111    /// the bridge IIFE to reference (not call) during snapshot creation.
1112    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    /// Create a BridgeCallContext with a byte writer and reader (wraps in WriterFrameSender
1131    /// and ReaderResponseReceiver). Convenient for tests that pre-serialize BridgeResponse bytes.
1132    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    /// Create a BridgeCallContext with a FrameSender, ResponseReceiver, call_id routing table,
1159    /// and shared call_id counter. All sessions sharing the same CallIdRouter must share
1160    /// the same counter to prevent call_id collisions in the routing table.
1161    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    /// Perform a sync-blocking bridge call.
1225    ///
1226    /// Generates a unique call_id, sends a BridgeCall message over IPC,
1227    /// blocks on read() for the BridgeResponse, and returns the result.
1228    /// Error responses from the host are returned as Err.
1229    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        // Optional diagnostic tracking. Correctness comes from the atomic
1248        // counter and recv_response(call_id) identity validation.
1249        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        // Send BridgeCall to host
1290        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        // Receive BridgeResponse directly (no re-serialization)
1315        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        // Validate and extract BridgeResponse
1395        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                // Registration already owns call/request/response capacity.
1470                // If supervision rejects the timer, retract the unpublished
1471                // route before returning so admission cannot leak permanently.
1472                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    /// Legacy one-step helper used by non-concurrent tests.
1537    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    /// Check if a call_id is currently pending.
1550    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    /// Number of pending calls.
1558    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    /// Shared writer that captures output for test inspection
1641    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    /// Serialize a BridgeResponse into length-prefixed binary frame bytes
1661    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        // Verify the BridgeCall was written correctly
1703        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        // Prepare two sequential responses
1752        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        // Empty reader = EOF error; call_id should be cleaned up
1770        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        // Response has call_id=99 but expected call_id=1
1784        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        // Response is not a BridgeResponse
1799        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        // Verify the BridgeCall was written correctly
1834        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        // Sync call uses call_id=1, async_send should get call_id=2
1870        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        // Verify the received event matches without any BinaryFrame hop.
1896        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        // Multiple senders can send concurrently without blocking each other
1904        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); // Drop original sender so rx closes when threads finish
1922
1923        for h in handles {
1924            h.join().expect("thread join");
1925        }
1926
1927        // All 40 frames should arrive and be decodable
1928        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        // Verify BridgeCallContext works with ChannelRuntimeEventSender end-to-end
1938        let (tx, rx) = crossbeam_channel::unbounded();
1939
1940        // Pre-serialize a BridgeResponse for the reader
1941        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        // Verify the BridgeCall went through the channel
1958        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(&registry),
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(&registry),
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(&registry),
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(&registry),
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(&registry),
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(&registry),
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(&registry),
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(&registry),
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        // Send multiple frames — buffer grows to high-water mark
3126        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        // Verify all events arrive with their payload intact.
3137        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        // Small follow-up events still go through the same sender.
3151        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}