Skip to main content

boson_core/models/
task_config.rs

1//! Task config model — priority, pool, retry, rate limits, and idempotency.
2//!
3//! [`TaskConfig`] is persisted per task. On enqueue, Boson merges descriptor defaults with
4//! any stored config to set [`Job`](crate::Job) priority, pool, and policies.
5
6use chrono::{DateTime, Utc};
7use serde::{Deserialize, Serialize};
8
9/// How backends enforce enqueue idempotency when a job key is present.
10///
11/// On Scylla, [`IdempotencyMode::Lwt`] uses a lightweight transaction; [`IdempotencyMode::None`]
12/// skips that path (at-least-once). Mem/SQL honor the same policy for check-then-reuse.
13#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
14#[serde(rename_all = "snake_case")]
15pub enum IdempotencyMode {
16    /// Exactly-once under concurrent enqueue (default).
17    #[default]
18    Lwt,
19    /// At-least-once; skip idempotency coordination (higher throughput).
20    None,
21}
22
23impl IdempotencyMode {
24    /// Parse from a config string (`lwt` / `none`).
25    #[must_use]
26    pub fn parse(s: &str) -> Option<Self> {
27        match s.trim().to_ascii_lowercase().as_str() {
28            "lwt" => Some(Self::Lwt),
29            "none" => Some(Self::None),
30            _ => None,
31        }
32    }
33}
34
35/// Retry policy for a task.
36///
37/// On handler failure, the worker reschedules while `job.attempt < max_attempts`.
38/// Delay before the next attempt (milliseconds):
39///
40/// `min(base_delay_ms × backoff_multiplier^(attempt - 1), max_delay_ms)`
41///
42/// where `attempt` is the 1-based attempt number on the job.
43#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize)]
44pub struct RetryPolicy {
45    /// Maximum retry attempts after the first failure (`0` = no retries).
46    pub max_attempts: u32,
47    /// Base delay in ms before the first retry.
48    pub base_delay_ms: u64,
49    /// Exponential backoff multiplier applied per attempt.
50    pub backoff_multiplier: f64,
51    /// Upper cap on retry delay in ms.
52    pub max_delay_ms: u64,
53}
54
55/// Rate limiting for enqueue backpressure.
56///
57/// `max_in_flight == 0` and `max_enqueue_per_second == 0` mean unlimited (default).
58///
59/// When a limit is exceeded, [`Boson::enqueue`](https://docs.rs/boson-runtime/latest/boson_runtime/struct.Boson.html#method.enqueue) returns
60/// [`BosonError::RateLimited`](crate::BosonError::RateLimited); callers should retry with backoff.
61///
62/// # Example
63///
64/// ```rust
65/// use boson_core::RateLimitPolicy;
66///
67/// // At most 10 active jobs and 5 enqueues per second for one task.
68/// let policy = RateLimitPolicy {
69///     max_in_flight: 10,
70///     max_enqueue_per_second: 5,
71/// };
72/// ```
73#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize)]
74pub struct RateLimitPolicy {
75    /// Max jobs for this task in `queued` + `running` at once. `0` = no limit.
76    pub max_in_flight: u32,
77    /// Max successful enqueues per wall-clock second per process. `0` = no limit.
78    pub max_enqueue_per_second: u32,
79}
80
81/// Per-task config persisted for admin UI and enqueue defaults.
82#[derive(Debug, Clone, Serialize, Deserialize)]
83pub struct TaskConfig {
84    /// Task name (unique key).
85    pub task_name: String,
86    /// Override priority (lower = higher priority).
87    pub priority: i32,
88    /// Override pool name.
89    pub pool: String,
90    /// Retry policy override.
91    pub retry_policy: RetryPolicy,
92    /// Enqueue rate limits (optional).
93    pub rate_limit_policy: RateLimitPolicy,
94    /// Per-task idempotency override (`None` = inherit runtime / builder default).
95    #[serde(default, skip_serializing_if = "Option::is_none")]
96    pub idempotency_mode: Option<IdempotencyMode>,
97    /// When created/updated.
98    pub updated_at: DateTime<Utc>,
99}
100
101impl TaskConfig {
102    /// Create default config for a task name.
103    ///
104    /// # Examples
105    ///
106    /// ```
107    /// use boson_core::TaskConfig;
108    ///
109    /// let config = TaskConfig::default_for("notify");
110    /// assert_eq!(config.task_name, "notify");
111    /// assert_eq!(config.pool, "global");
112    /// assert_eq!(config.priority, 1);
113    /// ```
114    #[must_use]
115    pub fn default_for(task_name: &str) -> Self {
116        Self {
117            task_name: task_name.to_string(),
118            priority: 1,
119            pool: "global".to_string(),
120            retry_policy: RetryPolicy::default(),
121            rate_limit_policy: RateLimitPolicy::default(),
122            idempotency_mode: None,
123            updated_at: Utc::now(),
124        }
125    }
126
127    /// Resolve effective idempotency mode (task override, else `runtime_default`).
128    #[must_use]
129    pub fn resolved_idempotency_mode(&self, runtime_default: IdempotencyMode) -> IdempotencyMode {
130        self.idempotency_mode.unwrap_or(runtime_default)
131    }
132
133    /// Build config from explicit policy defaults (typically from a task descriptor).
134    #[must_use]
135    pub fn from_policy_defaults(
136        task_name: &str,
137        priority: i32,
138        pool: impl Into<String>,
139        retry_policy: RetryPolicy,
140        rate_limit_policy: RateLimitPolicy,
141        idempotency_mode: Option<IdempotencyMode>,
142    ) -> Self {
143        Self {
144            task_name: task_name.to_string(),
145            priority,
146            pool: pool.into(),
147            retry_policy,
148            rate_limit_policy,
149            idempotency_mode,
150            updated_at: Utc::now(),
151        }
152    }
153
154    /// Fill [`idempotency_mode`](Self::idempotency_mode) from the runtime default when unset.
155    #[must_use]
156    pub const fn with_runtime_idempotency_fallback(
157        mut self,
158        runtime_default: IdempotencyMode,
159    ) -> Self {
160        if self.idempotency_mode.is_none() {
161            self.idempotency_mode = Some(runtime_default);
162        }
163        self
164    }
165}
166
167#[cfg(test)]
168mod tests {
169    use super::*;
170
171    #[test]
172    fn from_policy_defaults_sets_fields() {
173        let config = TaskConfig::from_policy_defaults(
174            "notify",
175            2,
176            "urgent",
177            RetryPolicy {
178                max_attempts: 5,
179                base_delay_ms: 100,
180                backoff_multiplier: 2.0,
181                max_delay_ms: 1000,
182            },
183            RateLimitPolicy {
184                max_in_flight: 10,
185                max_enqueue_per_second: 3,
186            },
187            Some(IdempotencyMode::None),
188        );
189        assert_eq!(config.task_name, "notify");
190        assert_eq!(config.priority, 2);
191        assert_eq!(config.pool, "urgent");
192        assert_eq!(config.retry_policy.max_attempts, 5);
193        assert_eq!(config.rate_limit_policy.max_in_flight, 10);
194        assert_eq!(config.idempotency_mode, Some(IdempotencyMode::None));
195    }
196
197    #[test]
198    fn runtime_idempotency_fallback_fills_none_only() {
199        let filled = TaskConfig::default_for("t").with_runtime_idempotency_fallback(IdempotencyMode::Lwt);
200        assert_eq!(filled.idempotency_mode, Some(IdempotencyMode::Lwt));
201
202        let kept = TaskConfig::default_for("t");
203        let mut kept = kept;
204        kept.idempotency_mode = Some(IdempotencyMode::None);
205        let kept = kept.with_runtime_idempotency_fallback(IdempotencyMode::Lwt);
206        assert_eq!(kept.idempotency_mode, Some(IdempotencyMode::None));
207    }
208}