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}