1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
// SPDX-License-Identifier: Apache-2.0
//! # velo-queue
//!
//! Named work queue abstraction for Velo distributed systems.
//!
//! A work queue is a named entity that you create or connect to, then get
//! typed sender/receiver handles for enqueuing and consuming work items.
//!
//! ## Backends
//!
//! | Backend | Feature | Description |
//! |---------|---------|-------------|
//! | [`InMemoryBackend`](backends::memory::InMemoryBackend) | *(always)* | `DashMap` + `flume` channels, for testing |
//! | [`MessengerQueueBackend`](backends::messenger::MessengerQueueBackend) | `messenger` | Actor on a velo instance via active messages |
//! | [`NatsQueueBackend`](backends::nats::NatsQueueBackend) | `nats` | NATS JetStream with WorkQueue retention |
//!
//! ## Quick Start
//!
//! ```rust,no_run
//! use serde::{Serialize, Deserialize};
//! use crate::queue::{sender, receiver, backends::memory::InMemoryBackend};
//!
//! #[derive(Serialize, Deserialize, Debug, PartialEq)]
//! struct Job { id: u64, payload: String }
//!
//! # async fn example() -> Result<(), Box<dyn std::error::Error>> {
//! let backend = InMemoryBackend::new(1024);
//!
//! let tx = sender::<Job>(&backend, "my-jobs").await?;
//! let rx = receiver::<Job>(&backend, "my-jobs").await?;
//!
//! tx.enqueue(&Job { id: 1, payload: "work".into() }).await?;
//! let job = rx.next().await?.unwrap();
//! assert_eq!(job.id, 1);
//! # Ok(())
//! # }
//! ```
//!
//! ## TODO: Acknowledgment Support
//!
//! Currently items are auto-acknowledged on receipt. A future iteration will add:
//! - `WorkItem<T>` wrapper with `ack()`, `nack(delay)`, `in_progress()`, `term()`
//! - `AckPolicy` config: `None` (auto-ack) vs `Manual` (explicit acknowledgment)
//! - Redelivery for nack'd or timed-out items
// Re-export primary types at crate root for convenience.
pub use ;
pub use ;
pub use NextOptions;
pub use WorkQueueReceiver;
pub use WorkQueueSender;
/// Create a typed sender for a named queue from the given backend.
pub async
/// Create a typed receiver for a named queue from the given backend.
pub async