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
61pub 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
101pub fn drop_scheduled_tasks_table() -> TableDropStatement {
103 Table::drop()
104 .table(scheduled_tasks_table())
105 .if_exists()
106 .to_owned()
107}
108
109pub 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
121pub 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#[derive(Clone, Debug, PartialEq, DeriveEntityModel)]
194#[sea_orm(table_name = "scheduled_tasks")]
195pub struct Model {
196 #[sea_orm(primary_key, auto_increment = false)]
198 pub task_id: String,
199 pub namespace: String,
201 pub task_name: String,
203 pub display_name: String,
205 pub next_run_at: DateTimeUtc,
207 pub claim_owner_id: Option<String>,
209 pub claim_expires_at: Option<DateTimeUtc>,
211 pub last_claimed_at: Option<DateTimeUtc>,
213 pub last_finished_at: Option<DateTimeUtc>,
215 pub created_at: DateTimeUtc,
217 pub updated_at: DateTimeUtc,
219}
220
221#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)]
222pub enum Relation {}
223
224impl ActiveModelBehavior for ActiveModel {}
225
226#[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 pub const fn new(db: DatabaseConnection) -> Self {
268 Self { db }
269 }
270
271 pub async fn ensure_task(&self, entry: ScheduledTaskCatalogEntry<'_>) -> crate::Result<Model> {
273 ensure_task(&self.db, entry).await
274 }
275
276 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 pub async fn renew_claim(&self, renewal: ScheduledTaskClaimRenewal<'_>) -> crate::Result<bool> {
286 renew_claim(&self.db, renewal).await
287 }
288
289 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 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 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 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 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 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}