Skip to main content

av_harness/
worker.rs

1//! Bounded asynchronous worker for loop analysis, event emission, and capture.
2
3use crate::session::{Session, Workflow};
4use av_bridge::{EventBus, PublishAck};
5use av_core::metrics::Registry;
6use av_events::{EventClass, EventMetrics, OcsfEventBuilder, StatusId, StopReason};
7use av_loopdetect::{BreakerVerdict, Embedder, NoopVectorSink, VectorSink};
8use serde_json::Value;
9use std::sync::atomic::Ordering;
10use std::sync::Arc;
11use tokio::sync::{mpsc, oneshot};
12use tracing::Instrument as _;
13
14/// Structured ATIF representation attached to an asynchronous event job.
15pub struct AtifCapture {
16    /// Step originator.
17    pub source: av_atif::Source,
18    /// Required dialog message.
19    pub message: Value,
20    /// Optional explicit reasoning.
21    pub reasoning_content: Option<String>,
22    /// Model used for this step.
23    pub model_name: Option<String>,
24    /// Structured tool calls.
25    pub tool_calls: Option<Vec<av_atif::ToolCall>>,
26    /// Structured environment feedback.
27    pub observation: Option<av_atif::Observation>,
28    /// Number of represented LLM calls.
29    pub llm_call_count: Option<u64>,
30}
31
32/// Work copied from the hot path for asynchronous processing.
33pub struct WorkerJob {
34    /// Session receiving the resulting event and trajectory step.
35    pub session: Arc<Session>,
36    /// Identity validated for the request that created this job.
37    pub identity: av_events::AgentIdentity,
38    /// Event class to emit when loop detection does not override it.
39    pub class: EventClass,
40    /// Class-specific event payload.
41    pub payload: Value,
42    /// Reasoning or response text used for loop detection and ATIF capture.
43    pub text: String,
44    /// Whether this job represents a reasoning step that should update the
45    /// semantic loop breaker.
46    pub analyze_loop: bool,
47    /// Event outcome.
48    pub status: StatusId,
49    /// Normalized stop reason for stop events.
50    pub stop_reason: Option<StopReason>,
51    /// Provider or source-native stop reason value.
52    pub native_stop_reason: Option<String>,
53    /// Token and compression metrics.
54    pub metrics: EventMetrics,
55    /// Cost attributed to this step, in micro-USD.
56    pub cost_usd_micros: u64,
57    /// ATIF step representation for unsigned workflows.
58    pub atif: Option<AtifCapture>,
59    /// Durable response-attempt marker cleared only after journal and broker commit.
60    pub response_marker: Option<String>,
61    /// Chat response attempt correlated across request and terminal capture records.
62    pub response_attempt: Option<ResponseAttempt>,
63}
64
65/// Durable request/terminal marker embedded in the active event journal.
66#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
67pub struct ResponseAttempt {
68    /// Stable ID shared by request admission and its terminal response event.
69    pub id: String,
70    /// False on request admission, true on response completion or dispatch failure.
71    pub terminal: bool,
72}
73
74#[derive(serde::Serialize, serde::Deserialize)]
75struct InFlightResponse {
76    session_id: String,
77    attempt_id: String,
78    request_digest: String,
79}
80
81#[derive(Debug, serde::Serialize, serde::Deserialize)]
82pub(crate) struct ActiveJournalRecord {
83    pub(crate) event: Value,
84    pub(crate) identity: av_events::AgentIdentity,
85    pub(crate) atif_step: Option<av_atif::Step>,
86    pub(crate) tool_calls: u64,
87    pub(crate) tool_allowed: u64,
88    pub(crate) tool_blocked: u64,
89    pub(crate) prompt_tokens: u64,
90    pub(crate) completion_tokens: u64,
91    pub(crate) cached_tokens: u64,
92    pub(crate) cost_usd_micros: u64,
93    pub(crate) stop_reason_id: Option<u8>,
94    pub(crate) response_attempt: Option<ResponseAttempt>,
95}
96
97#[derive(Debug, serde::Serialize, serde::Deserialize)]
98struct BrokerAckRecord {
99    session_id: String,
100    event_uid: String,
101    ack: PublishAck,
102}
103
104struct Envelope {
105    job: WorkerJob,
106    completion: Option<oneshot::Sender<Result<(), String>>>,
107    span: tracing::Span,
108    _capacity_permit: tokio::sync::OwnedSemaphorePermit,
109}
110
111/// Non-blocking submission error.
112#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
113pub enum SubmitError {
114    /// The bounded worker queue has no remaining capacity.
115    #[error("worker queue is full")]
116    Full,
117    /// The worker supervisor has stopped.
118    #[error("worker queue is closed")]
119    Closed,
120}
121
122/// Which admission stage a drop should be counted against. Each variant
123/// maps to a distinct `av_events_dropped_total{stage="..."}` counter,
124/// so operators can PromQL-alert on the actual bottleneck class.
125#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126enum DropStage {
127    /// Initial worker-side admission. `try_reserve` / `try_reserve_pair`'s
128    /// worker half. Bumps `av_events_dropped_total{stage="worker_queue"}`.
129    WorkerQueue,
130    /// Downstream response-capture submission (`ResponsePermit::submit`).
131    /// The permit itself only holds the response-capacity semaphore; the
132    /// mpsc slot is contested at submit time, so exhaustion at THIS
133    /// stage — which is a real observable class of failure — bumps
134    /// `av_events_dropped_total{stage="response_slot"}`.
135    ResponseSlot,
136}
137
138impl DropStage {
139    fn full_counter_key(self) -> &'static str {
140        match self {
141            Self::WorkerQueue => "av_events_dropped_total{stage=\"worker_queue\"}",
142            Self::ResponseSlot => "av_events_dropped_total{stage=\"response_slot\"}",
143        }
144    }
145
146    fn label(self) -> &'static str {
147        match self {
148            Self::WorkerQueue => "worker_queue",
149            Self::ResponseSlot => "response_slot",
150        }
151    }
152}
153
154/// Cloneable handle used by request handlers to submit worker jobs.
155#[derive(Clone)]
156pub struct WorkerHandle {
157    senders: Arc<Vec<mpsc::Sender<Envelope>>>,
158    metrics: Arc<Registry>,
159    pending: Arc<std::sync::atomic::AtomicU64>,
160    drained: Arc<tokio::sync::Notify>,
161    capacity: Arc<tokio::sync::Semaphore>,
162    /// Separate admission budget for downstream response jobs. Every
163    /// chat request reserves a worker slot AND a response slot; giving
164    /// them distinct semaphores + distinct
165    /// `av_events_dropped_total{stage=...}` counters lets operators see
166    /// which class of capacity ran out. Prior to this split both
167    /// reservations pulled from the same semaphore so admission was
168    /// silently halved and both failures aliased to `stage="worker_queue"`.
169    response_capacity: Arc<tokio::sync::Semaphore>,
170}
171
172/// Guaranteed slot in the bounded worker queue.
173pub struct WorkerPermit {
174    permit: mpsc::OwnedPermit<Envelope>,
175    pending: Arc<std::sync::atomic::AtomicU64>,
176    capacity_permit: tokio::sync::OwnedSemaphorePermit,
177}
178
179/// Reservation for the downstream response-capture worker job. Draws
180/// from a distinct capacity semaphore so operators can tell via
181/// `av_events_dropped_total{stage="response_slot"}` vs
182/// `stage="worker_queue"` which class of admission is the bottleneck.
183///
184/// Unlike [`WorkerPermit`], the response permit only holds the
185/// capacity semaphore — the mpsc queue slot is re-acquired at submit
186/// time via [`Self::submit`]. This keeps the initial admission cheap
187/// (one semaphore acquire instead of two mpsc reservations that would
188/// otherwise compete for the same shard's slot with the worker
189/// permit) and lets the response job land whenever worker capacity
190/// and the shard have room, which is the common case since response
191/// capture happens
192/// tens-of-seconds after the initial worker job has drained.
193///
194/// Held by the streaming-response wrapper in `routes.rs`
195/// (`AbortFinalizingStream`) for the lifetime of the forwarded response;
196/// drops on stream completion or client abort.
197pub struct ResponsePermit {
198    _capacity_permit: tokio::sync::OwnedSemaphorePermit,
199}
200
201impl ResponsePermit {
202    /// Commit a response-capture job. Consumes the response permit's
203    /// capacity slot and races for a shard's mpsc slot; if the shard is
204    /// momentarily full, returns a `SubmitError::Full` and bumps
205    /// `av_events_dropped_total{stage="response_slot"}` — NOT the
206    /// worker_queue counter. That distinction is the whole point of
207    /// the split: operators need to see which class of exhaustion is
208    /// producing drops.
209    pub fn submit(self, worker: &WorkerHandle, job: WorkerJob) -> Result<(), SubmitError> {
210        // Release the response-capacity permit before contending for
211        // the mpsc slot so a failed submit does not artificially pin
212        // the response budget. Dropping via NLL — the permit is not
213        // used by `try_submit_labeled`, which draws a fresh
214        // worker-capacity slot for the actual queue admission.
215        let _capacity_permit = self._capacity_permit;
216        worker.try_submit_labeled(job, DropStage::ResponseSlot)
217    }
218}
219
220/// Fused reservation covering a worker job AND its downstream response
221/// slot. Acquired atomically at request admission: on any failure the
222/// worker slot (if already held) is released via RAII before the error
223/// surfaces, so callers cannot end up with an orphaned half-permit.
224pub struct WorkerAndResponsePermit {
225    /// Permit for the initial worker job (dispatch / quota / receipt-sign).
226    pub worker: WorkerPermit,
227    /// Permit for the downstream response-capture job.
228    pub response: ResponsePermit,
229}
230
231impl WorkerPermit {
232    /// Commit a job into the previously reserved queue slot.
233    pub fn submit(self, job: WorkerJob) {
234        job.session.worker_job_started();
235        self.pending.fetch_add(1, Ordering::AcqRel);
236        let span = worker_span(&job);
237        self.permit.send(Envelope {
238            job,
239            completion: None,
240            span,
241            _capacity_permit: self.capacity_permit,
242        });
243    }
244}
245
246impl WorkerHandle {
247    /// Reserve queue capacity before committing quota or other state.
248    pub fn try_reserve(&self, session_id: &str) -> Result<WorkerPermit, SubmitError> {
249        let capacity_permit = self.try_capacity()?;
250        let Some(sender) = self.sender_for(session_id).cloned() else {
251            return Err(SubmitError::Closed);
252        };
253        match sender.try_reserve_owned() {
254            Ok(permit) => Ok(WorkerPermit {
255                permit,
256                pending: Arc::clone(&self.pending),
257                capacity_permit,
258            }),
259            Err(mpsc::error::TrySendError::Full(_sender)) => {
260                self.metrics
261                    .counter(
262                        "av_events_dropped_total{stage=\"worker_queue\"}",
263                        "Worker jobs dropped",
264                    )
265                    .inc();
266                Err(SubmitError::Full)
267            }
268            Err(mpsc::error::TrySendError::Closed(_sender)) => {
269                self.metrics
270                    .counter(
271                        "av_events_dropped_total{stage=\"worker_closed\"}",
272                        "Worker jobs dropped",
273                    )
274                    .inc();
275                Err(SubmitError::Closed)
276            }
277        }
278    }
279
280    /// Reserve the downstream response-capture slot. Distinct from
281    /// [`Self::try_reserve`] because response capture is a separate
282    /// stage with its own admission budget; keeping the counters
283    /// separate lets operators see which one ran out. The mpsc queue
284    /// slot is re-acquired at submit time via [`ResponsePermit::submit`].
285    fn try_reserve_response(&self) -> Result<ResponsePermit, SubmitError> {
286        Arc::clone(&self.response_capacity)
287            .try_acquire_owned()
288            .map(|p| ResponsePermit { _capacity_permit: p })
289            .map_err(|_| {
290                self.metrics
291                    .counter(
292                        "av_events_dropped_total{stage=\"response_slot\"}",
293                        "Response-slot reservations that failed",
294                    )
295                    .inc();
296                SubmitError::Full
297            })
298    }
299
300    /// Atomic acquire-both-or-fail for admission: obtains a worker
301    /// permit AND a response permit for the same request, distinct
302    /// counters on failure. On any error path (worker slot or response
303    /// slot exhaustion, closed channel), a partially-taken worker
304    /// permit drops via RAII before the error surfaces so callers
305    /// never observe an orphaned half-reservation.
306    pub fn try_reserve_pair(&self, session_id: &str) -> Result<WorkerAndResponsePermit, SubmitError> {
307        // Stage 1: worker slot. Bumps stage="worker_queue" on failure.
308        let worker = self.try_reserve(session_id)?;
309        // Stage 2: response slot. Bumps stage="response_slot" on
310        // failure. If this fails, `worker` drops via NLL — its
311        // OwnedSemaphorePermit + mpsc::OwnedPermit release cleanly
312        // and the counters agree with the caller's view (only the
313        // response-slot counter incremented).
314        let response = self.try_reserve_response()?;
315        Ok(WorkerAndResponsePermit { worker, response })
316    }
317
318    /// Submit without waiting. Used by rejection/failure paths; the hot path
319    /// reserves capacity up front via [`WorkerHandle::try_reserve`] and
320    /// submits through [`WorkerPermit::submit`].
321    pub fn try_submit(&self, job: WorkerJob) -> Result<(), SubmitError> {
322        self.try_submit_labeled(job, DropStage::WorkerQueue)
323    }
324
325    /// Same as [`Self::try_submit`], but the failure counters carry a
326    /// caller-supplied stage label so response-slot exhaustion mid-stream
327    /// can be distinguished from admission-side worker-queue exhaustion.
328    /// See [`ResponsePermit::submit`].
329    fn try_submit_labeled(&self, job: WorkerJob, stage: DropStage) -> Result<(), SubmitError> {
330        let capacity_permit = self.try_capacity_labeled(stage)?;
331        let Some(sender) = self.sender_for(&job.session.id).cloned() else {
332            return Err(SubmitError::Closed);
333        };
334        job.session.worker_job_started();
335        self.pending.fetch_add(1, Ordering::AcqRel);
336        let span = worker_span(&job);
337        match sender.try_send(Envelope {
338            job,
339            completion: None,
340            span,
341            _capacity_permit: capacity_permit,
342        }) {
343            Ok(()) => Ok(()),
344            Err(mpsc::error::TrySendError::Full(envelope)) => {
345                let session_id = envelope.job.session.id.clone();
346                envelope.job.session.worker_job_finished();
347                self.worker_job_finished();
348                self.metrics
349                    .counter(stage.full_counter_key(), "Worker jobs dropped")
350                    .inc();
351                tracing::warn!(
352                    session = %session_id,
353                    stage = %stage.label(),
354                    "worker queue full; request must fail closed"
355                );
356                Err(SubmitError::Full)
357            }
358            Err(mpsc::error::TrySendError::Closed(envelope)) => {
359                let session_id = envelope.job.session.id.clone();
360                envelope.job.session.worker_job_finished();
361                self.worker_job_finished();
362                self.metrics
363                    .counter(
364                        "av_events_dropped_total{stage=\"worker_closed\"}",
365                        "Worker jobs dropped",
366                    )
367                    .inc();
368                tracing::warn!(
369                    session = %session_id,
370                    stage = %stage.label(),
371                    "worker queue closed; request must fail closed"
372                );
373                Err(SubmitError::Closed)
374            }
375        }
376    }
377
378    /// Submit and wait for completion. Intended for lifecycle operations and
379    /// deterministic integration tests, never streaming hot paths.
380    pub async fn submit_and_wait(&self, job: WorkerJob) -> Result<(), String> {
381        let sender_for_session = self
382            .sender_for(&job.session.id)
383            .cloned()
384            .ok_or_else(|| "worker queue is closed".to_owned())?;
385        let (sender, receiver) = oneshot::channel();
386        let capacity_permit = Arc::clone(&self.capacity)
387            .acquire_owned()
388            .await
389            .map_err(|_| "worker capacity semaphore is closed".to_owned())?;
390        job.session.worker_job_started();
391        self.pending.fetch_add(1, Ordering::AcqRel);
392        let span = worker_span(&job);
393        sender_for_session
394            .send(Envelope {
395                job,
396                completion: Some(sender),
397                span,
398                _capacity_permit: capacity_permit,
399            })
400            .await
401            .map_err(|error| {
402                error.0.job.session.worker_job_finished();
403                self.worker_job_finished();
404                "worker queue is closed".to_owned()
405            })?;
406        receiver
407            .await
408            .map_err(|_| "worker completion channel closed".to_owned())?
409    }
410
411    /// Wait until every accepted job has completed processing.
412    pub async fn wait_idle(&self) {
413        loop {
414            // Same pinned Notified must span enable() and .await; a fresh
415            // notified() after enable is dropped would miss a notify_waiters()
416            // firing in the interval.
417            let notified = self.drained.notified();
418            let mut notified = std::pin::pin!(notified);
419            notified.as_mut().enable();
420            if self.pending.load(Ordering::Acquire) == 0 {
421                return;
422            }
423            notified.await;
424        }
425    }
426
427    fn worker_job_finished(&self) {
428        if self.pending.fetch_sub(1, Ordering::AcqRel) == 1 {
429            self.drained.notify_waiters();
430        }
431    }
432
433    fn try_capacity(&self) -> Result<tokio::sync::OwnedSemaphorePermit, SubmitError> {
434        self.try_capacity_labeled(DropStage::WorkerQueue)
435    }
436
437    /// Same as [`Self::try_capacity`] but the failure counter carries a
438    /// caller-supplied stage label. `WorkerQueue` accounts to the main
439    /// worker admission semaphore; `ResponseSlot` accounts to the mid-
440    /// stream response-capture submission path — so operators can tell
441    /// which class of exhaustion is producing drops.
442    fn try_capacity_labeled(
443        &self,
444        stage: DropStage,
445    ) -> Result<tokio::sync::OwnedSemaphorePermit, SubmitError> {
446        Arc::clone(&self.capacity).try_acquire_owned().map_err(|_| {
447            self.metrics
448                .counter(stage.full_counter_key(), "Worker jobs dropped")
449                .inc();
450            SubmitError::Full
451        })
452    }
453
454    fn sender_for(&self, session_id: &str) -> Option<&mpsc::Sender<Envelope>> {
455        let partitions = u32::try_from(self.senders.len()).ok()?;
456        let shard = av_bridge::bus::partition_for(session_id, partitions) as usize;
457        self.senders.get(shard)
458    }
459}
460
461/// Start the session-sharded worker pool (16 shards, each behind a
462/// bounded channel). Ordering is per session via hash sharding.
463pub fn spawn_worker(
464    capacity: usize,
465    bridge: Arc<dyn EventBus>,
466    embedder: Arc<dyn Embedder>,
467    metrics: Arc<Registry>,
468) -> WorkerHandle {
469    spawn_worker_with_sink(capacity, bridge, embedder, Arc::new(NoopVectorSink), metrics)
470}
471
472/// Start the worker pool with an explicit off-path vector sink.
473pub fn spawn_worker_with_sink(
474    capacity: usize,
475    bridge: Arc<dyn EventBus>,
476    embedder: Arc<dyn Embedder>,
477    vector_sink: Arc<dyn VectorSink>,
478    metrics: Arc<Registry>,
479) -> WorkerHandle {
480    spawn_worker_with_spool(capacity, bridge, embedder, vector_sink, None, metrics)
481}
482
483/// Start the worker pool with vector persistence and optional ATIF journal.
484pub fn spawn_worker_with_spool(
485    capacity: usize,
486    bridge: Arc<dyn EventBus>,
487    embedder: Arc<dyn Embedder>,
488    vector_sink: Arc<dyn VectorSink>,
489    spool_dir: Option<std::path::PathBuf>,
490    metrics: Arc<Registry>,
491) -> WorkerHandle {
492    spawn_worker_with_spool_authenticated(
493        capacity,
494        bridge,
495        embedder,
496        vector_sink,
497        spool_dir,
498        [0; 32],
499        metrics,
500    )
501}
502
503/// Start the worker pool with authenticated active-workflow journals.
504pub fn spawn_worker_with_spool_authenticated(
505    capacity: usize,
506    bridge: Arc<dyn EventBus>,
507    embedder: Arc<dyn Embedder>,
508    vector_sink: Arc<dyn VectorSink>,
509    spool_dir: Option<std::path::PathBuf>,
510    journal_key: [u8; 32],
511    metrics: Arc<Registry>,
512) -> WorkerHandle {
513    // Sharding is decoupled from `capacity`: routing is by
514    // `partition_for(session_id, MAX_SHARDS)`, so if we sized
515    // `senders.len() < MAX_SHARDS`, requests whose id hashes to a
516    // partition >= senders.len() would silently see
517    // `SubmitError::Closed` (via `sender_for` returning `None`).
518    // Under `capacity = 1`, 15/16 partitions were unroutable. Always
519    // spawn all MAX_SHARDS shards; the global semaphore continues to
520    // enforce the caller's admission cap so total in-flight work is
521    // still bounded by `capacity`.
522    const MAX_SHARDS: usize = 16;
523    let capacity = capacity.max(1);
524    let shard_count = MAX_SHARDS;
525    // Per-shard channel size: enough that a single shard can absorb
526    // the full admission burst if every request happens to hash to it,
527    // capped at `capacity` so we never over-allocate the caller's
528    // memory budget across shards.
529    let per_shard_capacity = capacity.min(capacity.div_ceil(shard_count).max(1) * shard_count);
530    let per_shard_capacity = per_shard_capacity.max(1);
531    let global_capacity = Arc::new(tokio::sync::Semaphore::new(capacity));
532    // Response-slot capacity mirrors worker capacity by default. Keeping
533    // them equal preserves the historical behaviour (both counters were
534    // silently drawn from the same pool) while distinct semaphores let
535    // operators observe which class exhausts first via
536    // `av_events_dropped_total{stage="response_slot"}` vs
537    // `stage="worker_queue"`. A follow-up can expose an independent
538    // `response_capacity` config field once operators have telemetry.
539    let response_capacity = Arc::new(tokio::sync::Semaphore::new(capacity));
540    let pending = Arc::new(std::sync::atomic::AtomicU64::new(0));
541    let drained = Arc::new(tokio::sync::Notify::new());
542    let mut senders = Vec::with_capacity(shard_count);
543    for _ in 0..shard_count {
544        let (sender, receiver) = mpsc::channel::<Envelope>(per_shard_capacity);
545        senders.push(sender);
546        spawn_worker_shard(
547            receiver,
548            Arc::clone(&bridge),
549            Arc::clone(&embedder),
550            Arc::clone(&vector_sink),
551            spool_dir.clone(),
552            journal_key,
553            Arc::clone(&metrics),
554            Arc::clone(&pending),
555            Arc::clone(&drained),
556        );
557    }
558    WorkerHandle {
559        senders: Arc::new(senders),
560        metrics,
561        pending,
562        drained,
563        capacity: global_capacity,
564        response_capacity,
565    }
566}
567
568#[allow(clippy::too_many_arguments)]
569fn spawn_worker_shard(
570    mut receiver: mpsc::Receiver<Envelope>,
571    bridge: Arc<dyn EventBus>,
572    embedder: Arc<dyn Embedder>,
573    vector_sink: Arc<dyn VectorSink>,
574    spool_dir: Option<std::path::PathBuf>,
575    journal_key: [u8; 32],
576    worker_metrics: Arc<Registry>,
577    worker_pending: Arc<std::sync::atomic::AtomicU64>,
578    worker_drained: Arc<tokio::sync::Notify>,
579) {
580    use futures::future::FutureExt as _;
581    tokio::spawn(async move {
582        loop {
583            let Some(envelope) = receiver.recv().await else {
584                break;
585            };
586            // Supervise the per-envelope routing code (Arc::clones,
587            // span instrumentation, capacity_permit drop, drained
588            // notify, worker_pending accounting). The INNER
589            // `tokio::spawn(process_job).await` already catches
590            // process_job panics via JoinError → av_worker_panics_total,
591            // but a panic in the OUTER routing (allocator failure
592            // inside a tracing::warn Display, `worker_pending
593            // .fetch_sub` accounting bug) would kill the whole shard
594            // driver task, and every future envelope routed to this
595            // shard would pile up forever — jamming 1/MAX_SHARDS of
596            // the session id space until process restart.
597            let bridge_ref = Arc::clone(&bridge);
598            let embedder_ref = Arc::clone(&embedder);
599            let vector_sink_ref = Arc::clone(&vector_sink);
600            let spool_dir_ref = spool_dir.clone();
601            let worker_metrics_ref = Arc::clone(&worker_metrics);
602            let worker_pending_ref = Arc::clone(&worker_pending);
603            let worker_drained_ref = Arc::clone(&worker_drained);
604            let outcome = std::panic::AssertUnwindSafe(process_envelope(
605                envelope,
606                bridge_ref,
607                embedder_ref,
608                vector_sink_ref,
609                spool_dir_ref,
610                journal_key,
611                worker_metrics_ref,
612                worker_pending_ref,
613                worker_drained_ref,
614            ))
615            .catch_unwind()
616            .await;
617            if let Err(panic) = outcome {
618                let msg = panic
619                    .downcast_ref::<&'static str>()
620                    .copied()
621                    .or_else(|| panic.downcast_ref::<String>().map(String::as_str))
622                    .unwrap_or("panic payload was not a string");
623                worker_metrics
624                    .counter(
625                        "av_worker_shard_panics_total",
626                        "Worker shard driver panicked outside a job; supervised via catch_unwind",
627                    )
628                    .inc();
629                tracing::error!(
630                    panic = %msg,
631                    "worker shard driver panicked during envelope routing; continuing"
632                );
633            }
634        }
635    });
636}
637
638/// One envelope's worth of shard-driver work, factored out so
639/// [`spawn_worker_shard`] can wrap the whole body in `catch_unwind`.
640#[allow(clippy::too_many_arguments)]
641async fn process_envelope(
642    envelope: Envelope,
643    bridge: Arc<dyn EventBus>,
644    embedder: Arc<dyn Embedder>,
645    vector_sink: Arc<dyn VectorSink>,
646    spool_dir: Option<std::path::PathBuf>,
647    journal_key: [u8; 32],
648    worker_metrics: Arc<Registry>,
649    worker_pending: Arc<std::sync::atomic::AtomicU64>,
650    worker_drained: Arc<tokio::sync::Notify>,
651) {
652    // Round-12 F3: `worker_pending` was previously decremented at the
653    // bottom of this function. Any panic between here and that line
654    // (tokio::spawn(...).await JoinError construction under runtime
655    // shutdown races, or a panic inside tracing::warn's Display
656    // formatter when the error string exceeds the allocator's
657    // fragmentation threshold) would leak the pending count — and
658    // `WorkerHandle::wait_idle()` at shutdown would spin on
659    // `notify.notified().await` forever, tripping the 30 s drain
660    // budget and terminating without capturing evidence. Encode the
661    // decrement in an RAII guard so unwinding is a valid release
662    // point.
663    let _pending_guard = PendingGuard::new(Arc::clone(&worker_pending), Arc::clone(&worker_drained));
664    let Envelope {
665        job,
666        completion,
667        span,
668        _capacity_permit: capacity_permit,
669    } = envelope;
670    let session = Arc::clone(&job.session);
671    // Round-33 F2: guard the session-level pending-jobs decrement in
672    // the same RAII shape as `PendingGuard` (round-12 F3). The bare
673    // `session.worker_job_finished()` after `tokio::spawn(...).await`
674    // used to leak the session-level counter on any panic or drop
675    // between here and line ~712. A stuck `session.pending_jobs` means
676    // `close_session_locked -> wait_for_worker_jobs().await` blocks
677    // forever on `jobs_drained.notified()`, holding the session's
678    // lifecycle lock and starving every subsequent close / promote /
679    // recovery-adopt on that id. Class round-12 F3 closed, one call
680    // frame up.
681    let _session_pending_guard = SessionPendingGuard {
682        session: Arc::clone(&session),
683    };
684    let result = tokio::spawn(
685        {
686            let job_metrics = Arc::clone(&worker_metrics);
687            async move {
688                process_job(
689                    job,
690                    bridge,
691                    embedder,
692                    vector_sink,
693                    spool_dir,
694                    journal_key,
695                    job_metrics,
696                )
697                .await
698            }
699        }
700        .instrument(span),
701    )
702    .await;
703    let outcome = match result {
704        Ok(result) => result,
705        Err(error) => {
706            worker_metrics
707                .counter(
708                    "av_worker_panics_total",
709                    "Worker job panics isolated by supervisor",
710                )
711                .inc();
712            Err(format!("worker job panicked: {error}"))
713        }
714    };
715    if let Err(error) = &outcome {
716        // The session flips fail-closed silently otherwise: name the
717        // session and the root cause (e.g. EACCES on the journal
718        // file) so operators can trace "capture is incomplete"
719        // refusals back to this event.
720        tracing::warn!(session = %session.id, %error, "capture job failed; session is fail-closed");
721        session.mark_capture_failed();
722        worker_metrics
723            .counter("av_worker_errors_total", "Worker jobs that failed")
724            .inc();
725    }
726    // Round-33 F2: `_session_pending_guard`'s Drop calls
727    // `worker_job_finished()` — replaces the bare call previously
728    // here so a panic between `spawn.await` and this point cannot
729    // leak the session pending counter.
730    drop(capacity_permit);
731    // Guard runs Drop here (or on the panic path). The completion
732    // channel below is fine to send after — the receiver only needs
733    // the outcome, not the pending accounting.
734    if let Some(completion) = completion {
735        let _ = completion.send(outcome);
736    }
737}
738
739/// RAII decrement for `worker_pending`. Runs on every path — normal
740/// return, error return, or panic — so `wait_idle` never sees a
741/// permanently non-zero pending count after a panic in
742/// [`process_envelope`].
743struct PendingGuard {
744    pending: Arc<std::sync::atomic::AtomicU64>,
745    drained: Arc<tokio::sync::Notify>,
746}
747
748impl PendingGuard {
749    fn new(pending: Arc<std::sync::atomic::AtomicU64>, drained: Arc<tokio::sync::Notify>) -> Self {
750        Self { pending, drained }
751    }
752}
753
754impl Drop for PendingGuard {
755    fn drop(&mut self) {
756        if self.pending.fetch_sub(1, Ordering::AcqRel) == 1 {
757            self.drained.notify_waiters();
758        }
759    }
760}
761
762/// Round-33 F2: RAII pair to [`PendingGuard`] for the session-level
763/// pending-jobs counter. `Session::worker_job_finished()` was
764/// previously called as a bare method after `tokio::spawn(...).await`,
765/// but a panic in the routing-side epilogue (Display-side allocator
766/// failure, catch_unwind of the wrapper future, runtime shutdown
767/// mid-envelope) between `.await` and that call left
768/// `session.pending_jobs` stuck. `close_session_locked` then blocked
769/// forever on `wait_for_worker_jobs().await`, holding the session's
770/// lifecycle lock — the exact class round-12 F3 closed on the
771/// worker-level counter.
772struct SessionPendingGuard {
773    session: Arc<crate::session::Session>,
774}
775
776impl Drop for SessionPendingGuard {
777    fn drop(&mut self) {
778        self.session.worker_job_finished();
779    }
780}
781
782fn worker_span(job: &WorkerJob) -> tracing::Span {
783    tracing::info_span!(
784        parent: &tracing::Span::current(),
785        "agentvisor.worker",
786        session.id = %job.session.id,
787        event.class = ?job.class,
788    )
789}
790
791#[allow(clippy::too_many_arguments)]
792async fn process_job(
793    mut job: WorkerJob,
794    bridge: Arc<dyn EventBus>,
795    embedder: Arc<dyn Embedder>,
796    vector_sink: Arc<dyn VectorSink>,
797    spool_dir: Option<std::path::PathBuf>,
798    journal_key: [u8; 32],
799    metrics: Arc<Registry>,
800) -> Result<(), String> {
801    // A queued envelope for a session that has already been poisoned must not
802    // consume a sequence number or write to the journal — otherwise the seq
803    // its OcsfEvent carries (from `next_seq`, advanced by the failed prior
804    // envelope) would not match the entry's position on disk, breaking
805    // recovery's `event.metadata.sequence != index` check.
806    if job.session.capture_failed() {
807        return Ok(());
808    }
809    let response_marker = job.response_marker.clone();
810    let is_llm_agent_response = job
811        .atif
812        .as_ref()
813        .is_some_and(|capture| capture.source == av_atif::Source::Agent && capture.llm_call_count != Some(0));
814    let step_tokens = job
815        .metrics
816        .prompt_tokens
817        .unwrap_or(0)
818        .saturating_add(job.metrics.completion_tokens.unwrap_or(0));
819    let breaker = if job.analyze_loop {
820        let text = job.text.clone();
821        let embedding = tokio::task::spawn_blocking(move || embedder.try_embed(&text))
822            .await
823            .map_err(|error| error.to_string())??;
824        let nearest_similarity = vector_sink
825            .nearest_similarity(&job.session.id, &embedding)
826            .await?;
827        let verdict = job.session.loop_state.observe_embedding_with_similarity(
828            embedding.clone(),
829            step_tokens,
830            nearest_similarity,
831        );
832        vector_sink.record(&job.session.id, &embedding).await?;
833        Some(verdict)
834    } else {
835        None
836    };
837    // Accounting must key on the class the job was *submitted* with: a
838    // breaker trip replaces `class` with StopReason below, and deciding the
839    // prompt/completion buckets from the replaced class silently dropped a
840    // tripped chat admission's prompt tokens from the journal record and the
841    // receipt totals — undercounting exactly the runaway sessions the
842    // breaker exists to attest.
843    let submitted_class = job.class;
844    let (class, status, stop_reason, payload) = match breaker {
845        Some(BreakerVerdict::Tripped {
846            delta,
847            streak,
848            tokens_consumed,
849            action,
850        }) => {
851            // The trip is recorded in the trajectory, but operators watching
852            // logs/metrics must also see why this session's next request
853            // will be rejected/aborted/injected.
854            tracing::warn!(
855                session = %job.session.id,
856                delta,
857                streak,
858                tokens_consumed,
859                action = ?action,
860                "semantic loop circuit breaker tripped"
861            );
862            metrics
863                .counter("av_breaker_trips_total", "Semantic loop breaker trips")
864                .inc();
865            (
866                EventClass::StopReason,
867                StatusId::Failure,
868                Some(StopReason::LoopDetected),
869                serde_json::json!({
870                    "delta": delta,
871                    "streak": streak,
872                    "tokens_consumed": tokens_consumed,
873                    "action": action,
874                }),
875            )
876        }
877        _ => (job.class, job.status, job.stop_reason, job.payload),
878    };
879
880    let mut builder = OcsfEventBuilder::new(
881        class,
882        job.session.id.clone(),
883        job.identity.clone(),
884        job.session.next_seq(),
885    )
886    .status(status)
887    .payload(payload)
888    .metrics(job.metrics);
889    if let Some(reason) = stop_reason {
890        job.session.record_stop_reason(reason);
891        builder = match job.native_stop_reason {
892            Some(native) => builder.stop_reason_native(reason, native),
893            None => builder.stop_reason(reason),
894        };
895    }
896    let event = builder.build().map_err(|error| error.to_string())?;
897    let event_uid = event.metadata.uid.clone();
898    let value = serde_json::to_value(&event).map_err(|error| error.to_string())?;
899
900    let atif_step = if job.session.workflow == Workflow::Unsigned {
901        let capture = job
902            .atif
903            .take()
904            .ok_or_else(|| "unsigned worker job has no ATIF capture".to_owned())?;
905        let is_llm_step = capture.source == av_atif::Source::Agent && capture.llm_call_count != Some(0);
906        Some(av_atif::Step {
907            step_id: 0,
908            timestamp: Some(av_core::time::now_iso8601()),
909            source: capture.source,
910            message: capture.message,
911            reasoning_effort: None,
912            reasoning_content: capture.reasoning_content,
913            model_name: capture.model_name,
914            tool_calls: capture.tool_calls,
915            observation: capture.observation,
916            metrics: is_llm_step.then(|| av_atif::Metrics {
917                prompt_tokens: Some(job.metrics.prompt_tokens.unwrap_or(0)),
918                completion_tokens: Some(job.metrics.completion_tokens.unwrap_or(0)),
919                cached_tokens: Some(job.metrics.cached_tokens.unwrap_or(0)),
920                cost_usd: Some(job.cost_usd_micros as f64 / av_core::units::USD_MICROS_PER_DOLLAR as f64),
921                logprobs: None,
922                completion_token_ids: None,
923                prompt_token_ids: None,
924                extra: None,
925            }),
926            is_copied_context: None,
927            llm_call_count: capture.llm_call_count,
928            extra: None,
929        })
930    } else {
931        None
932    };
933    let is_tool_call = submitted_class == EventClass::ToolCall;
934    let is_response_accounting = is_llm_agent_response && submitted_class != EventClass::Compression;
935    let record = ActiveJournalRecord {
936        event: value.clone(),
937        identity: job.identity.clone(),
938        atif_step: atif_step.clone(),
939        tool_calls: u64::from(is_tool_call),
940        tool_allowed: u64::from(is_tool_call && status == StatusId::Success),
941        tool_blocked: u64::from(is_tool_call && status == StatusId::Failure),
942        prompt_tokens: if submitted_class == EventClass::Compression {
943            job.metrics.prompt_tokens.unwrap_or(0)
944        } else {
945            0
946        },
947        completion_tokens: if is_response_accounting {
948            job.metrics.completion_tokens.unwrap_or(0)
949        } else {
950            0
951        },
952        cached_tokens: if is_response_accounting {
953            job.metrics.cached_tokens.unwrap_or(0)
954        } else {
955            0
956        },
957        cost_usd_micros: if is_response_accounting {
958            job.cost_usd_micros
959        } else {
960            0
961        },
962        stop_reason_id: event.stop_reason_id,
963        response_attempt: job.response_attempt.clone(),
964    };
965    if let Some(directory) = spool_dir.as_deref() {
966        append_active_event_journal(directory, &job.session, &job.identity, &record, &journal_key).await?;
967    }
968
969    match job.session.workflow {
970        Workflow::Signed => job
971            .session
972            .chain
973            .lock()
974            .append(&value)
975            .map_err(|error| error.to_string())?,
976        Workflow::Unsigned => {
977            job.session
978                .atif
979                .lock()
980                .push_step(atif_step.ok_or_else(|| "unsigned ATIF step disappeared".to_owned())?)
981                .map_err(|error| error.to_string())?;
982        }
983    }
984
985    let topic = class.topic().to_owned();
986    let key = job.identity.instance_uid.clone();
987    let publish_topic = topic.clone();
988    let publish_key = key.clone();
989    let publish_uid = event_uid.clone();
990    let ack = tokio::task::spawn_blocking(move || {
991        bridge.publish_idempotent(&publish_topic, &publish_key, &value, &publish_uid)
992    })
993    .await
994    .map_err(|error| error.to_string())?
995    .map_err(|error| error.to_string())?;
996    if let Some(directory) = spool_dir.as_deref() {
997        persist_broker_ack(directory, &job.session.id, &event_uid, &ack, &journal_key).await?;
998    }
999    if is_tool_call {
1000        checked_atomic_add(&job.session.totals.tool_calls, 1, "tool calls")?;
1001        match status {
1002            StatusId::Success => {
1003                checked_atomic_add(&job.session.totals.tool_allowed, 1, "allowed tools")?;
1004            }
1005            StatusId::Failure => {
1006                checked_atomic_add(&job.session.totals.tool_blocked, 1, "blocked tools")?;
1007            }
1008            StatusId::Unknown => {}
1009            _ => {}
1010        }
1011    }
1012    if submitted_class == EventClass::Compression {
1013        checked_atomic_add(
1014            &job.session.totals.prompt_tokens,
1015            job.metrics.prompt_tokens.unwrap_or(0),
1016            "prompt tokens",
1017        )?;
1018    } else if is_llm_agent_response {
1019        checked_atomic_add(
1020            &job.session.totals.completion_tokens,
1021            job.metrics.completion_tokens.unwrap_or(0),
1022            "completion tokens",
1023        )?;
1024        checked_atomic_add(
1025            &job.session.totals.cached_tokens,
1026            job.metrics.cached_tokens.unwrap_or(0),
1027            "cached tokens",
1028        )?;
1029        checked_atomic_add(&job.session.totals.cost_usd_micros, job.cost_usd_micros, "cost")?;
1030    }
1031    if let (Some(directory), Some(attempt_id)) = (spool_dir.as_deref(), response_marker.as_deref()) {
1032        clear_response_marker(directory, &journal_key, &job.session.id, attempt_id).await?;
1033    }
1034    Ok(())
1035}
1036
1037fn checked_atomic_add(counter: &std::sync::atomic::AtomicU64, value: u64, field: &str) -> Result<(), String> {
1038    counter
1039        .fetch_update(Ordering::AcqRel, Ordering::Acquire, |current| {
1040            current
1041                .checked_add(value)
1042                .filter(|next| *next <= av_core::error::JCS_SAFE_MAX)
1043        })
1044        .map(|_| ())
1045        .map_err(|_| format!("{field} accounting exceeds JCS-safe bounds"))
1046}
1047
1048pub(crate) async fn create_response_marker(
1049    spool_dir: &std::path::Path,
1050    journal_key: &[u8; 32],
1051    session_id: &str,
1052    request_digest: String,
1053) -> Result<String, String> {
1054    let attempt_id = av_core::new_event_uid();
1055    let marker = InFlightResponse {
1056        session_id: session_id.to_owned(),
1057        attempt_id: attempt_id.clone(),
1058        request_digest,
1059    };
1060    let sealed = crate::journal::seal(journal_key, "in-flight-response", 0, &marker)?;
1061    let path = response_marker_path(spool_dir, session_id, &attempt_id);
1062    tokio::task::spawn_blocking(move || write_atomic_control(&path, &sealed))
1063        .await
1064        .map_err(|error| error.to_string())??;
1065    Ok(attempt_id)
1066}
1067
1068#[cfg(test)]
1069pub(crate) async fn ensure_no_inflight_responses(
1070    spool_dir: &std::path::Path,
1071    journal_key: &[u8; 32],
1072) -> Result<(), String> {
1073    let sessions = inflight_response_sessions(spool_dir, journal_key).await?;
1074    if let Some(session_id) = sessions.into_iter().next() {
1075        Err(format!(
1076            "provider response for session {session_id} was not durably captured"
1077        ))
1078    } else {
1079        Ok(())
1080    }
1081}
1082
1083pub(crate) async fn inflight_response_sessions(
1084    spool_dir: &std::path::Path,
1085    journal_key: &[u8; 32],
1086) -> Result<std::collections::HashSet<String>, String> {
1087    let spool_dir = spool_dir.to_path_buf();
1088    let directory = spool_dir.join(crate::spool::INFLIGHT_RESPONSES);
1089    let journal_key = *journal_key;
1090    tokio::task::spawn_blocking(move || {
1091        let entries = match std::fs::read_dir(&directory) {
1092            Ok(entries) => entries,
1093            Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1094                return Ok(std::collections::HashSet::new());
1095            }
1096            Err(error) => return Err(error.to_string()),
1097        };
1098        let mut sessions = std::collections::HashSet::new();
1099        for entry in entries {
1100            let path = entry.map_err(|error| error.to_string())?.path();
1101            if path.extension().and_then(std::ffi::OsStr::to_str) != Some("json") {
1102                continue;
1103            }
1104            let marker: InFlightResponse = crate::journal::open(
1105                &journal_key,
1106                "in-flight-response",
1107                0,
1108                // Round-18: cap sealed marker read at MAX_CONTROL_BYTES.
1109                &av_core::fsutil::read_capped(&path, av_core::fsutil::MAX_CONTROL_BYTES)
1110                    .map_err(|error| error.to_string())?,
1111            )?;
1112            if path != response_marker_path(&spool_dir, &marker.session_id, &marker.attempt_id) {
1113                return Err("in-flight response marker path does not match its payload".to_owned());
1114            }
1115            sessions.insert(marker.session_id);
1116        }
1117        Ok(sessions)
1118    })
1119    .await
1120    .map_err(|error| error.to_string())?
1121}
1122
1123async fn clear_response_marker(
1124    spool_dir: &std::path::Path,
1125    journal_key: &[u8; 32],
1126    session_id: &str,
1127    attempt_id: &str,
1128) -> Result<(), String> {
1129    let path = response_marker_path(spool_dir, session_id, attempt_id);
1130    let journal_key = *journal_key;
1131    let session_id = session_id.to_owned();
1132    let attempt_id = attempt_id.to_owned();
1133    tokio::task::spawn_blocking(move || {
1134        let marker: InFlightResponse = crate::journal::open(
1135            &journal_key,
1136            "in-flight-response",
1137            0,
1138            // Round-18: cap sealed marker read at MAX_CONTROL_BYTES.
1139            &av_core::fsutil::read_capped(&path, av_core::fsutil::MAX_CONTROL_BYTES)
1140                .map_err(|error| error.to_string())?,
1141        )?;
1142        if marker.session_id != session_id || marker.attempt_id != attempt_id {
1143            return Err("in-flight response marker does not match completed job".to_owned());
1144        }
1145        std::fs::remove_file(&path).map_err(|error| error.to_string())?;
1146        let parent = path
1147            .parent()
1148            .ok_or_else(|| "in-flight response marker has no parent".to_owned())?;
1149        av_core::fsutil::sync_directory(parent).map_err(|error| error.to_string())
1150    })
1151    .await
1152    .map_err(|error| error.to_string())?
1153}
1154
1155fn response_marker_path(
1156    spool_dir: &std::path::Path,
1157    session_id: &str,
1158    attempt_id: &str,
1159) -> std::path::PathBuf {
1160    let digest = av_core::digest::sha256_hex(format!("{session_id}:{attempt_id}").as_bytes());
1161    spool_dir
1162        .join(crate::spool::INFLIGHT_RESPONSES)
1163        .join(format!("{}.json", &digest[..32]))
1164}
1165
1166fn write_atomic_control(path: &std::path::Path, bytes: &[u8]) -> Result<(), String> {
1167    av_core::fsutil::write_atomic(path, bytes).map_err(|error| error.to_string())
1168}
1169
1170async fn append_active_event_journal(
1171    directory: &std::path::Path,
1172    session: &Session,
1173    _identity: &av_events::AgentIdentity,
1174    record: &ActiveJournalRecord,
1175    journal_key: &[u8; 32],
1176) -> Result<(), String> {
1177    let index = session.journal_index();
1178    let domain = format!("{}:active", session.id);
1179    let line = crate::journal::seal(journal_key, &domain, index, record)?;
1180    append_journal(directory, session, line, *journal_key).await?;
1181    session.commit_journal_index(index)
1182}
1183
1184async fn append_journal(
1185    directory: &std::path::Path,
1186    session: &Session,
1187    line: Vec<u8>,
1188    journal_key: [u8; 32],
1189) -> Result<(), String> {
1190    let directory = directory.to_path_buf();
1191    let session_id = session.id.clone();
1192    let identity = session.identity.clone();
1193    let workflow = session.workflow.as_str();
1194    tokio::task::spawn_blocking(move || -> Result<(), String> {
1195        use std::io::Write as _;
1196        std::fs::create_dir_all(&directory).map_err(|error| error.to_string())?;
1197        let digest = av_core::digest::sha256_hex(session_id.as_bytes());
1198        let stem = digest.get(..32).unwrap_or(&digest);
1199        let metadata_path = directory.join(format!("{stem}.session.json"));
1200        let metadata_payload = serde_json::json!({
1201            "journal_version": 2,
1202            "session_id": session_id,
1203            "identity": identity,
1204            "workflow": workflow,
1205        });
1206        if !metadata_path.exists() {
1207            let metadata = crate::journal::seal(&journal_key, "metadata", 0, &metadata_payload)?;
1208            // Use the centralized atomic writer: it uses an RAII guard
1209            // so any intermediate failure (write_all / sync_all /
1210            // rename) cannot leak a zero-byte `.tmp` orphan. A
1211            // repeatedly-failing journal writer used to be able to
1212            // exhaust the ext4 inode table long before disk-full.
1213            //
1214            // Round-37 F2: basename the paths in error strings. These
1215            // errors bubble to `tracing::warn!(session = %session.id,
1216            // %error, "capture job failed; session is fail-closed")`
1217            // at process_envelope; that warn exports through
1218            // tracing_opentelemetry -> OTLP -> SIEM, so a single
1219            // disk-full incident used to emit one span per pending
1220            // event with the full absolute journal tree. Basename
1221            // preserves enough context for triage (the stem encodes
1222            // the session id) without leaking the deployment topology.
1223            av_core::fsutil::write_atomic(&metadata_path, &metadata).map_err(|error| {
1224                format!(
1225                    "write journal metadata {}: {error}",
1226                    av_core::fsutil::basename(&metadata_path)
1227                )
1228            })?;
1229        } else {
1230            let stored: serde_json::Value = crate::journal::open(
1231                &journal_key,
1232                "metadata",
1233                0,
1234                // Round-18: cap sealed journal-metadata read at
1235                // MAX_CONTROL_BYTES.
1236                &av_core::fsutil::read_capped(&metadata_path, av_core::fsutil::MAX_CONTROL_BYTES)
1237                    .map_err(|error| error.to_string())?,
1238            )?;
1239            if stored != metadata_payload {
1240                return Err("journal metadata does not match session workflow and identity".to_owned());
1241            }
1242        }
1243        let journal_path = directory.join(format!("{stem}.events.ndjson"));
1244        // Track whether the journal is being created by *this* append so
1245        // we can fsync the containing directory once the file exists —
1246        // the metadata fsync above only durably named the metadata
1247        // file, not this new journal file. Without the dirent fsync,
1248        // a power loss on POSIX-conformant filesystems (xfs, btrfs)
1249        // can lose the entry entirely — the file appears not to exist
1250        // on restart, `recover_signed_journals` treats the journal as
1251        // empty, deletes the metadata, and every already-acked event
1252        // becomes orphaned on the broker.
1253        let journal_created = !journal_path.exists();
1254        let mut journal = std::fs::OpenOptions::new()
1255            .create(true)
1256            .append(true)
1257            .open(&journal_path)
1258            .map_err(|error| {
1259                format!(
1260                    "open event journal {}: {error}",
1261                    av_core::fsutil::basename(&journal_path)
1262                )
1263            })?;
1264        journal.write_all(&line).map_err(|error| error.to_string())?;
1265        journal.write_all(b"\n").map_err(|error| error.to_string())?;
1266        journal.sync_data().map_err(|error| error.to_string())?;
1267        if journal_created {
1268            std::fs::File::open(&directory)
1269                .and_then(|dir| dir.sync_all())
1270                .map_err(|error| {
1271                    // Round-37 F2: drop the full path entirely; the
1272                    // sibling `session = %session.id` field on the
1273                    // downstream warn already scopes this to a
1274                    // specific session, and the containing directory
1275                    // is `atif_spool_dir` — same across the whole
1276                    // deployment.
1277                    format!("fsync journal directory: {error}")
1278                })?;
1279        }
1280        Ok(())
1281    })
1282    .await
1283    .map_err(|error| error.to_string())?
1284}
1285
1286pub(crate) async fn persist_broker_ack(
1287    directory: &std::path::Path,
1288    session_id: &str,
1289    event_uid: &str,
1290    ack: &PublishAck,
1291    journal_key: &[u8; 32],
1292) -> Result<(), String> {
1293    let path = broker_ack_path(directory, session_id, event_uid);
1294    let record = BrokerAckRecord {
1295        session_id: session_id.to_owned(),
1296        event_uid: event_uid.to_owned(),
1297        ack: ack.clone(),
1298    };
1299    let sealed = crate::journal::seal(journal_key, "broker-ack", 0, &record)?;
1300    tokio::task::spawn_blocking(move || write_atomic_control(&path, &sealed))
1301        .await
1302        .map_err(|error| error.to_string())?
1303}
1304
1305pub(crate) async fn read_broker_ack(
1306    directory: &std::path::Path,
1307    session_id: &str,
1308    event_uid: &str,
1309    journal_key: &[u8; 32],
1310) -> Result<Option<PublishAck>, String> {
1311    let path = broker_ack_path(directory, session_id, event_uid);
1312    let bytes = match tokio::fs::read(&path).await {
1313        Ok(bytes) => bytes,
1314        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
1315        Err(error) => return Err(error.to_string()),
1316    };
1317    let record: BrokerAckRecord = crate::journal::open(journal_key, "broker-ack", 0, &bytes)?;
1318    if record.session_id != session_id || record.event_uid != event_uid {
1319        return Err("broker acknowledgment does not match event".to_owned());
1320    }
1321    Ok(Some(record.ack))
1322}
1323
1324fn broker_ack_path(directory: &std::path::Path, session_id: &str, event_uid: &str) -> std::path::PathBuf {
1325    let session_digest = av_core::digest::sha256_hex(session_id.as_bytes());
1326    let event_digest = av_core::digest::sha256_hex(event_uid.as_bytes());
1327    directory
1328        .join("broker-acks")
1329        .join(&session_digest[..32])
1330        .join(format!("{}.json", &event_digest[..32]))
1331}
1332
1333#[cfg(test)]
1334mod tests {
1335    #![allow(
1336        clippy::expect_used,
1337        clippy::indexing_slicing,
1338        clippy::panic,
1339        clippy::unwrap_used
1340    )]
1341
1342    use super::*;
1343    use av_bridge::{BusError, PublishAck, StoredEvent};
1344    use av_events::AgentIdentity;
1345    use av_loopdetect::{BreakerConfig, HashEmbedder};
1346    use parking_lot::Mutex;
1347    use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering};
1348
1349    #[derive(Default)]
1350    struct RecordingBus {
1351        events: Mutex<Vec<(String, String, Value)>>,
1352    }
1353
1354    struct PanicOnceSink {
1355        panicked: AtomicBool,
1356    }
1357
1358    struct SlowEmbedder;
1359
1360    impl Embedder for SlowEmbedder {
1361        fn dim(&self) -> usize {
1362            1
1363        }
1364
1365        fn embed(&self, _text: &str) -> Vec<f32> {
1366            std::thread::sleep(std::time::Duration::from_millis(100));
1367            vec![1.0]
1368        }
1369    }
1370
1371    struct SlowBus;
1372
1373    impl EventBus for SlowBus {
1374        fn publish(&self, topic: &str, _key: &str, _value: &Value) -> Result<PublishAck, BusError> {
1375            std::thread::sleep(std::time::Duration::from_millis(100));
1376            Ok(PublishAck {
1377                topic: topic.to_owned(),
1378                partition: 0,
1379                offset: 0,
1380            })
1381        }
1382
1383        fn fetch(
1384            &self,
1385            _topic: &str,
1386            _partition: u32,
1387            _offset: u64,
1388            _max: usize,
1389        ) -> Result<Vec<StoredEvent>, BusError> {
1390            Ok(Vec::new())
1391        }
1392
1393        fn partitions(&self, _topic: &str) -> Result<u32, BusError> {
1394            Ok(1)
1395        }
1396
1397        fn topics(&self) -> Vec<String> {
1398            EventClass::all()
1399                .iter()
1400                .map(|class| class.topic().to_owned())
1401                .collect()
1402        }
1403    }
1404
1405    impl VectorSink for PanicOnceSink {
1406        fn record<'a>(
1407            &'a self,
1408            _session_id: &'a str,
1409            _vector: &'a [f32],
1410        ) -> av_loopdetect::VectorSinkFuture<'a> {
1411            Box::pin(async move {
1412                if !self.panicked.swap(true, AtomicOrdering::AcqRel) {
1413                    panic!("injected vector sink panic");
1414                }
1415                Ok(())
1416            })
1417        }
1418    }
1419
1420    impl EventBus for RecordingBus {
1421        fn publish(&self, topic: &str, key: &str, value: &Value) -> Result<PublishAck, BusError> {
1422            let mut events = self.events.lock();
1423            let offset = events.len() as u64;
1424            events.push((topic.to_owned(), key.to_owned(), value.clone()));
1425            Ok(PublishAck {
1426                topic: topic.to_owned(),
1427                partition: 0,
1428                offset,
1429            })
1430        }
1431
1432        fn publish_idempotent(
1433            &self,
1434            topic: &str,
1435            key: &str,
1436            value: &Value,
1437            event_uid: &str,
1438        ) -> Result<PublishAck, BusError> {
1439            if let Some((offset, _)) = self
1440                .events
1441                .lock()
1442                .iter()
1443                .enumerate()
1444                .find(|(_, (_, _, event))| event["metadata"]["uid"] == event_uid)
1445            {
1446                return Ok(PublishAck {
1447                    topic: topic.to_owned(),
1448                    partition: 0,
1449                    offset: offset as u64,
1450                });
1451            }
1452            self.publish(topic, key, value)
1453        }
1454
1455        fn fetch(
1456            &self,
1457            _topic: &str,
1458            _partition: u32,
1459            _offset: u64,
1460            _max: usize,
1461        ) -> Result<Vec<StoredEvent>, BusError> {
1462            Ok(Vec::new())
1463        }
1464
1465        fn partitions(&self, _topic: &str) -> Result<u32, BusError> {
1466            Ok(1)
1467        }
1468
1469        fn topics(&self) -> Vec<String> {
1470            EventClass::all()
1471                .iter()
1472                .map(|class| class.topic().to_owned())
1473                .collect()
1474        }
1475    }
1476
1477    fn session(workflow: Workflow) -> Arc<Session> {
1478        Arc::new(Session::new(
1479            "session-1".to_owned(),
1480            workflow,
1481            AgentIdentity {
1482                version: "1".to_owned(),
1483                charter: "test".into(),
1484                instance_uid: "instance-1".to_owned(),
1485                ttl_remaining_s: Some(600),
1486            },
1487            BreakerConfig {
1488                min_tokens: u64::MAX,
1489                ..BreakerConfig::default()
1490            },
1491        ))
1492    }
1493
1494    fn job(session: Arc<Session>) -> WorkerJob {
1495        WorkerJob {
1496            identity: session.current_identity(),
1497            session,
1498            class: EventClass::Compression,
1499            payload: serde_json::json!({"kind": "response"}),
1500            text: "a useful response".to_owned(),
1501            analyze_loop: true,
1502            status: StatusId::Success,
1503            stop_reason: None,
1504            native_stop_reason: None,
1505            metrics: EventMetrics {
1506                prompt_tokens: Some(100),
1507                completion_tokens: Some(20),
1508                cached_tokens: Some(40),
1509                pruned_tokens: Some(10),
1510                pruning_ratio_millis: Some(100),
1511            },
1512            cost_usd_micros: 250,
1513            atif: Some(AtifCapture {
1514                source: av_atif::Source::Agent,
1515                message: Value::String("a useful response".to_owned()),
1516                reasoning_content: None,
1517                model_name: None,
1518                tool_calls: None,
1519                observation: None,
1520                llm_call_count: Some(1),
1521            }),
1522            response_marker: None,
1523            response_attempt: None,
1524        }
1525    }
1526
1527    #[test]
1528    fn receipt_accounting_never_wraps_or_exceeds_jcs_bounds() {
1529        let counter = std::sync::atomic::AtomicU64::new(av_core::error::JCS_SAFE_MAX);
1530        assert!(checked_atomic_add(&counter, 1, "test").is_err());
1531        assert_eq!(counter.load(Ordering::Acquire), av_core::error::JCS_SAFE_MAX);
1532    }
1533
1534    #[tokio::test(flavor = "current_thread")]
1535    async fn synchronous_backends_do_not_block_tokio_worker_thread() {
1536        let worker = spawn_worker(
1537            4,
1538            Arc::new(SlowBus),
1539            Arc::new(SlowEmbedder),
1540            Arc::new(Registry::new()),
1541        );
1542        worker.try_submit(job(session(Workflow::Signed))).unwrap();
1543
1544        let started = std::time::Instant::now();
1545        tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1546        assert!(
1547            started.elapsed() < std::time::Duration::from_millis(75),
1548            "synchronous worker dependency blocked the Tokio reactor"
1549        );
1550        worker.wait_idle().await;
1551    }
1552
1553    #[tokio::test]
1554    async fn response_marker_is_authenticated_and_cleared_after_publish() {
1555        let directory = tempfile::tempdir().unwrap();
1556        let journal_key = [13; 32];
1557        let marker = create_response_marker(
1558            directory.path(),
1559            &journal_key,
1560            "session-1",
1561            "request-digest".to_owned(),
1562        )
1563        .await
1564        .unwrap();
1565        assert!(ensure_no_inflight_responses(directory.path(), &journal_key)
1566            .await
1567            .unwrap_err()
1568            .contains("was not durably captured"));
1569
1570        let worker = spawn_worker_with_spool_authenticated(
1571            4,
1572            Arc::new(RecordingBus::default()),
1573            Arc::new(HashEmbedder::default()),
1574            Arc::new(NoopVectorSink),
1575            Some(directory.path().to_path_buf()),
1576            journal_key,
1577            Arc::new(Registry::new()),
1578        );
1579        let session = session(Workflow::Signed);
1580        let mut response = job(session);
1581        response.response_marker = Some(marker);
1582        worker.submit_and_wait(response).await.unwrap();
1583        ensure_no_inflight_responses(directory.path(), &journal_key)
1584            .await
1585            .unwrap();
1586    }
1587
1588    #[tokio::test]
1589    async fn response_marker_rejects_payload_mutation() {
1590        let directory = tempfile::tempdir().unwrap();
1591        let journal_key = [17; 32];
1592        create_response_marker(
1593            directory.path(),
1594            &journal_key,
1595            "session-1",
1596            "request-digest".to_owned(),
1597        )
1598        .await
1599        .unwrap();
1600        let marker_path = std::fs::read_dir(directory.path().join(crate::spool::INFLIGHT_RESPONSES))
1601            .unwrap()
1602            .next()
1603            .unwrap()
1604            .unwrap()
1605            .path();
1606        let mut envelope: Value = serde_json::from_slice(&std::fs::read(&marker_path).unwrap()).unwrap();
1607        envelope["payload"]["request_digest"] = Value::String("forged".to_owned());
1608        std::fs::write(marker_path, serde_json::to_vec(&envelope).unwrap()).unwrap();
1609        assert!(ensure_no_inflight_responses(directory.path(), &journal_key)
1610            .await
1611            .unwrap_err()
1612            .contains("authentication failed"));
1613    }
1614
1615    #[tokio::test]
1616    async fn signed_job_publishes_and_extends_chain() {
1617        let bridge = Arc::new(RecordingBus::default());
1618        let metrics = Arc::new(Registry::new());
1619        let worker = spawn_worker(4, bridge.clone(), Arc::new(HashEmbedder::default()), metrics);
1620        let session = session(Workflow::Signed);
1621        worker.submit_and_wait(job(Arc::clone(&session))).await.unwrap();
1622
1623        assert_eq!(session.chain.lock().count(), 1);
1624        assert_eq!(session.totals.prompt_tokens.load(Ordering::Acquire), 100);
1625        let events = bridge.events.lock();
1626        assert_eq!(events.len(), 1);
1627        assert_eq!(events[0].0, "agent.compression");
1628        assert_eq!(events[0].1, "instance-1");
1629    }
1630
1631    #[tokio::test]
1632    async fn signed_job_journals_event_and_accounting() {
1633        let directory = tempfile::tempdir().unwrap();
1634        let journal_key = [9; 32];
1635        let worker = spawn_worker_with_spool_authenticated(
1636            4,
1637            Arc::new(RecordingBus::default()),
1638            Arc::new(HashEmbedder::default()),
1639            Arc::new(NoopVectorSink),
1640            Some(directory.path().to_path_buf()),
1641            journal_key,
1642            Arc::new(Registry::new()),
1643        );
1644        let session = session(Workflow::Signed);
1645        worker.submit_and_wait(job(Arc::clone(&session))).await.unwrap();
1646
1647        let digest = av_core::digest::sha256_hex(session.id.as_bytes());
1648        let stem = digest.get(..32).unwrap();
1649        let metadata: Value = crate::journal::open(
1650            &journal_key,
1651            "metadata",
1652            0,
1653            &std::fs::read(directory.path().join(format!("{stem}.session.json"))).unwrap(),
1654        )
1655        .unwrap();
1656        assert_eq!(metadata["workflow"], "signed");
1657        let journal =
1658            std::fs::read_to_string(directory.path().join(format!("{stem}.events.ndjson"))).unwrap();
1659        let record: ActiveJournalRecord =
1660            crate::journal::open(&journal_key, "session-1:active", 0, journal.trim().as_bytes()).unwrap();
1661        assert_eq!(record.prompt_tokens, 100);
1662        assert_eq!(record.completion_tokens, 0);
1663        assert_eq!(record.event["class_name"], "agent.compression");
1664    }
1665
1666    /// A session marked `capture_failed` (typically because a prior envelope's
1667    /// journal append failed) must not have subsequent queued envelopes advance
1668    /// `next_seq` and land in the journal — their `event.metadata.sequence`
1669    /// would exceed the entry's byte position, which recovery rejects with
1670    /// "signed event sequence does not match active journal index". Regression
1671    /// against the process_job pattern that treated capture_failed as a
1672    /// post-processing check only.
1673    #[tokio::test]
1674    async fn capture_failed_session_does_not_write_further_journal_entries() {
1675        let directory = tempfile::tempdir().unwrap();
1676        let journal_key = [23; 32];
1677        let worker = spawn_worker_with_spool_authenticated(
1678            4,
1679            Arc::new(RecordingBus::default()),
1680            Arc::new(HashEmbedder::default()),
1681            Arc::new(NoopVectorSink),
1682            Some(directory.path().to_path_buf()),
1683            journal_key,
1684            Arc::new(Registry::new()),
1685        );
1686        let session = session(Workflow::Signed);
1687        worker.submit_and_wait(job(Arc::clone(&session))).await.unwrap();
1688        let digest = av_core::digest::sha256_hex(session.id.as_bytes());
1689        let stem = digest.get(..32).unwrap();
1690        let journal_path = directory.path().join(format!("{stem}.events.ndjson"));
1691        let bytes_after_first = std::fs::read(&journal_path).unwrap();
1692        assert!(
1693            !bytes_after_first.is_empty(),
1694            "first envelope must journal one line"
1695        );
1696        assert_eq!(
1697            bytes_after_first.iter().filter(|byte| **byte == b'\n').count(),
1698            1,
1699            "first envelope must produce exactly one journal line",
1700        );
1701        session.mark_capture_failed();
1702        // The second envelope must silently short-circuit: no journal line, no
1703        // chain append, no next_seq increment for the poisoned session.
1704        let seq_before = session.next_seq();
1705        worker.submit_and_wait(job(Arc::clone(&session))).await.unwrap();
1706        let seq_after = session.next_seq();
1707        assert_eq!(
1708            seq_after,
1709            seq_before + 1,
1710            "process_job for a capture_failed session must not advance next_seq",
1711        );
1712        assert_eq!(
1713            std::fs::read(&journal_path).unwrap(),
1714            bytes_after_first,
1715            "process_job for a capture_failed session must not append to the journal",
1716        );
1717        assert_eq!(
1718            session.chain.lock().count(),
1719            1,
1720            "process_job for a capture_failed session must not extend the chain",
1721        );
1722    }
1723
1724    #[tokio::test]
1725    async fn signed_job_journal_recovers_replays_and_issues_one_receipt() {
1726        let directory = tempfile::tempdir().unwrap();
1727        let bridge = Arc::new(RecordingBus::default());
1728        let metrics = Arc::new(Registry::new());
1729        let signer = Arc::new(av_receipts::Ed25519Signer::from_seed(&[14; 32]));
1730        let journal_key = crate::journal::key_from_signer(signer.as_ref());
1731        let worker = spawn_worker_with_spool_authenticated(
1732            4,
1733            bridge.clone(),
1734            Arc::new(HashEmbedder::default()),
1735            Arc::new(NoopVectorSink),
1736            Some(directory.path().to_path_buf()),
1737            journal_key,
1738            Arc::clone(&metrics),
1739        );
1740        let active = session(Workflow::Signed);
1741        worker.submit_and_wait(job(Arc::clone(&active))).await.unwrap();
1742        let digest = av_core::digest::sha256_hex(active.id.as_bytes());
1743        let stem = digest.get(..32).unwrap();
1744        let ack_path = std::fs::read_dir(directory.path().join("broker-acks").join(stem))
1745            .unwrap()
1746            .next()
1747            .unwrap()
1748            .unwrap()
1749            .path();
1750        std::fs::remove_file(ack_path).unwrap();
1751
1752        let registry = crate::session::SessionRegistry::new();
1753        let finalizer = crate::reconciler::Finalizer::with_bridge(
1754            signer,
1755            directory.path().to_path_buf(),
1756            metrics,
1757            bridge.clone(),
1758        );
1759        assert_eq!(
1760            finalizer
1761                .recover_spooled_sessions(&registry, &Default::default())
1762                .await
1763                .unwrap(),
1764            1
1765        );
1766        let recovered = registry.get(&active.id).unwrap();
1767        let receipt = recovered.receipt.lock().clone().unwrap();
1768        receipt.verify_embedded().unwrap();
1769        assert!(matches!(
1770            receipt.body.subject,
1771            av_receipts::ReceiptSubject::EventChain { event_count: 1, .. }
1772        ));
1773        assert_eq!(receipt.body.cost.prompt_tokens, 100);
1774        let events = bridge.events.lock();
1775        assert_eq!(events.len(), 3, "acknowledged active event must not be replayed");
1776        assert_eq!(events[0].2["metadata"]["sequence"], 0);
1777        assert_eq!(events[1].2["metadata"]["sequence"], 1);
1778        assert_eq!(events[2].2["metadata"]["sequence"], 2);
1779        assert!(!directory.path().join(format!("{stem}.events.ndjson")).exists());
1780    }
1781
1782    #[tokio::test]
1783    async fn recovery_quarantines_only_the_incomplete_signed_session() {
1784        let directory = tempfile::tempdir().unwrap();
1785        let bridge = Arc::new(RecordingBus::default());
1786        let metrics = Arc::new(Registry::new());
1787        let signer = Arc::new(av_receipts::Ed25519Signer::from_seed(&[27; 32]));
1788        let journal_key = crate::journal::key_from_signer(signer.as_ref());
1789        let worker = spawn_worker_with_spool_authenticated(
1790            8,
1791            bridge.clone(),
1792            Arc::new(HashEmbedder::default()),
1793            Arc::new(NoopVectorSink),
1794            Some(directory.path().to_path_buf()),
1795            journal_key,
1796            Arc::clone(&metrics),
1797        );
1798        let incomplete = Arc::new(Session::new(
1799            "incomplete-signed".to_owned(),
1800            Workflow::Signed,
1801            session(Workflow::Signed).current_identity(),
1802            BreakerConfig::default(),
1803        ));
1804        let complete = Arc::new(Session::new(
1805            "complete-signed".to_owned(),
1806            Workflow::Signed,
1807            session(Workflow::Signed).current_identity(),
1808            BreakerConfig::default(),
1809        ));
1810        let mut incomplete_request = job(Arc::clone(&incomplete));
1811        incomplete_request.response_attempt = Some(ResponseAttempt {
1812            id: "incomplete-attempt".to_owned(),
1813            terminal: false,
1814        });
1815        worker.submit_and_wait(incomplete_request).await.unwrap();
1816        worker.submit_and_wait(job(Arc::clone(&complete))).await.unwrap();
1817
1818        let registry = crate::session::SessionRegistry::new();
1819        let finalizer = crate::reconciler::Finalizer::with_bridge(
1820            signer,
1821            directory.path().to_path_buf(),
1822            metrics,
1823            bridge,
1824        );
1825        assert_eq!(
1826            finalizer
1827                .recover_spooled_sessions(&registry, &Default::default())
1828                .await
1829                .unwrap(),
1830            2
1831        );
1832        let quarantined = registry.get("incomplete-signed").unwrap();
1833        assert!(quarantined.capture_failed());
1834        assert!(quarantined.receipt.lock().is_none());
1835        assert!(registry.get("complete-signed").unwrap().receipt.lock().is_some());
1836    }
1837
1838    #[tokio::test]
1839    async fn tampered_signed_journal_is_never_turned_into_a_receipt() {
1840        let directory = tempfile::tempdir().unwrap();
1841        let signer = Arc::new(av_receipts::Ed25519Signer::from_seed(&[18; 32]));
1842        let journal_key = crate::journal::key_from_signer(signer.as_ref());
1843        let worker = spawn_worker_with_spool_authenticated(
1844            4,
1845            Arc::new(RecordingBus::default()),
1846            Arc::new(HashEmbedder::default()),
1847            Arc::new(NoopVectorSink),
1848            Some(directory.path().to_path_buf()),
1849            journal_key,
1850            Arc::new(Registry::new()),
1851        );
1852        let active = session(Workflow::Signed);
1853        worker.submit_and_wait(job(Arc::clone(&active))).await.unwrap();
1854        let digest = av_core::digest::sha256_hex(active.id.as_bytes());
1855        let stem = digest.get(..32).unwrap();
1856        let journal_path = directory.path().join(format!("{stem}.events.ndjson"));
1857        let mut envelope: Value =
1858            serde_json::from_str(std::fs::read_to_string(&journal_path).unwrap().trim()).unwrap();
1859        envelope["payload"]["prompt_tokens"] = serde_json::json!(999_999);
1860        std::fs::write(
1861            &journal_path,
1862            format!("{}\n", serde_json::to_string(&envelope).unwrap()),
1863        )
1864        .unwrap();
1865        let registry = crate::session::SessionRegistry::new();
1866        let finalizer = crate::reconciler::Finalizer::new(
1867            signer,
1868            directory.path().to_path_buf(),
1869            Arc::new(Registry::new()),
1870        );
1871        // Round-41 F1: per-session errors during signed recovery no
1872        // longer propagate to the outer `recover_spooled_sessions`
1873        // Err (which used to head-of-line-block every other session
1874        // for the reconciler tick). They now warn+continue. The
1875        // security invariant this test locks in is unchanged: the
1876        // tampered journal MUST NOT be turned into a receipt AND
1877        // the corrupted session MUST NOT be installed into the
1878        // registry (both were previously ensured by the outer Err
1879        // short-circuit). Now they're ensured by the async block
1880        // returning Err BEFORE `try_insert_recovered` runs.
1881        let outcome = finalizer
1882            .recover_spooled_sessions(&registry, &Default::default())
1883            .await;
1884        assert!(
1885            outcome.is_ok(),
1886            "round-41 F1: per-session HMAC failures warn+continue instead of propagating, got {outcome:?}"
1887        );
1888        // The corrupted session must NOT be installed into the
1889        // registry — the security property this test was written to
1890        // enforce.
1891        assert!(registry.get(&active.id).is_none());
1892    }
1893
1894    #[tokio::test]
1895    async fn signed_restart_reuses_receipt_persisted_before_journal_cleanup() {
1896        let directory = tempfile::tempdir().unwrap();
1897        let signer = Arc::new(av_receipts::Ed25519Signer::from_seed(&[15; 32]));
1898        let journal_key = crate::journal::key_from_signer(signer.as_ref());
1899        let worker = spawn_worker_with_spool_authenticated(
1900            4,
1901            Arc::new(RecordingBus::default()),
1902            Arc::new(HashEmbedder::default()),
1903            Arc::new(NoopVectorSink),
1904            Some(directory.path().to_path_buf()),
1905            journal_key,
1906            Arc::new(Registry::new()),
1907        );
1908        let active = session(Workflow::Signed);
1909        worker.submit_and_wait(job(Arc::clone(&active))).await.unwrap();
1910        let digest = av_core::digest::sha256_hex(active.id.as_bytes());
1911        let stem = digest.get(..32).unwrap();
1912        let metadata_path = directory.path().join(format!("{stem}.session.json"));
1913        let journal_path = directory.path().join(format!("{stem}.events.ndjson"));
1914        let metadata = std::fs::read(&metadata_path).unwrap();
1915        let journal = std::fs::read(&journal_path).unwrap();
1916        let ack_path = std::fs::read_dir(directory.path().join("broker-acks").join(stem))
1917            .unwrap()
1918            .next()
1919            .unwrap()
1920            .unwrap()
1921            .path();
1922        let ack = std::fs::read(&ack_path).unwrap();
1923        let first_finalizer = crate::reconciler::Finalizer::new(
1924            signer.clone(),
1925            directory.path().to_path_buf(),
1926            Arc::new(Registry::new()),
1927        );
1928        let crate::reconciler::FinalizeOutcome::Receipt { receipt: first } = first_finalizer
1929            .close_session(active, av_events::StopReason::SessionClosed)
1930            .await
1931            .unwrap()
1932        else {
1933            panic!("expected signed receipt")
1934        };
1935
1936        std::fs::write(&metadata_path, metadata).unwrap();
1937        std::fs::write(&journal_path, journal).unwrap();
1938        std::fs::create_dir_all(ack_path.parent().unwrap()).unwrap();
1939        std::fs::write(&ack_path, ack).unwrap();
1940        let registry = crate::session::SessionRegistry::new();
1941        let after_restart = crate::reconciler::Finalizer::new(
1942            signer,
1943            directory.path().to_path_buf(),
1944            Arc::new(Registry::new()),
1945        );
1946        assert_eq!(
1947            after_restart
1948                .recover_spooled_sessions(&registry, &Default::default())
1949                .await
1950                .unwrap(),
1951            1
1952        );
1953        let recovered = registry.get("session-1").unwrap();
1954        assert_eq!(
1955            recovered.receipt.lock().as_ref().unwrap().body.receipt_id,
1956            first.body.receipt_id
1957        );
1958    }
1959
1960    #[tokio::test]
1961    async fn unsigned_job_preserves_required_atif_metrics() {
1962        let bridge = Arc::new(RecordingBus::default());
1963        let worker = spawn_worker(
1964            4,
1965            bridge,
1966            Arc::new(HashEmbedder::default()),
1967            Arc::new(Registry::new()),
1968        );
1969        let session = session(Workflow::Unsigned);
1970        worker.submit_and_wait(job(Arc::clone(&session))).await.unwrap();
1971
1972        let trajectory = session.take_trajectory();
1973        assert_eq!(trajectory.steps.len(), 1);
1974        assert_eq!(
1975            trajectory.steps[0].metrics.as_ref().unwrap().cached_tokens,
1976            Some(40)
1977        );
1978        assert_eq!(session.totals.cached_tokens.load(Ordering::Acquire), 0);
1979    }
1980
1981    #[tokio::test]
1982    async fn unsigned_job_journal_survives_restart_and_promotes() {
1983        let directory = tempfile::tempdir().unwrap();
1984        let bridge = Arc::new(RecordingBus::default());
1985        let metrics = Arc::new(Registry::new());
1986        let signer = Arc::new(av_receipts::Ed25519Signer::from_seed(&[13; 32]));
1987        let journal_key = crate::journal::key_from_signer(signer.as_ref());
1988        let worker = spawn_worker_with_spool_authenticated(
1989            4,
1990            bridge.clone(),
1991            Arc::new(HashEmbedder::default()),
1992            Arc::new(NoopVectorSink),
1993            Some(directory.path().to_path_buf()),
1994            journal_key,
1995            Arc::clone(&metrics),
1996        );
1997        let active = session(Workflow::Unsigned);
1998        worker.submit_and_wait(job(Arc::clone(&active))).await.unwrap();
1999        let mut tool = job(Arc::clone(&active));
2000        tool.class = EventClass::ToolCall;
2001        tool.status = StatusId::Failure;
2002        tool.stop_reason = Some(StopReason::PolicyBlocked);
2003        tool.metrics = EventMetrics::default();
2004        tool.cost_usd_micros = 0;
2005        tool.atif.as_mut().unwrap().llm_call_count = Some(0);
2006        worker.submit_and_wait(tool).await.unwrap();
2007        assert_eq!(bridge.events.lock().len(), 2);
2008
2009        let registry = crate::session::SessionRegistry::new();
2010        let finalizer = crate::reconciler::Finalizer::with_bridge(
2011            signer,
2012            directory.path().to_path_buf(),
2013            metrics,
2014            bridge.clone(),
2015        );
2016        assert_eq!(
2017            finalizer
2018                .recover_spooled_sessions(&registry, &Default::default())
2019                .await
2020                .unwrap(),
2021            1
2022        );
2023        let recovered = registry.get(&active.id).unwrap();
2024        assert!(recovered.atif_path.lock().as_ref().unwrap().exists());
2025        assert_eq!(
2026            bridge.events.lock().len(),
2027            2,
2028            "acknowledged unsigned events must not replay"
2029        );
2030        assert_eq!(recovered.current_identity().ttl_remaining_s, Some(600));
2031        let receipt = finalizer.promote(recovered).await.unwrap();
2032        receipt.verify_embedded().unwrap();
2033        assert_eq!(receipt.body.cost.prompt_tokens, 100);
2034        assert_eq!(receipt.body.tool_calls.total, 1);
2035        assert_eq!(receipt.body.tool_calls.blocked, 1);
2036        assert_eq!(receipt.body.stop_reason_id, StopReason::PolicyBlocked.id());
2037        assert!(matches!(
2038            receipt.body.subject,
2039            av_receipts::ReceiptSubject::AtifTrajectory {
2040                step_count: 2,
2041                retroactive: true,
2042                ..
2043            }
2044        ));
2045    }
2046
2047    #[tokio::test]
2048    async fn full_queue_is_counted_without_blocking() {
2049        let metrics = Arc::new(Registry::new());
2050        let (sender, _receiver) = mpsc::channel(1);
2051        let worker = WorkerHandle {
2052            senders: Arc::new(vec![sender]),
2053            metrics: Arc::clone(&metrics),
2054            pending: Arc::new(std::sync::atomic::AtomicU64::new(0)),
2055            drained: Arc::new(tokio::sync::Notify::new()),
2056            capacity: Arc::new(tokio::sync::Semaphore::new(1)),
2057            response_capacity: Arc::new(tokio::sync::Semaphore::new(1)),
2058        };
2059        worker.try_submit(job(session(Workflow::Signed))).unwrap();
2060        assert_eq!(
2061            worker.try_submit(job(session(Workflow::Signed))),
2062            Err(SubmitError::Full)
2063        );
2064        assert!(metrics
2065            .render()
2066            .contains("av_events_dropped_total{stage=\"worker_queue\"} 1"));
2067    }
2068
2069    /// The fused permit split (worker_queue vs response_slot) exists
2070    /// so operators can distinguish which class of admission ran out.
2071    /// Assert the counters are actually distinct at the metrics
2072    /// registry level under two adversarial saturations.
2073    #[tokio::test]
2074    async fn worker_queue_and_response_slot_counters_are_distinct() {
2075        let metrics = Arc::new(Registry::new());
2076        let (sender, _receiver) = mpsc::channel(1);
2077        let worker = WorkerHandle {
2078            senders: Arc::new(vec![sender]),
2079            metrics: Arc::clone(&metrics),
2080            pending: Arc::new(std::sync::atomic::AtomicU64::new(0)),
2081            drained: Arc::new(tokio::sync::Notify::new()),
2082            capacity: Arc::new(tokio::sync::Semaphore::new(1)),
2083            response_capacity: Arc::new(tokio::sync::Semaphore::new(1)),
2084        };
2085        // First pair succeeds — worker and response semaphores each go
2086        // to 0.
2087        let first = worker.try_reserve_pair("s").expect("first pair");
2088        // Second pair fails at the worker stage (first failure). Only
2089        // the worker_queue counter should be bumped.
2090        assert_eq!(worker.try_reserve_pair("s").err(), Some(SubmitError::Full));
2091        let rendered = metrics.render();
2092        assert!(
2093            rendered.contains("av_events_dropped_total{stage=\"worker_queue\"} 1"),
2094            "expected worker_queue counter increment, got:\n{rendered}"
2095        );
2096        assert!(
2097            !rendered.contains("av_events_dropped_total{stage=\"response_slot\"}"),
2098            "response_slot must not have been touched when worker_queue exhausts first, got:\n{rendered}"
2099        );
2100        // Free the worker slot but keep the response permit held. A
2101        // fresh pair acquire should succeed on the worker side then
2102        // fail at the response stage, bumping response_slot.
2103        drop(first.worker);
2104        let _still_holding_response = first.response;
2105        assert_eq!(worker.try_reserve_pair("s").err(), Some(SubmitError::Full));
2106        let rendered = metrics.render();
2107        assert!(
2108            rendered.contains("av_events_dropped_total{stage=\"response_slot\"} 1"),
2109            "expected response_slot counter increment, got:\n{rendered}"
2110        );
2111        assert!(
2112            rendered.contains("av_events_dropped_total{stage=\"worker_queue\"} 1"),
2113            "worker_queue must remain at 1 (only response stage exhausted this time), got:\n{rendered}"
2114        );
2115    }
2116
2117    /// The prior review caught that `ResponsePermit::submit` was going
2118    /// through the generic `try_submit` path, which bumps
2119    /// `stage="worker_queue"` on any capacity/mpsc-full failure — so a
2120    /// mid-stream response-capture job that raced with re-admitted
2121    /// worker jobs would be silently misattributed to the wrong
2122    /// counter, defeating the observability goal of the fused-permit
2123    /// split. Locking the correct routing: saturate worker capacity
2124    /// AFTER a response permit is issued, then commit the response job
2125    /// and assert it lands on `stage="response_slot"`.
2126    #[tokio::test]
2127    async fn response_permit_submit_bumps_response_slot_on_worker_capacity_exhaustion() {
2128        let bridge = Arc::new(RecordingBus::default());
2129        let metrics = Arc::new(Registry::new());
2130        let worker = spawn_worker(
2131            2,
2132            bridge.clone(),
2133            Arc::new(HashEmbedder::default()),
2134            Arc::clone(&metrics),
2135        );
2136        // Issue a pair permit (worker=1 held; response=1 held for the
2137        // rest of the test).
2138        let permits = worker
2139            .try_reserve_pair("mid-stream-response")
2140            .expect("initial pair");
2141        // Drop the worker permit so that half is free; the response
2142        // half stays held until we submit below.
2143        drop(permits.worker);
2144        // Now consume both remaining worker slots directly so any
2145        // further worker-side acquire fails. The response semaphore
2146        // is untouched — permits.response still holds its slot.
2147        let hog_a = worker.try_reserve("hog-a").expect("hog-a");
2148        let hog_b = worker.try_reserve("hog-b").expect("hog-b");
2149        // Submit the response-capture job. It draws from worker
2150        // capacity + mpsc — worker capacity is exhausted → `Full`.
2151        // The counter increment MUST be `stage="response_slot"`,
2152        // NOT `stage="worker_queue"`.
2153        let job = job(session(Workflow::Signed));
2154        let err = permits.response.submit(&worker, job).unwrap_err();
2155        assert_eq!(err, SubmitError::Full);
2156        let rendered = metrics.render();
2157        assert!(
2158            rendered.contains("av_events_dropped_total{stage=\"response_slot\"} 1"),
2159            "response_slot MUST have been bumped by ResponsePermit::submit; got:\n{rendered}"
2160        );
2161        assert!(
2162            !rendered.contains("av_events_dropped_total{stage=\"worker_queue\"} 1"),
2163            "worker_queue must NOT be bumped by ResponsePermit::submit failure — that would \
2164             mask which class of exhaustion produced the drop; got:\n{rendered}"
2165        );
2166        drop(hog_a);
2167        drop(hog_b);
2168    }
2169
2170    #[tokio::test]
2171    async fn one_shard_can_borrow_all_global_capacity() {
2172        const CAPACITY: usize = 16;
2173        let worker = spawn_worker(
2174            CAPACITY,
2175            Arc::new(RecordingBus::default()),
2176            Arc::new(HashEmbedder::default()),
2177            Arc::new(Registry::new()),
2178        );
2179        let partitions = u32::try_from(CAPACITY).unwrap();
2180        let target = av_bridge::bus::partition_for("target", partitions);
2181        let session_ids: Vec<String> = (0..10_000)
2182            .map(|index| format!("same-shard-{index}"))
2183            .filter(|id| av_bridge::bus::partition_for(id, partitions) == target)
2184            .take(CAPACITY)
2185            .collect();
2186        assert_eq!(session_ids.len(), CAPACITY);
2187        let permits: Vec<_> = session_ids
2188            .iter()
2189            .map(|id| worker.try_reserve(id).unwrap())
2190            .collect();
2191        assert!(matches!(
2192            worker.try_reserve("globally-full"),
2193            Err(SubmitError::Full)
2194        ));
2195        drop(permits);
2196        assert!(worker.try_reserve("capacity-released").is_ok());
2197    }
2198
2199    #[tokio::test]
2200    async fn panic_is_counted_and_supervisor_processes_next_job() {
2201        let bridge = Arc::new(RecordingBus::default());
2202        let metrics = Arc::new(Registry::new());
2203        let worker = spawn_worker_with_sink(
2204            4,
2205            bridge.clone(),
2206            Arc::new(HashEmbedder::default()),
2207            Arc::new(PanicOnceSink {
2208                panicked: AtomicBool::new(false),
2209            }),
2210            Arc::clone(&metrics),
2211        );
2212        let failed_session = session(Workflow::Signed);
2213        assert!(worker
2214            .submit_and_wait(job(Arc::clone(&failed_session)))
2215            .await
2216            .is_err());
2217        assert!(failed_session.capture_failed());
2218        worker
2219            .submit_and_wait(job(session(Workflow::Signed)))
2220            .await
2221            .unwrap();
2222        assert_eq!(bridge.events.lock().len(), 1);
2223        assert!(metrics.render().contains("av_worker_panics_total 1"));
2224    }
2225
2226    // Regression: subscribe (enable) and await must operate on the SAME pinned
2227    // Notified, or notify_waiters() firing between the two loses the wakeup.
2228    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
2229    async fn wait_idle_survives_notify_races_under_contention() {
2230        let bridge = Arc::new(RecordingBus::default());
2231        let worker = spawn_worker(
2232            64,
2233            bridge.clone(),
2234            Arc::new(HashEmbedder::default()),
2235            Arc::new(Registry::new()),
2236        );
2237        for _ in 0..64 {
2238            for _ in 0..4 {
2239                worker.try_submit(job(session(Workflow::Signed))).unwrap();
2240            }
2241            tokio::time::timeout(std::time::Duration::from_secs(5), worker.wait_idle())
2242                .await
2243                .expect("wait_idle deadlocked: notify_waiters was lost during subscribe/check race");
2244        }
2245    }
2246
2247    // Deterministic proof of the enable/await invariant: producer's
2248    // notify_waiters() fires between the consumer's check and its await.
2249    // If the consumer re-subscribes with a fresh notified() after check, the
2250    // notify_waiters() is lost and the consumer deadlocks. The correct pattern
2251    // keeps the same pinned Notified alive across the check.
2252    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2253    async fn notified_enable_check_await_pattern_is_race_free() {
2254        use std::sync::atomic::AtomicUsize;
2255        use tokio::sync::Notify;
2256
2257        for _ in 0..64 {
2258            let notify = Arc::new(Notify::new());
2259            let pending = Arc::new(AtomicUsize::new(1));
2260            let (check_done_tx, check_done_rx) = tokio::sync::oneshot::channel::<()>();
2261            let (notify_done_tx, notify_done_rx) = tokio::sync::oneshot::channel::<()>();
2262
2263            let n_c = Arc::clone(&notify);
2264            let p_c = Arc::clone(&pending);
2265            let consumer = tokio::spawn(async move {
2266                let notified = n_c.notified();
2267                let mut notified = std::pin::pin!(notified);
2268                notified.as_mut().enable();
2269                assert_ne!(p_c.load(AtomicOrdering::Acquire), 0, "producer ran early");
2270                check_done_tx.send(()).unwrap();
2271                notify_done_rx.await.unwrap();
2272                notified.await;
2273            });
2274
2275            let n_p = Arc::clone(&notify);
2276            let p_p = Arc::clone(&pending);
2277            let producer = tokio::spawn(async move {
2278                check_done_rx.await.unwrap();
2279                p_p.store(0, AtomicOrdering::Release);
2280                n_p.notify_waiters();
2281                notify_done_tx.send(()).unwrap();
2282            });
2283
2284            tokio::time::timeout(std::time::Duration::from_secs(2), async {
2285                producer.await.unwrap();
2286                consumer.await.unwrap();
2287            })
2288            .await
2289            .expect("consumer deadlocked: enable/await pattern violated");
2290        }
2291    }
2292
2293    // ------------------------------------------------------------------
2294    // Congestion & bottleneck stress tests.
2295    // ------------------------------------------------------------------
2296
2297    /// `try_submit` must never block the caller when the worker is saturated —
2298    /// it must return `SubmitError::Full` immediately. A regression that made
2299    /// the hot path await capacity would silently convert backpressure into
2300    /// head-of-line blocking on the request pipeline.
2301    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2302    async fn try_submit_returns_full_within_the_first_millisecond_when_saturated() {
2303        const CAPACITY: usize = 4;
2304        let worker = spawn_worker(
2305            CAPACITY,
2306            Arc::new(RecordingBus::default()),
2307            Arc::new(HashEmbedder::default()),
2308            Arc::new(Registry::new()),
2309        );
2310        // Hold every permit so try_submit sees a saturated semaphore.
2311        let session = session(Workflow::Signed);
2312        let _held: Vec<_> = (0..CAPACITY)
2313            .map(|i| worker.try_reserve(&format!("hold-{i}")).unwrap())
2314            .collect();
2315        let started = std::time::Instant::now();
2316        let result = worker.try_submit(job(Arc::clone(&session)));
2317        let elapsed = started.elapsed();
2318        assert!(matches!(result, Err(SubmitError::Full)));
2319        assert!(
2320            elapsed < std::time::Duration::from_millis(1),
2321            "try_submit blocked for {elapsed:?} when saturated \
2322             (backpressure should be synchronous)",
2323        );
2324    }
2325
2326    /// Metrics must record every drop; operators cannot see congestion
2327    /// otherwise. This locks the counter name and its increment on both
2328    /// the semaphore-exhaustion path (returning Full from `try_capacity`)
2329    /// and the closed-channel path.
2330    #[tokio::test]
2331    async fn every_dropped_job_increments_the_drop_metric() {
2332        const CAPACITY: usize = 2;
2333        let bridge = Arc::new(RecordingBus::default());
2334        let metrics = Arc::new(Registry::new());
2335        let worker = spawn_worker(
2336            CAPACITY,
2337            bridge,
2338            Arc::new(HashEmbedder::default()),
2339            Arc::clone(&metrics),
2340        );
2341        let _hold: Vec<_> = (0..CAPACITY)
2342            .map(|i| worker.try_reserve(&format!("s-{i}")).unwrap())
2343            .collect();
2344        for _ in 0..17 {
2345            assert!(matches!(worker.try_reserve("x"), Err(SubmitError::Full)));
2346        }
2347        let rendered = metrics.render();
2348        assert!(
2349            rendered.contains("av_events_dropped_total{stage=\"worker_queue\"} 17"),
2350            "expected 17 drops in metrics, got:\n{rendered}",
2351        );
2352    }
2353
2354    /// Capacity released by dropping a permit is immediately reusable.
2355    /// If the semaphore forgot to release on drop, a bursty workload would
2356    /// hit a false Full condition and reject every subsequent request
2357    /// until the worker was recycled.
2358    #[tokio::test]
2359    async fn capacity_is_reusable_after_permit_drop() {
2360        const CAPACITY: usize = 8;
2361        let worker = spawn_worker(
2362            CAPACITY,
2363            Arc::new(RecordingBus::default()),
2364            Arc::new(HashEmbedder::default()),
2365            Arc::new(Registry::new()),
2366        );
2367        for cycle in 0..1_000 {
2368            let permits: Vec<_> = (0..CAPACITY)
2369                .map(|i| worker.try_reserve(&format!("cycle-{cycle}-{i}")).unwrap())
2370                .collect();
2371            assert!(matches!(worker.try_reserve("overflow"), Err(SubmitError::Full)));
2372            drop(permits);
2373            // Full capacity must be immediately re-acquirable.
2374            let recycled: Vec<_> = (0..CAPACITY)
2375                .map(|i| worker.try_reserve(&format!("recycle-{cycle}-{i}")).unwrap())
2376                .collect();
2377            drop(recycled);
2378        }
2379    }
2380
2381    /// Shard partitioning is deterministic and roughly balanced. If one
2382    /// session hashed to the same shard as many others, that shard's queue
2383    /// would fill up first and every other shard would run under-utilized
2384    /// (see `one_shard_can_borrow_all_global_capacity` for the extreme).
2385    /// This test asserts the *balance* invariant across a wide id space.
2386    #[test]
2387    fn partition_for_distributes_uniformly_across_shards() {
2388        const SHARDS: u32 = 16;
2389        const SAMPLES: u32 = 100_000;
2390        let mut hits = [0u32; SHARDS as usize];
2391        for i in 0..SAMPLES {
2392            let id = format!("session-{i}");
2393            let s = av_bridge::bus::partition_for(&id, SHARDS) as usize;
2394            hits[s] += 1;
2395        }
2396        let expected = SAMPLES / SHARDS;
2397        for (shard, count) in hits.iter().enumerate() {
2398            let lower = expected * 9 / 10;
2399            let upper = expected * 11 / 10;
2400            assert!(
2401                (lower..=upper).contains(count),
2402                "shard {shard} got {count} hits, expected {expected} (10% band)",
2403            );
2404        }
2405    }
2406
2407    /// All 16 partitions must have a live shard channel, even when the
2408    /// caller passes a small `capacity`. An earlier version sized
2409    /// `shard_count = capacity.min(16)`, so with `capacity < 16` any
2410    /// session id hashing to a partition without a shard silently
2411    /// returned `SubmitError::Closed` — 15/16 of the id space was
2412    /// unroutable at capacity = 1.
2413    #[tokio::test]
2414    async fn every_shard_is_routable_at_capacity_one() {
2415        let worker = spawn_worker(
2416            1,
2417            Arc::new(RecordingBus::default()),
2418            Arc::new(HashEmbedder::default()),
2419            Arc::new(Registry::new()),
2420        );
2421        // Probe every one of the 16 possible partition indices with a
2422        // synthetic id crafted to hash there.
2423        for target in 0..16u32 {
2424            let mut candidate = None;
2425            for i in 0..10_000u64 {
2426                let id = format!("probe-{target}-{i}");
2427                if av_bridge::bus::partition_for(&id, 16) == target {
2428                    candidate = Some(id);
2429                    break;
2430                }
2431            }
2432            let id = candidate.unwrap_or_else(|| panic!("no id found for shard {target}"));
2433            let permit = worker
2434                .try_reserve(&id)
2435                .unwrap_or_else(|error| panic!("shard {target} not routable: {error:?}"));
2436            drop(permit);
2437        }
2438    }
2439
2440    /// Design invariant: `submit_and_wait` canceled at any point must not
2441    /// leave the `pending` counter incremented, because a leaked pending
2442    /// count would deadlock `wait_idle` on shutdown. The coupled
2443    /// semaphore + shard-channel capacity guarantees `send().await` never
2444    /// blocks while a permit is held, so cancellation before the envelope
2445    /// enters the channel can only happen while the counter has not yet
2446    /// been incremented. This test locks that invariant across many
2447    /// cancellation timings.
2448    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2449    async fn cancelling_submit_and_wait_never_leaks_pending() {
2450        const CAPACITY: usize = 4;
2451        let bridge = Arc::new(RecordingBus::default());
2452        let worker = Arc::new(spawn_worker(
2453            CAPACITY,
2454            bridge.clone(),
2455            Arc::new(HashEmbedder::default()),
2456            Arc::new(Registry::new()),
2457        ));
2458        // Saturate the semaphore so submit_and_wait always blocks on
2459        // acquire_owned; abort at randomized delays to hit every polling
2460        // window of the acquire future.
2461        let held: Vec<_> = (0..CAPACITY)
2462            .map(|i| worker.try_reserve(&format!("hold-{i}")).unwrap())
2463            .collect();
2464        for delay_us in [0u64, 50, 200, 800, 3_200] {
2465            let hostile_session = session(Workflow::Signed);
2466            let worker_for_task = Arc::clone(&worker);
2467            let task = tokio::spawn(async move {
2468                let _ = worker_for_task
2469                    .submit_and_wait(job(Arc::clone(&hostile_session)))
2470                    .await;
2471            });
2472            if delay_us > 0 {
2473                tokio::time::sleep(std::time::Duration::from_micros(delay_us)).await;
2474            }
2475            task.abort();
2476            let _ = task.await;
2477        }
2478        drop(held);
2479        tokio::time::timeout(std::time::Duration::from_secs(2), worker.wait_idle())
2480            .await
2481            .expect("wait_idle deadlocked — a canceled submit_and_wait leaked pending");
2482    }
2483}