use std::future::Future;
use std::pin::Pin;
use boson_core::{ExecutionContext, IdempotencyMode, RateLimitPolicy, Result, RetryPolicy};
use serde_json::Value;
pub type InvokeFn = fn(
Box<dyn ExecutionContext>,
Value,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send + 'static>>;
#[derive(Debug, Clone, Copy)]
pub struct TaskDefaults {
pub priority: i32,
pub pool: &'static str,
pub retry: RetryPolicy,
pub rate: RateLimitPolicy,
}
impl TaskDefaults {
#[must_use]
pub const fn standard() -> Self {
Self {
priority: 1,
pool: "global",
retry: RetryPolicy {
max_attempts: 3,
base_delay_ms: 1000,
backoff_multiplier: 2.0,
max_delay_ms: 30_000,
},
rate: RateLimitPolicy {
max_in_flight: 100,
max_enqueue_per_second: 50,
},
}
}
}
#[derive(Clone, Copy)]
pub struct TaskDescriptor {
pub name: &'static str,
pub invoke: InvokeFn,
pub signature_json: &'static str,
pub signature_hash: u64,
pub default_priority: i32,
pub default_pool: &'static str,
pub default_retry_max_attempts: u32,
pub default_retry_base_delay_ms: u64,
pub default_retry_backoff_multiplier: f64,
pub default_retry_max_delay_ms: u64,
pub default_rate_max_in_flight: u32,
pub default_rate_max_enqueue_per_second: u32,
pub default_idempotency_mode: Option<IdempotencyMode>,
}
impl TaskDescriptor {
pub const fn new(name: &'static str, invoke: InvokeFn) -> Self {
Self::with_defaults(name, invoke, "{}", 0, TaskDefaults::standard())
}
pub const fn with_defaults(
name: &'static str,
invoke: InvokeFn,
signature_json: &'static str,
signature_hash: u64,
defaults: TaskDefaults,
) -> Self {
Self {
name,
invoke,
signature_json,
signature_hash,
default_priority: defaults.priority,
default_pool: defaults.pool,
default_retry_max_attempts: defaults.retry.max_attempts,
default_retry_base_delay_ms: defaults.retry.base_delay_ms,
default_retry_backoff_multiplier: defaults.retry.backoff_multiplier,
default_retry_max_delay_ms: defaults.retry.max_delay_ms,
default_rate_max_in_flight: defaults.rate.max_in_flight,
default_rate_max_enqueue_per_second: defaults.rate.max_enqueue_per_second,
default_idempotency_mode: None,
}
}
#[allow(clippy::too_many_arguments)]
pub const fn with_policy(
name: &'static str,
invoke: InvokeFn,
signature_json: &'static str,
signature_hash: u64,
priority: i32,
pool: &'static str,
max_attempts: u32,
base_delay_ms: u64,
backoff_multiplier: f64,
max_delay_ms: u64,
max_in_flight: u32,
max_enqueue_per_second: u32,
idempotency_mode: Option<IdempotencyMode>,
) -> Self {
Self {
name,
invoke,
signature_json,
signature_hash,
default_priority: priority,
default_pool: pool,
default_retry_max_attempts: max_attempts,
default_retry_base_delay_ms: base_delay_ms,
default_retry_backoff_multiplier: backoff_multiplier,
default_retry_max_delay_ms: max_delay_ms,
default_rate_max_in_flight: max_in_flight,
default_rate_max_enqueue_per_second: max_enqueue_per_second,
default_idempotency_mode: idempotency_mode,
}
}
#[must_use]
pub fn to_task_config(&self) -> boson_core::TaskConfig {
boson_core::TaskConfig::from_policy_defaults(
self.name,
self.default_priority,
self.default_pool,
RetryPolicy {
max_attempts: self.default_retry_max_attempts,
base_delay_ms: self.default_retry_base_delay_ms,
backoff_multiplier: self.default_retry_backoff_multiplier,
max_delay_ms: self.default_retry_max_delay_ms,
},
RateLimitPolicy {
max_in_flight: self.default_rate_max_in_flight,
max_enqueue_per_second: self.default_rate_max_enqueue_per_second,
},
self.default_idempotency_mode,
)
}
}
impl std::fmt::Debug for TaskDescriptor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TaskDescriptor")
.field("name", &self.name)
.field("signature_json", &self.signature_json)
.field("signature_hash", &self.signature_hash)
.field("default_priority", &self.default_priority)
.field("default_pool", &self.default_pool)
.field("default_retry_max_attempts", &self.default_retry_max_attempts)
.field("default_retry_base_delay_ms", &self.default_retry_base_delay_ms)
.field(
"default_retry_backoff_multiplier",
&self.default_retry_backoff_multiplier,
)
.field("default_retry_max_delay_ms", &self.default_retry_max_delay_ms)
.field("default_rate_max_in_flight", &self.default_rate_max_in_flight)
.field(
"default_rate_max_enqueue_per_second",
&self.default_rate_max_enqueue_per_second,
)
.field("default_idempotency_mode", &self.default_idempotency_mode)
.field("invoke", &"<fn>")
.finish()
}
}
quark::inventory::collect!(TaskDescriptor);
impl quark::Registrable for TaskDescriptor {
fn registry_key(&self) -> &str {
self.name
}
}