Skip to main content

a3s_lane/
queue.rs

1//! Core queue implementation with lanes and priority scheduling
2
3#[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
26/// Lane identifier
27pub type LaneId = String;
28
29/// Command identifier
30pub type CommandId = String;
31
32/// Lane priority (lower number = higher priority)
33pub type Priority = u8;
34
35/// Lane priorities
36pub 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/// Command to be executed
48#[async_trait]
49pub trait Command: Send + Sync {
50    /// Execute the command
51    async fn execute(&self) -> Result<serde_json::Value>;
52
53    /// Get command type (for logging/debugging)
54    fn command_type(&self) -> &str;
55}
56
57/// A simple JSON-based command for data-driven Rust usage.
58///
59/// Returns the payload as-is when executed. Useful when commands are represented
60/// as JSON data but still run through the Rust queue API.
61pub struct JsonCommand {
62    command_type: String,
63    payload: serde_json::Value,
64}
65
66impl JsonCommand {
67    /// Create a new JSON command.
68    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
87/// Command wrapper
88struct 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    /// Keeps queue-level ownership alive from admission through terminal completion.
98    _work_permit: Option<WorkPermit>,
99    /// Submission time used by the priority booster to calculate deadline proximity
100    #[cfg(feature = "distributed")]
101    enqueue_time: chrono::DateTime<chrono::Utc>,
102}
103
104/// Lane state
105#[allow(dead_code)]
106struct LaneState {
107    /// Lane configuration
108    config: LaneConfig,
109
110    /// Priority
111    priority: Priority,
112
113    /// Pending commands (FIFO queue)
114    pending: VecDeque<CommandWrapper>,
115
116    /// Active command count
117    active: usize,
118
119    /// Commands waiting for their retry backoff to expire
120    delayed_retries: usize,
121
122    /// Semaphore for concurrency control
123    semaphore: Arc<Semaphore>,
124
125    /// True when the lane is currently considered under pressure
126    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
152/// Lane
153pub struct Lane {
154    id: LaneId,
155    state: Arc<Mutex<LaneState>>,
156    storage: Option<Arc<dyn Storage>>,
157    /// Rate limiter instantiated from LaneConfig.rate_limit (None = unlimited)
158    #[cfg(feature = "distributed")]
159    rate_limiter: RateLimiter,
160    /// Priority booster instantiated from LaneConfig.priority_boost
161    #[cfg(feature = "distributed")]
162    booster: Option<PriorityBooster>,
163}
164
165impl Lane {
166    /// Create a new lane
167    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    /// Create a new lane with storage
191    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    /// Get lane ID
220    pub fn id(&self) -> &str {
221        &self.id
222    }
223
224    /// Get lane priority
225    pub async fn priority(&self) -> Priority {
226        self.state.lock().await.priority
227    }
228
229    /// Get effective priority, applying boost based on the front command's age
230    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    /// Enqueue a command
243    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    /// Enqueue a command with queue-level lifecycle ownership.
251    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        // Persist to storage if available
279        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!({}), // Empty payload for now
285                retry_count: 0,
286                created_at: Utc::now(),
287                last_attempt_at: None,
288            };
289            // Ignore storage errors to not block command execution
290            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    /// Re-enqueue a command for retry (internal use)
300    async fn retry_command(&self, mut wrapper: CommandWrapper, delay: std::time::Duration) {
301        wrapper.attempt += 1;
302
303        // Count delayed retries as outstanding lane work before the active attempt is released.
304        {
305            let mut state = self.state.lock().await;
306            state.delayed_retries += 1;
307        }
308
309        // Spawn a task to re-enqueue after delay.
310        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    /// Try to dequeue a command for execution
320    async fn try_dequeue(&self) -> Option<CommandWrapper> {
321        // Check rate limiter before acquiring the state lock
322        #[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    /// Mark a command as completed
336    async fn mark_completed(&self) {
337        let mut state = self.state.lock().await;
338        state.active = state.active.saturating_sub(1);
339    }
340
341    /// Get lane status
342    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    /// Check for pressure state transitions.
358    ///
359    /// Returns the event key to emit if a transition occurred, or `None` if no change.
360    /// - Transitions to pressured when `pending >= threshold` and was not already pressured.
361    /// - Transitions to idle when `pending == 0` and was previously pressured.
362    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/// Lane status
383#[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
391/// Admission and outstanding-work state shared by all commands in a queue.
392struct QueueLifecycle {
393    /// Serializes the submit admission decision with shutdown.
394    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
443/// RAII ownership for one command accepted by [`CommandQueue::submit`].
444struct 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/// Command queue
456#[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    /// Create a new command queue
467    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    /// Create a new command queue with a dead letter queue
478    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    /// Create a new command queue with storage
489    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    /// Create a new command queue with both DLQ and storage
500    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    /// Get the storage backend
515    pub fn storage(&self) -> Option<&Arc<dyn Storage>> {
516        self.storage.as_ref()
517    }
518
519    /// Get the dead letter queue
520    pub fn dlq(&self) -> Option<&DeadLetterQueue> {
521        self.dlq.as_ref()
522    }
523
524    /// Check if shutdown is in progress
525    pub fn is_shutting_down(&self) -> bool {
526        self.lifecycle.is_shutting_down()
527    }
528
529    /// Initiate graceful shutdown - stop accepting new commands
530    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    /// Wait for all pending commands to complete (with timeout)
538    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    /// Register a lane
579    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    /// Submit a command to a lane
585    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    /// Start the scheduler
611    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    /// Schedule the next command
634    async fn schedule_next(&self) {
635        // Find the highest-priority lane with pending commands.
636        // effective_priority applies any deadline-based boost configured on the lane.
637        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        // Check pressure transitions for all lanes and emit events
646        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        // Sort by priority (lower number = higher priority)
662        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    /// Subscribe to all queue lifecycle events as an `EventStream` (implements `Stream`)
811    pub fn subscribe_stream(&self) -> EventStream {
812        self.event_emitter.subscribe_stream()
813    }
814
815    /// Subscribe to filtered queue lifecycle events as an `EventStream`
816    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    /// Get queue status for all lanes
824    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
842/// Built-in lane IDs
843pub 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    /// Test command implementation
897    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    /// Failing test command
933    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        // Verify priority ordering: system has highest priority (lowest number)
958        // Using const block to satisfy clippy assertions_on_constants
959        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        // Register multiple lanes
1098        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        // Start scheduler
1263        Arc::clone(&queue).start_scheduler().await;
1264
1265        // Submit a command that takes longer than timeout
1266        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        // Wait for result
1273        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        // Should be a timeout error
1279        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        // Submit a command with delay but no timeout configured
1299        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        // Should succeed
1311        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        // Submit a fast command with long timeout
1327        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        // Should succeed
1339        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        // Command that fails first 2 times, then succeeds
1348        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        // Wait for result (should succeed after retries)
1385        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); // Failed twice, succeeded on 3rd attempt
1394    }
1395
1396    #[tokio::test]
1397    async fn test_command_retry_exhausted() {
1398        // Command that always fails
1399        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        // Wait for result (should fail after exhausting retries)
1426        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        // Verify command was only executed once (no retries on success)
1484        assert_eq!(counter.load(Ordering::SeqCst), 1);
1485    }
1486
1487    #[tokio::test]
1488    async fn test_dlq_integration() {
1489        // Command that always fails
1490        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        // Submit a failing command
1514        let cmd = Box::new(FailCommand);
1515        let rx = queue.submit("test-lane", cmd).await.unwrap();
1516
1517        // Wait for result
1518        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        // Check DLQ
1526        let dlq = queue.dlq().expect("DLQ should exist");
1527        tokio::time::sleep(std::time::Duration::from_millis(50)).await; // Give time for DLQ push
1528
1529        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); // Initial + 2 retries
1535        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        // Initiate shutdown
1556        queue.shutdown().await;
1557        assert!(queue.is_shutting_down());
1558
1559        // Try to submit a command - should be rejected
1560        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            // Expected
1566        } 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        // Submit a slow command
1705        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        // Initiate shutdown
1712        queue.shutdown().await;
1713
1714        // Drain should wait for the command to complete
1715        let drain_result = queue.drain(std::time::Duration::from_secs(2)).await;
1716        assert!(drain_result.is_ok());
1717
1718        // Command should have completed
1719        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        // Submit a very slow command
1738        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        // Initiate shutdown
1745        queue.shutdown().await;
1746
1747        // Drain with short timeout should fail
1748        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            // Expected
1752        } 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    // ── Pressure event tests ───────────────────────────────────────────────────
1788
1789    #[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        // Subscribe before enqueuing to capture all events
1801        let mut stream = emitter.subscribe_filtered(|e| e.key == events::QUEUE_LANE_PRESSURE);
1802
1803        // Enqueue 2 commands (meets threshold=2)
1804        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        // Start scheduler — first tick calls check_pressure → pending=2 >= 2 → emit PRESSURE
1810        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        // Subscribe to both pressure and idle events
1832        let mut stream = emitter.subscribe_filtered(|e| {
1833            e.key == events::QUEUE_LANE_PRESSURE || e.key == events::QUEUE_LANE_IDLE
1834        });
1835
1836        // Enqueue 1 command (meets threshold=1)
1837        let cmd = Box::new(TestCommand::new(serde_json::json!({})));
1838        std::mem::drop(queue.submit("test-lane", cmd).await.unwrap());
1839
1840        // Start scheduler
1841        Arc::clone(&queue).start_scheduler().await;
1842
1843        // First event must be pressure
1844        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        // After dequeue, pending=0 → idle event on the next scheduler tick
1851        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        // No pressure threshold
1866        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        // Enqueue several commands
1875        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        // Allow time for scheduler ticks to run
1883        tokio::time::sleep(std::time::Duration::from_millis(200)).await;
1884
1885        // No pressure/idle events should have been emitted
1886        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    // ── Bug-fix tests ──────────────────────────────────────────────────────────
1896
1897    /// Bug fix: EventEmitter.emit() is now called on submit, start, complete, retry,
1898    /// dead-letter, fail, and shutdown.
1899    #[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        // Wait for command to complete
1942        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        // Token bucket starts full (1 token for per_second(1))
1977        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        // First dequeue consumes the single available token
1988        let first = lane.try_dequeue().await;
1989        assert!(first.is_some(), "first dequeue should succeed");
1990        // Return the slot so capacity isn't the limiting factor
1991        lane.mark_completed().await;
1992
1993        // No tokens left — rate limiter should block the second dequeue
1994        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        // Boost activates with 9 s remaining on a 10 s deadline.
2004        // A freshly-enqueued command has ~10 s left, so no boost yet.
2005        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        // Without any pending command the lane returns its base priority
2012        assert_eq!(lane.effective_priority().await, priorities::QUERY);
2013
2014        // A just-enqueued command has ~10 s left — no boost
2015        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        // Deadline already passed — effective priority should reach 0 (maximum)
2027        let config = LaneConfig::new(1, 4).with_priority_boost(PriorityBoostConfig::standard(
2028            std::time::Duration::from_millis(1), // 1 ms deadline
2029        ));
2030        let lane = Lane::new("test", config, priorities::PROMPT); // base = 5
2031
2032        let _rx = lane
2033            .enqueue(Box::new(TestCommand::new(serde_json::json!(1))))
2034            .await;
2035
2036        // Sleep past the deadline so the booster gives priority 0
2037        tokio::time::sleep(std::time::Duration::from_millis(10)).await;
2038
2039        assert_eq!(lane.effective_priority().await, 0);
2040    }
2041}