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(®istry),
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(®istry),
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}