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` (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}