Skip to main content

boson_runtime/registry/
descriptor.rs

1//! Task descriptor for auto-registration and manual registration.
2//!
3//! A [`TaskDescriptor`] is static metadata for one handler. Defaults seed the initial
4//! [`TaskConfig`](boson_core::TaskConfig) on first enqueue; admin updates can override
5//! priority, pool, retry, and rate limits at runtime.
6
7use std::future::Future;
8use std::pin::Pin;
9
10use boson_core::{ExecutionContext, IdempotencyMode, RateLimitPolicy, Result, RetryPolicy};
11use serde_json::Value;
12
13/// Invokes a registered task with execution context and JSON parameters.
14pub type InvokeFn = fn(
15    Box<dyn ExecutionContext>,
16    Value,
17) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
18
19/// Default retry, rate, priority, and pool settings for a registered task.
20#[derive(Debug, Clone, Copy)]
21pub struct TaskDefaults {
22    /// Default priority (lower = higher priority).
23    pub priority: i32,
24    /// Default pool name.
25    pub pool: &'static str,
26    /// Default retry policy.
27    pub retry: RetryPolicy,
28    /// Default enqueue rate limits.
29    pub rate: RateLimitPolicy,
30}
31
32impl TaskDefaults {
33    /// Standard production-like defaults for manual registration.
34    ///
35    /// - priority `1`, pool `"global"`
36    /// - retry: 3 attempts, 1000 ms base delay, `2.0×` multiplier, `30_000` ms cap
37    /// - rate: 100 max in-flight, 50 enqueues per second
38    #[must_use]
39    pub const fn standard() -> Self {
40        Self {
41            priority: 1,
42            pool: "global",
43            retry: RetryPolicy {
44                max_attempts: 3,
45                base_delay_ms: 1000,
46                backoff_multiplier: 2.0,
47                max_delay_ms: 30_000,
48            },
49            rate: RateLimitPolicy {
50                max_in_flight: 100,
51                max_enqueue_per_second: 50,
52            },
53        }
54    }
55}
56
57/// Descriptor for a registered task.
58///
59/// Registration defaults flow into [`TaskConfig`](boson_core::TaskConfig) on first enqueue.
60/// Use [`TaskDescriptor::with_defaults`] to set retry, rate, priority, and pool in one call.
61///
62/// ## Signature versioning
63///
64/// - [`signature_json`](Self::signature_json) — JSON schema string describing task parameters
65///   (convention for tooling; Boson does not validate against it at runtime).
66/// - [`signature_hash`](Self::signature_hash) — hash of the parameter schema/version. Stored on
67///   each [`Job`](boson_core::Job) at enqueue; if the registered task's hash changes while a job
68///   is still active, dispatch returns [`BosonError::SignatureMismatch`](boson_core::BosonError::SignatureMismatch).
69///   Bump `signature_hash` when you change parameter shape incompatibly.
70#[derive(Clone, Copy)]
71pub struct TaskDescriptor {
72    /// Unique task name (registry key and enqueue target).
73    pub name: &'static str,
74    /// Function to invoke the task.
75    pub invoke: InvokeFn,
76    /// JSON schema string for parameters (documentation / tooling; not validated at runtime).
77    pub signature_json: &'static str,
78    /// Version hash checked against enqueued jobs; change when parameters change incompatibly.
79    pub signature_hash: u64,
80    /// Default priority (lower = higher priority).
81    pub default_priority: i32,
82    /// Default pool name for worker assignment.
83    pub default_pool: &'static str,
84    /// Default retry max attempts (see [`RetryPolicy::max_attempts`]).
85    pub default_retry_max_attempts: u32,
86    /// Default retry base delay ms (see [`RetryPolicy::base_delay_ms`]).
87    pub default_retry_base_delay_ms: u64,
88    /// Default retry backoff multiplier (see [`RetryPolicy::backoff_multiplier`]).
89    pub default_retry_backoff_multiplier: f64,
90    /// Default retry max delay ms (see [`RetryPolicy::max_delay_ms`]).
91    pub default_retry_max_delay_ms: u64,
92    /// Default max in-flight jobs (see [`RateLimitPolicy::max_in_flight`]; `0` = unlimited).
93    pub default_rate_max_in_flight: u32,
94    /// Default max enqueues per second (see [`RateLimitPolicy::max_enqueue_per_second`]; `0` = unlimited).
95    pub default_rate_max_enqueue_per_second: u32,
96    /// Per-task idempotency override (`None` = inherit runtime default).
97    pub default_idempotency_mode: Option<IdempotencyMode>,
98}
99
100impl TaskDescriptor {
101    /// Minimal descriptor for tests (`signature_json` `"{}"`, `signature_hash` `0`).
102    pub const fn new(name: &'static str, invoke: InvokeFn) -> Self {
103        Self::with_defaults(name, invoke, "{}", 0, TaskDefaults::standard())
104    }
105
106    /// Descriptor with grouped policy defaults.
107    pub const fn with_defaults(
108        name: &'static str,
109        invoke: InvokeFn,
110        signature_json: &'static str,
111        signature_hash: u64,
112        defaults: TaskDefaults,
113    ) -> Self {
114        Self {
115            name,
116            invoke,
117            signature_json,
118            signature_hash,
119            default_priority: defaults.priority,
120            default_pool: defaults.pool,
121            default_retry_max_attempts: defaults.retry.max_attempts,
122            default_retry_base_delay_ms: defaults.retry.base_delay_ms,
123            default_retry_backoff_multiplier: defaults.retry.backoff_multiplier,
124            default_retry_max_delay_ms: defaults.retry.max_delay_ms,
125            default_rate_max_in_flight: defaults.rate.max_in_flight,
126            default_rate_max_enqueue_per_second: defaults.rate.max_enqueue_per_second,
127            default_idempotency_mode: None,
128        }
129    }
130
131    /// Descriptor with explicit per-field policy defaults.
132    ///
133    /// Used by the [`#[task]`](https://docs.rs/boson-macros/latest/boson_macros/attr.task.html) attribute when policy fields are set on the handler.
134    /// Prefer [`Self::with_defaults`] for manual registration in tests.
135    #[allow(clippy::too_many_arguments)]
136    pub const fn with_policy(
137        name: &'static str,
138        invoke: InvokeFn,
139        signature_json: &'static str,
140        signature_hash: u64,
141        priority: i32,
142        pool: &'static str,
143        max_attempts: u32,
144        base_delay_ms: u64,
145        backoff_multiplier: f64,
146        max_delay_ms: u64,
147        max_in_flight: u32,
148        max_enqueue_per_second: u32,
149        idempotency_mode: Option<IdempotencyMode>,
150    ) -> Self {
151        Self {
152            name,
153            invoke,
154            signature_json,
155            signature_hash,
156            default_priority: priority,
157            default_pool: pool,
158            default_retry_max_attempts: max_attempts,
159            default_retry_base_delay_ms: base_delay_ms,
160            default_retry_backoff_multiplier: backoff_multiplier,
161            default_retry_max_delay_ms: max_delay_ms,
162            default_rate_max_in_flight: max_in_flight,
163            default_rate_max_enqueue_per_second: max_enqueue_per_second,
164            default_idempotency_mode: idempotency_mode,
165        }
166    }
167
168    /// Materialize descriptor defaults into a [`TaskConfig`](boson_core::TaskConfig).
169    #[must_use]
170    pub fn to_task_config(&self) -> boson_core::TaskConfig {
171        boson_core::TaskConfig::from_policy_defaults(
172            self.name,
173            self.default_priority,
174            self.default_pool,
175            RetryPolicy {
176                max_attempts: self.default_retry_max_attempts,
177                base_delay_ms: self.default_retry_base_delay_ms,
178                backoff_multiplier: self.default_retry_backoff_multiplier,
179                max_delay_ms: self.default_retry_max_delay_ms,
180            },
181            RateLimitPolicy {
182                max_in_flight: self.default_rate_max_in_flight,
183                max_enqueue_per_second: self.default_rate_max_enqueue_per_second,
184            },
185            self.default_idempotency_mode,
186        )
187    }
188}
189
190impl std::fmt::Debug for TaskDescriptor {
191    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
192        f.debug_struct("TaskDescriptor")
193            .field("name", &self.name)
194            .field("signature_json", &self.signature_json)
195            .field("signature_hash", &self.signature_hash)
196            .field("default_priority", &self.default_priority)
197            .field("default_pool", &self.default_pool)
198            .field("default_retry_max_attempts", &self.default_retry_max_attempts)
199            .field("default_retry_base_delay_ms", &self.default_retry_base_delay_ms)
200            .field(
201                "default_retry_backoff_multiplier",
202                &self.default_retry_backoff_multiplier,
203            )
204            .field("default_retry_max_delay_ms", &self.default_retry_max_delay_ms)
205            .field("default_rate_max_in_flight", &self.default_rate_max_in_flight)
206            .field(
207                "default_rate_max_enqueue_per_second",
208                &self.default_rate_max_enqueue_per_second,
209            )
210            .field("default_idempotency_mode", &self.default_idempotency_mode)
211            .field("invoke", &"<fn>")
212            .finish()
213    }
214}
215
216quark::inventory::collect!(TaskDescriptor);
217
218impl quark::Registrable for TaskDescriptor {
219    fn registry_key(&self) -> &str {
220        self.name
221    }
222}