Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
86 changes: 54 additions & 32 deletions crates/taskito-core/src/pubsub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,13 @@
//! same fan-out semantics — in particular the idempotency-key salting, which
//! silently drops deliveries if a shell gets it wrong.

use std::collections::HashMap;

use crate::error::Result;
use crate::job::{Job, NewJob};
use crate::storage::models::SubscriptionRow;
use crate::storage::Storage;

/// Per-task delivery settings a subscriber's task declared at registration.
/// Shells pass what their task registry knows; unknown tasks fall back to
/// the queue-level defaults.
/// Queue-level fallback delivery settings, used when neither the publish call
/// nor the subscription row specifies a value.
#[derive(Clone, Copy)]
pub struct DeliveryDefaults {
pub priority: i32,
Expand All @@ -22,9 +19,9 @@ pub struct DeliveryDefaults {
}

/// Everything a publish shares across the jobs it fans out to. Per-subscriber
/// routing (task name, queue) comes from the subscription registry; delivery
/// settings resolve per field as explicit publish override, then the
/// subscriber task's own defaults, then the queue defaults.
/// routing (task name, queue) and delivery settings come from the subscription
/// registry; delivery settings resolve per field as explicit publish override,
/// then the subscription row's persisted setting, then the queue default.
pub struct PublishRequest {
pub topic: String,
/// Wire-envelope payload bytes; every subscriber receives the same body.
Expand All @@ -46,7 +43,6 @@ pub struct PublishRequest {
pub result_ttl_ms: Option<i64>,
pub namespace: Option<String>,
pub queue_defaults: DeliveryDefaults,
pub task_defaults: HashMap<String, DeliveryDefaults>,
}

/// Fan a message out to every active subscription of `topic` as ordinary
Expand Down Expand Up @@ -85,19 +81,28 @@ fn salted_unique_key(key: &str, topic: &str, subscription_name: &str) -> String
}

fn delivery_job(request: &PublishRequest, sub: &SubscriptionRow) -> NewJob {
let task = request
.task_defaults
.get(&sub.task_name)
.copied()
.unwrap_or(request.queue_defaults);
// Resolve each field independently: explicit publish override, then the
// subscription's persisted setting, then the queue default. Persisting on
// the row is what lets a producer-only process apply a subscriber's own
// retry/timeout/priority without ever loading its task code.
let defaults = request.queue_defaults;
NewJob {
queue: sub.queue.clone(),
task_name: sub.task_name.clone(),
payload: request.payload.clone(),
priority: request.priority.unwrap_or(task.priority),
priority: request
.priority
.or(sub.priority)
.unwrap_or(defaults.priority),
scheduled_at: request.scheduled_at,
max_retries: request.max_retries.unwrap_or(task.max_retries),
timeout_ms: request.timeout_ms.unwrap_or(task.timeout_ms),
max_retries: request
.max_retries
.or(sub.max_retries)
.unwrap_or(defaults.max_retries),
timeout_ms: request
.timeout_ms
.or(sub.timeout_ms)
.unwrap_or(defaults.timeout_ms),
unique_key: request
.idempotency_key
.as_ref()
Expand Down Expand Up @@ -158,11 +163,23 @@ mod tests {
max_retries: 3,
timeout_ms: 300_000,
},
task_defaults: HashMap::new(),
}
}

fn subscribe(storage: &SqliteStorage, topic: &str, name: &str, task: &str) {
subscribe_with(storage, topic, name, task, None, None, None);
}

#[allow(clippy::too_many_arguments)]
fn subscribe_with(
storage: &SqliteStorage,
topic: &str,
name: &str,
task: &str,
priority: Option<i32>,
max_retries: Option<i32>,
timeout_ms: Option<i64>,
) {
storage
.register_subscription(&NewSubscriptionRow {
topic,
Expand All @@ -173,6 +190,9 @@ mod tests {
durable: true,
owner_worker_id: None,
created_at: now_millis(),
priority,
max_retries,
timeout_ms,
})
.unwrap();
}
Expand Down Expand Up @@ -227,28 +247,30 @@ mod tests {
}

#[test]
fn delivery_settings_resolve_override_then_task_then_queue() {
fn delivery_settings_resolve_override_then_subscription_then_queue() {
use std::collections::HashMap;
let storage = SqliteStorage::in_memory().unwrap();
subscribe(&storage, "orders", "email", "send_email");
// send_email persisted its own settings on the subscription row;
// audit_log has none → queue defaults.
subscribe_with(
&storage,
"orders",
"email",
"send_email",
Some(5),
Some(0),
Some(60_000),
);
subscribe(&storage, "orders", "audit", "audit_log");

let mut req = request("orders", None);
req.task_defaults.insert(
"send_email".to_string(),
DeliveryDefaults {
priority: 5,
max_retries: 0,
timeout_ms: 60_000,
},
);
let jobs = publish_to_topic(&storage, &req).unwrap();
let jobs = publish_to_topic(&storage, &request("orders", None)).unwrap();
let by_task: HashMap<_, _> = jobs.iter().map(|j| (j.task_name.as_str(), j)).collect();
// send_email declared its own settings; audit_log falls back to queue defaults.
assert_eq!(by_task["send_email"].max_retries, 0);
assert_eq!(by_task["send_email"].priority, 5);
assert_eq!(by_task["send_email"].timeout_ms, 60_000);
assert_eq!(by_task["audit_log"].max_retries, 3);

// An explicit publish-level override beats both.
// An explicit publish-level override beats both the row and the default.
let mut req = request("orders", None);
req.max_retries = Some(9);
let jobs = publish_to_topic(&storage, &req).unwrap();
Expand Down
9 changes: 9 additions & 0 deletions crates/taskito-core/src/storage/diesel_common/migrations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -275,6 +275,9 @@ pub fn create_tables(d: &Dialect) -> Vec<String> {
durable {bool_true},
owner_worker_id TEXT,
created_at {bi} NOT NULL,
priority INTEGER,
max_retries INTEGER,
timeout_ms {bi},
PRIMARY KEY (topic, subscription_name)
)"
),
Expand Down Expand Up @@ -387,6 +390,12 @@ pub fn alter_statements(d: &Dialect) -> Vec<String> {
format!("ALTER TABLE jobs ADD COLUMN {ife}has_deps {bool_false}"),
// DLQ auto-retry counter: tracks how many times an entry was auto-retried
format!("ALTER TABLE dead_letter ADD COLUMN {ife}dlq_retry_count INTEGER NOT NULL DEFAULT 0"),
// Per-subscription delivery settings: let publish_to_topic apply each
// subscriber's own retry/timeout/priority even from a producer process
// that never loaded the subscriber task. NULL = fall back to queue defaults.
format!("ALTER TABLE topic_subscriptions ADD COLUMN {ife}priority INTEGER"),
format!("ALTER TABLE topic_subscriptions ADD COLUMN {ife}max_retries INTEGER"),
format!("ALTER TABLE topic_subscriptions ADD COLUMN {ife}timeout_ms {bi}"),
]
}

Expand Down
8 changes: 8 additions & 0 deletions crates/taskito-core/src/storage/models.rs
Original file line number Diff line number Diff line change
Expand Up @@ -421,6 +421,11 @@ pub struct SubscriptionRow {
pub durable: bool,
pub owner_worker_id: Option<String>,
pub created_at: i64,
/// Per-subscription delivery settings persisted at registration so
/// `publish_to_topic` applies them cross-process. `None` = queue default.
pub priority: Option<i32>,
pub max_retries: Option<i32>,
pub timeout_ms: Option<i64>,
}

/// Insertable/updatable struct for subscription registrations.
Expand All @@ -435,6 +440,9 @@ pub struct NewSubscriptionRow<'a> {
pub durable: bool,
pub owner_worker_id: Option<&'a str>,
pub created_at: i64,
pub priority: Option<i32>,
pub max_retries: Option<i32>,
pub timeout_ms: Option<i64>,
}

// ── Distributed Locks ───────────────────────────────────────────
Expand Down
3 changes: 3 additions & 0 deletions crates/taskito-core/src/storage/postgres/pubsub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,9 @@ impl PostgresStorage {
topic_subscriptions::queue.eq(sub.queue),
topic_subscriptions::durable.eq(sub.durable),
topic_subscriptions::owner_worker_id.eq(sub.owner_worker_id),
topic_subscriptions::priority.eq(sub.priority),
topic_subscriptions::max_retries.eq(sub.max_retries),
topic_subscriptions::timeout_ms.eq(sub.timeout_ms),
))
.execute(&mut conn)?;

Expand Down
14 changes: 14 additions & 0 deletions crates/taskito-core/src/storage/redis_backend/pubsub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,14 @@ struct SubEntry {
durable: bool,
owner_worker_id: Option<String>,
created_at: i64,
// Older blobs (registered before per-subscription settings) lack these;
// default to None so they resolve to queue defaults, matching prior behavior.
#[serde(default)]
priority: Option<i32>,
#[serde(default)]
max_retries: Option<i32>,
#[serde(default)]
timeout_ms: Option<i64>,
}

impl From<SubEntry> for SubscriptionRow {
Expand All @@ -33,6 +41,9 @@ impl From<SubEntry> for SubscriptionRow {
durable: e.durable,
owner_worker_id: e.owner_worker_id,
created_at: e.created_at,
priority: e.priority,
max_retries: e.max_retries,
timeout_ms: e.timeout_ms,
}
}
}
Expand Down Expand Up @@ -84,6 +95,9 @@ impl RedisStorage {
durable: sub.durable,
owner_worker_id: sub.owner_worker_id.map(|s| s.to_string()),
created_at: sub.created_at,
priority: sub.priority,
max_retries: sub.max_retries,
timeout_ms: sub.timeout_ms,
};
let blob_key = self.key(&["sub", sub.topic, sub.subscription_name]);
let by_topic = self.key(&["subs", "by_topic", sub.topic]);
Expand Down
3 changes: 3 additions & 0 deletions crates/taskito-core/src/storage/schema.rs
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,9 @@ diesel::table! {
durable -> Bool,
owner_worker_id -> Nullable<Text>,
created_at -> BigInt,
priority -> Nullable<Integer>,
max_retries -> Nullable<Integer>,
timeout_ms -> Nullable<BigInt>,
}
}

Expand Down
3 changes: 3 additions & 0 deletions crates/taskito-core/src/storage/sqlite/pubsub.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,9 @@ impl SqliteStorage {
topic_subscriptions::queue.eq(sub.queue),
topic_subscriptions::durable.eq(sub.durable),
topic_subscriptions::owner_worker_id.eq(sub.owner_worker_id),
topic_subscriptions::priority.eq(sub.priority),
topic_subscriptions::max_retries.eq(sub.max_retries),
topic_subscriptions::timeout_ms.eq(sub.timeout_ms),
))
.execute(&mut conn)?;

Expand Down
3 changes: 3 additions & 0 deletions crates/taskito-core/src/storage/sqlite/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1385,6 +1385,9 @@ fn make_sub<'a>(
durable: owner.is_none(),
owner_worker_id: owner,
created_at,
priority: None,
max_retries: None,
timeout_ms: None,
}
}

Expand Down
3 changes: 3 additions & 0 deletions crates/taskito-core/tests/rust/storage_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -870,6 +870,9 @@ fn test_topic_subscriptions_crud(s: &impl Storage) {
durable: owner.is_none(),
owner_worker_id: owner,
created_at,
priority: None,
max_retries: None,
timeout_ms: None,
};

// Upsert idempotency: re-registering (topic, name) updates in place.
Expand Down
37 changes: 8 additions & 29 deletions crates/taskito-java/src/convert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,6 @@
//! Option and filter structs cross as JSON strings (decoded here); opaque job
//! payloads cross as raw `byte[]` and are never interpreted by the core.

use std::collections::HashMap;

use serde::{Deserialize, Serialize};
use taskito_core::job::{now_millis, Job, NewJob};
use taskito_core::pubsub::{DeliveryDefaults, PublishRequest};
Expand Down Expand Up @@ -103,18 +101,6 @@ pub struct PublishOptions {
pub timeout_ms: Option<i64>,
pub expires_ms: Option<i64>,
pub result_ttl_ms: Option<i64>,
/// Per-task delivery settings from the SDK's task registry, keyed by task name.
pub task_defaults: Option<HashMap<String, TaskDeliveryDefaults>>,
}

/// A subscriber task's own delivery settings. Absent fields mirror
/// [`build_new_job`]'s zero defaults so a delivery resolves like a plain enqueue.
#[derive(Deserialize, Default, Clone, Copy)]
#[serde(rename_all = "camelCase", default)]
pub struct TaskDeliveryDefaults {
pub priority: i32,
pub max_retries: i32,
pub timeout_ms: i64,
}

/// Build a core [`PublishRequest`] from a publish call. The queue-level
Expand Down Expand Up @@ -146,21 +132,6 @@ pub fn build_publish_request(
max_retries: 0,
timeout_ms: 0,
},
task_defaults: options
.task_defaults
.unwrap_or_default()
.into_iter()
.map(|(name, task)| {
(
name,
DeliveryDefaults {
priority: task.priority,
max_retries: task.max_retries,
timeout_ms: task.timeout_ms,
},
)
})
.collect(),
}
}

Expand Down Expand Up @@ -326,6 +297,8 @@ pub struct WorkerOptions {
}

/// One topic subscription declared by the SDK, registered at worker start.
/// Delivery settings (from the subscriber task's own config) are persisted on
/// the row so a producer-only process applies them without loading the task.
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct SubscriptionSpec {
Expand All @@ -334,6 +307,12 @@ pub struct SubscriptionSpec {
pub task_name: String,
pub queue: String,
pub durable: bool,
#[serde(default)]
pub priority: Option<i32>,
#[serde(default)]
pub max_retries: Option<i32>,
#[serde(default)]
pub timeout_ms: Option<i64>,
}

/// A task's retry-backoff curve. Fields left unset fall back to the core's
Expand Down
Loading