1mod config;
78mod connect;
79mod enqueue_rate;
80mod fleet;
81pub mod keys;
82mod publish;
83mod workqueue;
84
85pub use config::{EnqueueMode, NatsEnqueueConfig};
86pub use fleet::connect_fleet_from_env;
87
88pub use workqueue::{connect_auto, NatsWorkQueueBackend};
89
90use std::sync::Arc;
91
92use async_nats::jetstream::kv::Store;
93use async_nats::jetstream;
94use async_trait::async_trait;
95use boson_core::{
96 BosonError, IdempotencyMode, Job, JobEnqueueDisposition, JobStatus, QueueBackend, Result, Run,
97 RunStatus, TaskConfig, TaskRunStats,
98};
99use chrono::{DateTime, Utc};
100use enqueue_rate::EnqueueRateLimiter;
101use futures::StreamExt;
102use serde::{Deserialize, Serialize};
103use uuid::Uuid;
104
105#[derive(Debug, Clone, Serialize, Deserialize)]
107struct LeaseRow {
108 lease_id: String,
109 job_id: String,
110 worker_id: String,
111 expires_at: DateTime<Utc>,
112}
113
114pub struct NatsQueueBackend {
119 kv: Store,
120 keys: keys::Keyspace,
121 enqueue_rate: EnqueueRateLimiter,
122}
123
124impl std::fmt::Debug for NatsQueueBackend {
125 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
126 f.debug_struct("NatsQueueBackend").finish_non_exhaustive()
127 }
128}
129
130impl NatsQueueBackend {
131 pub async fn connect(url: &str) -> Result<Self> {
148 Self::connect_with_keyspace(url, keys::Keyspace::from_env()).await
149 }
150
151 pub async fn connect_with_keyspace(url: &str, keyspace: keys::Keyspace) -> Result<Self> {
157 let client = connect::connect_nats(url).await.map_err(map_err)?;
158 let js = jetstream::new(client);
159 let bucket = keyspace.bucket();
160 let kv = match js.get_key_value(&bucket).await {
161 Ok(store) => store,
162 Err(_) => js
163 .create_key_value(async_nats::jetstream::kv::Config {
164 bucket,
165 ..Default::default()
166 })
167 .await
168 .map_err(map_err)?,
169 };
170 Ok(Self {
171 kv,
172 keys: keyspace,
173 enqueue_rate: EnqueueRateLimiter::new(),
174 })
175 }
176
177 #[must_use]
179 pub fn test_url() -> String {
180 std::env::var("BOSON_TEST_NATS_URL").unwrap_or_else(|_| "nats://127.0.0.1:4222".into())
181 }
182
183 pub async fn flush_namespace(&self) -> Result<()> {
189 let prefix = self.keys.namespace_prefix();
190 let keys = self.list_keys_prefixed(&prefix).await?;
191 for key in keys {
192 self.kv_delete(&key).await?;
193 }
194 Ok(())
195 }
196
197 async fn kv_get(&self, key: &str) -> Result<Option<Vec<u8>>> {
198 Ok(self
199 .kv
200 .get(key)
201 .await
202 .map_err(map_err)?
203 .map(|bytes| bytes.to_vec()))
204 }
205
206 async fn kv_put(&self, key: &str, value: &[u8]) -> Result<()> {
207 self.kv.put(key, value.to_vec().into()).await.map_err(map_err)?;
208 Ok(())
209 }
210
211 async fn kv_delete(&self, key: &str) -> Result<()> {
212 self.kv.delete(key).await.map_err(map_err)?;
213 Ok(())
214 }
215
216 async fn list_keys_prefixed(&self, prefix: &str) -> Result<Vec<String>> {
217 let mut keys = Vec::new();
218 let mut stream = self.kv.keys().await.map_err(map_err)?;
219 while let Some(key) = stream.next().await {
220 let key = key.map_err(map_err)?;
221 if key.starts_with(prefix) {
222 keys.push(key);
223 }
224 }
225 Ok(keys)
226 }
227
228 async fn load_job(&self, job_id: &str) -> Result<Option<Job>> {
229 let raw = self.kv_get(&self.keys.job(job_id)).await?;
230 raw.map_or(Ok(None), |bytes| {
231 serde_json::from_slice(&bytes).map_err(map_err).map(Some)
232 })
233 }
234
235 async fn save_job(&self, job: &Job) -> Result<()> {
236 let bytes = serde_json::to_vec(job).map_err(map_err)?;
237 self.kv_put(&self.keys.job(&job.job_id), &bytes).await
238 }
239
240 async fn add_ready(&self, job: &Job) -> Result<()> {
241 if job.status != JobStatus::Queued {
242 return Ok(());
243 }
244 let key = self.keys.ready(
245 &job.pool,
246 job.priority,
247 job.created_at.timestamp_millis(),
248 &job.job_id,
249 );
250 self.kv_put(&key, job.job_id.as_bytes()).await?;
251 self.kv_put(&self.keys.pool_marker(&job.pool), b"1").await?;
252 Ok(())
253 }
254
255 async fn remove_ready_for_job(&self, job: &Job) -> Result<()> {
256 let prefix = self.keys.ready_prefix(&job.pool);
257 let keys = self.list_keys_prefixed(&prefix).await?;
258 for key in keys {
259 if key.ends_with(&job.job_id) {
260 self.kv_delete(&key).await?;
261 }
262 }
263 Ok(())
264 }
265
266 async fn load_run(&self, run_id: &str) -> Result<Option<Run>> {
267 let raw = self.kv_get(&self.keys.run(run_id)).await?;
268 raw.map_or(Ok(None), |bytes| {
269 serde_json::from_slice(&bytes).map_err(map_err).map(Some)
270 })
271 }
272
273 async fn save_run(&self, run: &Run) -> Result<()> {
274 let bytes = serde_json::to_vec(run).map_err(map_err)?;
275 self.kv_put(&self.keys.run(&run.run_id), &bytes).await
276 }
277
278 async fn load_lease_row(&self, lease_id: &str) -> Result<Option<LeaseRow>> {
279 let raw = self.kv_get(&self.keys.lease(lease_id)).await?;
280 raw.map_or(Ok(None), |bytes| {
281 serde_json::from_slice(&bytes).map_err(map_err).map(Some)
282 })
283 }
284}
285
286fn map_err(e: impl std::fmt::Display) -> BosonError {
287 BosonError::Backend(e.to_string())
288}
289
290#[async_trait]
291impl QueueBackend for NatsQueueBackend {
292 async fn upsert_job(&self, job: &Job) -> Result<()> {
293 let existing = self.load_job(&job.job_id).await?;
294 if let Some(ref old) = existing {
295 if old.status == JobStatus::Queued && job.status != JobStatus::Queued {
296 self.remove_ready_for_job(old).await?;
297 } else if job.status == JobStatus::Queued {
298 self.remove_ready_for_job(old).await?;
299 self.add_ready(job).await?;
300 }
301 } else if job.status == JobStatus::Queued {
302 self.add_ready(job).await?;
303 }
304 self.save_job(job).await
305 }
306
307 async fn enqueue_with_policies(
308 &self,
309 job: Job,
310 task_config: &TaskConfig,
311 ) -> Result<(String, JobEnqueueDisposition)> {
312 let idempotency = task_config.resolved_idempotency_mode(IdempotencyMode::Lwt);
313 let mut job = job;
314 if idempotency == IdempotencyMode::Lwt {
315 if let Some(ref key) = job.idempotency_key {
316 if !key.is_empty() {
317 let idem_key = self.keys.idempotency(key);
318 let inserted = self.kv_get(&idem_key).await?.is_none();
319 if inserted {
320 self.kv_put(&idem_key, job.job_id.as_bytes()).await?;
321 } else if let Some(bytes) = self.kv_get(&idem_key).await? {
322 let prior_id = String::from_utf8_lossy(&bytes).into_owned();
323 if let Some(prior) = self.load_job(&prior_id).await? {
324 if matches!(prior.status, JobStatus::Queued | JobStatus::Running) {
325 return Ok((
326 prior_id,
327 JobEnqueueDisposition::ReusedIdempotent,
328 ));
329 }
330 }
331 self.kv_put(&idem_key, job.job_id.as_bytes()).await?;
332 }
333 }
334 }
335 } else {
336 job.idempotency_key = None;
337 }
338
339 let policy = &task_config.rate_limit_policy;
340 if policy.max_in_flight > 0 {
341 let count = self.count_active_jobs_for_task(&job.task_name).await?;
342 if count >= policy.max_in_flight {
343 return Err(BosonError::RateLimited(job.task_name.clone()));
344 }
345 }
346 if policy.max_enqueue_per_second > 0
347 && !self
348 .enqueue_rate
349 .try_record(&job.task_name, policy.max_enqueue_per_second)
350 {
351 return Err(BosonError::RateLimited(job.task_name.clone()));
352 }
353
354 let job_id = job.job_id.clone();
355 self.save_job(&job).await?;
356 self.add_ready(&job).await?;
357 Ok((job_id, JobEnqueueDisposition::InsertedNew))
358 }
359
360 async fn get_job(&self, job_id: &str) -> Result<Option<Job>> {
361 self.load_job(job_id).await
362 }
363
364 async fn list_jobs(
365 &self,
366 status_filter: Option<JobStatus>,
367 offset: usize,
368 limit: usize,
369 ) -> Result<Vec<Job>> {
370 let prefix = self.keys.job_prefix();
371 let keys = self.list_keys_prefixed(&prefix).await?;
372 let mut jobs = Vec::new();
373 for key in keys {
374 if let Some(bytes) = self.kv_get(&key).await? {
375 if let Ok(job) = serde_json::from_slice::<Job>(&bytes) {
376 if status_filter.is_none_or(|st| job.status == st) {
377 jobs.push(job);
378 }
379 }
380 }
381 }
382 jobs.sort_by_key(|j| j.created_at);
383 Ok(jobs.into_iter().skip(offset).take(limit).collect())
384 }
385
386 async fn cancel_job_if_active(&self, job_id: &str) -> Result<()> {
387 let Some(mut job) = self.load_job(job_id).await? else {
388 return Err(BosonError::JobNotFound(job_id.to_string()));
389 };
390 if !matches!(job.status, JobStatus::Queued | JobStatus::Running) {
391 return Ok(());
392 }
393 if job.status == JobStatus::Queued {
394 self.remove_ready_for_job(&job).await?;
395 }
396 job.status = JobStatus::Canceled;
397 self.save_job(&job).await
398 }
399
400 async fn try_claim_job(&self, job_id: &str) -> Result<Option<Job>> {
401 let Some(mut job) = self.load_job(job_id).await? else {
402 return Ok(None);
403 };
404 if job.status != JobStatus::Queued {
405 return Ok(None);
406 }
407 job.status = JobStatus::Running;
408 self.save_job(&job).await?;
409 self.remove_ready_for_job(&job).await?;
410 Ok(Some(job))
411 }
412
413 async fn revert_job_to_queued(&self, job_id: &str) -> Result<()> {
414 let Some(mut job) = self.load_job(job_id).await? else {
415 return Ok(());
416 };
417 if job.status != JobStatus::Running {
418 return Ok(());
419 }
420 job.status = JobStatus::Queued;
421 self.save_job(&job).await?;
422 self.add_ready(&job).await
423 }
424
425 async fn distinct_pools_queued(&self) -> Result<Vec<String>> {
426 let prefix = self.keys.pool_prefix();
427 let keys = self.list_keys_prefixed(&prefix).await?;
428 let mut out: Vec<String> = keys
429 .iter()
430 .filter_map(|key| key.strip_prefix(&prefix).map(str::to_string))
431 .collect();
432 out.sort();
433 Ok(out)
434 }
435
436 async fn list_queued_for_pool_sorted(&self, pool: &str, limit: usize) -> Result<Vec<Job>> {
437 let limit = limit.max(1);
438 let prefix = self.keys.ready_prefix(pool);
439 let mut keys = self.list_keys_prefixed(&prefix).await?;
440 keys.sort();
441 let mut jobs = Vec::new();
442 for key in keys.into_iter().take(limit) {
443 let Some(job_id) = key.rsplit('.').next() else {
444 continue;
445 };
446 if let Some(job) = self.load_job(job_id).await? {
447 if job.status == JobStatus::Queued && job.pool == pool {
448 jobs.push(job);
449 }
450 }
451 }
452 Ok(jobs)
453 }
454
455 async fn count_jobs(&self, status_filter: Option<JobStatus>) -> Result<u64> {
456 let jobs = self.list_jobs(status_filter, 0, usize::MAX).await?;
457 Ok(u64::try_from(jobs.len()).unwrap_or(u64::MAX))
458 }
459
460 async fn count_jobs_for_task(
461 &self,
462 task_name: &str,
463 status: Option<JobStatus>,
464 ) -> Result<u64> {
465 let jobs = self.list_jobs(status, 0, usize::MAX).await?;
466 let count = jobs.iter().filter(|j| j.task_name == task_name).count();
467 Ok(u64::try_from(count).unwrap_or(u64::MAX))
468 }
469
470 async fn count_active_jobs_for_task(&self, task_name: &str) -> Result<u32> {
471 let jobs = self.list_jobs(None, 0, usize::MAX).await?;
472 let count = jobs
473 .iter()
474 .filter(|j| {
475 j.task_name == task_name
476 && matches!(j.status, JobStatus::Queued | JobStatus::Running)
477 })
478 .count();
479 Ok(u32::try_from(count).unwrap_or(u32::MAX))
480 }
481
482 async fn find_nonterminal_by_idempotency_key(&self, key: &str) -> Result<Option<String>> {
483 if key.is_empty() {
484 return Ok(None);
485 }
486 let Some(bytes) = self.kv_get(&self.keys.idempotency(key)).await? else {
487 return Ok(None);
488 };
489 let job_id = String::from_utf8_lossy(&bytes).into_owned();
490 if let Some(job) = self.load_job(&job_id).await? {
491 if matches!(job.status, JobStatus::Queued | JobStatus::Running) {
492 return Ok(Some(job_id));
493 }
494 }
495 Ok(None)
496 }
497
498 async fn upsert_run(&self, run: &Run) -> Result<()> {
499 self.save_run(run).await
500 }
501
502 async fn get_run(&self, run_id: &str) -> Result<Option<Run>> {
503 self.load_run(run_id).await
504 }
505
506 async fn list_runs(
507 &self,
508 job_id_filter: Option<&str>,
509 offset: usize,
510 limit: usize,
511 ) -> Result<Vec<Run>> {
512 let prefix = self.keys.run_prefix();
513 let keys = self.list_keys_prefixed(&prefix).await?;
514 let mut runs = Vec::new();
515 for key in keys {
516 if let Some(bytes) = self.kv_get(&key).await? {
517 if let Ok(run) = serde_json::from_slice::<Run>(&bytes) {
518 if job_id_filter.is_none_or(|id| run.job_id == id) {
519 runs.push(run);
520 }
521 }
522 }
523 }
524 runs.sort_by_key(|r| r.started_at);
525 Ok(runs.into_iter().skip(offset).take(limit).collect())
526 }
527
528 async fn finish_run(
529 &self,
530 run_id: &str,
531 status: RunStatus,
532 duration_ms: Option<i64>,
533 error_message: Option<String>,
534 ) -> Result<()> {
535 let Some(mut run) = self.load_run(run_id).await? else {
536 return Ok(());
537 };
538 run.status = status;
539 run.finished_at = Some(Utc::now());
540 run.duration_ms = duration_ms;
541 run.error_message = error_message;
542 self.save_run(&run).await
543 }
544
545 async fn count_runs(&self, job_id_filter: Option<&str>) -> Result<u64> {
546 let runs = self.list_runs(job_id_filter, 0, usize::MAX).await?;
547 Ok(u64::try_from(runs.len()).unwrap_or(u64::MAX))
548 }
549
550 async fn count_runs_since(&self, since: DateTime<Utc>) -> Result<u64> {
551 let runs = self.list_runs(None, 0, usize::MAX).await?;
552 let count = runs.iter().filter(|r| r.started_at >= since).count();
553 Ok(u64::try_from(count).unwrap_or(u64::MAX))
554 }
555
556 async fn task_run_stats(&self, task_name: &str) -> Result<TaskRunStats> {
557 let runs = self.list_runs(None, 0, usize::MAX).await?;
558 let filtered: Vec<_> = runs.iter().filter(|r| r.task_name == task_name).collect();
559 let runs_total = u32::try_from(filtered.len()).unwrap_or(u32::MAX);
560 let success_count = u32::try_from(
561 filtered
562 .iter()
563 .filter(|r| r.status == RunStatus::Success)
564 .count(),
565 )
566 .unwrap_or(u32::MAX);
567 Ok(TaskRunStats {
568 runs_total,
569 success_count,
570 })
571 }
572
573 async fn get_task_config(&self, task_name: &str) -> Result<Option<TaskConfig>> {
574 let raw = self.kv_get(&self.keys.task_config(task_name)).await?;
575 raw.map_or(Ok(None), |bytes| {
576 serde_json::from_slice(&bytes).map_err(map_err).map(Some)
577 })
578 }
579
580 async fn upsert_task_config(&self, config: &TaskConfig) -> Result<()> {
581 let bytes = serde_json::to_vec(config).map_err(map_err)?;
582 self.kv_put(&self.keys.task_config(&config.task_name), &bytes)
583 .await
584 }
585
586 async fn try_claim_run_lease(
587 &self,
588 job_id: &str,
589 worker_id: &str,
590 ttl_secs: i64,
591 ) -> Result<Option<String>> {
592 if let Some(bytes) = self.kv_get(&self.keys.lease_by_job(job_id)).await? {
593 let lid = String::from_utf8_lossy(&bytes).into_owned();
594 if let Some(row) = self.load_lease_row(&lid).await? {
595 if row.expires_at > Utc::now() {
596 return Ok(None);
597 }
598 }
599 }
600 if self.kv_get(&self.keys.lease_by_job(job_id)).await?.is_some() {
601 return Ok(None);
602 }
603 let lease_id = Uuid::new_v4().to_string();
604 let row = LeaseRow {
605 lease_id: lease_id.clone(),
606 job_id: job_id.to_string(),
607 worker_id: worker_id.to_string(),
608 expires_at: Utc::now() + chrono::Duration::seconds(ttl_secs),
609 };
610 let json = serde_json::to_vec(&row).map_err(map_err)?;
611 self.kv_put(&self.keys.lease_by_job(job_id), lease_id.as_bytes())
612 .await?;
613 self.kv_put(&self.keys.lease(&lease_id), &json).await?;
614 Ok(Some(lease_id))
615 }
616
617 async fn extend_lease(&self, lease_id: &str, ttl_secs: i64) -> Result<()> {
618 let Some(mut row) = self.load_lease_row(lease_id).await? else {
619 return Ok(());
620 };
621 row.expires_at = Utc::now() + chrono::Duration::seconds(ttl_secs);
622 let json = serde_json::to_vec(&row).map_err(map_err)?;
623 self.kv_put(&self.keys.lease(lease_id), &json).await
624 }
625
626 async fn release_lease(&self, lease_id: &str) -> Result<()> {
627 let Some(row) = self.load_lease_row(lease_id).await? else {
628 return Ok(());
629 };
630 self.kv_delete(&self.keys.lease(lease_id)).await?;
631 self.kv_delete(&self.keys.lease_by_job(&row.job_id)).await
632 }
633
634 async fn expired_lease_job_pairs(&self) -> Result<Vec<(String, String)>> {
635 let prefix = self.keys.lease_prefix();
636 let keys = self.list_keys_prefixed(&prefix).await?;
637 let now = Utc::now();
638 let mut out = Vec::new();
639 for key in keys {
640 if key.contains(".lease_by_job.") {
641 continue;
642 }
643 if let Some(bytes) = self.kv_get(&key).await? {
644 if let Ok(row) = serde_json::from_slice::<LeaseRow>(&bytes) {
645 if row.expires_at <= now {
646 out.push((row.lease_id, row.job_id));
647 }
648 }
649 }
650 }
651 Ok(out)
652 }
653}
654
655pub async fn install_default_nats_backend(url: &str) -> Result<Arc<NatsQueueBackend>> {
661 let backend = Arc::new(NatsQueueBackend::connect(url).await?);
662 boson_core::QueueRouter::set_global(boson_core::QueueRouter::with_default(
663 Arc::clone(&backend) as Arc<dyn QueueBackend>,
664 ));
665 Ok(backend)
666}