Skip to main content

boson_runtime/
bootstrap.rs

1//! Resolve builder dependencies and construct [`Boson`] / [`ManualWorker`].
2
3use std::sync::Arc;
4
5use boson_core::{
6    default_backend_from_global, BosonError, QueueBackend, Result,
7};
8use boson_telemetry::{install_ops_log, NoOpsLog, OpsLog};
9
10use crate::registry::TaskRegistry;
11use crate::telemetry::log_runtime_ready;
12use crate::worker::{spawn_worker, ManualWorker, WorkerSettings};
13use crate::{Boson, BosonBuilder};
14
15impl BosonBuilder {
16    pub(crate) fn resolve_backend(&self) -> Result<Arc<dyn QueueBackend>> {
17        if let Some(ref b) = self.queue_backend {
18            return Ok(Arc::clone(b));
19        }
20        default_backend_from_global()
21    }
22
23    pub(crate) fn resolve_registry(&self) -> Arc<TaskRegistry> {
24        if let Some(ref r) = self.registry {
25            return Arc::clone(r);
26        }
27        if self.use_auto_registry {
28            Arc::new(TaskRegistry::auto_discover())
29        } else {
30            Arc::new(TaskRegistry::new())
31        }
32    }
33
34    pub(crate) fn resolve_worker_settings(&self) -> WorkerSettings {
35        WorkerSettings::resolve(
36            self.worker_id.clone(),
37            self.lease_ttl_secs,
38            self.runtime_label.clone(),
39            self.worker_pools.clone(),
40            self.worker_poll_interval_ms,
41        )
42    }
43
44    pub(crate) fn install_ops_log(&self) {
45        let ops = self
46            .ops_log
47            .clone()
48            .unwrap_or_else(|| Arc::new(NoOpsLog) as Arc<dyn OpsLog>);
49        install_ops_log(ops);
50    }
51
52    /// Build [`Boson`] and optionally spawn the background worker loop.
53    ///
54    /// When [`without_worker`](Self::without_worker) was **not** set (default), a Tokio task polls
55    /// queued jobs, claims them, dispatches registered handlers (from
56    /// [`auto_registry`](Self::auto_registry) or [`registry`](Self::registry)), and applies retry
57    /// policy — Mode 1 embedded or Mode 2 worker.
58    ///
59    /// When [`without_worker`](Self::without_worker) **was** set, no claim loop starts; use this
60    /// for Mode 2 enqueue hosts, then [`configure`](crate::configure) + `send_with`.
61    ///
62    /// Enqueue with [`Boson::enqueue`] or macro `send_with` after [`configure`](crate::configure).
63    /// For step-driven tests, use [`build_manual`](Self::build_manual) instead.
64    ///
65    /// Getting started:
66    /// [Mode 1](https://docs.rs/uf-boson/latest/boson/index.html#mode-1--embedded-one-binary) /
67    /// [Mode 2](https://docs.rs/uf-boson/latest/boson/index.html#mode-2--remote-worker-two-binaries).
68    ///
69    /// # Example — Mode 1 embedded
70    ///
71    /// ```rust,no_run
72    /// use std::sync::Arc;
73    ///
74    /// use boson_backend_mem::MemQueueBackend;
75    /// use boson_core::JsonExecutionContextFactory;
76    /// use boson_runtime::{configure, Boson};
77    ///
78    /// # fn main() -> boson_core::Result<()> {
79    /// let boson = Boson::builder()
80    ///     .queue_backend(Arc::new(MemQueueBackend::new()))
81    ///     .execution_context_factory(JsonExecutionContextFactory)
82    ///     .auto_registry()
83    ///     .build()?; // worker loop runs in the background
84    /// configure(boson);
85    /// # Ok(())
86    /// # }
87    /// ```
88    ///
89    /// # Errors
90    ///
91    /// Returns [`BosonError::InvalidConfig`] if no execution context factory was set, or an error
92    /// if the queue backend cannot be resolved.
93    pub fn build(self) -> Result<Boson> {
94        self.install_ops_log();
95        let identity = self.execution_context_factory.clone().ok_or_else(|| {
96            BosonError::InvalidConfig("missing required execution_context_factory".into())
97        })?;
98        let backend = self.resolve_backend()?;
99        let registry = self.resolve_registry();
100        let worker = self.resolve_worker_settings();
101        let spawn_worker_flag = self.spawn_worker;
102
103        let boson = Boson::from_parts_with_idempotency(
104            Arc::clone(&backend),
105            Arc::clone(&registry),
106            worker.clone(),
107            self.idempotency_mode,
108        );
109        log_runtime_ready(&worker.runtime_label);
110
111        if spawn_worker_flag {
112            spawn_worker(backend, registry, identity, worker);
113        }
114
115        Ok(boson)
116    }
117
118    /// Build without a background worker; returns [`ManualWorker`] for step-driven execution.
119    ///
120    /// Prefer this in **tests** when you want to enqueue then call
121    /// [`ManualWorker::try_run_next`](crate::ManualWorker::try_run_next). For Mode 2 enqueue-only
122    /// production hosts that never drain locally, prefer [`without_worker`](Self::without_worker)
123    /// + [`build`](Self::build) instead (no `ManualWorker` handle needed).
124    ///
125    /// # Example — enqueue then drain one job
126    ///
127    /// ```rust,no_run
128    /// use std::sync::Arc;
129    ///
130    /// use boson_backend_mem::MemQueueBackend;
131    /// use boson_core::{ExecutionContext, JsonExecutionContextFactory};
132    /// use boson_macros::task;
133    /// use boson_runtime::{configure, Boson, ManualWorker};
134    ///
135    /// #[task(name = "ping")]
136    /// async fn ping(_ctx: Box<dyn ExecutionContext>, n: u32) -> boson_core::Result<()> {
137    ///     let _ = n;
138    ///     Ok(())
139    /// }
140    ///
141    /// # async fn demo() -> boson_core::Result<()> {
142    /// let (boson, manual): (_, ManualWorker) = Boson::builder()
143    ///     .queue_backend(Arc::new(MemQueueBackend::new()))
144    ///     .execution_context_factory(JsonExecutionContextFactory)
145    ///     .auto_registry()
146    ///     .without_worker()
147    ///     .build_manual()?;
148    /// configure(boson);
149    /// Ping::send_with(serde_json::json!({}), PingParams { n: 1 }).await?;
150    /// assert!(manual.try_run_next().await);
151    /// # Ok(())
152    /// # }
153    /// ```
154    ///
155    /// # Errors
156    ///
157    /// Returns [`BosonError::InvalidConfig`] if no execution context factory was set, or an error
158    /// if the queue backend cannot be resolved.
159    pub fn build_manual(self) -> Result<(Boson, ManualWorker)> {
160        self.install_ops_log();
161        let identity = self.execution_context_factory.clone().ok_or_else(|| {
162            BosonError::InvalidConfig("missing required execution_context_factory".into())
163        })?;
164        let backend = self.resolve_backend()?;
165        let registry = self.resolve_registry();
166        let worker = self.resolve_worker_settings();
167        let boson = Boson::from_parts_with_idempotency(
168            Arc::clone(&backend),
169            Arc::clone(&registry),
170            worker.clone(),
171            self.idempotency_mode,
172        );
173        let manual = ManualWorker::new(backend, registry, identity, worker);
174        Ok((boson, manual))
175    }
176}