@@ -1486,44 +1486,78 @@ async fn repeated_step_names_are_numbered() -> Result<()> {
14861486
14871487#[ tokio:: test]
14881488#[ ignore = "requires a Postgres database initialized with Absurd SQL" ]
1489- async fn sleep_until_suspends_then_resumes_from_checkpoint ( ) -> Result < ( ) > {
1490- let ( queue, client) = test_client ( ) . await ?;
1489+ async fn sleep_for_suspends_until_duration_elapses ( ) -> Result < ( ) > {
1490+ let ( queue, client, pool) = test_client_with_single_connection_pool ( ) . await ?;
1491+ let base = utc ( "2024-05-05T10:00:00Z" ) ?;
1492+ set_fake_now ( & pool, Some ( base) ) . await ?;
1493+
1494+ let task = Task :: < ( ) , Value > :: new ( "sleep-for" ) . queue ( & queue) ;
1495+ client. register ( & task, |( ) , mut ctx| async move {
1496+ ctx. sleep_for ( "wait-for" , Duration :: from_secs ( 60 ) ) . await ?;
1497+ Ok ( json ! ( { "resumed" : true } ) )
1498+ } ) ?;
1499+
1500+ let spawned = client. spawn ( & task, ( ) , Default :: default ( ) ) . await ?;
1501+ assert_eq ! ( client. work_batch( WorkBatchOptions :: new( ) ) . await ?, 1 ) ;
1502+
1503+ let task = fetch_task ( & queue, spawned. task_id ) . await ?;
1504+ let run = fetch_run ( & queue, spawned. run_id ) . await ?;
1505+ let wake_at = base + ChronoDuration :: seconds ( 60 ) ;
1506+ assert_eq ! ( task. state, "sleeping" ) ;
1507+ assert_eq ! ( run. state, "sleeping" ) ;
1508+ assert_eq ! ( run. available_at, Some ( wake_at) ) ;
1509+
1510+ set_fake_now ( & pool, Some ( wake_at + ChronoDuration :: seconds ( 5 ) ) ) . await ?;
1511+ assert_eq ! ( client. work_batch( WorkBatchOptions :: new( ) ) . await ?, 1 ) ;
1512+
1513+ let task = fetch_task ( & queue, spawned. task_id ) . await ?;
1514+ assert_eq ! ( task. state, "completed" ) ;
1515+ assert_eq ! ( task. completed_payload, Some ( json!( { "resumed" : true } ) ) ) ;
1516+
1517+ client. drop_queue ( ) . await ?;
1518+ Ok ( ( ) )
1519+ }
1520+
1521+ #[ tokio:: test]
1522+ #[ ignore = "requires a Postgres database initialized with Absurd SQL" ]
1523+ async fn sleep_until_checkpoint_prevents_rescheduling ( ) -> Result < ( ) > {
1524+ let ( queue, client, pool) = test_client_with_single_connection_pool ( ) . await ?;
1525+ let base = utc ( "2024-05-06T09:00:00Z" ) ?;
1526+ set_fake_now ( & pool, Some ( base) ) . await ?;
1527+ let wake_at = base + ChronoDuration :: minutes ( 5 ) ;
14911528 let executions = Arc :: new ( AtomicUsize :: new ( 0 ) ) ;
14921529
1493- let task = Task :: < ( ) , Value > :: new ( "sleepy " ) . queue ( & queue) ;
1530+ let task = Task :: < ( ) , Value > :: new ( "sleep-until " ) . queue ( & queue) ;
14941531 client. register ( & task, {
14951532 let executions = Arc :: clone ( & executions) ;
14961533 move |( ) , mut ctx| {
14971534 let executions = Arc :: clone ( & executions) ;
14981535 async move {
14991536 let execution = executions. fetch_add ( 1 , Ordering :: SeqCst ) + 1 ;
1500- let wake_at = Utc :: now ( ) + ChronoDuration :: seconds ( 1 ) ;
1501- ctx. sleep_until ( "pause" , wake_at) . await ?;
1537+ ctx. sleep_until ( "sleep-step" , wake_at) . await ?;
15021538 Ok ( json ! ( { "executions" : execution } ) )
15031539 }
15041540 }
15051541 } ) ?;
15061542
15071543 let spawned = client. spawn ( & task, ( ) , Default :: default ( ) ) . await ?;
1508-
15091544 assert_eq ! ( client. work_batch( WorkBatchOptions :: new( ) ) . await ?, 1 ) ;
15101545
1546+ let checkpoints = fetch_checkpoints ( & queue, spawned. task_id ) . await ?;
1547+ assert_eq ! (
1548+ checkpoints,
1549+ vec![ ( "sleep-step" . to_string( ) , serde_json:: to_value( wake_at) ?) ] ,
1550+ ) ;
1551+
15111552 let task = fetch_task ( & queue, spawned. task_id ) . await ?;
15121553 let run = fetch_run ( & queue, spawned. run_id ) . await ?;
15131554 assert_eq ! ( task. state, "sleeping" ) ;
15141555 assert_eq ! ( run. state, "sleeping" ) ;
15151556 assert ! ( run. wake_event. is_none( ) ) ;
15161557 assert ! ( run. failure_reason. is_none( ) ) ;
1558+ assert_eq ! ( run. available_at, Some ( wake_at) ) ;
15171559
1518- let checkpoints = fetch_checkpoints ( & queue, spawned. task_id ) . await ?;
1519- assert_eq ! ( checkpoints. len( ) , 1 ) ;
1520- assert_eq ! ( checkpoints[ 0 ] . 0 , "pause" ) ;
1521- let checkpoint_wake: DateTime < Utc > = serde_json:: from_value ( checkpoints[ 0 ] . 1 . clone ( ) ) ?;
1522- assert_eq ! ( run. available_at, Some ( checkpoint_wake) ) ;
1523-
1524- assert_eq ! ( client. work_batch( WorkBatchOptions :: new( ) ) . await ?, 0 ) ;
1525-
1526- tokio:: time:: sleep ( Duration :: from_millis ( 1_100 ) ) . await ;
1560+ set_fake_now ( & pool, Some ( wake_at) ) . await ?;
15271561 assert_eq ! ( client. work_batch( WorkBatchOptions :: new( ) ) . await ?, 1 ) ;
15281562
15291563 let task = fetch_task ( & queue, spawned. task_id ) . await ?;
0 commit comments