reinhardt-tasks 0.3.2

Background task execution and scheduling
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
//! Settings fragments for task queues, workers, webhooks, and broker backends.
//!
//! These fragments are the settings-first configuration entry points for the
//! task system. Each maps to a `[tasks_*]` TOML section and can be composed
//! into a project's settings with the `#[settings]` macro. Conversions into the
//! deprecated compatibility `XxxConfig` types are provided for the migration
//! window; new code should prefer the fragments and the
//! `create_*_from_settings` constructors.

#![allow(deprecated)] // Conversions target legacy config types during the compatibility window.

use std::collections::HashMap;
use std::time::Duration;

use reinhardt_core::macros::settings;
use serde::{Deserialize, Serialize};

use crate::webhook::{HttpWebhookSender, RetryConfig, WebhookConfig};
use crate::worker::{Worker, WorkerConfig};

// --- defaults -------------------------------------------------------------

fn default_queue_name() -> String {
	"default".to_string()
}
fn default_max_retries() -> u32 {
	3
}
fn default_worker_name() -> String {
	"worker".to_string()
}
fn default_concurrency() -> usize {
	4
}
fn default_poll_interval_ms() -> u64 {
	1000
}
fn default_webhook_method() -> String {
	"POST".to_string()
}
fn default_webhook_timeout_secs() -> u64 {
	5
}
fn default_retry_max_retries() -> u32 {
	3
}
fn default_retry_initial_backoff_ms() -> u64 {
	100
}
fn default_retry_max_backoff_ms() -> u64 {
	30_000
}
fn default_retry_backoff_multiplier() -> f64 {
	2.0
}

// --- queue ----------------------------------------------------------------

/// Task queue settings fragment.
///
/// Maps to the `[tasks_queue]` section. This fragment defines the queue
/// configuration section; applying `name` / `max_retries` to a running
/// queue is deferred to the post-deprecation queue model, because the
/// current `TaskQueue` is a stateless, zero-sized delegator. See
/// reinhardt-web#5068 for the rationale.
#[settings(fragment = true, section = "tasks_queue")]
#[non_exhaustive]
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct QueueSettings {
	/// The name of the queue.
	#[serde(default = "default_queue_name")]
	pub name: String,
	/// Maximum number of retry attempts for failed tasks.
	#[serde(default = "default_max_retries")]
	pub max_retries: u32,
}

impl Default for QueueSettings {
	fn default() -> Self {
		Self {
			name: default_queue_name(),
			max_retries: default_max_retries(),
		}
	}
}

// --- webhook --------------------------------------------------------------

/// Retry policy value object embedded in [`WebhookSettings`].
///
/// This is not an independently loadable section; it is nested under
/// `[tasks_webhook.retry]`.
#[settings(fragment = true)]
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct WebhookRetrySettings {
	/// Maximum number of retry attempts.
	#[serde(default = "default_retry_max_retries")]
	pub max_retries: u32,
	/// Initial backoff between retries, in milliseconds.
	#[serde(default = "default_retry_initial_backoff_ms")]
	pub initial_backoff_ms: u64,
	/// Maximum backoff between retries, in milliseconds.
	#[serde(default = "default_retry_max_backoff_ms")]
	pub max_backoff_ms: u64,
	/// Backoff multiplier for exponential backoff.
	#[serde(default = "default_retry_backoff_multiplier")]
	pub backoff_multiplier: f64,
}

impl Default for WebhookRetrySettings {
	fn default() -> Self {
		Self {
			max_retries: default_retry_max_retries(),
			initial_backoff_ms: default_retry_initial_backoff_ms(),
			max_backoff_ms: default_retry_max_backoff_ms(),
			backoff_multiplier: default_retry_backoff_multiplier(),
		}
	}
}

impl From<&WebhookRetrySettings> for RetryConfig {
	fn from(settings: &WebhookRetrySettings) -> Self {
		Self {
			max_retries: settings.max_retries,
			initial_backoff: Duration::from_millis(settings.initial_backoff_ms),
			max_backoff: Duration::from_millis(settings.max_backoff_ms),
			backoff_multiplier: settings.backoff_multiplier,
		}
	}
}

/// Webhook delivery settings fragment.
///
/// Maps to the `[tasks_webhook]` section.
#[settings(fragment = true, section = "tasks_webhook")]
#[non_exhaustive]
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct WebhookSettings {
	/// Target URL for webhook delivery.
	#[serde(default)]
	pub url: String,
	/// HTTP method to use.
	#[serde(default = "default_webhook_method")]
	pub method: String,
	/// Additional headers to include with each request.
	#[serde(default)]
	pub headers: HashMap<String, String>,
	/// Request timeout, in seconds.
	#[serde(default = "default_webhook_timeout_secs")]
	pub timeout_secs: u64,
	/// Retry policy.
	#[setting(node)]
	#[serde(default)]
	pub retry: WebhookRetrySettings,
}

impl Default for WebhookSettings {
	fn default() -> Self {
		Self {
			url: String::new(),
			method: default_webhook_method(),
			headers: HashMap::new(),
			timeout_secs: default_webhook_timeout_secs(),
			retry: WebhookRetrySettings::default(),
		}
	}
}

impl From<&WebhookSettings> for WebhookConfig {
	fn from(settings: &WebhookSettings) -> Self {
		Self {
			url: settings.url.clone(),
			method: settings.method.clone(),
			headers: settings.headers.clone(),
			timeout: Duration::from_secs(settings.timeout_secs),
			retry_config: RetryConfig::from(&settings.retry),
		}
	}
}

/// Build an [`HttpWebhookSender`] from a [`WebhookSettings`] fragment.
pub fn create_webhook_sender_from_settings(settings: &WebhookSettings) -> HttpWebhookSender {
	HttpWebhookSender::new(WebhookConfig::from(settings))
}

// --- worker ---------------------------------------------------------------

/// Task worker settings fragment.
///
/// Maps to the `[tasks_worker]` section.
#[settings(fragment = true, section = "tasks_worker")]
#[non_exhaustive]
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct WorkerSettings {
	/// Name of this worker instance.
	#[serde(default = "default_worker_name")]
	pub name: String,
	/// Number of concurrent task handlers.
	#[serde(default = "default_concurrency")]
	pub concurrency: usize,
	/// How long to wait between queue polls, in milliseconds.
	#[serde(default = "default_poll_interval_ms")]
	pub poll_interval_ms: u64,
	/// Webhook delivery targets for task completion notifications.
	#[setting(node)]
	#[serde(default)]
	pub webhooks: Vec<WebhookSettings>,
}

impl Default for WorkerSettings {
	fn default() -> Self {
		Self {
			name: default_worker_name(),
			concurrency: default_concurrency(),
			poll_interval_ms: default_poll_interval_ms(),
			webhooks: Vec::new(),
		}
	}
}

impl From<&WorkerSettings> for WorkerConfig {
	fn from(settings: &WorkerSettings) -> Self {
		Self {
			name: settings.name.clone(),
			concurrency: settings.concurrency,
			poll_interval: Duration::from_millis(settings.poll_interval_ms),
			webhook_configs: settings.webhooks.iter().map(WebhookConfig::from).collect(),
		}
	}
}

/// Build a [`Worker`] from a [`WorkerSettings`] fragment.
pub fn create_worker_from_settings(settings: &WorkerSettings) -> Worker {
	Worker::new(WorkerConfig::from(settings))
}

// --- sqs backend ----------------------------------------------------------

#[cfg(feature = "sqs-backend")]
fn default_sqs_visibility_timeout() -> i32 {
	30
}
#[cfg(feature = "sqs-backend")]
fn default_sqs_max_messages() -> i32 {
	1
}
#[cfg(feature = "sqs-backend")]
fn default_sqs_wait_time_seconds() -> i32 {
	0
}

/// Amazon SQS backend settings fragment.
///
/// Maps to the `[tasks_sqs]` section. Available with the `sqs-backend` feature.
#[cfg(feature = "sqs-backend")]
#[settings(fragment = true, section = "tasks_sqs")]
#[non_exhaustive]
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct SqsSettings {
	/// The SQS queue URL.
	#[serde(default)]
	pub queue_url: String,
	/// Message visibility timeout, in seconds.
	#[serde(default = "default_sqs_visibility_timeout")]
	pub visibility_timeout: i32,
	/// Maximum number of messages to receive per poll (capped at 10 by SQS).
	#[serde(default = "default_sqs_max_messages")]
	pub max_messages: i32,
	/// Wait time for long polling, in seconds.
	#[serde(default = "default_sqs_wait_time_seconds")]
	pub wait_time_seconds: i32,
}

#[cfg(feature = "sqs-backend")]
impl Default for SqsSettings {
	fn default() -> Self {
		Self {
			queue_url: String::new(),
			visibility_timeout: default_sqs_visibility_timeout(),
			max_messages: default_sqs_max_messages(),
			wait_time_seconds: default_sqs_wait_time_seconds(),
		}
	}
}

#[cfg(feature = "sqs-backend")]
impl From<&SqsSettings> for crate::backends::sqs::SqsConfig {
	fn from(settings: &SqsSettings) -> Self {
		// SqsConfig fields are private; rebuild through the builder API.
		crate::backends::sqs::SqsConfig::new(settings.queue_url.clone())
			.with_visibility_timeout(settings.visibility_timeout)
			.with_max_messages(settings.max_messages)
			.with_wait_time_seconds(settings.wait_time_seconds)
	}
}

/// Build an [`SqsBackend`](crate::backends::sqs::SqsBackend) from an
/// [`SqsSettings`] fragment.
#[cfg(feature = "sqs-backend")]
pub async fn create_sqs_backend_from_settings(
	settings: &SqsSettings,
) -> Result<crate::backends::sqs::SqsBackend, crate::TaskExecutionError> {
	crate::backends::sqs::SqsBackend::new(crate::backends::sqs::SqsConfig::from(settings)).await
}

// --- rabbitmq backend -----------------------------------------------------

#[cfg(feature = "rabbitmq-backend")]
fn default_rabbitmq_url() -> String {
	"amqp://localhost:5672/%2f".to_string()
}
#[cfg(feature = "rabbitmq-backend")]
fn default_rabbitmq_queue_name() -> String {
	"reinhardt_tasks".to_string()
}
#[cfg(feature = "rabbitmq-backend")]
fn default_rabbitmq_routing_key() -> String {
	"reinhardt_tasks".to_string()
}

/// RabbitMQ backend settings fragment.
///
/// Maps to the `[tasks_rabbitmq]` section. Available with the
/// `rabbitmq-backend` feature.
#[cfg(feature = "rabbitmq-backend")]
#[settings(fragment = true, section = "tasks_rabbitmq")]
#[non_exhaustive]
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct RabbitMQSettings {
	/// The AMQP connection URL.
	#[serde(default = "default_rabbitmq_url")]
	pub url: String,
	/// The queue name to publish to.
	#[serde(default = "default_rabbitmq_queue_name")]
	pub queue_name: String,
	/// The exchange name (empty string for the default exchange).
	#[serde(default)]
	pub exchange_name: String,
	/// The routing key.
	#[serde(default = "default_rabbitmq_routing_key")]
	pub routing_key: String,
}

#[cfg(feature = "rabbitmq-backend")]
impl Default for RabbitMQSettings {
	fn default() -> Self {
		Self {
			url: default_rabbitmq_url(),
			queue_name: default_rabbitmq_queue_name(),
			exchange_name: String::new(),
			routing_key: default_rabbitmq_routing_key(),
		}
	}
}

#[cfg(feature = "rabbitmq-backend")]
impl From<&RabbitMQSettings> for crate::backends::rabbitmq::RabbitMQConfig {
	fn from(settings: &RabbitMQSettings) -> Self {
		Self {
			url: settings.url.clone(),
			queue_name: settings.queue_name.clone(),
			exchange_name: settings.exchange_name.clone(),
			routing_key: settings.routing_key.clone(),
		}
	}
}

/// Build a [`RabbitMQBackend`](crate::backends::rabbitmq::RabbitMQBackend) from
/// a [`RabbitMQSettings`] fragment.
#[cfg(feature = "rabbitmq-backend")]
pub async fn create_rabbitmq_backend_from_settings(
	settings: &RabbitMQSettings,
) -> Result<crate::backends::rabbitmq::RabbitMQBackend, lapin::Error> {
	crate::backends::rabbitmq::RabbitMQBackend::new(
		crate::backends::rabbitmq::RabbitMQConfig::from(settings),
	)
	.await
}

#[cfg(test)]
mod tests {
	use super::*;
	use reinhardt_conf::settings::fragment::SettingsFragment;

	#[rstest::rstest]
	fn section_names_are_crate_prefixed() {
		// Arrange / Act / Assert
		assert_eq!(QueueSettings::section(), "tasks_queue");
		assert_eq!(WorkerSettings::section(), "tasks_worker");
		assert_eq!(WebhookSettings::section(), "tasks_webhook");
	}

	#[rstest::rstest]
	fn queue_settings_default_has_expected_values() {
		// Arrange
		let settings = QueueSettings::default();

		// Act / Assert
		assert_eq!(settings.name, "default");
		assert_eq!(settings.max_retries, 3);
	}

	#[rstest::rstest]
	fn worker_settings_convert_milliseconds_to_duration() {
		// Arrange
		let settings = WorkerSettings {
			name: "ingest".to_string(),
			concurrency: 8,
			poll_interval_ms: 2500,
			webhooks: Vec::new(),
		};

		// Act
		let config = WorkerConfig::from(&settings);

		// Assert
		assert_eq!(config.name, "ingest");
		assert_eq!(config.concurrency, 8);
		assert_eq!(config.poll_interval, Duration::from_millis(2500));
		assert!(config.webhook_configs.is_empty());
	}

	#[rstest::rstest]
	fn webhook_settings_convert_seconds_and_nested_retry() {
		// Arrange
		let settings = WebhookSettings::default();

		// Act
		let config = WebhookConfig::from(&settings);

		// Assert
		assert_eq!(config.method, "POST");
		assert_eq!(config.timeout, Duration::from_secs(5));
		assert_eq!(config.retry_config.max_retries, 3);
		assert_eq!(
			config.retry_config.initial_backoff,
			Duration::from_millis(100)
		);
		assert_eq!(config.retry_config.max_backoff, Duration::from_secs(30));
		assert_eq!(config.retry_config.backoff_multiplier, 2.0);
	}

	#[rstest::rstest]
	fn worker_settings_map_nested_webhooks() {
		// Arrange
		let settings = WorkerSettings {
			name: "w".to_string(),
			concurrency: 1,
			poll_interval_ms: 100,
			webhooks: vec![WebhookSettings {
				url: "https://example.com/hook".to_string(),
				..WebhookSettings::default()
			}],
		};

		// Act
		let config = WorkerConfig::from(&settings);

		// Assert
		assert_eq!(config.webhook_configs.len(), 1);
		assert_eq!(config.webhook_configs[0].url, "https://example.com/hook");
	}

	#[rstest::rstest]
	fn webhook_settings_deserialize_with_defaults() {
		// Arrange — only `url` is provided; everything else falls back to defaults.
		let json = r#"{ "url": "https://example.com/hook", "timeout_secs": 10 }"#;

		// Act
		let settings: WebhookSettings = serde_json::from_str(json).unwrap();
		let config = WebhookConfig::from(&settings);

		// Assert
		assert_eq!(config.url, "https://example.com/hook");
		assert_eq!(config.method, "POST");
		assert_eq!(config.timeout, Duration::from_secs(10));
		assert_eq!(config.retry_config.max_retries, 3);
	}

	#[cfg(feature = "sqs-backend")]
	#[rstest::rstest]
	fn sqs_settings_default_converts_to_config() {
		// Arrange
		let settings = SqsSettings {
			queue_url: "https://sqs.example.com/q".to_string(),
			..SqsSettings::default()
		};

		// Act
		let config = crate::backends::sqs::SqsConfig::from(&settings);

		// Assert — round-trips through the builder, which caps max_messages at 10.
		assert!(format!("{config:?}").contains("https://sqs.example.com/q"));
	}

	#[cfg(feature = "rabbitmq-backend")]
	#[rstest::rstest]
	fn rabbitmq_settings_default_converts_to_config() {
		// Arrange
		let settings = RabbitMQSettings::default();

		// Act
		let config = crate::backends::rabbitmq::RabbitMQConfig::from(&settings);

		// Assert
		assert_eq!(config.queue_name, "reinhardt_tasks");
		assert_eq!(config.routing_key, "reinhardt_tasks");
	}
}