1use 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
25pub const SCHEDULED_TASKS_TABLE: &str = "scheduled_tasks";
27pub const SCHEDULED_TASK_ID_COLUMN: &str = "task_id";
29pub const SCHEDULED_TASK_NAMESPACE_COLUMN: &str = "namespace";
31pub const SCHEDULED_TASK_NAME_COLUMN: &str = "task_name";
33pub const SCHEDULED_TASK_DISPLAY_NAME_COLUMN: &str = "display_name";
35pub const SCHEDULED_TASK_NEXT_RUN_AT_COLUMN: &str = "next_run_at";
37pub const SCHEDULED_TASK_CLAIM_OWNER_ID_COLUMN: &str = "claim_owner_id";
39pub const SCHEDULED_TASK_CLAIM_EXPIRES_AT_COLUMN: &str = "claim_expires_at";
41pub const SCHEDULED_TASK_LAST_CLAIMED_AT_COLUMN: &str = "last_claimed_at";
43pub const SCHEDULED_TASK_LAST_FINISHED_AT_COLUMN: &str = "last_finished_at";
45pub const SCHEDULED_TASK_CREATED_AT_COLUMN: &str = "created_at";
47pub const SCHEDULED_TASK_UPDATED_AT_COLUMN: &str = "updated_at";
49pub const SCHEDULED_TASK_NAMESPACE_NAME_UNIQUE_INDEX: &str =
51 "idx_scheduled_tasks_namespace_name_unique";
52pub 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#[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#[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#[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#[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#[derive(Clone, Debug, PartialEq, DeriveEntityModel)]
198#[sea_orm(table_name = "scheduled_tasks")]
199pub struct Model {
200 #[sea_orm(primary_key, auto_increment = false)]
202 pub task_id: String,
203 pub namespace: String,
205 pub task_name: String,
207 pub display_name: String,
209 pub next_run_at: DateTimeUtc,
211 pub claim_owner_id: Option<String>,
213 pub claim_expires_at: Option<DateTimeUtc>,
215 pub last_claimed_at: Option<DateTimeUtc>,
217 pub last_finished_at: Option<DateTimeUtc>,
219 pub created_at: DateTimeUtc,
221 pub updated_at: DateTimeUtc,
223}
224
225#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
226pub enum Relation {}
227
228impl ActiveModelBehavior for ActiveModel {}
229
230#[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 #[must_use]
272 pub const fn new(db: DatabaseConnection) -> Self {
273 Self { db }
274 }
275
276 pub async fn ensure_task(&self, entry: ScheduledTaskCatalogEntry<'_>) -> crate::Result<Model> {
282 ensure_task(&self.db, entry).await
283 }
284
285 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 pub async fn renew_claim(&self, renewal: ScheduledTaskClaimRenewal<'_>) -> crate::Result<bool> {
303 renew_claim(&self.db, renewal).await
304 }
305
306 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 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 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 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 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 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}