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}