fix(workflow): snapshot definition on run for approval resume

Store definition JSON + hash on workflow_runs at trigger time so resume
rejects live NIP-33 drift instead of silently executing post-gate edits.

Signed-off-by: Joshua Belke <joshua@innovationhub-act.org>
This commit is contained in:
Josh Belke 2026-08-09 22:00:45 -04:00
commit d98dc1301b
8 changed files with 282 additions and 31 deletions

View file

@ -3969,12 +3969,17 @@ impl Db {
}
/// Create a new workflow run.
///
/// Pass `definition_snapshot` / `definition_hash` from the live workflow at
/// trigger time so approval resume can reject definition drift.
pub async fn create_workflow_run(
&self,
community_id: CommunityId,
workflow_id: Uuid,
trigger_event_id: Option<&[u8]>,
trigger_context: Option<&serde_json::Value>,
definition_snapshot: Option<&serde_json::Value>,
definition_hash: Option<&[u8]>,
) -> Result<Uuid> {
workflow::create_workflow_run(
&self.pool,
@ -3982,6 +3987,8 @@ impl Db {
workflow_id,
trigger_event_id,
trigger_context,
definition_snapshot,
definition_hash,
)
.await
}

View file

@ -212,6 +212,19 @@ pub struct WorkflowRunRecord {
/// Serialized `TriggerContext` captured at workflow start.
/// NULL for runs created before this column was added (backwards-compatible).
pub trigger_context: Option<serde_json::Value>,
/// Definition JSON copied from the workflow at trigger time.
///
/// Approval resume must execute this snapshot (or reject) rather than
/// re-reading the live `workflows.definition`, which is NIP-33 LWW and can
/// change during a multi-day approval wait (`meridian-lyw4`).
/// NULL for runs created before this column was added.
pub definition_snapshot: Option<serde_json::Value>,
/// SHA-256 of the snapshotted definition at trigger time.
///
/// Compared to the live workflow's `definition_hash` on resume; a mismatch
/// is an explicit rejection, never a silent step swap.
/// NULL for runs created before this column was added.
pub definition_hash: Option<Vec<u8>>,
/// When execution began.
pub started_at: Option<DateTime<Utc>>,
/// When execution finished (success or failure).
@ -848,20 +861,28 @@ pub async fn delete_workflow_for_owner(
/// `trigger_context` is the serialized `TriggerContext` for this run. It is stored
/// so that post-approval resume steps can restore the original trigger data and
/// correctly resolve `{{trigger.*}}` template variables.
///
/// `definition_snapshot` / `definition_hash` freeze the workflow body that this
/// run is authorized to execute. Callers must pass the live definition at
/// trigger time; approval resume compares the stored hash to the live row and
/// rejects (or freezes to the snapshot) on drift (`meridian-lyw4`).
pub async fn create_workflow_run(
pool: &PgPool,
community_id: CommunityId,
workflow_id: Uuid,
trigger_event_id: Option<&[u8]>,
trigger_context: Option<&serde_json::Value>,
definition_snapshot: Option<&serde_json::Value>,
definition_hash: Option<&[u8]>,
) -> Result<Uuid> {
let id = Uuid::new_v4();
sqlx::query(
r#"
INSERT INTO workflow_runs
(community_id, id, workflow_id, status, trigger_event_id, current_step, execution_trace, trigger_context)
VALUES ($1, $2, $3, 'pending', $4, 0, '[]', $5)
(community_id, id, workflow_id, status, trigger_event_id, current_step,
execution_trace, trigger_context, definition_snapshot, definition_hash)
VALUES ($1, $2, $3, 'pending', $4, 0, '[]', $5, $6, $7)
"#,
)
.bind(community_id.as_uuid())
@ -869,6 +890,8 @@ pub async fn create_workflow_run(
.bind(workflow_id)
.bind(trigger_event_id)
.bind(trigger_context)
.bind(definition_snapshot)
.bind(definition_hash)
.execute(pool)
.await?;
@ -884,7 +907,8 @@ pub async fn get_workflow_run(
let row = sqlx::query(
r#"
SELECT community_id, id, workflow_id, status::text AS status, trigger_event_id, current_step,
execution_trace, trigger_context, started_at, completed_at, error_message, created_at
execution_trace, trigger_context, definition_snapshot, definition_hash,
started_at, completed_at, error_message, created_at
FROM workflow_runs
WHERE community_id = $1 AND id = $2
"#,
@ -909,7 +933,8 @@ pub async fn list_workflow_runs(
let rows = sqlx::query(
r#"
SELECT community_id, id, workflow_id, status::text AS status, trigger_event_id, current_step,
execution_trace, trigger_context, started_at, completed_at, error_message, created_at
execution_trace, trigger_context, definition_snapshot, definition_hash,
started_at, completed_at, error_message, created_at
FROM workflow_runs
WHERE community_id = $1 AND workflow_id = $2
ORDER BY created_at DESC
@ -1219,6 +1244,8 @@ fn row_to_run_record(row: sqlx::postgres::PgRow) -> Result<WorkflowRunRecord> {
current_step: row.try_get("current_step")?,
execution_trace: row.try_get("execution_trace")?,
trigger_context: row.try_get("trigger_context")?,
definition_snapshot: row.try_get("definition_snapshot")?,
definition_hash: row.try_get("definition_hash")?,
started_at: row.try_get("started_at")?,
completed_at: row.try_get("completed_at")?,
error_message: row.try_get("error_message")?,
@ -1523,6 +1550,8 @@ mod tests {
{ "step": "s1", "status": "completed" }
]),
trigger_context: None,
definition_snapshot: None,
definition_hash: None,
started_at: Some(now),
completed_at: None,
error_message: None,
@ -1533,12 +1562,46 @@ mod tests {
assert_eq!(record.workflow_id, workflow_id);
assert_eq!(record.status, RunStatus::Running);
assert_eq!(record.trigger_event_id, Some(trigger_event_id));
assert_eq!(record.trigger_context, None);
assert_eq!(record.definition_snapshot, None);
assert_eq!(record.definition_hash, None);
assert_eq!(record.current_step, 2);
assert!(record.started_at.is_some());
assert!(record.completed_at.is_none());
assert!(record.error_message.is_none());
}
#[test]
fn workflow_run_record_carries_definition_snapshot() {
let now = Utc::now();
let snapshot = serde_json::json!({
"name": "snap",
"trigger": { "on": "webhook" },
"steps": []
});
let hash = vec![0xab; 32];
let record = WorkflowRunRecord {
id: Uuid::new_v4(),
community_id: CommunityId::from_uuid(Uuid::new_v4()),
workflow_id: Uuid::new_v4(),
status: RunStatus::WaitingApproval,
trigger_event_id: None,
current_step: 1,
execution_trace: serde_json::json!([]),
trigger_context: None,
definition_snapshot: Some(snapshot.clone()),
definition_hash: Some(hash.clone()),
started_at: Some(now),
completed_at: None,
error_message: None,
created_at: now,
};
assert_eq!(record.definition_snapshot, Some(snapshot));
assert_eq!(record.definition_hash, Some(hash));
assert_eq!(record.status, RunStatus::WaitingApproval);
}
#[test]
fn workflow_run_record_no_trigger_event() {
let now = Utc::now();
@ -1551,6 +1614,8 @@ mod tests {
current_step: 0,
execution_trace: serde_json::json!([]),
trigger_context: None,
definition_snapshot: None,
definition_hash: None,
started_at: None,
completed_at: None,
error_message: None,
@ -1574,6 +1639,8 @@ mod tests {
current_step: 1,
execution_trace: serde_json::json!([]),
trigger_context: None,
definition_snapshot: None,
definition_hash: None,
started_at: Some(now),
completed_at: Some(now),
error_message: Some("step timeout exceeded".to_owned()),
@ -1605,6 +1672,8 @@ mod tests {
current_step: 2,
execution_trace: trace.clone(),
trigger_context: None,
definition_snapshot: None,
definition_hash: None,
started_at: Some(now),
completed_at: Some(now),
error_message: None,
@ -1627,6 +1696,8 @@ mod tests {
current_step: 0,
execution_trace: serde_json::json!([]),
trigger_context: None,
definition_snapshot: None,
definition_hash: None,
started_at: None,
completed_at: None,
error_message: None,
@ -1965,7 +2036,7 @@ mod tests {
.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)
let run_id = create_workflow_run(&pool, community, workflow_id, None, None, None, None)
.await
.expect("create run ok");
@ -1994,7 +2065,7 @@ mod tests {
// 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)
let other_run = create_workflow_run(&pool, community, workflow_id, None, None, None, None)
.await
.expect("create second run ok");
let reattached =
@ -2007,6 +2078,49 @@ mod tests {
);
}
/// `create_workflow_run` persists the definition snapshot + hash so approval
/// resume can reject live drift (`meridian-lyw4`). RED without migration
/// `0032` columns / SELECT projection.
#[tokio::test]
#[ignore = "requires Postgres"]
async fn create_workflow_run_persists_definition_snapshot() {
let pool = setup_pool().await;
let community = make_community(&pool).await;
let (workflow_id, _) = make_workflow_in(&pool, community).await;
let snapshot = serde_json::json!({
"name": "snap-run",
"trigger": { "on": "webhook" },
"steps": [
{
"id": "gate",
"action": "request_approval",
"from": "@lead",
"message": "ok?"
}
]
});
let hash = vec![0x11u8; 32];
let run_id = create_workflow_run(
&pool,
community,
workflow_id,
None,
None,
Some(&snapshot),
Some(&hash),
)
.await
.expect("create run with snapshot");
let run = get_workflow_run(&pool, community, run_id)
.await
.expect("fetch run");
assert_eq!(run.definition_snapshot.as_ref(), Some(&snapshot));
assert_eq!(run.definition_hash.as_deref(), Some(hash.as_slice()));
}
/// 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
@ -2275,10 +2389,10 @@ mod tests {
insert_workflow_with_ids(&pool, community_a, workflow_id, channel_id, "wf-A").await;
insert_workflow_with_ids(&pool, community_b, workflow_id, Uuid::new_v4(), "wf-B").await;
let run_a = create_workflow_run(&pool, community_a, workflow_id, None, None)
let run_a = create_workflow_run(&pool, community_a, workflow_id, None, None, None, None)
.await
.expect("run A");
let run_b = create_workflow_run(&pool, community_b, workflow_id, None, None)
let run_b = create_workflow_run(&pool, community_b, workflow_id, None, None, None, None)
.await
.expect("run B");

View file

@ -2099,7 +2099,14 @@ pub async fn workflow_webhook(
let run_id = state
.db
.create_workflow_run(community_id, id, None, trigger_ctx_json.as_ref())
.create_workflow_run(
community_id,
id,
None,
trigger_ctx_json.as_ref(),
Some(&workflow.definition),
Some(workflow.definition_hash.as_slice()),
)
.await
.map_err(|e| super::internal_error(&format!("db error: {e}")))?;

View file

@ -1015,6 +1015,8 @@ async fn handle_workflow_trigger(
workflow_id,
Some(&event_id_bytes),
trigger_ctx_json.as_ref(),
Some(&workflow.definition),
Some(workflow.definition_hash.as_slice()),
)
.await
.map_err(|e| IngestError::Internal(format!("error: db create_workflow_run: {e}")))?;
@ -1394,27 +1396,37 @@ async fn resume_workflow_after_approval(
}
};
let def: meridian_workflow::WorkflowDef =
match serde_json::from_value(workflow.definition.clone()) {
Ok(d) => d,
Err(e) => {
tracing::error!("resume_workflow: failed to parse workflow definition: {e}");
if let Err(db_err) = db
.update_workflow_run(
community_id,
run_id,
RunStatus::Failed,
run.current_step,
&run.execution_trace,
Some(&format!("definition parse error: {e}")),
)
.await
{
tracing::error!("resume_workflow: failed to mark run as failed: {db_err}");
}
return;
// Freeze to the definition snapshotted at trigger time. A live edit during
// the approval wait must not silently change post-gate steps — reject with
// an explicit run failure when the live hash drifts (`meridian-lyw4`).
let def = match meridian_workflow::definition_for_approval_resume(
run.definition_snapshot.as_ref(),
run.definition_hash.as_deref(),
&workflow.definition_hash,
) {
Ok(d) => d,
Err(e) => {
tracing::warn!(
run_id = %run_id,
workflow_id = %workflow_id,
"resume_workflow: {e}"
);
if let Err(db_err) = db
.update_workflow_run(
community_id,
run_id,
RunStatus::Failed,
run.current_step,
&run.execution_trace,
Some(&e.to_string()),
)
.await
{
tracing::error!("resume_workflow: failed to mark run as failed: {db_err}");
}
};
return;
}
};
// Reconstruct step_outputs from execution trace for template resolution
let mut initial_outputs: std::collections::HashMap<String, serde_json::Value> =

View file

@ -867,6 +867,46 @@ async fn resolve_webhook_secret_headers(
Ok(Some(resolved))
}
/// Resolve the workflow definition an approval-resume may execute.
///
/// A run stores the definition JSON and its hash at trigger time. Resume must
/// not re-read the live `workflows` row for step bodies: that row is NIP-33
/// LWW and can change during a multi-day approval wait (`meridian-lyw4`).
///
/// - Missing snapshot/hash → reject with an explicit record (fail closed).
/// - Live hash ≠ snapshotted hash → reject with an explicit mismatch record.
/// - Match → parse and return the **snapshotted** definition (freeze), never
/// the live JSON, so a hash collision cannot still swap step bodies.
pub fn definition_for_approval_resume(
snapshot: Option<&serde_json::Value>,
snapshotted_hash: Option<&[u8]>,
live_hash: &[u8],
) -> Result<WorkflowDef, WorkflowError> {
let Some(snapshot) = snapshot else {
return Err(WorkflowError::InvalidDefinition(
"approval resume rejected: run has no definition snapshot".into(),
));
};
let Some(snap_hash) = snapshotted_hash else {
return Err(WorkflowError::InvalidDefinition(
"approval resume rejected: run has no definition hash".into(),
));
};
if snap_hash != live_hash {
return Err(WorkflowError::InvalidDefinition(format!(
"approval resume rejected: workflow definition changed during approval wait \
(snapshotted hash {}, live hash {})",
hex::encode(snap_hash),
hex::encode(live_hash),
)));
}
serde_json::from_value(snapshot.clone()).map_err(|e| {
WorkflowError::InvalidDefinition(format!(
"approval resume rejected: snapshotted definition is corrupt: {e}"
))
})
}
/// Resolve a `request_approval` step's absolute expiry.
///
/// `now` is a parameter rather than read inside, so the bound is testable
@ -1552,6 +1592,66 @@ mod tests {
assert!(resolve_approval_expiry(Some("soon"), fixed_now()).is_err());
}
fn minimal_workflow_json() -> serde_json::Value {
serde_json::json!({
"name": "freeze-test",
"trigger": { "on": "webhook" },
"steps": [
{
"id": "gate",
"action": "request_approval",
"from": "@lead",
"message": "ok?"
},
{
"id": "after",
"action": "add_reaction",
"emoji": "thumbsup"
}
]
})
}
#[test]
fn approval_resume_rejects_when_snapshot_missing() {
let err = definition_for_approval_resume(None, Some(&[0x11]), &[0x11])
.expect_err("missing snapshot must reject");
let msg = err.to_string();
assert!(
msg.contains("no definition snapshot"),
"expected explicit missing-snapshot record, got: {msg}"
);
}
#[test]
fn approval_resume_rejects_when_live_definition_hash_drifts() {
let snapshot = minimal_workflow_json();
let snap_hash = [0xAA; 32];
let live_hash = [0xBB; 32];
let err = definition_for_approval_resume(Some(&snapshot), Some(&snap_hash), &live_hash)
.expect_err("hash drift must reject");
let msg = err.to_string();
assert!(
msg.contains("definition changed during approval wait"),
"expected explicit mismatch record, got: {msg}"
);
assert!(msg.contains(&hex::encode(snap_hash)));
assert!(msg.contains(&hex::encode(live_hash)));
}
#[test]
fn approval_resume_freezes_to_snapshotted_definition_on_hash_match() {
let snapshot = minimal_workflow_json();
let hash = [0xCC; 32];
// Live body could differ textually; hash match is the gate, and the
// returned def must come from the snapshot — not from a live re-read.
let def = definition_for_approval_resume(Some(&snapshot), Some(&hash), &hash)
.expect("matching hash resumes from snapshot");
assert_eq!(def.name, "freeze-test");
assert_eq!(def.steps.len(), 2);
assert_eq!(def.steps[1].id, "after");
}
#[test]
fn pending_approval_debug_never_leaks_the_token() {
// The token is a bearer credential: whoever holds it can act on the

View file

@ -37,7 +37,7 @@ pub mod schema;
pub use action_sink::{ActionSink, ActionSinkError};
pub use error::{PartialProgress, WorkflowError};
pub use executor::{ExecutionResult, PendingApproval};
pub use executor::{definition_for_approval_resume, ExecutionResult, PendingApproval};
pub use schema::{ActionDef, Step, TriggerDef, WorkflowDef};
use std::collections::HashMap;
@ -549,6 +549,8 @@ impl WorkflowEngine {
workflow.id,
Some(&trigger_event_id_bytes),
Some(&trigger_ctx_json),
Some(&workflow.definition),
Some(workflow.definition_hash.as_slice()),
)
.await
{
@ -811,6 +813,8 @@ impl WorkflowEngine {
workflow.id,
None, // no trigger event for cron
trigger_ctx_json.as_ref(),
Some(&workflow.definition),
Some(workflow.definition_hash.as_slice()),
)
.await
{

View file

@ -0,0 +1,7 @@
-- Snapshot the workflow definition onto each run at trigger time so an
-- approval-resume cannot silently pick up a later edit of the live definition
-- (meridian-lyw4). Both columns are nullable for rows created before this
-- migration; the resume path fails closed when they are absent.
ALTER TABLE workflow_runs
ADD COLUMN definition_snapshot JSONB,
ADD COLUMN definition_hash BYTEA;

View file

@ -8,7 +8,7 @@ files are compiled into the relay binary and applied automatically at startup.
## Ownership
`NNNN_snake_case.sql`, zero-padded four-digit sequence, currently `0001` →
`0031`. `0001_initial_schema.sql` is the consolidated multi-tenant baseline; every
`0032`. `0001_initial_schema.sql` is the consolidated multi-tenant baseline; every
later file is an incremental change. This upper bound is hand-maintained and
nothing checks it — it read `0026` for five migrations, including the two
brand-transition files carrying the fence GUC compatibility, in the document