1use 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
14pub struct AtifCapture {
16 pub source: av_atif::Source,
18 pub message: Value,
20 pub reasoning_content: Option<String>,
22 pub model_name: Option<String>,
24 pub tool_calls: Option<Vec<av_atif::ToolCall>>,
26 pub observation: Option<av_atif::Observation>,
28 pub llm_call_count: Option<u64>,
30}
31
32pub struct WorkerJob {
34 pub session: Arc<Session>,
36 pub identity: av_events::AgentIdentity,
38 pub class: EventClass,
40 pub payload: Value,
42 pub text: String,
44 pub analyze_loop: bool,
47 pub status: StatusId,
49 pub stop_reason: Option<StopReason>,
51 pub native_stop_reason: Option<String>,
53 pub metrics: EventMetrics,
55 pub cost_usd_micros: u64,
57 pub atif: Option<AtifCapture>,
59 pub response_marker: Option<String>,
61 pub response_attempt: Option<ResponseAttempt>,
63}
64
65#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
67pub struct ResponseAttempt {
68 pub id: String,
70 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#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
113pub enum SubmitError {
114 #[error("worker queue is full")]
116 Full,
117 #[error("worker queue is closed")]
119 Closed,
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126enum DropStage {
127 WorkerQueue,
130 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#[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 response_capacity: Arc<tokio::sync::Semaphore>,
170}
171
172pub struct WorkerPermit {
174 permit: mpsc::OwnedPermit<Envelope>,
175 pending: Arc<std::sync::atomic::AtomicU64>,
176 capacity_permit: tokio::sync::OwnedSemaphorePermit,
177}
178
179pub struct ResponsePermit {
198 _capacity_permit: tokio::sync::OwnedSemaphorePermit,
199}
200
201impl ResponsePermit {
202 pub fn submit(self, worker: &WorkerHandle, job: WorkerJob) -> Result<(), SubmitError> {
210 let _capacity_permit = self._capacity_permit;
216 worker.try_submit_labeled(job, DropStage::ResponseSlot)
217 }
218}
219
220pub struct WorkerAndResponsePermit {
225 pub worker: WorkerPermit,
227 pub response: ResponsePermit,
229}
230
231impl WorkerPermit {
232 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 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 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 pub fn try_reserve_pair(&self, session_id: &str) -> Result<WorkerAndResponsePermit, SubmitError> {
307 let worker = self.try_reserve(session_id)?;
309 let response = self.try_reserve_response()?;
315 Ok(WorkerAndResponsePermit { worker, response })
316 }
317
318 pub fn try_submit(&self, job: WorkerJob) -> Result<(), SubmitError> {
322 self.try_submit_labeled(job, DropStage::WorkerQueue)
323 }
324
325 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 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 pub async fn wait_idle(&self) {
413 loop {
414 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 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
461pub 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
472pub 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
483pub 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
503pub 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 const MAX_SHARDS: usize = 16;
523 let capacity = capacity.max(1);
524 let shard_count = MAX_SHARDS;
525 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 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 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#[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 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 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 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 drop(capacity_permit);
731 if let Some(completion) = completion {
735 let _ = completion.send(outcome);
736 }
737}
738
739struct 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
762struct 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 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 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 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 &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 &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 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 &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 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 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 #[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 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(®istry, &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(®istry, &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 let outcome = finalizer
1882 .recover_spooled_sessions(®istry, &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 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(®istry, &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(®istry, &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 #[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 let first = worker.try_reserve_pair("s").expect("first pair");
2088 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 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 #[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 let permits = worker
2139 .try_reserve_pair("mid-stream-response")
2140 .expect("initial pair");
2141 drop(permits.worker);
2144 let hog_a = worker.try_reserve("hog-a").expect("hog-a");
2148 let hog_b = worker.try_reserve("hog-b").expect("hog-b");
2149 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 #[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 #[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(¬ify);
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(¬ify);
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 #[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 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 #[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 #[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 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 #[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 #[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 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 #[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 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}