Skip to main content

boson_backend_nats/
lib.rs

1//! `NATS` `JetStream` [`QueueBackend`] for fleet-scale deployments (Mode 2 remote / multi-host).
2//!
3//! **When to use:** broker-backed fleets with NATS `JetStream` (KV and/or workqueue). Not a
4//! `boson` facade feature — depend on this crate directly. Mode 2 workers need unique
5//! `worker_id` and `lease_ttl_secs > 0`.
6//!
7//! Getting started:
8//! [Mode 2](https://docs.rs/uf-boson/latest/boson/index.html#mode-2--remote-worker-two-binaries).
9//! Full Compose / KV vs `WorkQueue` / env: [crate README](https://github.com/unified-field-dev/boson/blob/main/boson-backend-nats/README.md).
10//!
11//! Fleet URL precedence: `BOSON_NATS_POOL_ROUTING` over `BOSON_NATS_URLS`
12//! (see [`connect_fleet_from_env`]).
13//!
14//! ## Mode 2 — Enqueue binary
15//!
16//! Shared NATS with the worker. No claim loop in this process:
17//!
18//! ```rust,ignore
19//! use std::sync::Arc;
20//!
21//! use boson_backend_nats::NatsQueueBackend;
22//! use boson_core::JsonExecutionContextFactory;
23//! use boson_runtime::{configure, Boson};
24//!
25//! # async fn boot_enqueue() -> boson_core::Result<()> {
26//! let url = std::env::var("BOSON_NATS_URL")
27//!     .unwrap_or_else(|_| "nats://127.0.0.1:4222".into());
28//! let backend = NatsQueueBackend::connect(&url).await?;
29//! let boson = Boson::builder()
30//!     .queue_backend(Arc::new(backend))
31//!     .execution_context_factory(JsonExecutionContextFactory)
32//!     .auto_registry()
33//!     .without_worker()
34//!     .build()?;
35//! configure(boson);
36//! // MyTask::send_with(...).await?;
37//! # Ok(())
38//! # }
39//! ```
40//!
41//! Also [`connect_auto`] / [`connect_fleet_from_env`] with the same `without_worker` + `configure`
42//! pattern.
43//!
44//! ## Mode 2 — Worker binary
45//!
46//! Same NATS URL / fleet, unique `worker_id`, and `lease_ttl_secs > 0`:
47//!
48//! ```rust,ignore
49//! use std::sync::Arc;
50//!
51//! use boson_backend_nats::NatsQueueBackend;
52//! use boson_core::JsonExecutionContextFactory;
53//! use boson_runtime::Boson;
54//!
55//! # async fn boot_worker() -> boson_core::Result<()> {
56//! let url = std::env::var("BOSON_NATS_URL")
57//!     .unwrap_or_else(|_| "nats://127.0.0.1:4222".into());
58//! let backend = NatsQueueBackend::connect(&url).await?;
59//! let _boson = Boson::builder()
60//!     .queue_backend(Arc::new(backend))
61//!     .execution_context_factory(JsonExecutionContextFactory)
62//!     .worker_id(std::env::var("BOSON_WORKER_ID").unwrap_or_else(|_| "worker-1".into()))
63//!     .lease_ttl_secs(30)
64//!     .auto_registry()
65//!     .build()?;
66//! # Ok(())
67//! # }
68//! ```
69//!
70//! Other Mode 2 backends:
71//! [`SQLite`](../boson_backend_sqlite/index.html#mode-2--enqueue-binary),
72//! [Postgres](../boson_backend_postgres/index.html#mode-2--enqueue-binary),
73//! [Redis](../boson_backend_redis/index.html#mode-2--enqueue-binary).
74//!
75//! Custom adapters: **How to implement** on [`QueueBackend`].
76
77mod 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/// Lease row persisted in KV.
106#[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
114/// `NATS` `JetStream` KV queue backend.
115///
116/// Mode 2 examples: [enqueue](index.html#mode-2--enqueue-binary) /
117/// [worker](index.html#mode-2--worker-binary).
118pub 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    /// Connect to `NATS` at `url` and open the KV bucket.
132    ///
133    /// # Examples
134    ///
135    /// ```rust,no_run
136    /// use boson_backend_nats::NatsQueueBackend;
137    ///
138    /// # async fn connect() -> boson_core::Result<()> {
139    /// let backend = NatsQueueBackend::connect("nats://127.0.0.1:4222").await?;
140    /// # Ok(())
141    /// # }
142    /// ```
143    ///
144    /// # Errors
145    ///
146    /// Returns an error when `NATS` or KV setup fails.
147    pub async fn connect(url: &str) -> Result<Self> {
148        Self::connect_with_keyspace(url, keys::Keyspace::from_env()).await
149    }
150
151    /// Connect with explicit key namespace.
152    ///
153    /// # Errors
154    ///
155    /// Returns an error when `NATS` or KV setup fails.
156    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    /// `NATS` URL for tests.
178    #[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    /// Delete all keys in this namespace (test isolation).
184    ///
185    /// # Errors
186    ///
187    /// Returns an error when KV list/delete fails.
188    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
655/// Install default `NATS` backend on global router (tests).
656///
657/// # Errors
658///
659/// Returns an error when `NATS` is unreachable.
660pub 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}