aster_forge_db/
scheduled_task.rs

1//! Database-backed scheduled task catalog and claim store.
2//!
3//! Scheduled task rows persist product runtime schedules across process restarts and coordinate
4//! due-work claims across service instances. `aster_forge_tasks` owns the public scheduling DTOs
5//! and runner trait; this module only supplies the `SeaORM` table contract and store implementation.
6//! Product crates still own task names, intervals, execution bodies, and outcome records.
7
8use std::time::Duration;
9
10use sea_orm::entity::prelude::*;
11use sea_orm::sea_query::{
12    Alias, ColumnDef, Index, IndexCreateStatement, Table, TableCreateStatement, TableDropStatement,
13};
14use sea_orm::{
15    ActiveModelTrait, ColumnTrait, Condition, DatabaseBackend, DatabaseConnection, EntityTrait,
16    QueryFilter, Set, sea_query::Expr,
17};
18
19use crate::DbError;
20use aster_forge_tasks::{
21    ScheduledTaskCatalogEntry, ScheduledTaskClaim, ScheduledTaskClaimRenewal,
22    ScheduledTaskClaimRequest, ScheduledTaskCompletion,
23};
24
25/// Scheduled task table name.
26pub const SCHEDULED_TASKS_TABLE: &str = "scheduled_tasks";
27/// Stable row identifier column.
28pub const SCHEDULED_TASK_ID_COLUMN: &str = "task_id";
29/// Product namespace column.
30pub const SCHEDULED_TASK_NAMESPACE_COLUMN: &str = "namespace";
31/// Product task name column.
32pub const SCHEDULED_TASK_NAME_COLUMN: &str = "task_name";
33/// Operator-facing display name column.
34pub const SCHEDULED_TASK_DISPLAY_NAME_COLUMN: &str = "display_name";
35/// Next due timestamp column.
36pub const SCHEDULED_TASK_NEXT_RUN_AT_COLUMN: &str = "next_run_at";
37/// Current claim owner column.
38pub const SCHEDULED_TASK_CLAIM_OWNER_ID_COLUMN: &str = "claim_owner_id";
39/// Current claim expiry column.
40pub const SCHEDULED_TASK_CLAIM_EXPIRES_AT_COLUMN: &str = "claim_expires_at";
41/// Last claim timestamp column.
42pub const SCHEDULED_TASK_LAST_CLAIMED_AT_COLUMN: &str = "last_claimed_at";
43/// Last completion timestamp column.
44pub const SCHEDULED_TASK_LAST_FINISHED_AT_COLUMN: &str = "last_finished_at";
45/// Row creation timestamp column.
46pub const SCHEDULED_TASK_CREATED_AT_COLUMN: &str = "created_at";
47/// Row update timestamp column.
48pub const SCHEDULED_TASK_UPDATED_AT_COLUMN: &str = "updated_at";
49/// Unique index name for one task per namespace/name pair.
50pub const SCHEDULED_TASK_NAMESPACE_NAME_UNIQUE_INDEX: &str =
51    "idx_scheduled_tasks_namespace_name_unique";
52/// Index name for due-time claim scans.
53pub const SCHEDULED_TASK_NEXT_RUN_INDEX: &str = "idx_scheduled_tasks_next_run";
54
55const SCHEDULED_TASK_ID_MAX_LEN: usize = 191;
56const SCHEDULED_TASK_NAMESPACE_MAX_LEN: usize = 64;
57const SCHEDULED_TASK_NAME_MAX_LEN: usize = 128;
58const SCHEDULED_TASK_DISPLAY_NAME_MAX_LEN: usize = 191;
59const SCHEDULED_TASK_OWNER_ID_MAX_LEN: usize = 191;
60
61/// Builds the shared `scheduled_tasks` table creation statement.
62#[must_use]
63pub fn create_scheduled_tasks_table(backend: DatabaseBackend) -> TableCreateStatement {
64    Table::create()
65        .table(scheduled_tasks_table())
66        .if_not_exists()
67        .col(
68            ColumnDef::new(scheduled_task_id())
69                .string_len(191)
70                .not_null()
71                .primary_key(),
72        )
73        .col(
74            ColumnDef::new(scheduled_task_namespace())
75                .string_len(64)
76                .not_null(),
77        )
78        .col(
79            ColumnDef::new(scheduled_task_name())
80                .string_len(128)
81                .not_null(),
82        )
83        .col(
84            ColumnDef::new(scheduled_task_display_name())
85                .string_len(191)
86                .not_null(),
87        )
88        .col(utc_datetime_column(backend, scheduled_task_next_run_at()).not_null())
89        .col(
90            ColumnDef::new(scheduled_task_claim_owner_id())
91                .string_len(191)
92                .null(),
93        )
94        .col(utc_datetime_column(backend, scheduled_task_claim_expires_at()).null())
95        .col(utc_datetime_column(backend, scheduled_task_last_claimed_at()).null())
96        .col(utc_datetime_column(backend, scheduled_task_last_finished_at()).null())
97        .col(utc_datetime_column(backend, scheduled_task_created_at()).not_null())
98        .col(utc_datetime_column(backend, scheduled_task_updated_at()).not_null())
99        .to_owned()
100}
101
102/// Builds the shared `scheduled_tasks` table drop statement.
103#[must_use]
104pub fn drop_scheduled_tasks_table() -> TableDropStatement {
105    Table::drop()
106        .table(scheduled_tasks_table())
107        .if_exists()
108        .to_owned()
109}
110
111/// Builds the unique index for one scheduled task per namespace/name pair.
112#[must_use]
113pub fn create_scheduled_tasks_namespace_name_unique_index() -> IndexCreateStatement {
114    Index::create()
115        .name(SCHEDULED_TASK_NAMESPACE_NAME_UNIQUE_INDEX)
116        .table(scheduled_tasks_table())
117        .col(scheduled_task_namespace())
118        .col(scheduled_task_name())
119        .unique()
120        .if_not_exists()
121        .to_owned()
122}
123
124/// Builds the due-time index used by scheduled task claim checks.
125#[must_use]
126pub fn create_scheduled_tasks_next_run_index() -> IndexCreateStatement {
127    Index::create()
128        .name(SCHEDULED_TASK_NEXT_RUN_INDEX)
129        .table(scheduled_tasks_table())
130        .col(scheduled_task_next_run_at())
131        .if_not_exists()
132        .to_owned()
133}
134
135fn scheduled_tasks_table() -> Alias {
136    Alias::new(SCHEDULED_TASKS_TABLE)
137}
138
139fn scheduled_task_id() -> Alias {
140    Alias::new(SCHEDULED_TASK_ID_COLUMN)
141}
142
143fn scheduled_task_namespace() -> Alias {
144    Alias::new(SCHEDULED_TASK_NAMESPACE_COLUMN)
145}
146
147fn scheduled_task_name() -> Alias {
148    Alias::new(SCHEDULED_TASK_NAME_COLUMN)
149}
150
151fn scheduled_task_display_name() -> Alias {
152    Alias::new(SCHEDULED_TASK_DISPLAY_NAME_COLUMN)
153}
154
155fn scheduled_task_next_run_at() -> Alias {
156    Alias::new(SCHEDULED_TASK_NEXT_RUN_AT_COLUMN)
157}
158
159fn scheduled_task_claim_owner_id() -> Alias {
160    Alias::new(SCHEDULED_TASK_CLAIM_OWNER_ID_COLUMN)
161}
162
163fn scheduled_task_claim_expires_at() -> Alias {
164    Alias::new(SCHEDULED_TASK_CLAIM_EXPIRES_AT_COLUMN)
165}
166
167fn scheduled_task_last_claimed_at() -> Alias {
168    Alias::new(SCHEDULED_TASK_LAST_CLAIMED_AT_COLUMN)
169}
170
171fn scheduled_task_last_finished_at() -> Alias {
172    Alias::new(SCHEDULED_TASK_LAST_FINISHED_AT_COLUMN)
173}
174
175fn scheduled_task_created_at() -> Alias {
176    Alias::new(SCHEDULED_TASK_CREATED_AT_COLUMN)
177}
178
179fn scheduled_task_updated_at() -> Alias {
180    Alias::new(SCHEDULED_TASK_UPDATED_AT_COLUMN)
181}
182
183fn utc_datetime_column(backend: DatabaseBackend, column: Alias) -> ColumnDef {
184    let mut definition = ColumnDef::new(column);
185    match backend {
186        DatabaseBackend::MySql => {
187            definition.custom(Alias::new("datetime(6)"));
188        }
189        _ => {
190            definition.timestamp_with_time_zone();
191        }
192    }
193    definition
194}
195
196/// `SeaORM` model for `scheduled_tasks`.
197#[derive(Clone, Debug, PartialEq, DeriveEntityModel)]
198#[sea_orm(table_name = "scheduled_tasks")]
199pub struct Model {
200    /// Stable row identifier built from namespace and task name.
201    #[sea_orm(primary_key, auto_increment = false)]
202    pub task_id: String,
203    /// Product namespace.
204    pub namespace: String,
205    /// Stable product task name.
206    pub task_name: String,
207    /// Operator-facing display name.
208    pub display_name: String,
209    /// Next due timestamp.
210    pub next_run_at: DateTimeUtc,
211    /// Runtime owner currently claiming this due run.
212    pub claim_owner_id: Option<String>,
213    /// Timestamp after which another runtime may reclaim this due run.
214    pub claim_expires_at: Option<DateTimeUtc>,
215    /// Timestamp of the last successful claim.
216    pub last_claimed_at: Option<DateTimeUtc>,
217    /// Timestamp of the last successful completion.
218    pub last_finished_at: Option<DateTimeUtc>,
219    /// Row creation timestamp.
220    pub created_at: DateTimeUtc,
221    /// Row update timestamp.
222    pub updated_at: DateTimeUtc,
223}
224
225#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
226pub enum Relation {}
227
228impl ActiveModelBehavior for ActiveModel {}
229
230/// SeaORM-backed scheduled task store.
231#[derive(Clone)]
232pub struct ScheduledTaskDbStore {
233    db: DatabaseConnection,
234}
235
236#[async_trait::async_trait]
237impl aster_forge_tasks::ScheduledTaskStore for ScheduledTaskDbStore {
238    type Error = DbError;
239
240    async fn ensure_scheduled_task(
241        &self,
242        entry: ScheduledTaskCatalogEntry<'_>,
243    ) -> std::result::Result<(), Self::Error> {
244        self.ensure_task(entry).await.map(|_| ())
245    }
246
247    async fn claim_scheduled_task(
248        &self,
249        request: ScheduledTaskClaimRequest<'_>,
250    ) -> std::result::Result<Option<ScheduledTaskClaim>, Self::Error> {
251        self.claim_due(request).await
252    }
253
254    async fn renew_scheduled_task_claim(
255        &self,
256        renewal: ScheduledTaskClaimRenewal<'_>,
257    ) -> std::result::Result<bool, Self::Error> {
258        self.renew_claim(renewal).await
259    }
260
261    async fn complete_scheduled_task(
262        &self,
263        completion: ScheduledTaskCompletion,
264    ) -> std::result::Result<bool, Self::Error> {
265        self.complete_claim(completion).await
266    }
267}
268
269impl ScheduledTaskDbStore {
270    /// Creates a scheduled task store from a `SeaORM` database connection.
271    #[must_use]
272    pub const fn new(db: DatabaseConnection) -> Self {
273        Self { db }
274    }
275
276    /// Ensures one product scheduled task is present in the catalog.
277    ///
278    /// # Errors
279    ///
280    /// Returns an error when the database operation fails.
281    pub async fn ensure_task(&self, entry: ScheduledTaskCatalogEntry<'_>) -> crate::Result<Model> {
282        ensure_task(&self.db, entry).await
283    }
284
285    /// Attempts to claim one due scheduled task firing.
286    ///
287    /// # Errors
288    ///
289    /// Returns an error when the database operation fails.
290    pub async fn claim_due(
291        &self,
292        request: ScheduledTaskClaimRequest<'_>,
293    ) -> crate::Result<Option<ScheduledTaskClaim>> {
294        claim_due(&self.db, request).await
295    }
296
297    /// Renews an owned claim while the task body is still running.
298    ///
299    /// # Errors
300    ///
301    /// Returns an error when the database operation fails.
302    pub async fn renew_claim(&self, renewal: ScheduledTaskClaimRenewal<'_>) -> crate::Result<bool> {
303        renew_claim(&self.db, renewal).await
304    }
305
306    /// Completes a claimed firing and advances the next due timestamp.
307    ///
308    /// # Errors
309    ///
310    /// Returns an error when the database operation fails.
311    pub async fn complete_claim(&self, completion: ScheduledTaskCompletion) -> crate::Result<bool> {
312        complete_claim(&self.db, completion).await
313    }
314}
315
316async fn ensure_task(
317    db: &DatabaseConnection,
318    entry: ScheduledTaskCatalogEntry<'_>,
319) -> crate::Result<Model> {
320    validate_catalog_entry(entry)?;
321    let task_id = scheduled_task_row_id(entry.namespace, entry.task_name)?;
322    let insert_result = ActiveModel {
323        task_id: Set(task_id.clone()),
324        namespace: Set(entry.namespace.to_string()),
325        task_name: Set(entry.task_name.to_string()),
326        display_name: Set(entry.display_name.to_string()),
327        next_run_at: Set(entry.first_run_at),
328        claim_owner_id: Set(None),
329        claim_expires_at: Set(None),
330        last_claimed_at: Set(None),
331        last_finished_at: Set(None),
332        created_at: Set(entry.first_run_at),
333        updated_at: Set(entry.first_run_at),
334    }
335    .insert(db)
336    .await;
337
338    match insert_result {
339        Ok(model) => Ok(model),
340        Err(insert_error) => refresh_existing_task(db, task_id, entry, insert_error).await,
341    }
342}
343
344async fn refresh_existing_task(
345    db: &DatabaseConnection,
346    task_id: String,
347    entry: ScheduledTaskCatalogEntry<'_>,
348    insert_error: sea_orm::DbErr,
349) -> crate::Result<Model> {
350    let existing = Entity::find_by_id(task_id.clone())
351        .one(db)
352        .await
353        .map_err(DbError::from)?;
354    let Some(existing) = existing else {
355        return Err(DbError::from(insert_error));
356    };
357    if existing.display_name == entry.display_name {
358        return Ok(existing);
359    }
360
361    Entity::update_many()
362        .col_expr(
363            Column::DisplayName,
364            Expr::value(entry.display_name.to_string()),
365        )
366        .col_expr(Column::UpdatedAt, Expr::value(entry.first_run_at))
367        .filter(Column::TaskId.eq(task_id.clone()))
368        .exec(db)
369        .await
370        .map_err(DbError::from)?;
371
372    Entity::find_by_id(task_id)
373        .one(db)
374        .await
375        .map_err(DbError::from)?
376        .ok_or_else(|| DbError::database_operation("scheduled task disappeared after update"))
377}
378
379async fn claim_due(
380    db: &DatabaseConnection,
381    request: ScheduledTaskClaimRequest<'_>,
382) -> crate::Result<Option<ScheduledTaskClaim>> {
383    validate_claim_request(request)?;
384    let task_id = scheduled_task_row_id(request.namespace, request.task_name)?;
385    let Some(existing) = Entity::find_by_id(task_id.clone())
386        .one(db)
387        .await
388        .map_err(DbError::from)?
389    else {
390        return Ok(None);
391    };
392
393    if existing.next_run_at > request.now {
394        return Ok(None);
395    }
396    if is_claim_fresh(&existing, request.now) {
397        return Ok(None);
398    }
399
400    let claim_expires_at = request
401        .now
402        .checked_add_signed(chrono_duration_from_std(request.claim_ttl)?)
403        .ok_or_else(|| DbError::non_retryable("scheduled task claim expiry overflow"))?;
404    let claim_available = Condition::any()
405        .add(Column::ClaimOwnerId.is_null())
406        .add(Column::ClaimExpiresAt.is_null())
407        .add(Column::ClaimExpiresAt.lte(request.now));
408    let update = Entity::update_many()
409        .col_expr(
410            Column::ClaimOwnerId,
411            Expr::value(Some(request.owner_id.to_string())),
412        )
413        .col_expr(Column::ClaimExpiresAt, Expr::value(Some(claim_expires_at)))
414        .col_expr(Column::LastClaimedAt, Expr::value(Some(request.now)))
415        .col_expr(Column::UpdatedAt, Expr::value(request.now))
416        .filter(Column::TaskId.eq(task_id.clone()))
417        .filter(Column::NextRunAt.eq(existing.next_run_at))
418        .filter(Column::NextRunAt.lte(request.now))
419        .filter(claim_available)
420        .exec(db)
421        .await
422        .map_err(DbError::from)?;
423
424    if update.rows_affected != 1 {
425        return Ok(None);
426    }
427
428    Ok(Some(ScheduledTaskClaim {
429        task_id,
430        namespace: existing.namespace,
431        task_name: existing.task_name,
432        owner_id: request.owner_id.to_string(),
433        scheduled_at: existing.next_run_at,
434        claimed_at: request.now,
435        claim_expires_at,
436    }))
437}
438
439async fn renew_claim(
440    db: &DatabaseConnection,
441    renewal: ScheduledTaskClaimRenewal<'_>,
442) -> crate::Result<bool> {
443    validate_renewal(&renewal)?;
444    let claim_expires_at = renewal
445        .now
446        .checked_add_signed(chrono_duration_from_std(renewal.claim_ttl)?)
447        .ok_or_else(|| DbError::non_retryable("scheduled task claim expiry overflow"))?;
448    // The owner + claim-timestamp predicate matches completion: a firing another
449    // runtime reclaimed (new owner or newer last_claimed_at) can never be revived.
450    let update = Entity::update_many()
451        .col_expr(Column::ClaimExpiresAt, Expr::value(Some(claim_expires_at)))
452        .col_expr(Column::UpdatedAt, Expr::value(renewal.now))
453        .filter(Column::TaskId.eq(renewal.claim.task_id.clone()))
454        .filter(Column::ClaimOwnerId.eq(renewal.claim.owner_id.clone()))
455        .filter(Column::LastClaimedAt.eq(renewal.claim.claimed_at))
456        .exec(db)
457        .await
458        .map_err(DbError::from)?;
459
460    Ok(update.rows_affected == 1)
461}
462
463async fn complete_claim(
464    db: &DatabaseConnection,
465    completion: ScheduledTaskCompletion,
466) -> crate::Result<bool> {
467    validate_completion(&completion)?;
468    let update = Entity::update_many()
469        .col_expr(Column::NextRunAt, Expr::value(completion.next_run_at))
470        .col_expr(Column::ClaimOwnerId, Expr::value(Option::<String>::None))
471        .col_expr(
472            Column::ClaimExpiresAt,
473            Expr::value(Option::<chrono::DateTime<chrono::Utc>>::None),
474        )
475        .col_expr(
476            Column::LastFinishedAt,
477            Expr::value(Some(completion.finished_at)),
478        )
479        .col_expr(Column::UpdatedAt, Expr::value(completion.finished_at))
480        .filter(Column::TaskId.eq(completion.claim.task_id))
481        .filter(Column::ClaimOwnerId.eq(completion.claim.owner_id))
482        .filter(Column::LastClaimedAt.eq(completion.claim.claimed_at))
483        .exec(db)
484        .await
485        .map_err(DbError::from)?;
486
487    Ok(update.rows_affected == 1)
488}
489
490fn is_claim_fresh(model: &Model, now: chrono::DateTime<chrono::Utc>) -> bool {
491    matches!(
492        (&model.claim_owner_id, model.claim_expires_at),
493        (Some(_), Some(expires_at)) if expires_at > now
494    )
495}
496
497fn scheduled_task_row_id(namespace: &str, task_name: &str) -> crate::Result<String> {
498    let task_id = format!("{namespace}:{task_name}");
499    if task_id.len() > SCHEDULED_TASK_ID_MAX_LEN {
500        return Err(DbError::non_retryable(format!(
501            "scheduled task id must be at most {SCHEDULED_TASK_ID_MAX_LEN} bytes"
502        )));
503    }
504    Ok(task_id)
505}
506
507fn validate_catalog_entry(entry: ScheduledTaskCatalogEntry<'_>) -> crate::Result<()> {
508    validate_non_empty("scheduled task namespace", entry.namespace)?;
509    validate_non_empty("scheduled task name", entry.task_name)?;
510    validate_non_empty("scheduled task display name", entry.display_name)?;
511    validate_max_len(
512        "scheduled task namespace",
513        entry.namespace,
514        SCHEDULED_TASK_NAMESPACE_MAX_LEN,
515    )?;
516    validate_max_len(
517        "scheduled task name",
518        entry.task_name,
519        SCHEDULED_TASK_NAME_MAX_LEN,
520    )?;
521    validate_max_len(
522        "scheduled task display name",
523        entry.display_name,
524        SCHEDULED_TASK_DISPLAY_NAME_MAX_LEN,
525    )
526}
527
528fn validate_claim_request(request: ScheduledTaskClaimRequest<'_>) -> crate::Result<()> {
529    validate_non_empty("scheduled task owner id", request.owner_id)?;
530    validate_max_len(
531        "scheduled task owner id",
532        request.owner_id,
533        SCHEDULED_TASK_OWNER_ID_MAX_LEN,
534    )?;
535    validate_max_len(
536        "scheduled task namespace",
537        request.namespace,
538        SCHEDULED_TASK_NAMESPACE_MAX_LEN,
539    )?;
540    validate_max_len(
541        "scheduled task name",
542        request.task_name,
543        SCHEDULED_TASK_NAME_MAX_LEN,
544    )?;
545    if request.claim_ttl.is_zero() {
546        return Err(DbError::non_retryable(
547            "scheduled task claim TTL must not be zero",
548        ));
549    }
550    Ok(())
551}
552
553fn validate_completion(completion: &ScheduledTaskCompletion) -> crate::Result<()> {
554    if completion.next_run_at <= completion.claim.scheduled_at {
555        return Err(DbError::non_retryable(
556            "scheduled task next run must be after the claimed scheduled time",
557        ));
558    }
559    Ok(())
560}
561
562fn validate_renewal(renewal: &ScheduledTaskClaimRenewal<'_>) -> crate::Result<()> {
563    validate_non_empty("scheduled task owner id", &renewal.claim.owner_id)?;
564    validate_max_len(
565        "scheduled task owner id",
566        &renewal.claim.owner_id,
567        SCHEDULED_TASK_OWNER_ID_MAX_LEN,
568    )?;
569    if renewal.claim_ttl.is_zero() {
570        return Err(DbError::non_retryable(
571            "scheduled task claim TTL must not be zero",
572        ));
573    }
574    Ok(())
575}
576
577fn validate_non_empty(name: &str, value: &str) -> crate::Result<()> {
578    if value.trim().is_empty() {
579        return Err(DbError::non_retryable(format!("{name} must not be empty")));
580    }
581    Ok(())
582}
583
584fn validate_max_len(name: &str, value: &str, max_len: usize) -> crate::Result<()> {
585    if value.len() > max_len {
586        return Err(DbError::non_retryable(format!(
587            "{name} must be at most {max_len} bytes"
588        )));
589    }
590    Ok(())
591}
592
593fn chrono_duration_from_std(duration: Duration) -> crate::Result<chrono::Duration> {
594    chrono::Duration::from_std(duration)
595        .map_err(|_| DbError::non_retryable("duration is too large for chrono"))
596}
597
598#[cfg(test)]
599mod tests {
600    use chrono::{Duration as ChronoDuration, TimeZone, Utc};
601    use sea_orm::sea_query::{MysqlQueryBuilder, PostgresQueryBuilder, SqliteQueryBuilder};
602    use sea_orm::{ConnectionTrait, Database, DatabaseBackend, Schema};
603
604    use super::{
605        Entity, ScheduledTaskCatalogEntry, ScheduledTaskClaimRenewal, ScheduledTaskClaimRequest,
606        ScheduledTaskCompletion, ScheduledTaskDbStore,
607        create_scheduled_tasks_namespace_name_unique_index, create_scheduled_tasks_next_run_index,
608        create_scheduled_tasks_table,
609    };
610
611    async fn sqlite_store() -> ScheduledTaskDbStore {
612        let db = Database::connect("sqlite::memory:")
613            .await
614            .expect("sqlite memory database should connect");
615        let schema = Schema::new(db.get_database_backend());
616        let statement = schema.create_table_from_entity(Entity);
617        db.execute(&statement)
618            .await
619            .expect("scheduled tasks table should be created");
620        ScheduledTaskDbStore::new(db)
621    }
622
623    async fn sqlite_store_from_builders() -> ScheduledTaskDbStore {
624        let db = Database::connect("sqlite::memory:")
625            .await
626            .expect("sqlite memory database should connect");
627        let backend = db.get_database_backend();
628        db.execute(&create_scheduled_tasks_table(backend))
629            .await
630            .expect("scheduled tasks table builder should execute");
631        db.execute(&create_scheduled_tasks_namespace_name_unique_index())
632            .await
633            .expect("scheduled tasks unique index builder should execute");
634        db.execute(&create_scheduled_tasks_next_run_index())
635            .await
636            .expect("scheduled tasks due index builder should execute");
637        ScheduledTaskDbStore::new(db)
638    }
639
640    fn create_table_sql(backend: DatabaseBackend) -> String {
641        let table = create_scheduled_tasks_table(backend);
642        match backend {
643            DatabaseBackend::MySql => table.to_string(MysqlQueryBuilder),
644            DatabaseBackend::Postgres => table.to_string(PostgresQueryBuilder),
645            DatabaseBackend::Sqlite => table.to_string(SqliteQueryBuilder),
646            _ => unreachable!("unsupported backend in scheduled task table test"),
647        }
648    }
649
650    fn entry(first_run_at: chrono::DateTime<Utc>) -> ScheduledTaskCatalogEntry<'static> {
651        ScheduledTaskCatalogEntry {
652            namespace: "aster_yggdrasil",
653            task_name: "audit-cleanup",
654            display_name: "Audit cleanup",
655            first_run_at,
656        }
657    }
658
659    #[test]
660    fn create_scheduled_tasks_table_uses_stable_shape() {
661        let sqlite_sql = create_table_sql(DatabaseBackend::Sqlite);
662        assert!(sqlite_sql.contains("CREATE TABLE IF NOT EXISTS \"scheduled_tasks\""));
663        assert!(sqlite_sql.contains("\"task_id\" varchar(191) NOT NULL PRIMARY KEY"));
664        assert!(sqlite_sql.contains("\"namespace\" varchar(64) NOT NULL"));
665        assert!(sqlite_sql.contains("\"next_run_at\" timestamp_with_timezone_text NOT NULL"));
666        let namespace_index =
667            create_scheduled_tasks_namespace_name_unique_index().to_string(SqliteQueryBuilder);
668        assert!(namespace_index.contains("idx_scheduled_tasks_namespace_name_unique"));
669        assert!(namespace_index.contains("\"namespace\", \"task_name\""));
670        let next_run_index = create_scheduled_tasks_next_run_index().to_string(SqliteQueryBuilder);
671        assert!(next_run_index.contains("idx_scheduled_tasks_next_run"));
672
673        let mysql_sql = create_table_sql(DatabaseBackend::MySql);
674        assert!(mysql_sql.contains("`next_run_at` datetime(6) NOT NULL"));
675
676        let postgres_sql = create_table_sql(DatabaseBackend::Postgres);
677        assert!(postgres_sql.contains("\"next_run_at\" timestamp with time zone NOT NULL"));
678    }
679
680    #[tokio::test]
681    async fn scheduled_tasks_builders_execute_on_sqlite() {
682        let store = sqlite_store_from_builders().await;
683        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
684
685        let inserted = store
686            .ensure_task(entry(now))
687            .await
688            .expect("scheduled task should insert with builder-created schema");
689
690        assert_eq!(inserted.task_id, "aster_yggdrasil:audit-cleanup");
691    }
692
693    #[tokio::test]
694    async fn ensure_task_rejects_invalid_catalog_values() {
695        let store = sqlite_store().await;
696        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
697
698        let empty = store
699            .ensure_task(ScheduledTaskCatalogEntry {
700                namespace: " ",
701                ..entry(now)
702            })
703            .await
704            .expect_err("empty namespace should be rejected");
705        assert!(empty.to_string().contains("must not be empty"));
706
707        let too_long = store
708            .ensure_task(ScheduledTaskCatalogEntry {
709                task_name: "x".repeat(129).as_str(),
710                ..entry(now)
711            })
712            .await
713            .expect_err("long task name should be rejected");
714        assert!(too_long.to_string().contains("at most 128 bytes"));
715    }
716
717    #[tokio::test]
718    async fn ensure_task_inserts_and_refreshes_display_name_without_resetting_schedule() {
719        let store = sqlite_store().await;
720        let first_run_at = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
721
722        let inserted = store
723            .ensure_task(entry(first_run_at))
724            .await
725            .expect("scheduled task should insert");
726        assert_eq!(inserted.task_id, "aster_yggdrasil:audit-cleanup");
727        assert_eq!(inserted.next_run_at, first_run_at);
728
729        let refreshed = store
730            .ensure_task(ScheduledTaskCatalogEntry {
731                display_name: "Audit cleanup v2",
732                first_run_at: first_run_at + ChronoDuration::hours(1),
733                ..entry(first_run_at)
734            })
735            .await
736            .expect("scheduled task should refresh");
737        assert_eq!(refreshed.display_name, "Audit cleanup v2");
738        assert_eq!(refreshed.next_run_at, first_run_at);
739    }
740
741    #[tokio::test]
742    async fn claim_due_claims_once_until_completion_or_expiry() {
743        let store = sqlite_store().await;
744        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
745        store
746            .ensure_task(entry(now))
747            .await
748            .expect("scheduled task should insert");
749
750        let claim = store
751            .claim_due(ScheduledTaskClaimRequest {
752                namespace: "aster_yggdrasil",
753                task_name: "audit-cleanup",
754                owner_id: "runtime-a",
755                now,
756                claim_ttl: std::time::Duration::from_secs(30),
757            })
758            .await
759            .expect("claim should succeed")
760            .expect("task should be due");
761        assert_eq!(claim.scheduled_at, now);
762
763        let blocked = store
764            .claim_due(ScheduledTaskClaimRequest {
765                namespace: "aster_yggdrasil",
766                task_name: "audit-cleanup",
767                owner_id: "runtime-b",
768                now: now + ChronoDuration::seconds(1),
769                claim_ttl: std::time::Duration::from_secs(30),
770            })
771            .await
772            .expect("standby claim should succeed");
773        assert!(blocked.is_none());
774
775        let reclaimed = store
776            .claim_due(ScheduledTaskClaimRequest {
777                namespace: "aster_yggdrasil",
778                task_name: "audit-cleanup",
779                owner_id: "runtime-b",
780                now: now + ChronoDuration::seconds(31),
781                claim_ttl: std::time::Duration::from_secs(30),
782            })
783            .await
784            .expect("expired claim should be reclaimable")
785            .expect("task should still be due");
786        assert_eq!(reclaimed.owner_id, "runtime-b");
787    }
788
789    #[tokio::test]
790    async fn claim_due_blocks_duplicate_fresh_claim_from_same_owner() {
791        let store = sqlite_store().await;
792        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
793        store
794            .ensure_task(entry(now))
795            .await
796            .expect("scheduled task should insert");
797
798        let first = store
799            .claim_due(ScheduledTaskClaimRequest {
800                namespace: "aster_yggdrasil",
801                task_name: "audit-cleanup",
802                owner_id: "runtime-a",
803                now,
804                claim_ttl: std::time::Duration::from_secs(30),
805            })
806            .await
807            .expect("first claim should query")
808            .expect("task should be due");
809        assert_eq!(first.owner_id, "runtime-a");
810
811        let duplicate = store
812            .claim_due(ScheduledTaskClaimRequest {
813                namespace: "aster_yggdrasil",
814                task_name: "audit-cleanup",
815                owner_id: "runtime-a",
816                now: now + ChronoDuration::seconds(1),
817                claim_ttl: std::time::Duration::from_secs(30),
818            })
819            .await
820            .expect("duplicate claim should query");
821        assert!(duplicate.is_none());
822    }
823
824    #[tokio::test]
825    async fn claim_due_skips_not_due_and_rejects_zero_ttl() {
826        let store = sqlite_store().await;
827        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
828        store
829            .ensure_task(entry(now + ChronoDuration::minutes(5)))
830            .await
831            .expect("scheduled task should insert");
832
833        let not_due = store
834            .claim_due(ScheduledTaskClaimRequest {
835                namespace: "aster_yggdrasil",
836                task_name: "audit-cleanup",
837                owner_id: "runtime-a",
838                now,
839                claim_ttl: std::time::Duration::from_secs(30),
840            })
841            .await
842            .expect("not-due claim check should succeed");
843        assert!(not_due.is_none());
844
845        let zero_ttl = store
846            .claim_due(ScheduledTaskClaimRequest {
847                namespace: "aster_yggdrasil",
848                task_name: "audit-cleanup",
849                owner_id: "runtime-a",
850                now,
851                claim_ttl: std::time::Duration::ZERO,
852            })
853            .await
854            .expect_err("zero TTL should be rejected");
855        assert!(zero_ttl.to_string().contains("must not be zero"));
856    }
857
858    #[tokio::test]
859    async fn completing_claim_advances_next_run_and_clears_claim() {
860        let store = sqlite_store().await;
861        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
862        store
863            .ensure_task(entry(now))
864            .await
865            .expect("scheduled task should insert");
866        let claim = store
867            .claim_due(ScheduledTaskClaimRequest {
868                namespace: "aster_yggdrasil",
869                task_name: "audit-cleanup",
870                owner_id: "runtime-a",
871                now,
872                claim_ttl: std::time::Duration::from_secs(30),
873            })
874            .await
875            .expect("claim should succeed")
876            .expect("task should be due");
877        let next_run_at = now + ChronoDuration::hours(1);
878
879        assert!(
880            store
881                .complete_claim(ScheduledTaskCompletion {
882                    claim,
883                    finished_at: now + ChronoDuration::seconds(5),
884                    next_run_at,
885                })
886                .await
887                .expect("completion should succeed")
888        );
889
890        let second = store
891            .claim_due(ScheduledTaskClaimRequest {
892                namespace: "aster_yggdrasil",
893                task_name: "audit-cleanup",
894                owner_id: "runtime-a",
895                now: now + ChronoDuration::minutes(30),
896                claim_ttl: std::time::Duration::from_secs(30),
897            })
898            .await
899            .expect("claim check should succeed");
900        assert!(second.is_none());
901    }
902
903    #[tokio::test]
904    async fn complete_claim_requires_matching_owner_and_claim_timestamp() {
905        let store = sqlite_store().await;
906        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
907        store
908            .ensure_task(entry(now))
909            .await
910            .expect("scheduled task should insert");
911        let claim = store
912            .claim_due(ScheduledTaskClaimRequest {
913                namespace: "aster_yggdrasil",
914                task_name: "audit-cleanup",
915                owner_id: "runtime-a",
916                now,
917                claim_ttl: std::time::Duration::from_secs(30),
918            })
919            .await
920            .expect("claim should succeed")
921            .expect("task should be due");
922
923        let mut wrong_owner = claim.clone();
924        wrong_owner.owner_id = "runtime-b".to_string();
925        assert!(
926            !store
927                .complete_claim(ScheduledTaskCompletion {
928                    claim: wrong_owner,
929                    finished_at: now + ChronoDuration::seconds(5),
930                    next_run_at: now + ChronoDuration::hours(1),
931                })
932                .await
933                .expect("wrong owner completion should query")
934        );
935
936        let mut wrong_claim_time = claim;
937        wrong_claim_time.claimed_at += ChronoDuration::seconds(1);
938        assert!(
939            !store
940                .complete_claim(ScheduledTaskCompletion {
941                    claim: wrong_claim_time,
942                    finished_at: now + ChronoDuration::seconds(5),
943                    next_run_at: now + ChronoDuration::hours(1),
944                })
945                .await
946                .expect("wrong claim timestamp completion should query")
947        );
948    }
949
950    #[tokio::test]
951    async fn complete_claim_rejects_next_run_that_does_not_advance_schedule() {
952        let store = sqlite_store().await;
953        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
954        store
955            .ensure_task(entry(now))
956            .await
957            .expect("scheduled task should insert");
958        let claim = store
959            .claim_due(ScheduledTaskClaimRequest {
960                namespace: "aster_yggdrasil",
961                task_name: "audit-cleanup",
962                owner_id: "runtime-a",
963                now,
964                claim_ttl: std::time::Duration::from_secs(30),
965            })
966            .await
967            .expect("claim should succeed")
968            .expect("task should be due");
969
970        let error = store
971            .complete_claim(ScheduledTaskCompletion {
972                claim,
973                finished_at: now + ChronoDuration::seconds(5),
974                next_run_at: now,
975            })
976            .await
977            .expect_err("non-advancing next run should be rejected");
978
979        assert!(error.to_string().contains("must be after"));
980    }
981
982    #[tokio::test]
983    async fn renew_claim_extends_expiry_and_blocks_reclaim_while_fresh() {
984        let store = sqlite_store().await;
985        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
986        store
987            .ensure_task(entry(now))
988            .await
989            .expect("scheduled task should insert");
990        let claim = store
991            .claim_due(ScheduledTaskClaimRequest {
992                namespace: "aster_yggdrasil",
993                task_name: "audit-cleanup",
994                owner_id: "runtime-a",
995                now,
996                claim_ttl: std::time::Duration::from_secs(30),
997            })
998            .await
999            .expect("claim should succeed")
1000            .expect("task should be due");
1001
1002        // Renew near the end of the original window: expiry moves to now+55.
1003        assert!(
1004            store
1005                .renew_claim(ScheduledTaskClaimRenewal {
1006                    claim: &claim,
1007                    now: now + ChronoDuration::seconds(25),
1008                    claim_ttl: std::time::Duration::from_secs(30),
1009                })
1010                .await
1011                .expect("renewal should succeed")
1012        );
1013
1014        // Past the original expiry, the renewed claim still blocks other runtimes.
1015        let blocked = store
1016            .claim_due(ScheduledTaskClaimRequest {
1017                namespace: "aster_yggdrasil",
1018                task_name: "audit-cleanup",
1019                owner_id: "runtime-b",
1020                now: now + ChronoDuration::seconds(40),
1021                claim_ttl: std::time::Duration::from_secs(30),
1022            })
1023            .await
1024            .expect("claim check should succeed");
1025        assert!(blocked.is_none());
1026
1027        // Past the renewed expiry, the firing is reclaimable again.
1028        let reclaimed = store
1029            .claim_due(ScheduledTaskClaimRequest {
1030                namespace: "aster_yggdrasil",
1031                task_name: "audit-cleanup",
1032                owner_id: "runtime-b",
1033                now: now + ChronoDuration::seconds(56),
1034                claim_ttl: std::time::Duration::from_secs(30),
1035            })
1036            .await
1037            .expect("expired claim should be reclaimable")
1038            .expect("task should still be due");
1039        assert_eq!(reclaimed.owner_id, "runtime-b");
1040    }
1041
1042    #[tokio::test]
1043    async fn renew_claim_requires_matching_owner_and_claim_timestamp() {
1044        let store = sqlite_store().await;
1045        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
1046        store
1047            .ensure_task(entry(now))
1048            .await
1049            .expect("scheduled task should insert");
1050        let claim = store
1051            .claim_due(ScheduledTaskClaimRequest {
1052                namespace: "aster_yggdrasil",
1053                task_name: "audit-cleanup",
1054                owner_id: "runtime-a",
1055                now,
1056                claim_ttl: std::time::Duration::from_secs(30),
1057            })
1058            .await
1059            .expect("claim should succeed")
1060            .expect("task should be due");
1061
1062        let mut wrong_owner = claim.clone();
1063        wrong_owner.owner_id = "runtime-b".to_string();
1064        assert!(
1065            !store
1066                .renew_claim(ScheduledTaskClaimRenewal {
1067                    claim: &wrong_owner,
1068                    now: now + ChronoDuration::seconds(10),
1069                    claim_ttl: std::time::Duration::from_secs(30),
1070                })
1071                .await
1072                .expect("wrong owner renewal should query")
1073        );
1074
1075        let mut wrong_claim_time = claim.clone();
1076        wrong_claim_time.claimed_at += ChronoDuration::seconds(1);
1077        assert!(
1078            !store
1079                .renew_claim(ScheduledTaskClaimRenewal {
1080                    claim: &wrong_claim_time,
1081                    now: now + ChronoDuration::seconds(10),
1082                    claim_ttl: std::time::Duration::from_secs(30),
1083                })
1084                .await
1085                .expect("wrong claim timestamp renewal should query")
1086        );
1087    }
1088
1089    #[tokio::test]
1090    async fn renew_claim_after_expiry_without_contestation_revives_claim() {
1091        let store = sqlite_store().await;
1092        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
1093        store
1094            .ensure_task(entry(now))
1095            .await
1096            .expect("scheduled task should insert");
1097        let claim = store
1098            .claim_due(ScheduledTaskClaimRequest {
1099                namespace: "aster_yggdrasil",
1100                task_name: "audit-cleanup",
1101                owner_id: "runtime-a",
1102                now,
1103                claim_ttl: std::time::Duration::from_secs(30),
1104            })
1105            .await
1106            .expect("claim should succeed")
1107            .expect("task should be due");
1108
1109        // The ownership row is untouched past expiry, so a stalled worker that
1110        // recovers before anyone reclaims resumes its own claim instead of
1111        // abandoning the firing to a duplicate execution.
1112        assert!(
1113            store
1114                .renew_claim(ScheduledTaskClaimRenewal {
1115                    claim: &claim,
1116                    now: now + ChronoDuration::seconds(40),
1117                    claim_ttl: std::time::Duration::from_secs(30),
1118                })
1119                .await
1120                .expect("late renewal should succeed while the row is uncontested")
1121        );
1122    }
1123
1124    #[tokio::test]
1125    async fn renew_claim_rejects_zero_ttl() {
1126        let store = sqlite_store().await;
1127        let now = Utc.with_ymd_and_hms(2026, 6, 26, 1, 0, 0).unwrap();
1128        store
1129            .ensure_task(entry(now))
1130            .await
1131            .expect("scheduled task should insert");
1132        let claim = store
1133            .claim_due(ScheduledTaskClaimRequest {
1134                namespace: "aster_yggdrasil",
1135                task_name: "audit-cleanup",
1136                owner_id: "runtime-a",
1137                now,
1138                claim_ttl: std::time::Duration::from_secs(30),
1139            })
1140            .await
1141            .expect("claim should succeed")
1142            .expect("task should be due");
1143
1144        let error = store
1145            .renew_claim(ScheduledTaskClaimRenewal {
1146                claim: &claim,
1147                now: now + ChronoDuration::seconds(10),
1148                claim_ttl: std::time::Duration::ZERO,
1149            })
1150            .await
1151            .expect_err("zero TTL renewal should be rejected");
1152        assert!(error.to_string().contains("must not be zero"));
1153    }
1154}