1use std::future::Future;
8use std::pin::Pin;
9
10use boson_core::{ExecutionContext, IdempotencyMode, RateLimitPolicy, Result, RetryPolicy};
11use serde_json::Value;
12
13pub type InvokeFn = fn(
15 Box<dyn ExecutionContext>,
16 Value,
17) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
18
19#[derive(Debug, Clone, Copy)]
21pub struct TaskDefaults {
22 pub priority: i32,
24 pub pool: &'static str,
26 pub retry: RetryPolicy,
28 pub rate: RateLimitPolicy,
30}
31
32impl TaskDefaults {
33 #[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#[derive(Clone, Copy)]
71pub struct TaskDescriptor {
72 pub name: &'static str,
74 pub invoke: InvokeFn,
76 pub signature_json: &'static str,
78 pub signature_hash: u64,
80 pub default_priority: i32,
82 pub default_pool: &'static str,
84 pub default_retry_max_attempts: u32,
86 pub default_retry_base_delay_ms: u64,
88 pub default_retry_backoff_multiplier: f64,
90 pub default_retry_max_delay_ms: u64,
92 pub default_rate_max_in_flight: u32,
94 pub default_rate_max_enqueue_per_second: u32,
96 pub default_idempotency_mode: Option<IdempotencyMode>,
98}
99
100impl TaskDescriptor {
101 pub const fn new(name: &'static str, invoke: InvokeFn) -> Self {
103 Self::with_defaults(name, invoke, "{}", 0, TaskDefaults::standard())
104 }
105
106 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 #[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 #[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}