- 1
//! Asynchronous event and work queue engines. - 2
//! - 3
//! Provides the core traits (`EventPublisher`, `EventSubscriber`, `WorkQueue`), - 4
//! the real production NATS Core & JetStream engine (`NatsBus`), and the - 5
//! high-throughput standalone concurrent engine (`InMemoryBus`). - 6
- 7
use std::collections::{HashMap, VecDeque}; - 8
use std::sync::{Arc, Mutex}; - 9
use std::time::Duration; - 10
- 11
use async_trait::async_trait; - 12
use thiserror::Error; - 13
use tokio::sync::{broadcast, mpsc}; - 14
- 15
use crate::envelope::MessageEnvelope; - 16
use crate::telemetry::{BusMetrics, DeadLetterEvent}; - 17
- 18
#[derive(Debug, Error)] - 19
pub enum BusError { - 20
#[error("connection error: {0}")] - 21
Connection(String), - 22
#[error("publish error on subject '{subject}': {reason}")] - 23
Publish { subject: String, reason: String }, - 24
#[error("subscription error for pattern '{pattern}': {reason}")] - 25
Subscribe { pattern: String, reason: String }, - 26
#[error("work queue error for queue '{queue}': {reason}")] - 27
WorkQueue { queue: String, reason: String }, - 28
#[error("serialization/deserialization error: {0}")] - 29
Serialization(String), - 30
#[error("task acknowledgement error: {0}")] - 31
AckError(String), - 32
#[error("maximum retry attempts exceeded (poison pill routed to DLQ)")] - 33
MaxRetriesExceeded, - 34
#[error("internal bus error: {0}")] - 35
Internal(String), - 36
} - 37
- 38
/// Acknowledge handle for a claimed work queue task. - 39
#[async_trait] - 40
pub trait TaskAckHandle: Send + Sync { - 41
/// Acknowledge successful task completion. - 42
async fn ack(&self) -> Result<(), BusError>; - 43
/// Negative acknowledge (re-queue task for another attempt or route to DLQ). - 44
async fn nack(&self, retry: bool) -> Result<(), BusError>; - 45
} - 46
- 47
/// A claimed task from a durable work queue. - 48
pub struct ClaimedTask { - 49
pub task_id: String, - 50
pub envelope: MessageEnvelope, - 51
pub attempts: u32, - 52
pub ack_handle: Box<dyn TaskAckHandle>, - 53
} - 54
- 55
/// Publisher contract for streaming ephemeral events. - 56
#[async_trait] - 57
pub trait EventPublisher: Send + Sync { - 58
async fn publish(&self, subject: &str, envelope: MessageEnvelope) -> Result<(), BusError>; - 59
} - 60
- 61
/// Subscriber contract for streaming ephemeral events. - 62
#[async_trait] - 63
pub trait EventSubscriber: Send + Sync { - 64
async fn subscribe( - 65
&self, - 66
subject_pattern: &str, - 67
) -> Result<mpsc::Receiver<MessageEnvelope>, BusError>; - 68
} - 69
- 70
/// Durable distributed work queue contract for competing consumer agent pools. - 71
#[async_trait] - 72
pub trait WorkQueue: Send + Sync { - 73
/// Enqueue a durable task into a named work queue. - 74
async fn enqueue(&self, queue: &str, envelope: MessageEnvelope) -> Result<String, BusError>; - 75
- 76
/// Claim a pending task from the work queue. - 77
async fn claim( - 78
&self, - 79
queue: &str, - 80
worker_id: &str, - 81
timeout: Duration, - 82
) -> Result<Option<ClaimedTask>, BusError>; - 83
} - 84
- 85
// --------------------------------------------------------------------------- - 86
// Real Production NATS Core + JetStream Engine - 87
// --------------------------------------------------------------------------- - 88
- 89
/// Configuration for connecting to a distributed NATS cluster. - 90
#[derive(Debug, Clone)] - 91
pub struct NatsConfig { - 92
pub url: String, - 93
pub credentials_jwt: Option<String>, - 94
pub nkey_seed: Option<String>, - 95
pub connect_timeout: Duration, - 96
} - 97
- 98
impl Default for NatsConfig { - 99
fn default() -> Self { - 100
Self { - 101
url: "nats://127.0.0.1:4222".to_string(), - 102
credentials_jwt: None, - 103
nkey_seed: None, - 104
connect_timeout: Duration::from_secs(5), - 105
} - 106
} - 107
} - 108
- 109
/// Production distributed messaging engine powered by NATS Core & JetStream. - 110
pub struct NatsBus { - 111
client: async_nats::Client, - 112
js: async_nats::jetstream::Context, - 113
metrics: Arc<BusMetrics>, - 114
} - 115
- 116
impl NatsBus { - 117
/// Connect to the NATS cluster and initialize JetStream contexts. - 118
pub async fn connect(config: NatsConfig) -> Result<Self, BusError> { - 119
let mut connect_opts = async_nats::ConnectOptions::new(); - 120
connect_opts = connect_opts.connection_timeout(config.connect_timeout); - 121
- 122
if let Some(creds) = config.credentials_jwt { - 123
connect_opts = connect_opts - 124
.credentials(&creds) - 125
.map_err(|e| BusError::Connection(e.to_string()))?; - 126
} else if let Some(nkey) = config.nkey_seed { - 127
connect_opts = connect_opts.nkey(nkey); - 128
} - 129
- 130
let client = connect_opts - 131
.connect(&config.url) - 132
.await - 133
.map_err(|e| BusError::Connection(e.to_string()))?; - 134
- 135
let js = async_nats::jetstream::new(client.clone()); - 136
let metrics = BusMetrics::new(); - 137
- 138
Ok(Self { - 139
client, - 140
js, - 141
metrics, - 142
}) - 143
} - 144
- 145
/// Access live operational metrics. - 146
pub fn metrics(&self) -> Arc<BusMetrics> { - 147
self.metrics.clone() - 148
} - 149
} - 150
- 151
#[async_trait] - 152
impl EventPublisher for NatsBus { - 153
async fn publish(&self, subject: &str, envelope: MessageEnvelope) -> Result<(), BusError> { - 154
let payload = - 155
serde_json::to_vec(&envelope).map_err(|e| BusError::Serialization(e.to_string()))?; - 156
let payload_len = payload.len(); - 157
- 158
self.client - 159
.publish(subject.to_string(), payload.into()) - 160
.await - 161
.map_err(|e| BusError::Publish { - 162
subject: subject.to_string(), - 163
reason: e.to_string(), - 164
})?; - 165
- 166
self.metrics.record_published(payload_len); - 167
Ok(()) - 168
} - 169
} - 170
- 171
#[async_trait] - 172
impl EventSubscriber for NatsBus { - 173
async fn subscribe( - 174
&self, - 175
subject_pattern: &str, - 176
) -> Result<mpsc::Receiver<MessageEnvelope>, BusError> { - 177
use futures::StreamExt; - 178
- 179
let mut sub = self - 180
.client - 181
.subscribe(subject_pattern.to_string()) - 182
.await - 183
.map_err(|e| BusError::Subscribe { - 184
pattern: subject_pattern.to_string(), - 185
reason: e.to_string(), - 186
})?; - 187
- 188
let (tx, rx) = mpsc::channel(2048); - 189
let metrics = self.metrics.clone(); - 190
- 191
tokio::spawn(async move { - 192
while let Some(msg) = sub.next().await { - 193
metrics.record_received(msg.payload.len()); - 194
if let Ok(env) = serde_json::from_slice::<MessageEnvelope>(&msg.payload) - 195
&& tx.send(env).await.is_err() - 196
{ - 197
break; - 198
} - 199
} - 200
}); - 201
- 202
Ok(rx) - 203
} - 204
} - 205
- 206
struct NatsJetStreamAckHandle { - 207
msg: Arc<tokio::sync::Mutex<Option<async_nats::jetstream::Message>>>, - 208
} - 209
- 210
#[async_trait] - 211
impl TaskAckHandle for NatsJetStreamAckHandle { - 212
async fn ack(&self) -> Result<(), BusError> { - 213
let mut guard = self.msg.lock().await; - 214
if let Some(msg) = guard.take() { - 215
msg.ack() - 216
.await - 217
.map_err(|e| BusError::AckError(e.to_string()))?; - 218
} - 219
Ok(()) - 220
} - 221
- 222
async fn nack(&self, retry: bool) -> Result<(), BusError> { - 223
let mut guard = self.msg.lock().await; - 224
if let Some(msg) = guard.take() { - 225
if retry { - 226
msg.ack_with(async_nats::jetstream::message::AckKind::Nak(None)) - 227
.await - 228
.map_err(|e| BusError::AckError(e.to_string()))?; - 229
} else { - 230
msg.ack() - 231
.await - 232
.map_err(|e| BusError::AckError(e.to_string()))?; - 233
} - 234
} - 235
Ok(()) - 236
} - 237
} - 238
- 239
#[async_trait] - 240
impl WorkQueue for NatsBus { - 241
async fn enqueue(&self, queue: &str, envelope: MessageEnvelope) -> Result<String, BusError> { - 242
let subject = format!("vak.work.{queue}.task"); - 243
let payload = - 244
serde_json::to_vec(&envelope).map_err(|e| BusError::Serialization(e.to_string()))?; - 245
let payload_len = payload.len(); - 246
- 247
let ack_future = self - 248
.js - 249
.publish(subject.clone(), payload.into()) - 250
.await - 251
.map_err(|e| BusError::WorkQueue { - 252
queue: queue.to_string(), - 253
reason: e.to_string(), - 254
})?; - 255
- 256
let ack = ack_future.await.map_err(|e| BusError::WorkQueue { - 257
queue: queue.to_string(), - 258
reason: e.to_string(), - 259
})?; - 260
- 261
self.metrics.record_published(payload_len); - 262
Ok(format!("{}:{}", ack.stream, ack.sequence)) - 263
} - 264
- 265
async fn claim( - 266
&self, - 267
queue: &str, - 268
_worker_id: &str, - 269
timeout: Duration, - 270
) -> Result<Option<ClaimedTask>, BusError> { - 271
use futures::StreamExt; - 272
- 273
let stream_name = format!("VAK_WORK_{}", queue.to_uppercase()); - 274
let consumer_name = format!("worker_{queue}"); - 275
- 276
let stream = match self.js.get_stream(&stream_name).await { - 277
Ok(s) => s, - 278
Err(_) => return Ok(None), - 279
}; - 280
- 281
let consumer = match stream.get_consumer(&consumer_name).await { - 282
Ok(c) => c, - 283
Err(_) => return Ok(None), - 284
}; - 285
- 286
let mut messages = match consumer.fetch().max_messages(1).messages().await { - 287
Ok(m) => m, - 288
Err(_) => return Ok(None), - 289
}; - 290
- 291
match tokio::time::timeout(timeout, messages.next()).await { - 292
Ok(Some(Ok(msg))) => { - 293
self.metrics.record_received(msg.payload.len()); - 294
let envelope: MessageEnvelope = serde_json::from_slice(&msg.payload) - 295
.map_err(|e| BusError::Serialization(e.to_string()))?; - 296
let attempts = 1; - 297
let task_id = envelope.id.clone(); - 298
let ack_handle = Box::new(NatsJetStreamAckHandle { - 299
msg: Arc::new(tokio::sync::Mutex::new(Some(msg))), - 300
}); - 301
- 302
Ok(Some(ClaimedTask { - 303
task_id, - 304
envelope, - 305
attempts, - 306
ack_handle, - 307
})) - 308
} - 309
_ => Ok(None), - 310
} - 311
} - 312
} - 313
- 314
// --------------------------------------------------------------------------- - 315
// Standalone Concurrency Engine (Zero-Dependency Local/Test Engine) - 316
// --------------------------------------------------------------------------- - 317
- 318
struct InMemoryTask { - 319
id: String, - 320
envelope: MessageEnvelope, - 321
attempts: u32, - 322
} - 323
- 324
/// Standalone, high-throughput concurrent engine matching NATS Core & JetStream semantics. - 325
pub struct InMemoryBus { - 326
events_tx: broadcast::Sender<(String, MessageEnvelope)>, - 327
queues: Arc<Mutex<HashMap<String, VecDeque<InMemoryTask>>>>, - 328
dlq: Arc<Mutex<Vec<DeadLetterEvent>>>, - 329
metrics: Arc<BusMetrics>, - 330
} - 331
- 332
impl Default for InMemoryBus { - 333
fn default() -> Self { - 334
Self::new() - 335
} - 336
} - 337
- 338
impl InMemoryBus { - 339
pub fn new() -> Self { - 340
let (events_tx, _) = broadcast::channel(4096); - 341
Self { - 342
events_tx, - 343
queues: Arc::new(Mutex::new(HashMap::new())), - 344
dlq: Arc::new(Mutex::new(Vec::new())), - 345
metrics: BusMetrics::new(), - 346
} - 347
} - 348
- 349
pub fn metrics(&self) -> Arc<BusMetrics> { - 350
self.metrics.clone() - 351
} - 352
- 353
pub fn dead_letters(&self) -> Vec<DeadLetterEvent> { - 354
self.dlq.lock().map(|d| d.clone()).unwrap_or_default() - 355
} - 356
} - 357
- 358
#[async_trait] - 359
impl EventPublisher for InMemoryBus { - 360
async fn publish(&self, subject: &str, envelope: MessageEnvelope) -> Result<(), BusError> { - 361
let size = serde_json::to_vec(&envelope) - 362
.map(|v| v.len()) - 363
.unwrap_or(128); - 364
self.metrics.record_published(size); - 365
let _ = self.events_tx.send((subject.to_string(), envelope)); - 366
Ok(()) - 367
} - 368
} - 369
- 370
#[async_trait] - 371
impl EventSubscriber for InMemoryBus { - 372
async fn subscribe( - 373
&self, - 374
subject_pattern: &str, - 375
) -> Result<mpsc::Receiver<MessageEnvelope>, BusError> { - 376
let mut b_rx = self.events_tx.subscribe(); - 377
let (tx, rx) = mpsc::channel(1024); - 378
let pattern = subject_pattern.to_string(); - 379
let metrics = self.metrics.clone(); - 380
- 381
tokio::spawn(async move { - 382
while let Ok((subj, env)) = b_rx.recv().await { - 383
if crate::subjects::matches_pattern(&pattern, &subj) { - 384
metrics.record_received(128); - 385
if tx.send(env).await.is_err() { - 386
break; - 387
} - 388
} - 389
} - 390
}); - 391
- 392
Ok(rx) - 393
} - 394
} - 395
- 396
struct InMemoryAckHandle { - 397
queue: String, - 398
task_id: String, - 399
envelope: MessageEnvelope, - 400
attempts: u32, - 401
max_retries: u32, - 402
queues: Arc<Mutex<HashMap<String, VecDeque<InMemoryTask>>>>, - 403
dlq: Arc<Mutex<Vec<DeadLetterEvent>>>, - 404
metrics: Arc<BusMetrics>, - 405
} - 406
- 407
#[async_trait] - 408
impl TaskAckHandle for InMemoryAckHandle { - 409
async fn ack(&self) -> Result<(), BusError> { - 410
// Task is already removed on claim; ack finalizes it. - 411
Ok(()) - 412
} - 413
- 414
async fn nack(&self, retry: bool) -> Result<(), BusError> { - 415
if retry && self.attempts < self.max_retries { - 416
if let Ok(mut map) = self.queues.lock() { - 417
let q = map.entry(self.queue.clone()).or_default(); - 418
q.push_front(InMemoryTask { - 419
id: self.task_id.clone(), - 420
envelope: self.envelope.clone(), - 421
attempts: self.attempts + 1, - 422
}); - 423
} - 424
} else { - 425
// Divert to Dead-Letter Queue (DLQ) - 426
self.metrics.record_dead_letter(); - 427
if let Ok(mut dlq) = self.dlq.lock() { - 428
dlq.push(DeadLetterEvent::new( - 429
self.envelope.clone(), - 430
&self.queue, - 431
"worker", - 432
self.attempts, - 433
"retry threshold exhausted", - 434
)); - 435
} - 436
} - 437
Ok(()) - 438
} - 439
} - 440
- 441
#[async_trait] - 442
impl WorkQueue for InMemoryBus { - 443
async fn enqueue(&self, queue: &str, envelope: MessageEnvelope) -> Result<String, BusError> { - 444
let task_id = envelope.id.clone(); - 445
let size = serde_json::to_vec(&envelope) - 446
.map(|v| v.len()) - 447
.unwrap_or(128); - 448
self.metrics.record_published(size); - 449
- 450
if let Ok(mut map) = self.queues.lock() { - 451
let q = map.entry(queue.to_string()).or_default(); - 452
q.push_back(InMemoryTask { - 453
id: task_id.clone(), - 454
envelope, - 455
attempts: 1, - 456
}); - 457
self.metrics.update_queue_lag(q.len() as u64); - 458
} - 459
- 460
Ok(task_id) - 461
} - 462
- 463
async fn claim( - 464
&self, - 465
queue: &str, - 466
_worker_id: &str, - 467
_timeout: Duration, - 468
) -> Result<Option<ClaimedTask>, BusError> { - 469
let task_opt = { - 470
let mut map = match self.queues.lock() { - 471
Ok(m) => m, - 472
Err(_) => return Ok(None), - 473
}; - 474
let q = map.entry(queue.to_string()).or_default(); - 475
let item = q.pop_front(); - 476
self.metrics.update_queue_lag(q.len() as u64); - 477
item - 478
}; - 479
- 480
if let Some(task) = task_opt { - 481
self.metrics.record_received(128); - 482
let ack_handle = Box::new(InMemoryAckHandle { - 483
queue: queue.to_string(), - 484
task_id: task.id.clone(), - 485
envelope: task.envelope.clone(), - 486
attempts: task.attempts, - 487
max_retries: 3, - 488
queues: self.queues.clone(), - 489
dlq: self.dlq.clone(), - 490
metrics: self.metrics.clone(), - 491
}); - 492
- 493
Ok(Some(ClaimedTask { - 494
task_id: task.id, - 495
envelope: task.envelope, - 496
attempts: task.attempts, - 497
ack_handle, - 498
})) - 499
} else { - 500
Ok(None) - 501
} - 502
} - 503
} - 504
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.