Skip to main content

boson_runtime/worker/
config.rs

1//! Worker lease and identity settings composed at build time.
2
3/// Resolved worker settings for claim, lease, and telemetry labels.
4///
5/// Usually constructed by [`BosonBuilder`](crate::BosonBuilder) via [`worker_id`](crate::BosonBuilder::worker_id)
6/// and [`lease_ttl_secs`](crate::BosonBuilder::lease_ttl_secs). Defaults: worker id from
7/// `INSTANCE_ID` / `BOSON_WORKER_ID` / `boson-worker-1`, lease TTL `0` (single-process, no
8/// distributed leases).
9///
10/// # Example — multiple worker processes
11///
12/// Run one Boson instance per process against shared persistence. Each process needs a unique
13/// [`worker_id`](Self::worker_id) and a positive [`lease_ttl_secs`](Self::lease_ttl_secs) so
14/// [`QueueBackend::try_claim_run_lease`](boson_core::QueueBackend::try_claim_run_lease) prevents
15/// double execution:
16///
17/// ```rust,no_run
18/// use std::sync::Arc;
19///
20/// use boson_backend_mem::MemQueueBackend;
21/// use boson_core::JsonExecutionContextFactory;
22/// use boson_runtime::Boson;
23///
24/// # fn main() -> boson_core::Result<()> {
25/// let _boson = Boson::builder()
26///     .queue_backend(Arc::new(MemQueueBackend::new()))
27///     .execution_context_factory(JsonExecutionContextFactory)
28///     .worker_id("worker-a")
29///     .lease_ttl_secs(30)
30///     .auto_registry()
31///     .build()?; // background loop claims and runs jobs for this worker id
32/// # Ok(())
33/// # }
34/// ```
35#[derive(Debug, Clone)]
36pub struct WorkerSettings {
37    /// Worker identity for distributed lease claims.
38    pub worker_id: String,
39    /// Run lease TTL in seconds; `0` disables lease coordination.
40    pub lease_ttl_secs: i64,
41    /// Telemetry/runtime label (topology slug or host-provided).
42    pub runtime_label: String,
43    /// When set, poll only these pools (shared-nothing fleet pinning).
44    pub worker_pools: Option<Vec<String>>,
45    /// Delay between worker poll ticks in milliseconds (default 50).
46    pub worker_poll_interval_ms: u64,
47    /// Skip persisting run rows on the hot path (bench / throughput mode).
48    pub skip_run_persistence: bool,
49}
50
51impl WorkerSettings {
52    /// Default embedded monolith: no leases, label `embedded`.
53    #[must_use]
54    pub fn embedded() -> Self {
55        Self {
56            worker_id: resolve_worker_id_from_env(),
57            lease_ttl_secs: 0,
58            runtime_label: "embedded".into(),
59            worker_pools: resolve_worker_pools_from_env(),
60            worker_poll_interval_ms: resolve_worker_poll_interval_from_env(),
61            skip_run_persistence: resolve_skip_run_persistence_from_env(),
62        }
63    }
64
65    /// Build settings from optional builder overrides.
66    pub fn resolve(
67        worker_id: Option<String>,
68        lease_ttl_secs: Option<i64>,
69        runtime_label: Option<String>,
70        worker_pools: Option<Vec<String>>,
71        worker_poll_interval_ms: Option<u64>,
72    ) -> Self {
73        Self {
74            worker_id: worker_id.unwrap_or_else(resolve_worker_id_from_env),
75            lease_ttl_secs: lease_ttl_secs.unwrap_or_else(resolve_lease_ttl_from_env),
76            runtime_label: runtime_label.unwrap_or_else(|| "embedded".to_string()),
77            worker_pools: worker_pools.or_else(resolve_worker_pools_from_env),
78            worker_poll_interval_ms: worker_poll_interval_ms
79                .unwrap_or_else(resolve_worker_poll_interval_from_env),
80            skip_run_persistence: resolve_skip_run_persistence_from_env(),
81        }
82    }
83
84    /// Pools this worker should poll. When pinned, uses [`Self::worker_pools`]; otherwise backend discovery.
85    #[must_use]
86    pub fn pools_to_poll(&self, discovered: Vec<String>) -> Vec<String> {
87        match &self.worker_pools {
88            Some(pools) if !pools.is_empty() => pools.clone(),
89            _ => discovered,
90        }
91    }
92}
93
94fn resolve_worker_id_from_env() -> String {
95    std::env::var("INSTANCE_ID")
96        .or_else(|_| std::env::var("BOSON_WORKER_ID"))
97        .unwrap_or_else(|_| "boson-worker-1".to_string())
98}
99
100fn resolve_lease_ttl_from_env() -> i64 {
101    std::env::var("BOSON_LEASE_TTL_SECS")
102        .ok()
103        .and_then(|s| s.parse().ok())
104        .unwrap_or(0)
105}
106
107fn resolve_worker_pools_from_env() -> Option<Vec<String>> {
108    std::env::var("BOSON_WORKER_POOLS")
109        .ok()
110        .map(|s| {
111            s.split(',')
112                .map(str::trim)
113                .filter(|p| !p.is_empty())
114                .map(ToString::to_string)
115                .collect()
116        })
117        .filter(|pools: &Vec<String>| !pools.is_empty())
118}
119
120fn resolve_worker_poll_interval_from_env() -> u64 {
121    std::env::var("BOSON_WORKER_POLL_MS")
122        .ok()
123        .and_then(|s| s.parse().ok())
124        .unwrap_or(50)
125}
126
127fn resolve_skip_run_persistence_from_env() -> bool {
128    std::env::var("BOSON_SKIP_RUN_ROWS")
129        .ok()
130        .is_some_and(|v| matches!(v.to_ascii_lowercase().as_str(), "1" | "true" | "yes"))
131}