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