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