From c60357ea6168001a164cd4637ac688730086adf1 Mon Sep 17 00:00:00 2001 From: Shubham Mishra Date: Thu, 13 Aug 2026 17:31:07 +0530 Subject: [PATCH 1/3] make scheduled function behaviour as Skip --- crates/core/src/host/scheduler.rs | 123 ++++++++++++++++++++++++------ 1 file changed, 100 insertions(+), 23 deletions(-) diff --git a/crates/core/src/host/scheduler.rs b/crates/core/src/host/scheduler.rs index c0bda92f0a5..009b408cc04 100644 --- a/crates/core/src/host/scheduler.rs +++ b/crates/core/src/host/scheduler.rs @@ -16,7 +16,7 @@ use spacetimedb_datastore::locking_tx_datastore::MutTxId; use spacetimedb_datastore::system_tables::{StScheduledFields, ST_SCHEDULED_ID}; use spacetimedb_datastore::traits::IsolationLevel; use spacetimedb_lib::scheduler::ScheduleAt; -use spacetimedb_lib::Timestamp; +use spacetimedb_lib::{TimeDuration, Timestamp}; use spacetimedb_primitives::{ColId, TableId}; use spacetimedb_sats::bsatn::ToBsatn as _; use spacetimedb_sats::AlgebraicValue; @@ -447,6 +447,12 @@ struct Reschedule { at_real: Instant, } +#[derive(Clone)] +struct ScheduledFunctionTiming { + function_name: Arc, + intended_at: Timestamp, +} + enum ScheduledProcedureStep { Done(CallScheduledFunctionResult, bool), Procedure { @@ -516,13 +522,13 @@ fn prepare_scheduled_procedure_call( inst: &mut impl WasmInstance, ) -> ScheduledProcedureStep { let ScheduledFunctionParams(item) = params; - let delay = scheduled_function_delay_context_for_item(&item); + let timing = scheduled_function_timing_for_item(&item); let id = scheduled_item_id(&item); let db = &**module_info.relational_db(); let tx = db.begin_mut_tx(IsolationLevel::Serializable, Workload::Internal); let params = procedure_call_params_for_queued_item(module_info, db, &tx, item); - let (timestamp, instant, params) = match params { + let (timestamp, _instant, params) = match params { // If the function was already deleted, leave the `ScheduledFunction` // in the database for when the module restarts. Ok(None) => return ScheduledProcedureStep::Done(CallScheduledFunctionResult { reschedule: None }, false), @@ -531,7 +537,10 @@ fn prepare_scheduled_procedure_call( // All we can do here is log an error. log::error!("could not determine scheduled procedure or its parameters: {err:#}"); let reschedule = id.and_then(|id| { - let reschedule_from = (Timestamp::now(), Instant::now()); + let reschedule_from = timing + .as_ref() + .map(|timing| timing.intended_at) + .unwrap_or_else(Timestamp::now); delete_scheduled_function_row(module_info, db, id, Some(tx), reschedule_from, inst_common, inst) }); return ScheduledProcedureStep::Done(CallScheduledFunctionResult { reschedule }, false); @@ -541,10 +550,15 @@ fn prepare_scheduled_procedure_call( // For scheduled procedures, it's incorrect to retry them if execution aborts midway, // so we must remove the schedule row before executing. let reschedule = id.and_then(|id| { - delete_scheduled_function_row(module_info, db, id, Some(tx), (timestamp, instant), inst_common, inst) + let reschedule_from = timing.as_ref().map(|timing| timing.intended_at).unwrap_or(timestamp); + delete_scheduled_function_row(module_info, db, id, Some(tx), reschedule_from, inst_common, inst) + }); + let delay = timing.map(|timing| { + ( + timing.function_name, + scheduled_function_delay(timestamp, timing.intended_at), + ) }); - let delay = - delay.map(|(function_name, requested_at)| (function_name, scheduled_function_delay(timestamp, requested_at))); ScheduledProcedureStep::Procedure { params: Box::new(params), reschedule, @@ -559,13 +573,13 @@ fn call_scheduled_reducer_until_done( inst: &mut impl WasmInstance, ) -> (CallScheduledFunctionResult, bool) { let ScheduledFunctionParams(item) = params; - let delay = scheduled_function_delay_context_for_item(&item); + let timing = scheduled_function_timing_for_item(&item); let id = scheduled_item_id(&item); let db = &**module_info.relational_db(); let tx = db.begin_mut_tx(IsolationLevel::Serializable, Workload::Internal); let params = reducer_call_params_for_queued_item(module_info, db, &tx, item); - let (timestamp, instant, params) = match params { + let (timestamp, _instant, params) = match params { // If the function was already deleted, leave the `ScheduledFunction` // in the database for when the module restarts. Ok(None) => return (CallScheduledFunctionResult { reschedule: None }, false), @@ -574,18 +588,22 @@ fn call_scheduled_reducer_until_done( // All we can do here is log an error. log::error!("could not determine scheduled reducer or its parameters: {err:#}"); let reschedule = id.and_then(|id| { - let reschedule_from = (Timestamp::now(), Instant::now()); + let reschedule_from = timing + .as_ref() + .map(|timing| timing.intended_at) + .unwrap_or_else(Timestamp::now); delete_scheduled_function_row(module_info, db, id, Some(tx), reschedule_from, inst_common, inst) }); return (CallScheduledFunctionResult { reschedule }, false); } }; - if let Some((function_name, requested_at)) = delay { - let delay = scheduled_function_delay(timestamp, requested_at); - record_scheduled_function_delay(module_info, &function_name, delay); + if let Some(timing) = timing.as_ref() { + let delay = scheduled_function_delay(timestamp, timing.intended_at); + record_scheduled_function_delay(module_info, &timing.function_name, delay); } - call_scheduled_reducer_with_tx(module_info, db, id, tx, (timestamp, instant), params, inst_common, inst) + let reschedule_from = timing.as_ref().map(|timing| timing.intended_at).unwrap_or(timestamp); + call_scheduled_reducer_with_tx(module_info, db, id, tx, reschedule_from, params, inst_common, inst) } fn scheduled_item_id(item: &QueueItem) -> Option { @@ -595,9 +613,12 @@ fn scheduled_item_id(item: &QueueItem) -> Option { } } -fn scheduled_function_delay_context_for_item(item: &QueueItem) -> Option<(Arc, Timestamp)> { +fn scheduled_function_timing_for_item(item: &QueueItem) -> Option { match item { - QueueItem::Id { function_name, at, .. } => Some((function_name.clone(), *at)), + QueueItem::Id { function_name, at, .. } => Some(ScheduledFunctionTiming { + function_name: function_name.clone(), + intended_at: *at, + }), QueueItem::VolatileNonatomicImmediate { .. } => None, } } @@ -630,7 +651,7 @@ fn call_scheduled_reducer_with_tx( db: &RelationalDB, id: Option, mut tx: MutTxId, - reschedule_from: (Timestamp, Instant), + reschedule_from: Timestamp, params: CallReducerParams, inst_common: &mut InstanceCommon, inst: &mut impl WasmInstance, @@ -682,20 +703,49 @@ fn delete_scheduled_function_row( db: &RelationalDB, id: ScheduledFunctionId, tx: Option, - reschedule_from: (Timestamp, Instant), + reschedule_from: Timestamp, inst_common: &mut InstanceCommon, inst: &mut impl WasmInstance, ) -> Option { - let (timestamp, instant) = reschedule_from; let tx = tx.unwrap_or_else(|| db.begin_mut_tx(IsolationLevel::Serializable, Workload::Internal)); let schedule_at = delete_scheduled_function_row_with_tx(module_info, db, tx, id, inst_common, inst)?; let ScheduleAt::Interval(dur) = schedule_at else { return None; }; - Some(Reschedule { - at_ts: schedule_at.to_timestamp_from(timestamp), - at_real: instant + dur.to_duration_abs(), - }) + Some(next_interval_reschedule(reschedule_from, dur)) +} + +fn next_interval_reschedule(last_intended_at: Timestamp, interval: TimeDuration) -> Reschedule { + let now_ts = Timestamp::now(); + let at_ts = next_interval_tick_after(last_intended_at, interval, now_ts); + let delay = at_ts.duration_since(now_ts).unwrap_or(Duration::ZERO); + Reschedule { + at_ts, + at_real: Instant::now() + delay, + } +} + +fn next_interval_tick_after(last_intended_at: Timestamp, interval: TimeDuration, now: Timestamp) -> Timestamp { + let interval = interval.abs(); + let interval_micros = interval.to_micros(); + if interval_micros == 0 { + return last_intended_at + interval; + } + + let elapsed_micros = now + .time_duration_since(last_intended_at) + .map(|duration| duration.to_micros()) + .unwrap_or(0); + let ticks = if elapsed_micros <= 0 { + 1 + } else { + (elapsed_micros / interval_micros).saturating_add(1) + }; + let offset = interval_micros.checked_mul(ticks).unwrap_or(i64::MAX); + last_intended_at + .checked_add(TimeDuration::from_micros(offset)) + .or_else(|| now.checked_add(interval)) + .unwrap_or(now) } /// Deletes a scheduled-row entry inside an existing mutable transaction. @@ -902,3 +952,30 @@ fn read_schedule_at(row: &RowRef<'_>, at_column: ColId) -> anyhow::Result Timestamp { + Timestamp::from_micros_since_unix_epoch(micros) + } + + #[test] + fn next_interval_tick_uses_last_intended_tick_when_current() { + let next = next_interval_tick_after(ts(1_000), TimeDuration::from_micros(100), ts(1_000)); + assert_eq!(next, ts(1_100)); + } + + #[test] + fn next_interval_tick_skips_missed_ticks() { + let next = next_interval_tick_after(ts(1_000), TimeDuration::from_micros(100), ts(1_350)); + assert_eq!(next, ts(1_400)); + } + + #[test] + fn next_interval_tick_is_strictly_after_now_on_boundary() { + let next = next_interval_tick_after(ts(1_000), TimeDuration::from_micros(100), ts(1_300)); + assert_eq!(next, ts(1_400)); + } +} From 98f55fea4119eada9ca21736599959b3956f5f45 Mon Sep 17 00:00:00 2001 From: Shubham Mishra Date: Thu, 13 Aug 2026 17:35:19 +0530 Subject: [PATCH 2/3] naming --- crates/core/src/host/scheduler.rs | 31 +++++++++++++------------------ 1 file changed, 13 insertions(+), 18 deletions(-) diff --git a/crates/core/src/host/scheduler.rs b/crates/core/src/host/scheduler.rs index 009b408cc04..315fec27d73 100644 --- a/crates/core/src/host/scheduler.rs +++ b/crates/core/src/host/scheduler.rs @@ -448,7 +448,7 @@ struct Reschedule { } #[derive(Clone)] -struct ScheduledFunctionTiming { +struct ScheduledInvocation { function_name: Arc, intended_at: Timestamp, } @@ -522,7 +522,7 @@ fn prepare_scheduled_procedure_call( inst: &mut impl WasmInstance, ) -> ScheduledProcedureStep { let ScheduledFunctionParams(item) = params; - let timing = scheduled_function_timing_for_item(&item); + let timing = scheduled_invocation_for_item(&item); let id = scheduled_item_id(&item); let db = &**module_info.relational_db(); let tx = db.begin_mut_tx(IsolationLevel::Serializable, Workload::Internal); @@ -547,18 +547,13 @@ fn prepare_scheduled_procedure_call( } }; + let intended_at = timing.as_ref().map_or(timestamp, |timing| timing.intended_at); + // For scheduled procedures, it's incorrect to retry them if execution aborts midway, // so we must remove the schedule row before executing. - let reschedule = id.and_then(|id| { - let reschedule_from = timing.as_ref().map(|timing| timing.intended_at).unwrap_or(timestamp); - delete_scheduled_function_row(module_info, db, id, Some(tx), reschedule_from, inst_common, inst) - }); - let delay = timing.map(|timing| { - ( - timing.function_name, - scheduled_function_delay(timestamp, timing.intended_at), - ) - }); + let reschedule = + id.and_then(|id| delete_scheduled_function_row(module_info, db, id, Some(tx), intended_at, inst_common, inst)); + let delay = timing.map(|timing| (timing.function_name, scheduled_function_delay(timestamp, intended_at))); ScheduledProcedureStep::Procedure { params: Box::new(params), reschedule, @@ -573,7 +568,7 @@ fn call_scheduled_reducer_until_done( inst: &mut impl WasmInstance, ) -> (CallScheduledFunctionResult, bool) { let ScheduledFunctionParams(item) = params; - let timing = scheduled_function_timing_for_item(&item); + let timing = scheduled_invocation_for_item(&item); let id = scheduled_item_id(&item); let db = &**module_info.relational_db(); let tx = db.begin_mut_tx(IsolationLevel::Serializable, Workload::Internal); @@ -598,12 +593,12 @@ fn call_scheduled_reducer_until_done( } }; + let intended_at = timing.as_ref().map_or(timestamp, |timing| timing.intended_at); if let Some(timing) = timing.as_ref() { - let delay = scheduled_function_delay(timestamp, timing.intended_at); + let delay = scheduled_function_delay(timestamp, intended_at); record_scheduled_function_delay(module_info, &timing.function_name, delay); } - let reschedule_from = timing.as_ref().map(|timing| timing.intended_at).unwrap_or(timestamp); - call_scheduled_reducer_with_tx(module_info, db, id, tx, reschedule_from, params, inst_common, inst) + call_scheduled_reducer_with_tx(module_info, db, id, tx, intended_at, params, inst_common, inst) } fn scheduled_item_id(item: &QueueItem) -> Option { @@ -613,9 +608,9 @@ fn scheduled_item_id(item: &QueueItem) -> Option { } } -fn scheduled_function_timing_for_item(item: &QueueItem) -> Option { +fn scheduled_invocation_for_item(item: &QueueItem) -> Option { match item { - QueueItem::Id { function_name, at, .. } => Some(ScheduledFunctionTiming { + QueueItem::Id { function_name, at, .. } => Some(ScheduledInvocation { function_name: function_name.clone(), intended_at: *at, }), From 3c42c5a8c9c6f6b71f2055d57a1869e418496c00 Mon Sep 17 00:00:00 2001 From: Shubham Mishra Date: Thu, 13 Aug 2026 18:01:15 +0530 Subject: [PATCH 3/3] comment --- crates/core/src/host/scheduler.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/crates/core/src/host/scheduler.rs b/crates/core/src/host/scheduler.rs index 315fec27d73..29ec9c75834 100644 --- a/crates/core/src/host/scheduler.rs +++ b/crates/core/src/host/scheduler.rs @@ -720,6 +720,8 @@ fn next_interval_reschedule(last_intended_at: Timestamp, interval: TimeDuration) } } +/// Returns the first interval tick after `now`, anchored at `last_intended_at`, +/// skipping any ticks that have already been missed. fn next_interval_tick_after(last_intended_at: Timestamp, interval: TimeDuration, now: Timestamp) -> Timestamp { let interval = interval.abs(); let interval_micros = interval.to_micros();