mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(workflow): skip missed interval fires
Co-authored-by: npub1mprnacetjua2xx3p5eddmhxyk6wv929ymm5py8kd2xfxurxahspqqlgyta <d8473ee32b973aa31a21a65adddcc4b69cc2a8a4dee8121ecd51926e0cddbc02@sprout-oss.stage.blox.sqprod.co> Signed-off-by: npub1mprnacetjua2xx3p5eddmhxyk6wv929ymm5py8kd2xfxurxahspqqlgyta <d8473ee32b973aa31a21a65adddcc4b69cc2a8a4dee8121ecd51926e0cddbc02@sprout-oss.stage.blox.sqprod.co>
This commit is contained in:
parent
308181fe59
commit
71da65e513
@@ -549,13 +549,14 @@ fn cron_should_fire(
|
||||
}
|
||||
}
|
||||
|
||||
/// Return the next interval schedule instant due at or before `now`.
|
||||
/// Return the most recent interval schedule instant due at or before `now`.
|
||||
///
|
||||
/// The anchor is DB-authoritative: the last claimed schedule instant when one
|
||||
/// exists, otherwise the workflow row's `created_at`. Returning exactly
|
||||
/// `anchor + interval` makes all pods compute the same `(workflow_id,
|
||||
/// exists, otherwise the workflow row's `created_at`. Returning a canonical
|
||||
/// interval boundary makes all pods compute the same `(workflow_id,
|
||||
/// scheduled_for)` claim key even when their local clocks are in different parts
|
||||
/// of the same interval window; catch-up happens at most one interval per tick.
|
||||
/// of the same interval window. If multiple intervals elapsed while the engine
|
||||
/// was down, old intervals are skipped instead of replayed one per tick.
|
||||
///
|
||||
/// Returns `None` (and logs a warning) if the duration string is invalid.
|
||||
fn interval_should_fire(
|
||||
@@ -584,8 +585,13 @@ fn interval_should_fire(
|
||||
return None;
|
||||
}
|
||||
|
||||
let next_fire = anchor + chrono::Duration::seconds(interval_secs);
|
||||
(next_fire <= now).then_some(next_fire)
|
||||
let elapsed_secs = (now - anchor).num_seconds();
|
||||
if elapsed_secs < interval_secs {
|
||||
return None;
|
||||
}
|
||||
|
||||
let elapsed_intervals = elapsed_secs / interval_secs;
|
||||
Some(anchor + chrono::Duration::seconds(interval_secs * elapsed_intervals))
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
@@ -878,7 +884,21 @@ mod tests {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn interval_should_fire_returns_next_due_boundary_after_elapsed_intervals() {
|
||||
fn interval_should_fire_skips_missed_intervals_from_created_at() {
|
||||
let wf_id = Uuid::new_v4();
|
||||
let created_at = chrono::DateTime::parse_from_rfc3339("2026-06-15T09:00:00Z")
|
||||
.unwrap()
|
||||
.with_timezone(&Utc);
|
||||
let now = created_at + chrono::Duration::hours(3) + chrono::Duration::minutes(5);
|
||||
assert_eq!(
|
||||
interval_should_fire("1h", None, created_at, now, wf_id),
|
||||
Some(created_at + chrono::Duration::hours(3)),
|
||||
"should claim the most recent due boundary, not replay old missed intervals"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn interval_should_fire_skips_missed_intervals_after_latest_claim() {
|
||||
let wf_id = Uuid::new_v4();
|
||||
let created_at = chrono::DateTime::parse_from_rfc3339("2026-06-15T09:00:00Z")
|
||||
.unwrap()
|
||||
@@ -887,8 +907,8 @@ mod tests {
|
||||
let now = latest + chrono::Duration::hours(2) + chrono::Duration::minutes(5);
|
||||
assert_eq!(
|
||||
interval_should_fire("1h", Some(latest), created_at, now, wf_id),
|
||||
Some(latest + chrono::Duration::hours(1)),
|
||||
"should return the canonical next due boundary, not local now"
|
||||
Some(latest + chrono::Duration::hours(2)),
|
||||
"should return the most recent due boundary, not catch up one interval per tick"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user