Skip to main content

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