1#[cfg(feature = "distributed")]
4use crate::boost::PriorityBooster;
5use crate::config::LaneConfig;
6use crate::dlq::{DeadLetter, DeadLetterQueue};
7use crate::error::{LaneError, Result};
8use crate::event::{events, EventEmitter, EventStream, LaneEvent};
9#[cfg(feature = "distributed")]
10use crate::ratelimit::RateLimiter;
11use crate::retry::RetryPolicy;
12use crate::storage::{Storage, StoredCommand};
13#[cfg(feature = "telemetry")]
14use crate::telemetry;
15use async_trait::async_trait;
16use chrono::Utc;
17use serde::{Deserialize, Serialize};
18use std::collections::{HashMap, VecDeque};
19use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
20use std::sync::Arc;
21#[cfg(feature = "telemetry")]
22use std::time::Instant;
23use tokio::sync::{Mutex, Semaphore};
24use uuid::Uuid;
25
26pub type LaneId = String;
28
29pub type CommandId = String;
31
32pub type Priority = u8;
34
35pub mod priorities {
37 use super::Priority;
38
39 pub const SYSTEM: Priority = 0;
40 pub const CONTROL: Priority = 1;
41 pub const QUERY: Priority = 2;
42 pub const SESSION: Priority = 3;
43 pub const SKILL: Priority = 4;
44 pub const PROMPT: Priority = 5;
45}
46
47#[async_trait]
49pub trait Command: Send + Sync {
50 async fn execute(&self) -> Result<serde_json::Value>;
52
53 fn command_type(&self) -> &str;
55}
56
57pub struct JsonCommand {
62 command_type: String,
63 payload: serde_json::Value,
64}
65
66impl JsonCommand {
67 pub fn new(command_type: impl Into<String>, payload: serde_json::Value) -> Self {
69 Self {
70 command_type: command_type.into(),
71 payload,
72 }
73 }
74}
75
76#[async_trait]
77impl Command for JsonCommand {
78 async fn execute(&self) -> Result<serde_json::Value> {
79 Ok(self.payload.clone())
80 }
81
82 fn command_type(&self) -> &str {
83 &self.command_type
84 }
85}
86
87struct CommandWrapper {
89 id: CommandId,
90 command: Arc<dyn Command>,
91 result_tx: Option<tokio::sync::oneshot::Sender<Result<serde_json::Value>>>,
92 timeout: Option<std::time::Duration>,
93 retry_policy: RetryPolicy,
94 attempt: u32,
95 lane_id: LaneId,
96 command_type: String,
97 _work_permit: Option<WorkPermit>,
99 #[cfg(feature = "distributed")]
101 enqueue_time: chrono::DateTime<chrono::Utc>,
102}
103
104#[allow(dead_code)]
106struct LaneState {
107 config: LaneConfig,
109
110 priority: Priority,
112
113 pending: VecDeque<CommandWrapper>,
115
116 active: usize,
118
119 delayed_retries: usize,
121
122 semaphore: Arc<Semaphore>,
124
125 is_pressured: bool,
127}
128
129impl LaneState {
130 fn new(config: LaneConfig, priority: Priority) -> Self {
131 let semaphore = Arc::new(Semaphore::new(config.max_concurrency));
132 Self {
133 config,
134 priority,
135 pending: VecDeque::new(),
136 active: 0,
137 delayed_retries: 0,
138 semaphore,
139 is_pressured: false,
140 }
141 }
142
143 fn has_capacity(&self) -> bool {
144 self.active < self.config.max_concurrency
145 }
146
147 fn has_pending(&self) -> bool {
148 !self.pending.is_empty()
149 }
150}
151
152pub struct Lane {
154 id: LaneId,
155 state: Arc<Mutex<LaneState>>,
156 storage: Option<Arc<dyn Storage>>,
157 #[cfg(feature = "distributed")]
159 rate_limiter: RateLimiter,
160 #[cfg(feature = "distributed")]
162 booster: Option<PriorityBooster>,
163}
164
165impl Lane {
166 pub fn new(id: impl Into<String>, config: LaneConfig, priority: Priority) -> Self {
168 #[cfg(feature = "distributed")]
169 let rate_limiter = config
170 .rate_limit
171 .as_ref()
172 .map(RateLimiter::token_bucket)
173 .unwrap_or_default();
174 #[cfg(feature = "distributed")]
175 let booster = config
176 .priority_boost
177 .as_ref()
178 .map(|pb| PriorityBooster::new(pb.clone()));
179 Self {
180 id: id.into(),
181 state: Arc::new(Mutex::new(LaneState::new(config, priority))),
182 storage: None,
183 #[cfg(feature = "distributed")]
184 rate_limiter,
185 #[cfg(feature = "distributed")]
186 booster,
187 }
188 }
189
190 pub fn with_storage(
192 id: impl Into<String>,
193 config: LaneConfig,
194 priority: Priority,
195 storage: Arc<dyn Storage>,
196 ) -> Self {
197 #[cfg(feature = "distributed")]
198 let rate_limiter = config
199 .rate_limit
200 .as_ref()
201 .map(RateLimiter::token_bucket)
202 .unwrap_or_default();
203 #[cfg(feature = "distributed")]
204 let booster = config
205 .priority_boost
206 .as_ref()
207 .map(|pb| PriorityBooster::new(pb.clone()));
208 Self {
209 id: id.into(),
210 state: Arc::new(Mutex::new(LaneState::new(config, priority))),
211 storage: Some(storage),
212 #[cfg(feature = "distributed")]
213 rate_limiter,
214 #[cfg(feature = "distributed")]
215 booster,
216 }
217 }
218
219 pub fn id(&self) -> &str {
221 &self.id
222 }
223
224 pub async fn priority(&self) -> Priority {
226 self.state.lock().await.priority
227 }
228
229 pub async fn effective_priority(&self) -> Priority {
231 let state = self.state.lock().await;
232 let base = state.priority;
233 #[cfg(feature = "distributed")]
234 if let Some(booster) = &self.booster {
235 if let Some(front) = state.pending.front() {
236 return booster.calculate_priority(base, front.enqueue_time);
237 }
238 }
239 base
240 }
241
242 pub async fn enqueue(
244 &self,
245 command: Box<dyn Command>,
246 ) -> tokio::sync::oneshot::Receiver<Result<serde_json::Value>> {
247 self.enqueue_with_permit(command, None).await
248 }
249
250 async fn enqueue_with_permit(
252 &self,
253 command: Box<dyn Command>,
254 work_permit: Option<WorkPermit>,
255 ) -> tokio::sync::oneshot::Receiver<Result<serde_json::Value>> {
256 let (tx, rx) = tokio::sync::oneshot::channel();
257 let state = self.state.lock().await;
258 let timeout = state.config.default_timeout;
259 let retry_policy = state.config.retry_policy.clone();
260 drop(state);
261
262 let command_id = Uuid::new_v4().to_string();
263 let command_type = command.command_type().to_string();
264 let wrapper = CommandWrapper {
265 id: command_id.clone(),
266 command: Arc::from(command),
267 result_tx: Some(tx),
268 timeout,
269 retry_policy,
270 attempt: 0,
271 lane_id: self.id.clone(),
272 command_type: command_type.clone(),
273 _work_permit: work_permit,
274 #[cfg(feature = "distributed")]
275 enqueue_time: Utc::now(),
276 };
277
278 if let Some(storage) = &self.storage {
280 let stored_cmd = StoredCommand {
281 id: command_id,
282 command_type,
283 lane_id: self.id.clone(),
284 payload: serde_json::json!({}), retry_count: 0,
286 created_at: Utc::now(),
287 last_attempt_at: None,
288 };
289 let _ = storage.save_command(stored_cmd).await;
291 }
292
293 let mut state = self.state.lock().await;
294 state.pending.push_back(wrapper);
295
296 rx
297 }
298
299 async fn retry_command(&self, mut wrapper: CommandWrapper, delay: std::time::Duration) {
301 wrapper.attempt += 1;
302
303 {
305 let mut state = self.state.lock().await;
306 state.delayed_retries += 1;
307 }
308
309 let state_clone = Arc::clone(&self.state);
311 tokio::spawn(async move {
312 tokio::time::sleep(delay).await;
313 let mut state = state_clone.lock().await;
314 state.delayed_retries = state.delayed_retries.saturating_sub(1);
315 state.pending.push_back(wrapper);
316 });
317 }
318
319 async fn try_dequeue(&self) -> Option<CommandWrapper> {
321 #[cfg(feature = "distributed")]
323 if !self.rate_limiter.try_acquire().await {
324 return None;
325 }
326 let mut state = self.state.lock().await;
327 if state.has_capacity() && state.has_pending() {
328 state.active += 1;
329 state.pending.pop_front()
330 } else {
331 None
332 }
333 }
334
335 async fn mark_completed(&self) {
337 let mut state = self.state.lock().await;
338 state.active = state.active.saturating_sub(1);
339 }
340
341 pub async fn status(&self) -> LaneStatus {
343 let state = self.state.lock().await;
344 LaneStatus {
345 pending: state.pending.len() + state.delayed_retries,
346 active: state.active,
347 min: state.config.min_concurrency,
348 max: state.config.max_concurrency,
349 }
350 }
351
352 async fn is_idle(&self) -> bool {
353 let state = self.state.lock().await;
354 state.pending.is_empty() && state.active == 0 && state.delayed_retries == 0
355 }
356
357 async fn check_pressure(&self) -> Option<&'static str> {
363 let mut state = self.state.lock().await;
364 let threshold = match state.config.pressure_threshold {
365 Some(t) => t,
366 None => return None,
367 };
368 let pending = state.pending.len();
369 let was_pressured = state.is_pressured;
370 if pending >= threshold && !was_pressured {
371 state.is_pressured = true;
372 Some(events::QUEUE_LANE_PRESSURE)
373 } else if pending == 0 && was_pressured {
374 state.is_pressured = false;
375 Some(events::QUEUE_LANE_IDLE)
376 } else {
377 None
378 }
379 }
380}
381
382#[derive(Debug, Clone, Serialize, Deserialize)]
384pub struct LaneStatus {
385 pub pending: usize,
386 pub active: usize,
387 pub min: usize,
388 pub max: usize,
389}
390
391struct QueueLifecycle {
393 admission_gate: Mutex<()>,
395 is_shutting_down: AtomicBool,
396 outstanding: AtomicUsize,
397 changed: tokio::sync::Notify,
398}
399
400impl QueueLifecycle {
401 fn new() -> Self {
402 Self {
403 admission_gate: Mutex::new(()),
404 is_shutting_down: AtomicBool::new(false),
405 outstanding: AtomicUsize::new(0),
406 changed: tokio::sync::Notify::new(),
407 }
408 }
409
410 async fn admit(self: &Arc<Self>) -> Result<WorkPermit> {
411 let _gate = self.admission_gate.lock().await;
412 if self.is_shutting_down.load(Ordering::Acquire) {
413 return Err(LaneError::ShutdownInProgress);
414 }
415
416 self.outstanding.fetch_add(1, Ordering::AcqRel);
417 Ok(WorkPermit {
418 lifecycle: Arc::clone(self),
419 })
420 }
421
422 async fn begin_shutdown(&self) -> bool {
423 let _gate = self.admission_gate.lock().await;
424 let started = !self.is_shutting_down.swap(true, Ordering::AcqRel);
425 drop(_gate);
426 self.changed.notify_waiters();
427 started
428 }
429
430 fn is_shutting_down(&self) -> bool {
431 self.is_shutting_down.load(Ordering::Acquire)
432 }
433
434 fn is_idle(&self) -> bool {
435 self.outstanding.load(Ordering::Acquire) == 0
436 }
437
438 fn notify_changed(&self) {
439 self.changed.notify_waiters();
440 }
441}
442
443struct WorkPermit {
445 lifecycle: Arc<QueueLifecycle>,
446}
447
448impl Drop for WorkPermit {
449 fn drop(&mut self) {
450 self.lifecycle.outstanding.fetch_sub(1, Ordering::AcqRel);
451 self.lifecycle.notify_changed();
452 }
453}
454
455#[allow(dead_code)]
457pub struct CommandQueue {
458 lanes: Arc<Mutex<HashMap<LaneId, Arc<Lane>>>>,
459 event_emitter: EventEmitter,
460 dlq: Option<DeadLetterQueue>,
461 storage: Option<Arc<dyn Storage>>,
462 lifecycle: Arc<QueueLifecycle>,
463}
464
465impl CommandQueue {
466 pub fn new(event_emitter: EventEmitter) -> Self {
468 Self {
469 lanes: Arc::new(Mutex::new(HashMap::new())),
470 event_emitter,
471 dlq: None,
472 storage: None,
473 lifecycle: Arc::new(QueueLifecycle::new()),
474 }
475 }
476
477 pub fn with_dlq(event_emitter: EventEmitter, dlq_size: usize) -> Self {
479 Self {
480 lanes: Arc::new(Mutex::new(HashMap::new())),
481 event_emitter,
482 dlq: Some(DeadLetterQueue::new(dlq_size)),
483 storage: None,
484 lifecycle: Arc::new(QueueLifecycle::new()),
485 }
486 }
487
488 pub fn with_storage(event_emitter: EventEmitter, storage: Arc<dyn Storage>) -> Self {
490 Self {
491 lanes: Arc::new(Mutex::new(HashMap::new())),
492 event_emitter,
493 dlq: None,
494 storage: Some(storage),
495 lifecycle: Arc::new(QueueLifecycle::new()),
496 }
497 }
498
499 pub fn with_dlq_and_storage(
501 event_emitter: EventEmitter,
502 dlq_size: usize,
503 storage: Arc<dyn Storage>,
504 ) -> Self {
505 Self {
506 lanes: Arc::new(Mutex::new(HashMap::new())),
507 event_emitter,
508 dlq: Some(DeadLetterQueue::new(dlq_size)),
509 storage: Some(storage),
510 lifecycle: Arc::new(QueueLifecycle::new()),
511 }
512 }
513
514 pub fn storage(&self) -> Option<&Arc<dyn Storage>> {
516 self.storage.as_ref()
517 }
518
519 pub fn dlq(&self) -> Option<&DeadLetterQueue> {
521 self.dlq.as_ref()
522 }
523
524 pub fn is_shutting_down(&self) -> bool {
526 self.lifecycle.is_shutting_down()
527 }
528
529 pub async fn shutdown(&self) {
531 if self.lifecycle.begin_shutdown().await {
532 self.event_emitter
533 .emit(LaneEvent::empty(events::QUEUE_SHUTDOWN_STARTED));
534 }
535 }
536
537 pub async fn drain(&self, timeout: std::time::Duration) -> Result<()> {
539 match tokio::time::timeout(timeout, self.wait_for_idle()).await {
540 Ok(()) => Ok(()),
541 Err(_) => Err(LaneError::Timeout(timeout)),
542 }
543 }
544
545 pub(crate) async fn wait_for_idle(&self) {
546 loop {
547 let changed = self.lifecycle.changed.notified();
548 if self.is_idle().await {
549 return;
550 }
551
552 tokio::select! {
553 _ = changed => {}
554 _ = tokio::time::sleep(std::time::Duration::from_millis(10)) => {}
555 }
556 }
557 }
558
559 async fn is_idle(&self) -> bool {
560 if !self.lifecycle.is_idle() {
561 return false;
562 }
563
564 let lanes = {
565 let lanes = self.lanes.lock().await;
566 lanes.values().cloned().collect::<Vec<_>>()
567 };
568
569 for lane in lanes {
570 if !lane.is_idle().await {
571 return false;
572 }
573 }
574
575 true
576 }
577
578 pub async fn register_lane(&self, lane: Arc<Lane>) {
580 let mut lanes = self.lanes.lock().await;
581 lanes.insert(lane.id().to_string(), lane);
582 }
583
584 pub async fn submit(
586 &self,
587 lane_id: &str,
588 command: Box<dyn Command>,
589 ) -> Result<tokio::sync::oneshot::Receiver<Result<serde_json::Value>>> {
590 let work_permit = self.lifecycle.admit().await?;
591 let lane = {
592 let lanes = self.lanes.lock().await;
593 Arc::clone(
594 lanes
595 .get(lane_id)
596 .ok_or_else(|| LaneError::LaneNotFound(lane_id.to_string()))?,
597 )
598 };
599 let rx = lane.enqueue_with_permit(command, Some(work_permit)).await;
600 self.lifecycle.notify_changed();
601
602 self.event_emitter.emit(LaneEvent::with_map(
603 events::QUEUE_COMMAND_SUBMITTED,
604 HashMap::from([("lane_id".to_string(), serde_json::json!(lane_id))]),
605 ));
606
607 Ok(rx)
608 }
609
610 pub async fn start_scheduler(self: Arc<Self>) {
612 std::mem::drop(self.spawn_scheduler());
613 }
614
615 pub(crate) fn spawn_scheduler(self: Arc<Self>) -> tokio::task::JoinHandle<()> {
616 tokio::spawn(async move {
617 loop {
618 self.schedule_next().await;
619 if self.is_shutting_down() && self.is_idle().await {
620 self.event_emitter
621 .emit(LaneEvent::empty(events::QUEUE_SHUTDOWN_COMPLETE));
622 return;
623 }
624
625 tokio::select! {
626 _ = self.lifecycle.changed.notified() => {}
627 _ = tokio::time::sleep(tokio::time::Duration::from_millis(10)) => {}
628 }
629 }
630 })
631 }
632
633 async fn schedule_next(&self) {
635 let lanes = {
638 let lanes = self.lanes.lock().await;
639 lanes
640 .iter()
641 .map(|(id, lane)| (id.clone(), Arc::clone(lane)))
642 .collect::<Vec<_>>()
643 };
644
645 for (lane_id, lane) in &lanes {
647 if let Some(event_key) = lane.check_pressure().await {
648 self.event_emitter.emit(LaneEvent::with_map(
649 event_key,
650 HashMap::from([("lane_id".to_string(), serde_json::json!(lane_id))]),
651 ));
652 }
653 }
654
655 let mut lane_priorities = Vec::new();
656 for (_, lane) in &lanes {
657 let priority = lane.effective_priority().await;
658 lane_priorities.push((priority, Arc::clone(lane)));
659 }
660
661 lane_priorities.sort_by_key(|(priority, _)| *priority);
663
664 for (_, lane) in lane_priorities {
665 if let Some(mut wrapper) = lane.try_dequeue().await {
666 let lane_clone = Arc::clone(&lane);
667 let timeout = wrapper.timeout;
668 let retry_policy = wrapper.retry_policy.clone();
669 let attempt = wrapper.attempt;
670 let dlq = self.dlq.clone();
671 let command_id = wrapper.id.clone();
672 let command_type = wrapper.command_type.clone();
673 let lane_id = wrapper.lane_id.clone();
674 let storage = lane.storage.clone();
675 let event_emitter = self.event_emitter.clone();
676
677 event_emitter.emit(LaneEvent::with_map(
678 events::QUEUE_COMMAND_STARTED,
679 HashMap::from([
680 ("lane_id".to_string(), serde_json::json!(lane_id)),
681 ("command_id".to_string(), serde_json::json!(command_id)),
682 ("command_type".to_string(), serde_json::json!(command_type)),
683 ]),
684 ));
685
686 tokio::spawn(async move {
687 #[cfg(feature = "telemetry")]
688 let exec_start = Instant::now();
689
690 let result = match timeout {
691 Some(dur) => {
692 match tokio::time::timeout(dur, wrapper.command.execute()).await {
693 Ok(r) => r,
694 Err(_) => Err(LaneError::Timeout(dur)),
695 }
696 }
697 None => wrapper.command.execute().await,
698 };
699
700 match result {
701 Ok(value) => {
702 if let Some(storage) = &storage {
703 let _ = storage.remove_command(&command_id).await;
704 }
705
706 #[cfg(feature = "telemetry")]
707 telemetry::record_complete(
708 &lane_id,
709 exec_start.elapsed().as_secs_f64(),
710 );
711
712 event_emitter.emit(LaneEvent::with_map(
713 events::QUEUE_COMMAND_COMPLETED,
714 HashMap::from([
715 ("lane_id".to_string(), serde_json::json!(lane_id)),
716 ("command_id".to_string(), serde_json::json!(command_id)),
717 ]),
718 ));
719
720 if let Some(tx) = wrapper.result_tx.take() {
721 let _ = tx.send(Ok(value));
722 }
723 lane_clone.mark_completed().await;
724 }
725 Err(err) => {
726 if retry_policy.should_retry(attempt) {
727 let delay = retry_policy.delay_for_attempt(attempt + 1);
728
729 tracing::info!(
730 command_id = %command_id,
731 retry_attempt = attempt + 1,
732 "a3s.lane.retry: retrying command"
733 );
734
735 lane_clone.retry_command(wrapper, delay).await;
736
737 event_emitter.emit(LaneEvent::with_map(
738 events::QUEUE_COMMAND_RETRY,
739 HashMap::from([
740 ("lane_id".to_string(), serde_json::json!(lane_id)),
741 ("command_id".to_string(), serde_json::json!(command_id)),
742 ("attempt".to_string(), serde_json::json!(attempt + 1)),
743 ]),
744 ));
745 lane_clone.mark_completed().await;
746 } else {
747 #[cfg(feature = "telemetry")]
748 telemetry::record_failure(&lane_id);
749
750 if let Some(storage) = &storage {
751 let _ = storage.remove_command(&command_id).await;
752 }
753
754 let error_msg = err.to_string();
755 let is_timeout = matches!(err, LaneError::Timeout(_));
756
757 if let Some(dlq) = dlq {
758 let dead_letter = DeadLetter {
759 command_id: command_id.clone(),
760 command_type: command_type.clone(),
761 lane_id: lane_id.clone(),
762 error: error_msg.clone(),
763 attempts: attempt + 1,
764 failed_at: Utc::now(),
765 };
766 dlq.push(dead_letter).await;
767
768 event_emitter.emit(LaneEvent::with_map(
769 events::QUEUE_COMMAND_DEAD_LETTERED,
770 HashMap::from([
771 ("lane_id".to_string(), serde_json::json!(lane_id)),
772 (
773 "command_id".to_string(),
774 serde_json::json!(command_id),
775 ),
776 (
777 "command_type".to_string(),
778 serde_json::json!(command_type),
779 ),
780 ]),
781 ));
782 }
783
784 event_emitter.emit(LaneEvent::with_map(
785 if is_timeout {
786 events::QUEUE_COMMAND_TIMEOUT
787 } else {
788 events::QUEUE_COMMAND_FAILED
789 },
790 HashMap::from([
791 ("lane_id".to_string(), serde_json::json!(lane_id)),
792 ("command_id".to_string(), serde_json::json!(command_id)),
793 ("error".to_string(), serde_json::json!(error_msg)),
794 ]),
795 ));
796
797 if let Some(tx) = wrapper.result_tx.take() {
798 let _ = tx.send(Err(err));
799 }
800 lane_clone.mark_completed().await;
801 }
802 }
803 }
804 });
805 break;
806 }
807 }
808 }
809
810 pub fn subscribe_stream(&self) -> EventStream {
812 self.event_emitter.subscribe_stream()
813 }
814
815 pub fn subscribe_filtered(
817 &self,
818 filter: impl Fn(&LaneEvent) -> bool + Send + Sync + 'static,
819 ) -> EventStream {
820 self.event_emitter.subscribe_filtered(filter)
821 }
822
823 pub async fn status(&self) -> HashMap<LaneId, LaneStatus> {
825 let lanes = {
826 let lanes = self.lanes.lock().await;
827 lanes
828 .iter()
829 .map(|(id, lane)| (id.clone(), Arc::clone(lane)))
830 .collect::<Vec<_>>()
831 };
832 let mut status = HashMap::new();
833
834 for (id, lane) in lanes {
835 status.insert(id, lane.status().await);
836 }
837
838 status
839 }
840}
841
842pub mod lane_ids {
844 pub const SYSTEM: &str = "system";
845 pub const CONTROL: &str = "control";
846 pub const QUERY: &str = "query";
847 pub const SESSION: &str = "session";
848 pub const SKILL: &str = "skill";
849 pub const PROMPT: &str = "prompt";
850}
851
852#[cfg(test)]
853mod tests {
854 use super::*;
855 use crate::storage::StoredDeadLetter;
856
857 #[derive(Default)]
858 struct BlockingStorage {
859 save_started: tokio::sync::Notify,
860 release_save: tokio::sync::Notify,
861 }
862
863 #[async_trait]
864 impl Storage for BlockingStorage {
865 async fn save_command(&self, _command: StoredCommand) -> Result<()> {
866 self.save_started.notify_one();
867 self.release_save.notified().await;
868 Ok(())
869 }
870
871 async fn load_commands(&self) -> Result<Vec<StoredCommand>> {
872 Ok(Vec::new())
873 }
874
875 async fn remove_command(&self, _id: &str) -> Result<()> {
876 Ok(())
877 }
878
879 async fn save_dead_letter(&self, _letter: StoredDeadLetter) -> Result<()> {
880 Ok(())
881 }
882
883 async fn load_dead_letters(&self) -> Result<Vec<StoredDeadLetter>> {
884 Ok(Vec::new())
885 }
886
887 async fn clear_dead_letters(&self) -> Result<()> {
888 Ok(())
889 }
890
891 async fn clear_all(&self) -> Result<()> {
892 Ok(())
893 }
894 }
895
896 struct TestCommand {
898 result: serde_json::Value,
899 delay_ms: Option<u64>,
900 }
901
902 impl TestCommand {
903 fn new(result: serde_json::Value) -> Self {
904 Self {
905 result,
906 delay_ms: None,
907 }
908 }
909
910 fn with_delay(result: serde_json::Value, delay_ms: u64) -> Self {
911 Self {
912 result,
913 delay_ms: Some(delay_ms),
914 }
915 }
916 }
917
918 #[async_trait]
919 impl Command for TestCommand {
920 async fn execute(&self) -> Result<serde_json::Value> {
921 if let Some(delay) = self.delay_ms {
922 tokio::time::sleep(tokio::time::Duration::from_millis(delay)).await;
923 }
924 Ok(self.result.clone())
925 }
926
927 fn command_type(&self) -> &str {
928 "test"
929 }
930 }
931
932 struct FailingCommand {
934 message: String,
935 }
936
937 #[async_trait]
938 impl Command for FailingCommand {
939 async fn execute(&self) -> Result<serde_json::Value> {
940 Err(LaneError::Other(self.message.clone()))
941 }
942
943 fn command_type(&self) -> &str {
944 "failing"
945 }
946 }
947
948 #[test]
949 fn test_priorities() {
950 assert_eq!(priorities::SYSTEM, 0);
951 assert_eq!(priorities::CONTROL, 1);
952 assert_eq!(priorities::QUERY, 2);
953 assert_eq!(priorities::SESSION, 3);
954 assert_eq!(priorities::SKILL, 4);
955 assert_eq!(priorities::PROMPT, 5);
956
957 const _: () = {
960 assert!(priorities::SYSTEM < priorities::CONTROL);
961 assert!(priorities::CONTROL < priorities::QUERY);
962 assert!(priorities::QUERY < priorities::SESSION);
963 assert!(priorities::SESSION < priorities::SKILL);
964 assert!(priorities::SKILL < priorities::PROMPT);
965 };
966 }
967
968 #[test]
969 fn test_lane_ids() {
970 assert_eq!(lane_ids::SYSTEM, "system");
971 assert_eq!(lane_ids::CONTROL, "control");
972 assert_eq!(lane_ids::QUERY, "query");
973 assert_eq!(lane_ids::SESSION, "session");
974 assert_eq!(lane_ids::SKILL, "skill");
975 assert_eq!(lane_ids::PROMPT, "prompt");
976 }
977
978 #[test]
979 fn test_lane_new() {
980 let config = LaneConfig::new(1, 4);
981 let lane = Lane::new("test-lane", config, priorities::QUERY);
982
983 assert_eq!(lane.id(), "test-lane");
984 }
985
986 #[tokio::test]
987 async fn test_lane_priority() {
988 let config = LaneConfig::new(1, 4);
989 let lane = Lane::new("test", config, priorities::SESSION);
990
991 assert_eq!(lane.priority().await, priorities::SESSION);
992 }
993
994 #[tokio::test]
995 async fn test_lane_status_initial() {
996 let config = LaneConfig::new(2, 8);
997 let lane = Lane::new("test", config, priorities::QUERY);
998
999 let status = lane.status().await;
1000 assert_eq!(status.pending, 0);
1001 assert_eq!(status.active, 0);
1002 assert_eq!(status.min, 2);
1003 assert_eq!(status.max, 8);
1004 }
1005
1006 #[tokio::test]
1007 async fn test_lane_enqueue() {
1008 let config = LaneConfig::new(1, 4);
1009 let lane = Lane::new("test", config, priorities::QUERY);
1010
1011 let cmd = Box::new(TestCommand::new(serde_json::json!({"result": "ok"})));
1012 let _rx = lane.enqueue(cmd).await;
1013
1014 let status = lane.status().await;
1015 assert_eq!(status.pending, 1);
1016 }
1017
1018 #[tokio::test]
1019 async fn test_lane_status_serialization() {
1020 let status = LaneStatus {
1021 pending: 5,
1022 active: 2,
1023 min: 1,
1024 max: 8,
1025 };
1026
1027 let json = serde_json::to_string(&status).unwrap();
1028 assert!(json.contains("\"pending\":5"));
1029 assert!(json.contains("\"active\":2"));
1030 assert!(json.contains("\"min\":1"));
1031 assert!(json.contains("\"max\":8"));
1032
1033 let parsed: LaneStatus = serde_json::from_str(&json).unwrap();
1034 assert_eq!(parsed.pending, 5);
1035 assert_eq!(parsed.active, 2);
1036 }
1037
1038 #[tokio::test]
1039 async fn test_command_queue_new() {
1040 let emitter = EventEmitter::new(100);
1041 let queue = CommandQueue::new(emitter);
1042
1043 let status = queue.status().await;
1044 assert!(status.is_empty());
1045 }
1046
1047 #[tokio::test]
1048 async fn test_command_queue_register_lane() {
1049 let emitter = EventEmitter::new(100);
1050 let queue = CommandQueue::new(emitter);
1051
1052 let config = LaneConfig::new(1, 4);
1053 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1054
1055 queue.register_lane(lane).await;
1056
1057 let status = queue.status().await;
1058 assert!(status.contains_key("test-lane"));
1059 }
1060
1061 #[tokio::test]
1062 async fn test_command_queue_submit() {
1063 let emitter = EventEmitter::new(100);
1064 let queue = CommandQueue::new(emitter);
1065
1066 let config = LaneConfig::new(1, 4);
1067 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1068 queue.register_lane(lane).await;
1069
1070 let cmd = Box::new(TestCommand::new(serde_json::json!({"status": "ok"})));
1071 let result = queue.submit("test-lane", cmd).await;
1072
1073 assert!(result.is_ok());
1074 }
1075
1076 #[tokio::test]
1077 async fn test_command_queue_submit_unknown_lane() {
1078 let emitter = EventEmitter::new(100);
1079 let queue = CommandQueue::new(emitter);
1080
1081 let cmd = Box::new(TestCommand::new(serde_json::json!({})));
1082 let result = queue.submit("nonexistent", cmd).await;
1083
1084 assert!(result.is_err());
1085 if let Err(LaneError::LaneNotFound(id)) = result {
1086 assert_eq!(id, "nonexistent");
1087 } else {
1088 panic!("Expected LaneNotFound error");
1089 }
1090 }
1091
1092 #[tokio::test]
1093 async fn test_command_queue_multiple_lanes() {
1094 let emitter = EventEmitter::new(100);
1095 let queue = CommandQueue::new(emitter);
1096
1097 let configs = vec![
1099 ("system", priorities::SYSTEM, 1),
1100 ("control", priorities::CONTROL, 8),
1101 ("query", priorities::QUERY, 4),
1102 ];
1103
1104 for (id, priority, max) in configs {
1105 let config = LaneConfig::new(1, max);
1106 let lane = Arc::new(Lane::new(id, config, priority));
1107 queue.register_lane(lane).await;
1108 }
1109
1110 let status = queue.status().await;
1111 assert_eq!(status.len(), 3);
1112 assert!(status.contains_key("system"));
1113 assert!(status.contains_key("control"));
1114 assert!(status.contains_key("query"));
1115 }
1116
1117 #[tokio::test]
1118 async fn test_command_queue_status() {
1119 let emitter = EventEmitter::new(100);
1120 let queue = CommandQueue::new(emitter);
1121
1122 let config = LaneConfig::new(2, 16);
1123 let lane = Arc::new(Lane::new("query", config, priorities::QUERY));
1124 queue.register_lane(lane).await;
1125
1126 let status = queue.status().await;
1127 let lane_status = status.get("query").unwrap();
1128
1129 assert_eq!(lane_status.min, 2);
1130 assert_eq!(lane_status.max, 16);
1131 assert_eq!(lane_status.pending, 0);
1132 assert_eq!(lane_status.active, 0);
1133 }
1134
1135 #[test]
1136 fn test_lane_state_has_capacity() {
1137 let config = LaneConfig::new(1, 2);
1138 let mut state = LaneState::new(config, priorities::QUERY);
1139
1140 assert!(state.has_capacity());
1141 state.active = 1;
1142 assert!(state.has_capacity());
1143 state.active = 2;
1144 assert!(!state.has_capacity());
1145 }
1146
1147 #[tokio::test]
1148 async fn test_lane_state_has_pending() {
1149 let config = LaneConfig::new(1, 4);
1150 let state = LaneState::new(config, priorities::QUERY);
1151
1152 assert!(!state.has_pending());
1153 assert_eq!(state.pending.len(), 0);
1154 }
1155
1156 #[test]
1157 fn test_lane_status_debug() {
1158 let status = LaneStatus {
1159 pending: 3,
1160 active: 1,
1161 min: 1,
1162 max: 4,
1163 };
1164
1165 let debug_str = format!("{:?}", status);
1166 assert!(debug_str.contains("LaneStatus"));
1167 assert!(debug_str.contains("pending"));
1168 assert!(debug_str.contains("active"));
1169 }
1170
1171 #[test]
1172 fn test_lane_status_clone() {
1173 let status = LaneStatus {
1174 pending: 5,
1175 active: 2,
1176 min: 1,
1177 max: 8,
1178 };
1179
1180 let cloned = status.clone();
1181 assert_eq!(cloned.pending, 5);
1182 assert_eq!(cloned.active, 2);
1183 assert_eq!(cloned.min, 1);
1184 assert_eq!(cloned.max, 8);
1185 }
1186
1187 #[tokio::test]
1188 async fn test_command_execution() {
1189 let cmd = TestCommand::new(serde_json::json!({"value": 42}));
1190 let result = cmd.execute().await;
1191
1192 assert!(result.is_ok());
1193 assert_eq!(result.unwrap(), serde_json::json!({"value": 42}));
1194 }
1195
1196 #[tokio::test]
1197 async fn test_command_type() {
1198 let cmd = TestCommand::new(serde_json::json!({}));
1199 assert_eq!(cmd.command_type(), "test");
1200
1201 let failing = FailingCommand {
1202 message: "error".to_string(),
1203 };
1204 assert_eq!(failing.command_type(), "failing");
1205 }
1206
1207 #[tokio::test]
1208 async fn test_failing_command() {
1209 let cmd = FailingCommand {
1210 message: "Something went wrong".to_string(),
1211 };
1212
1213 let result = cmd.execute().await;
1214 assert!(result.is_err());
1215
1216 if let Err(LaneError::Other(msg)) = result {
1217 assert_eq!(msg, "Something went wrong");
1218 } else {
1219 panic!("Expected Other error");
1220 }
1221 }
1222
1223 #[tokio::test]
1224 async fn test_command_with_delay() {
1225 let cmd = TestCommand::with_delay(serde_json::json!({"delayed": true}), 10);
1226
1227 let start = std::time::Instant::now();
1228 let result = cmd.execute().await;
1229 let elapsed = start.elapsed();
1230
1231 assert!(result.is_ok());
1232 assert!(elapsed.as_millis() >= 10);
1233 }
1234
1235 #[test]
1236 fn test_priority_type() {
1237 let p: Priority = 5;
1238 assert_eq!(p, 5u8);
1239 }
1240
1241 #[test]
1242 fn test_lane_id_type() {
1243 let id: LaneId = "test-lane".to_string();
1244 assert_eq!(id, "test-lane");
1245 }
1246
1247 #[test]
1248 fn test_command_id_type() {
1249 let id: CommandId = "cmd-123".to_string();
1250 assert_eq!(id, "cmd-123");
1251 }
1252
1253 #[tokio::test]
1254 async fn test_command_timeout() {
1255 let emitter = EventEmitter::new(100);
1256 let queue = Arc::new(CommandQueue::new(emitter));
1257
1258 let config = LaneConfig::new(1, 4).with_timeout(std::time::Duration::from_millis(50));
1259 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1260 queue.register_lane(lane).await;
1261
1262 Arc::clone(&queue).start_scheduler().await;
1264
1265 let cmd = Box::new(TestCommand::with_delay(
1267 serde_json::json!({"result": "ok"}),
1268 200,
1269 ));
1270 let rx = queue.submit("test-lane", cmd).await.unwrap();
1271
1272 let result = tokio::time::timeout(std::time::Duration::from_secs(1), rx)
1274 .await
1275 .expect("Timeout waiting for result")
1276 .expect("Channel closed");
1277
1278 assert!(result.is_err());
1280 if let Err(LaneError::Timeout(dur)) = result {
1281 assert_eq!(dur, std::time::Duration::from_millis(50));
1282 } else {
1283 panic!("Expected Timeout error");
1284 }
1285 }
1286
1287 #[tokio::test]
1288 async fn test_command_no_timeout() {
1289 let emitter = EventEmitter::new(100);
1290 let queue = Arc::new(CommandQueue::new(emitter));
1291
1292 let config = LaneConfig::new(1, 4);
1293 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1294 queue.register_lane(lane).await;
1295
1296 Arc::clone(&queue).start_scheduler().await;
1297
1298 let cmd = Box::new(TestCommand::with_delay(
1300 serde_json::json!({"result": "ok"}),
1301 50,
1302 ));
1303 let rx = queue.submit("test-lane", cmd).await.unwrap();
1304
1305 let result = tokio::time::timeout(std::time::Duration::from_secs(1), rx)
1306 .await
1307 .expect("Timeout waiting for result")
1308 .expect("Channel closed");
1309
1310 assert!(result.is_ok());
1312 assert_eq!(result.unwrap()["result"], "ok");
1313 }
1314
1315 #[tokio::test]
1316 async fn test_command_completes_before_timeout() {
1317 let emitter = EventEmitter::new(100);
1318 let queue = Arc::new(CommandQueue::new(emitter));
1319
1320 let config = LaneConfig::new(1, 4).with_timeout(std::time::Duration::from_secs(5));
1321 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1322 queue.register_lane(lane).await;
1323
1324 Arc::clone(&queue).start_scheduler().await;
1325
1326 let cmd = Box::new(TestCommand::with_delay(
1328 serde_json::json!({"result": "fast"}),
1329 10,
1330 ));
1331 let rx = queue.submit("test-lane", cmd).await.unwrap();
1332
1333 let result = tokio::time::timeout(std::time::Duration::from_secs(1), rx)
1334 .await
1335 .expect("Timeout waiting for result")
1336 .expect("Channel closed");
1337
1338 assert!(result.is_ok());
1340 assert_eq!(result.unwrap()["result"], "fast");
1341 }
1342
1343 #[tokio::test]
1344 async fn test_command_retry_on_failure() {
1345 use std::sync::atomic::{AtomicU32, Ordering};
1346
1347 struct RetryableCommand {
1349 attempts: Arc<AtomicU32>,
1350 }
1351
1352 #[async_trait]
1353 impl Command for RetryableCommand {
1354 async fn execute(&self) -> Result<serde_json::Value> {
1355 let attempt = self.attempts.fetch_add(1, Ordering::SeqCst);
1356 if attempt < 2 {
1357 Err(LaneError::Other(format!("Attempt {} failed", attempt)))
1358 } else {
1359 Ok(serde_json::json!({"success": true, "attempts": attempt + 1}))
1360 }
1361 }
1362
1363 fn command_type(&self) -> &str {
1364 "retryable"
1365 }
1366 }
1367
1368 let emitter = EventEmitter::new(100);
1369 let queue = Arc::new(CommandQueue::new(emitter));
1370
1371 let retry_policy = RetryPolicy::fixed(3, std::time::Duration::from_millis(10));
1372 let config = LaneConfig::new(1, 4).with_retry_policy(retry_policy);
1373 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1374 queue.register_lane(lane).await;
1375
1376 Arc::clone(&queue).start_scheduler().await;
1377
1378 let attempts = Arc::new(AtomicU32::new(0));
1379 let cmd = Box::new(RetryableCommand {
1380 attempts: Arc::clone(&attempts),
1381 });
1382 let rx = queue.submit("test-lane", cmd).await.unwrap();
1383
1384 let result = tokio::time::timeout(std::time::Duration::from_secs(2), rx)
1386 .await
1387 .expect("Timeout waiting for result")
1388 .expect("Channel closed");
1389
1390 assert!(result.is_ok());
1391 let value = result.unwrap();
1392 assert_eq!(value["success"], true);
1393 assert_eq!(value["attempts"], 3); }
1395
1396 #[tokio::test]
1397 async fn test_command_retry_exhausted() {
1398 struct AlwaysFailCommand;
1400
1401 #[async_trait]
1402 impl Command for AlwaysFailCommand {
1403 async fn execute(&self) -> Result<serde_json::Value> {
1404 Err(LaneError::Other("Always fails".to_string()))
1405 }
1406
1407 fn command_type(&self) -> &str {
1408 "always_fail"
1409 }
1410 }
1411
1412 let emitter = EventEmitter::new(100);
1413 let queue = Arc::new(CommandQueue::new(emitter));
1414
1415 let retry_policy = RetryPolicy::fixed(2, std::time::Duration::from_millis(10));
1416 let config = LaneConfig::new(1, 4).with_retry_policy(retry_policy);
1417 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1418 queue.register_lane(lane).await;
1419
1420 Arc::clone(&queue).start_scheduler().await;
1421
1422 let cmd = Box::new(AlwaysFailCommand);
1423 let rx = queue.submit("test-lane", cmd).await.unwrap();
1424
1425 let result = tokio::time::timeout(std::time::Duration::from_secs(2), rx)
1427 .await
1428 .expect("Timeout waiting for result")
1429 .expect("Channel closed");
1430
1431 assert!(result.is_err());
1432 if let Err(LaneError::Other(msg)) = result {
1433 assert_eq!(msg, "Always fails");
1434 } else {
1435 panic!("Expected Other error");
1436 }
1437 }
1438
1439 #[tokio::test]
1440 async fn test_command_no_retry_on_success() {
1441 use std::sync::atomic::{AtomicU32, Ordering};
1442
1443 struct CountingCommand {
1444 counter: Arc<AtomicU32>,
1445 }
1446
1447 #[async_trait]
1448 impl Command for CountingCommand {
1449 async fn execute(&self) -> Result<serde_json::Value> {
1450 let count = self.counter.fetch_add(1, Ordering::SeqCst);
1451 Ok(serde_json::json!({"count": count + 1}))
1452 }
1453
1454 fn command_type(&self) -> &str {
1455 "counting"
1456 }
1457 }
1458
1459 let emitter = EventEmitter::new(100);
1460 let queue = Arc::new(CommandQueue::new(emitter));
1461
1462 let retry_policy = RetryPolicy::exponential(3);
1463 let config = LaneConfig::new(1, 4).with_retry_policy(retry_policy);
1464 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1465 queue.register_lane(lane).await;
1466
1467 Arc::clone(&queue).start_scheduler().await;
1468
1469 let counter = Arc::new(AtomicU32::new(0));
1470 let cmd = Box::new(CountingCommand {
1471 counter: Arc::clone(&counter),
1472 });
1473 let rx = queue.submit("test-lane", cmd).await.unwrap();
1474
1475 let result = tokio::time::timeout(std::time::Duration::from_secs(1), rx)
1476 .await
1477 .expect("Timeout waiting for result")
1478 .expect("Channel closed");
1479
1480 assert!(result.is_ok());
1481 assert_eq!(result.unwrap()["count"], 1);
1482
1483 assert_eq!(counter.load(Ordering::SeqCst), 1);
1485 }
1486
1487 #[tokio::test]
1488 async fn test_dlq_integration() {
1489 struct FailCommand;
1491
1492 #[async_trait]
1493 impl Command for FailCommand {
1494 async fn execute(&self) -> Result<serde_json::Value> {
1495 Err(LaneError::Other("Permanent failure".to_string()))
1496 }
1497
1498 fn command_type(&self) -> &str {
1499 "fail_command"
1500 }
1501 }
1502
1503 let emitter = EventEmitter::new(100);
1504 let queue = Arc::new(CommandQueue::with_dlq(emitter, 100));
1505
1506 let retry_policy = RetryPolicy::fixed(2, std::time::Duration::from_millis(10));
1507 let config = LaneConfig::new(1, 4).with_retry_policy(retry_policy);
1508 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1509 queue.register_lane(lane).await;
1510
1511 Arc::clone(&queue).start_scheduler().await;
1512
1513 let cmd = Box::new(FailCommand);
1515 let rx = queue.submit("test-lane", cmd).await.unwrap();
1516
1517 let result = tokio::time::timeout(std::time::Duration::from_secs(2), rx)
1519 .await
1520 .expect("Timeout waiting for result")
1521 .expect("Channel closed");
1522
1523 assert!(result.is_err());
1524
1525 let dlq = queue.dlq().expect("DLQ should exist");
1527 tokio::time::sleep(std::time::Duration::from_millis(50)).await; assert_eq!(dlq.len().await, 1);
1530
1531 let letters = dlq.list().await;
1532 assert_eq!(letters[0].command_type, "fail_command");
1533 assert_eq!(letters[0].lane_id, "test-lane");
1534 assert_eq!(letters[0].attempts, 3); assert!(letters[0].error.contains("Permanent failure"));
1536 }
1537
1538 #[tokio::test]
1539 async fn test_no_dlq_without_configuration() {
1540 let emitter = EventEmitter::new(100);
1541 let queue = Arc::new(CommandQueue::new(emitter));
1542
1543 assert!(queue.dlq().is_none());
1544 }
1545
1546 #[tokio::test]
1547 async fn test_shutdown_rejects_new_commands() {
1548 let emitter = EventEmitter::new(100);
1549 let queue = Arc::new(CommandQueue::new(emitter));
1550
1551 let config = LaneConfig::new(1, 4);
1552 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1553 queue.register_lane(lane).await;
1554
1555 queue.shutdown().await;
1557 assert!(queue.is_shutting_down());
1558
1559 let cmd = Box::new(TestCommand::new(serde_json::json!({"test": "data"})));
1561 let result = queue.submit("test-lane", cmd).await;
1562
1563 assert!(result.is_err());
1564 if let Err(LaneError::ShutdownInProgress) = result {
1565 } else {
1567 panic!("Expected ShutdownInProgress error");
1568 }
1569 }
1570
1571 #[tokio::test]
1572 async fn test_shutdown_preserves_submit_admitted_before_close() {
1573 let emitter = EventEmitter::new(100);
1574 let storage = Arc::new(BlockingStorage::default());
1575 let queue = Arc::new(CommandQueue::with_storage(
1576 emitter,
1577 Arc::clone(&storage) as Arc<dyn Storage>,
1578 ));
1579 let lane = Arc::new(Lane::with_storage(
1580 "test-lane",
1581 LaneConfig::new(1, 4),
1582 priorities::QUERY,
1583 Arc::clone(&storage) as Arc<dyn Storage>,
1584 ));
1585 queue.register_lane(lane).await;
1586 let scheduler = Arc::clone(&queue).spawn_scheduler();
1587
1588 let submit_queue = Arc::clone(&queue);
1589 let submit = tokio::spawn(async move {
1590 submit_queue
1591 .submit(
1592 "test-lane",
1593 Box::new(TestCommand::new(serde_json::json!({"accepted": true}))),
1594 )
1595 .await
1596 });
1597
1598 storage.save_started.notified().await;
1599 assert!(
1600 queue.lanes.try_lock().is_ok(),
1601 "submit must not retain the lane map lock during storage I/O"
1602 );
1603
1604 queue.shutdown().await;
1605 let rejected = queue
1606 .submit(
1607 "test-lane",
1608 Box::new(TestCommand::new(serde_json::json!({}))),
1609 )
1610 .await;
1611 assert!(matches!(rejected, Err(LaneError::ShutdownInProgress)));
1612 assert_eq!(queue.lifecycle.outstanding.load(Ordering::Acquire), 1);
1613
1614 storage.release_save.notify_one();
1615 let receiver = submit.await.expect("submit task should join").unwrap();
1616 let result = receiver.await.expect("result channel should remain open");
1617 assert_eq!(result.unwrap()["accepted"], true);
1618
1619 queue
1620 .drain(std::time::Duration::from_secs(1))
1621 .await
1622 .unwrap();
1623 scheduler.await.expect("scheduler should exit after drain");
1624 }
1625
1626 #[tokio::test(start_paused = true)]
1627 async fn test_drain_waits_for_delayed_retry() {
1628 use std::sync::atomic::{AtomicU32, Ordering as AtomicOrdering};
1629
1630 struct FailOnceCommand {
1631 attempts: Arc<AtomicU32>,
1632 }
1633
1634 #[async_trait]
1635 impl Command for FailOnceCommand {
1636 async fn execute(&self) -> Result<serde_json::Value> {
1637 if self.attempts.fetch_add(1, AtomicOrdering::SeqCst) == 0 {
1638 Err(LaneError::Other("retry".to_string()))
1639 } else {
1640 Ok(serde_json::json!({"retried": true}))
1641 }
1642 }
1643
1644 fn command_type(&self) -> &str {
1645 "fail_once"
1646 }
1647 }
1648
1649 let emitter = EventEmitter::new(100);
1650 let mut retries =
1651 emitter.subscribe_filtered(|event| event.key == events::QUEUE_COMMAND_RETRY);
1652 let queue = Arc::new(CommandQueue::new(emitter));
1653 let retry_delay = std::time::Duration::from_secs(60);
1654 let lane = Arc::new(Lane::new(
1655 "test-lane",
1656 LaneConfig::new(1, 4).with_retry_policy(RetryPolicy::fixed(1, retry_delay)),
1657 priorities::QUERY,
1658 ));
1659 queue.register_lane(Arc::clone(&lane)).await;
1660 let scheduler = Arc::clone(&queue).spawn_scheduler();
1661
1662 let receiver = queue
1663 .submit(
1664 "test-lane",
1665 Box::new(FailOnceCommand {
1666 attempts: Arc::new(AtomicU32::new(0)),
1667 }),
1668 )
1669 .await
1670 .unwrap();
1671 retries.recv().await.expect("retry event should be emitted");
1672 assert_eq!(lane.status().await.pending, 1);
1673
1674 queue.shutdown().await;
1675 let drain_queue = Arc::clone(&queue);
1676 let drain =
1677 tokio::spawn(
1678 async move { drain_queue.drain(std::time::Duration::from_secs(120)).await },
1679 );
1680 tokio::task::yield_now().await;
1681 assert!(
1682 !drain.is_finished(),
1683 "retry backoff must remain owned by drain"
1684 );
1685
1686 tokio::time::advance(retry_delay).await;
1687 let result = receiver.await.expect("result channel should remain open");
1688 assert_eq!(result.unwrap()["retried"], true);
1689 drain.await.expect("drain task should join").unwrap();
1690 scheduler.await.expect("scheduler should exit after drain");
1691 }
1692
1693 #[tokio::test]
1694 async fn test_drain_waits_for_completion() {
1695 let emitter = EventEmitter::new(100);
1696 let queue = Arc::new(CommandQueue::new(emitter));
1697
1698 let config = LaneConfig::new(1, 4);
1699 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1700 queue.register_lane(lane).await;
1701
1702 Arc::clone(&queue).start_scheduler().await;
1703
1704 let cmd = Box::new(TestCommand::with_delay(
1706 serde_json::json!({"result": "ok"}),
1707 100,
1708 ));
1709 let rx = queue.submit("test-lane", cmd).await.unwrap();
1710
1711 queue.shutdown().await;
1713
1714 let drain_result = queue.drain(std::time::Duration::from_secs(2)).await;
1716 assert!(drain_result.is_ok());
1717
1718 let result = tokio::time::timeout(std::time::Duration::from_millis(100), rx)
1720 .await
1721 .expect("Timeout")
1722 .expect("Channel closed");
1723 assert!(result.is_ok());
1724 }
1725
1726 #[tokio::test]
1727 async fn test_drain_timeout() {
1728 let emitter = EventEmitter::new(100);
1729 let queue = Arc::new(CommandQueue::new(emitter));
1730
1731 let config = LaneConfig::new(1, 4);
1732 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1733 queue.register_lane(lane).await;
1734
1735 Arc::clone(&queue).start_scheduler().await;
1736
1737 let cmd = Box::new(TestCommand::with_delay(
1739 serde_json::json!({"result": "ok"}),
1740 5000,
1741 ));
1742 let _rx = queue.submit("test-lane", cmd).await.unwrap();
1743
1744 queue.shutdown().await;
1746
1747 let drain_result = queue.drain(std::time::Duration::from_millis(50)).await;
1749 assert!(drain_result.is_err());
1750 if let Err(LaneError::Timeout(_)) = drain_result {
1751 } else {
1753 panic!("Expected Timeout error");
1754 }
1755 }
1756
1757 #[tokio::test(start_paused = true)]
1758 async fn test_drain_deadline_includes_lane_map_lock_wait() {
1759 let queue = Arc::new(CommandQueue::new(EventEmitter::new(100)));
1760 let lanes_guard = queue.lanes.lock().await;
1761 let timeout = std::time::Duration::from_secs(5);
1762 let drain_queue = Arc::clone(&queue);
1763 let drain = tokio::spawn(async move { drain_queue.drain(timeout).await });
1764
1765 tokio::task::yield_now().await;
1766 assert!(!drain.is_finished());
1767 tokio::time::advance(timeout).await;
1768
1769 assert!(matches!(
1770 drain.await.expect("drain task should join"),
1771 Err(LaneError::Timeout(duration)) if duration == timeout
1772 ));
1773 drop(lanes_guard);
1774 }
1775
1776 #[tokio::test]
1777 async fn test_is_shutting_down() {
1778 let emitter = EventEmitter::new(100);
1779 let queue = Arc::new(CommandQueue::new(emitter));
1780
1781 assert!(!queue.is_shutting_down());
1782
1783 queue.shutdown().await;
1784 assert!(queue.is_shutting_down());
1785 }
1786
1787 #[tokio::test]
1790 async fn test_lane_pressure_emits_on_threshold() {
1791 use crate::event::events;
1792
1793 let emitter = EventEmitter::new(100);
1794 let queue = Arc::new(CommandQueue::new(emitter.clone()));
1795
1796 let config = LaneConfig::new(1, 4).with_pressure_threshold(2);
1797 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1798 queue.register_lane(lane).await;
1799
1800 let mut stream = emitter.subscribe_filtered(|e| e.key == events::QUEUE_LANE_PRESSURE);
1802
1803 for _ in 0..2 {
1805 let cmd = Box::new(TestCommand::new(serde_json::json!({})));
1806 std::mem::drop(queue.submit("test-lane", cmd).await.unwrap());
1807 }
1808
1809 Arc::clone(&queue).start_scheduler().await;
1811
1812 let event = tokio::time::timeout(std::time::Duration::from_secs(1), stream.recv())
1813 .await
1814 .expect("No pressure event received within timeout")
1815 .expect("Stream ended");
1816
1817 assert_eq!(event.key, events::QUEUE_LANE_PRESSURE);
1818 }
1819
1820 #[tokio::test]
1821 async fn test_lane_idle_emits_when_drained() {
1822 use crate::event::events;
1823
1824 let emitter = EventEmitter::new(100);
1825 let queue = Arc::new(CommandQueue::new(emitter.clone()));
1826
1827 let config = LaneConfig::new(1, 4).with_pressure_threshold(1);
1828 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1829 queue.register_lane(lane).await;
1830
1831 let mut stream = emitter.subscribe_filtered(|e| {
1833 e.key == events::QUEUE_LANE_PRESSURE || e.key == events::QUEUE_LANE_IDLE
1834 });
1835
1836 let cmd = Box::new(TestCommand::new(serde_json::json!({})));
1838 std::mem::drop(queue.submit("test-lane", cmd).await.unwrap());
1839
1840 Arc::clone(&queue).start_scheduler().await;
1842
1843 let pressure = tokio::time::timeout(std::time::Duration::from_secs(1), stream.recv())
1845 .await
1846 .expect("No pressure event")
1847 .expect("Stream ended");
1848 assert_eq!(pressure.key, events::QUEUE_LANE_PRESSURE);
1849
1850 let idle = tokio::time::timeout(std::time::Duration::from_secs(1), stream.recv())
1852 .await
1853 .expect("No idle event")
1854 .expect("Stream ended");
1855 assert_eq!(idle.key, events::QUEUE_LANE_IDLE);
1856 }
1857
1858 #[tokio::test]
1859 async fn test_lane_no_pressure_without_threshold() {
1860 use crate::event::events;
1861
1862 let emitter = EventEmitter::new(100);
1863 let queue = Arc::new(CommandQueue::new(emitter.clone()));
1864
1865 let config = LaneConfig::new(1, 4);
1867 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1868 queue.register_lane(lane).await;
1869
1870 let mut stream = emitter.subscribe_filtered(|e| {
1871 e.key == events::QUEUE_LANE_PRESSURE || e.key == events::QUEUE_LANE_IDLE
1872 });
1873
1874 for _ in 0..5 {
1876 let cmd = Box::new(TestCommand::new(serde_json::json!({})));
1877 std::mem::drop(queue.submit("test-lane", cmd).await.unwrap());
1878 }
1879
1880 Arc::clone(&queue).start_scheduler().await;
1881
1882 tokio::time::sleep(std::time::Duration::from_millis(200)).await;
1884
1885 let result =
1887 tokio::time::timeout(std::time::Duration::from_millis(50), stream.recv()).await;
1888
1889 assert!(
1890 result.is_err(),
1891 "Should not receive pressure/idle events without threshold"
1892 );
1893 }
1894
1895 #[tokio::test]
1900 async fn test_event_emitted_on_submit() {
1901 use crate::event::events;
1902
1903 let emitter = EventEmitter::new(100);
1904 let mut rx = emitter.subscribe();
1905 let queue = CommandQueue::new(emitter);
1906
1907 let config = LaneConfig::new(1, 4);
1908 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1909 queue.register_lane(lane).await;
1910
1911 let cmd = Box::new(TestCommand::new(serde_json::json!({"result": "ok"})));
1912 std::mem::drop(queue.submit("test-lane", cmd).await.unwrap());
1913
1914 let event = tokio::time::timeout(std::time::Duration::from_millis(200), async {
1915 rx.recv().await.unwrap()
1916 })
1917 .await
1918 .expect("QUEUE_COMMAND_SUBMITTED event not received");
1919
1920 assert_eq!(event.key, events::QUEUE_COMMAND_SUBMITTED);
1921 }
1922
1923 #[tokio::test]
1924 async fn test_event_emitted_on_complete() {
1925 use crate::event::{events, EventStream};
1926
1927 let emitter = EventEmitter::new(100);
1928 let queue = Arc::new(CommandQueue::new(emitter.clone()));
1929
1930 let config = LaneConfig::new(1, 4);
1931 let lane = Arc::new(Lane::new("test-lane", config, priorities::QUERY));
1932 queue.register_lane(lane).await;
1933 Arc::clone(&queue).start_scheduler().await;
1934
1935 let mut stream: EventStream =
1936 emitter.subscribe_filtered(|e| e.key == events::QUEUE_COMMAND_COMPLETED);
1937
1938 let cmd = Box::new(TestCommand::new(serde_json::json!({"result": "ok"})));
1939 let rx = queue.submit("test-lane", cmd).await.unwrap();
1940
1941 let _ = tokio::time::timeout(std::time::Duration::from_secs(1), rx).await;
1943
1944 let event = tokio::time::timeout(std::time::Duration::from_millis(200), stream.recv())
1945 .await
1946 .expect("QUEUE_COMMAND_COMPLETED event not received")
1947 .expect("stream closed");
1948
1949 assert_eq!(event.key, events::QUEUE_COMMAND_COMPLETED);
1950 }
1951
1952 #[tokio::test]
1953 async fn test_event_emitted_on_shutdown() {
1954 use crate::event::events;
1955
1956 let emitter = EventEmitter::new(100);
1957 let mut rx = emitter.subscribe();
1958 let queue = CommandQueue::new(emitter);
1959
1960 queue.shutdown().await;
1961
1962 let event = tokio::time::timeout(std::time::Duration::from_millis(100), async {
1963 rx.recv().await.unwrap()
1964 })
1965 .await
1966 .expect("QUEUE_SHUTDOWN_STARTED event not received");
1967
1968 assert_eq!(event.key, events::QUEUE_SHUTDOWN_STARTED);
1969 }
1970
1971 #[cfg(feature = "distributed")]
1972 #[tokio::test]
1973 async fn test_rate_limit_blocks_dequeue() {
1974 use crate::ratelimit::RateLimitConfig;
1975
1976 let config = LaneConfig::new(1, 10).with_rate_limit(RateLimitConfig::per_second(1));
1978 let lane = Lane::new("test", config, priorities::QUERY);
1979
1980 let _rx1 = lane
1981 .enqueue(Box::new(TestCommand::new(serde_json::json!(1))))
1982 .await;
1983 let _rx2 = lane
1984 .enqueue(Box::new(TestCommand::new(serde_json::json!(2))))
1985 .await;
1986
1987 let first = lane.try_dequeue().await;
1989 assert!(first.is_some(), "first dequeue should succeed");
1990 lane.mark_completed().await;
1992
1993 let second = lane.try_dequeue().await;
1995 assert!(second.is_none(), "second dequeue should be rate-limited");
1996 }
1997
1998 #[cfg(feature = "distributed")]
1999 #[tokio::test]
2000 async fn test_effective_priority_no_boost_when_fresh() {
2001 use crate::boost::PriorityBoostConfig;
2002
2003 let config = LaneConfig::new(1, 4).with_priority_boost(
2006 PriorityBoostConfig::new(std::time::Duration::from_secs(10))
2007 .with_boost(std::time::Duration::from_secs(9), 2),
2008 );
2009 let lane = Lane::new("test", config, priorities::QUERY);
2010
2011 assert_eq!(lane.effective_priority().await, priorities::QUERY);
2013
2014 let _rx = lane
2016 .enqueue(Box::new(TestCommand::new(serde_json::json!(1))))
2017 .await;
2018 assert_eq!(lane.effective_priority().await, priorities::QUERY);
2019 }
2020
2021 #[cfg(feature = "distributed")]
2022 #[tokio::test]
2023 async fn test_effective_priority_boosted_when_past_deadline() {
2024 use crate::boost::PriorityBoostConfig;
2025
2026 let config = LaneConfig::new(1, 4).with_priority_boost(PriorityBoostConfig::standard(
2028 std::time::Duration::from_millis(1), ));
2030 let lane = Lane::new("test", config, priorities::PROMPT); let _rx = lane
2033 .enqueue(Box::new(TestCommand::new(serde_json::json!(1))))
2034 .await;
2035
2036 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
2038
2039 assert_eq!(lane.effective_priority().await, 0);
2040 }
2041}