mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(workflow): add scheduled_workflow_fires.workflow_run_id for the attach audit link
The scheduler's post-run `attach_scheduled_workflow_run` does `SET workflow_run_id`, but neither schema/schema.sql nor the initial migration declared that column, so every won scheduled claim would create+run, then the best-effort attach would hit `column "workflow_run_id" does not exist` (PG 42703) and warn forever — the audit link could never populate. Found by Max in cold review of 778a5d28c. - Add nullable `workflow_run_id UUID` to scheduled_workflow_fires (schema + migration) with a composite FK `(community_id, workflow_run_id) REFERENCES workflow_runs (community_id, id)`. The FK uses ON DELETE NO ACTION, not SET NULL: community_id is shared with the claim PK and is NOT NULL, so SET NULL is unimplementable (verified against live PG: it raises a NOT NULL violation mid-cascade). NO ACTION blocks a delete of a still-linked run cleanly; workflow_runs are not pruned today regardless. - Rewrite the stale scheduled-fires schema comment that still claimed community is 'resolved server-side from workflow_id, never a caller-supplied claim parameter' — contradicted by the S1 reconciliation: community is server provenance from list_all_enabled_workflows(), passed explicitly, since id is not globally unique. - Surface the interval-anchor read failure with a warn! instead of unwrap_or(None) swallowing it (still fail-closed: a missing anchor suppresses the tick and retries). - Add attach_links_run_to_claim_and_is_idempotent: proves the column populates on attach and the IS NULL guard makes a second attach a no-op. Proven RED against the pre-migration schema (the exact 42703 error), GREEN after — the regression that would have caught this gap. Verified vs live Postgres, serial: buzz-db 110 / buzz-workflow 145 / buzz-relay 414, 0 failed; clippy -D warnings clean on the trio. Co-authored-by: Tyler Longwell <tlongwell@block.xyz> Signed-off-by: Tyler Longwell <tlongwell@block.xyz>
This commit is contained in:
co-authored by
Tyler Longwell
parent
4dc65b6ee5
commit
00f75bcfc9
@@ -1753,6 +1753,72 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// `attach_scheduled_workflow_run` links a won claim to the run it created.
|
||||
/// This is the regression for the missing `scheduled_workflow_fires.
|
||||
/// workflow_run_id` column: before the schema added it, the UPDATE failed at
|
||||
/// runtime with `column "workflow_run_id" does not exist`, so the audit link
|
||||
/// silently never populated and the scheduler warned on every fire. This test
|
||||
/// proves the column is present, the attach writes it, and the
|
||||
/// `workflow_run_id IS NULL` guard makes a second attach a no-op. It is RED
|
||||
/// without the migration column.
|
||||
#[tokio::test]
|
||||
#[ignore = "requires Postgres"]
|
||||
async fn attach_links_run_to_claim_and_is_idempotent() {
|
||||
let pool = setup_pool().await;
|
||||
|
||||
let community = make_community(&pool).await;
|
||||
let (workflow_id, _) = make_workflow_in(&pool, community).await;
|
||||
let scheduled_for = Utc.with_ymd_and_hms(2026, 6, 27, 0, 2, 0).unwrap();
|
||||
|
||||
// Win the claim for this instant.
|
||||
claim_scheduled_workflow_fire(&pool, community, workflow_id, scheduled_for)
|
||||
.await
|
||||
.expect("claim ok")
|
||||
.expect("claim wins");
|
||||
|
||||
// Create the run the won claim is responsible for, then attach it.
|
||||
let run_id = create_workflow_run(&pool, community, workflow_id, None, None)
|
||||
.await
|
||||
.expect("create run ok");
|
||||
|
||||
let attached =
|
||||
attach_scheduled_workflow_run(&pool, community, workflow_id, scheduled_for, run_id)
|
||||
.await
|
||||
.expect("attach ok");
|
||||
assert!(attached, "first attach must update the claim row");
|
||||
|
||||
// The column is populated with the run id.
|
||||
let linked: Option<Uuid> = sqlx::query_scalar(
|
||||
"SELECT workflow_run_id FROM scheduled_workflow_fires \
|
||||
WHERE community_id = $1 AND workflow_id = $2 AND scheduled_for = $3",
|
||||
)
|
||||
.bind(community.as_uuid())
|
||||
.bind(workflow_id)
|
||||
.bind(scheduled_for)
|
||||
.fetch_one(&pool)
|
||||
.await
|
||||
.expect("row exists");
|
||||
assert_eq!(
|
||||
linked,
|
||||
Some(run_id),
|
||||
"the claim row must now point at the run it created"
|
||||
);
|
||||
|
||||
// A second attach is a no-op: the `workflow_run_id IS NULL` guard means
|
||||
// an already-linked claim is never re-pointed to a different run.
|
||||
let other_run = create_workflow_run(&pool, community, workflow_id, None, None)
|
||||
.await
|
||||
.expect("create second run ok");
|
||||
let reattached =
|
||||
attach_scheduled_workflow_run(&pool, community, workflow_id, scheduled_for, other_run)
|
||||
.await
|
||||
.expect("second attach ok");
|
||||
assert!(
|
||||
!reattached,
|
||||
"attach must not overwrite an already-linked claim row"
|
||||
);
|
||||
}
|
||||
|
||||
/// Documents the retention-vs-interval coupling Sami flagged for §5c:
|
||||
/// pruning every claim below the workflow's interval makes
|
||||
/// `latest_scheduled_workflow_fire` return `None`, which re-introduces the
|
||||
|
||||
@@ -429,11 +429,28 @@ impl WorkflowEngine {
|
||||
// a process bounce can't double-fire within an interval.
|
||||
let last = match self.last_fired.get(&(community_id, workflow.id)) {
|
||||
Some(t) => Some(*t),
|
||||
None => self
|
||||
None => match self
|
||||
.db
|
||||
.latest_scheduled_workflow_fire(community_id, workflow.id)
|
||||
.await
|
||||
.unwrap_or(None),
|
||||
{
|
||||
Ok(anchor) => anchor,
|
||||
Err(e) => {
|
||||
// Fail closed: a missing anchor reads as
|
||||
// last_fired = now in interval_should_fire,
|
||||
// so this tick is suppressed and the next
|
||||
// tick retries. Surface the read failure so
|
||||
// a persistently-unreadable anchor is visible
|
||||
// rather than silently stalling the schedule.
|
||||
tracing::warn!(
|
||||
community_id = %community_id,
|
||||
workflow_id = %workflow.id,
|
||||
"Cron tick: failed to read interval restart anchor, \
|
||||
suppressing this tick: {e}"
|
||||
);
|
||||
None
|
||||
}
|
||||
},
|
||||
};
|
||||
if !interval_should_fire(dur, last, now, workflow.id) {
|
||||
continue;
|
||||
|
||||
@@ -433,17 +433,28 @@ CREATE INDEX idx_workflow_approvals_status ON workflow_approvals (community_id,
|
||||
-- ── Scheduled workflow fires (cron claim) ─────────────────────────────────────
|
||||
-- Plan §5: the at-most-once cron fire claim. UNIQUE (community_id, workflow_id,
|
||||
-- scheduled_for) — only the pod that wins the claim insert creates the run.
|
||||
-- Restart-safe (DB-durable). community resolved server-side from workflow_id,
|
||||
-- never a caller-supplied claim parameter (S1 tenant binding).
|
||||
-- Restart-safe (DB-durable). community is server provenance: the scheduler passes
|
||||
-- workflow.community_id from list_all_enabled_workflows(), never a client input.
|
||||
-- workflow_id is NOT globally unique under the (community_id, id) workflow key, so
|
||||
-- the claim binds both community and id explicitly rather than resolving from id.
|
||||
-- workflow_run_id links the won claim to the run it created (audit; NULL until the
|
||||
-- post-insert attach, and stays NULL if run creation failed after a won claim).
|
||||
-- The FK to workflow_runs uses NO ACTION (not SET NULL): community_id is shared
|
||||
-- with the claim PK and is NOT NULL, so SET NULL is unimplementable here; a future
|
||||
-- delete of a still-linked run is blocked rather than orphaning the at-most-once
|
||||
-- claim row. workflow_runs are not pruned today, so this is a guardrail, not a path.
|
||||
|
||||
CREATE TABLE scheduled_workflow_fires (
|
||||
community_id UUID NOT NULL REFERENCES communities(id),
|
||||
workflow_id UUID NOT NULL,
|
||||
scheduled_for TIMESTAMPTZ NOT NULL,
|
||||
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
workflow_run_id UUID,
|
||||
PRIMARY KEY (community_id, workflow_id, scheduled_for),
|
||||
FOREIGN KEY (community_id, workflow_id)
|
||||
REFERENCES workflows (community_id, id) ON DELETE CASCADE
|
||||
REFERENCES workflows (community_id, id) ON DELETE CASCADE,
|
||||
FOREIGN KEY (community_id, workflow_run_id)
|
||||
REFERENCES workflow_runs (community_id, id) ON DELETE NO ACTION
|
||||
);
|
||||
|
||||
-- The interval anchor reads MAX(scheduled_for) per workflow; the janitor prunes
|
||||
|
||||
+14
-3
@@ -417,17 +417,28 @@ CREATE INDEX idx_workflow_approvals_status ON workflow_approvals (community_id,
|
||||
-- ── Scheduled workflow fires (cron claim) ─────────────────────────────────────
|
||||
-- Plan §5: the at-most-once cron fire claim. UNIQUE (community_id, workflow_id,
|
||||
-- scheduled_for) — only the pod that wins the claim insert creates the run.
|
||||
-- Restart-safe (DB-durable). community resolved server-side from workflow_id,
|
||||
-- never a caller-supplied claim parameter (S1 tenant binding).
|
||||
-- Restart-safe (DB-durable). community is server provenance: the scheduler passes
|
||||
-- workflow.community_id from list_all_enabled_workflows(), never a client input.
|
||||
-- workflow_id is NOT globally unique under the (community_id, id) workflow key, so
|
||||
-- the claim binds both community and id explicitly rather than resolving from id.
|
||||
-- workflow_run_id links the won claim to the run it created (audit; NULL until the
|
||||
-- post-insert attach, and stays NULL if run creation failed after a won claim).
|
||||
-- The FK to workflow_runs uses NO ACTION (not SET NULL): community_id is shared
|
||||
-- with the claim PK and is NOT NULL, so SET NULL is unimplementable here; a future
|
||||
-- delete of a still-linked run is blocked rather than orphaning the at-most-once
|
||||
-- claim row. workflow_runs are not pruned today, so this is a guardrail, not a path.
|
||||
|
||||
CREATE TABLE scheduled_workflow_fires (
|
||||
community_id UUID NOT NULL REFERENCES communities(id),
|
||||
workflow_id UUID NOT NULL,
|
||||
scheduled_for TIMESTAMPTZ NOT NULL,
|
||||
claimed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
workflow_run_id UUID,
|
||||
PRIMARY KEY (community_id, workflow_id, scheduled_for),
|
||||
FOREIGN KEY (community_id, workflow_id)
|
||||
REFERENCES workflows (community_id, id) ON DELETE CASCADE
|
||||
REFERENCES workflows (community_id, id) ON DELETE CASCADE,
|
||||
FOREIGN KEY (community_id, workflow_run_id)
|
||||
REFERENCES workflow_runs (community_id, id) ON DELETE NO ACTION
|
||||
);
|
||||
|
||||
-- The interval anchor reads MAX(scheduled_for) per workflow; the janitor prunes
|
||||
|
||||
Reference in New Issue
Block a user