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(®istry),
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(®istry),
170 worker.clone(),
171 self.idempotency_mode,
172 );
173 let manual = ManualWorker::new(backend, registry, identity, worker);
174 Ok((boson, manual))
175 }
176}