From 827a439710bc26293dc22f3f00afa0f6a17387c0 Mon Sep 17 00:00:00 2001 From: Gustavo Bertoi Date: Thu, 4 Jun 2026 02:02:59 -0300 Subject: [PATCH] fix(ha): claim workflows on spawn + persist running to stop the untriggered steal (#89) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit On the multi-node HA stack a freshly triggered workflow intermittently stayed in "untriggered" forever: the trigger path never set claimed_by and never persisted the untriggered->running transition (only journaled it), so the row sat at claimed_by=NULL, state=untriggered for the whole run. The NOTIFY-driven claim sweep on other nodes then stole and re-triggered it, spawning a duplicate handler and stalling the run; recovery (running/sleeping only) never healed it. Fix: - ClaimRepository.ClaimWorkflow(nodeID, workflowID) — atomically claim a single workflow for the owning node (memory no-op returns true; postgres UPDATE ... RETURNING rows-affected). - WorkflowHandler.Init claims the workflow for this node on both the new-trigger and replay paths; on loss to another live node it bows out (tells the instance supervisor it's done here) so no duplicate runs. It then persists the running/sleeping state immediately so the DB row reflects reality. - Claim sweep lease-expired branch and recoverWorkflows now include 'untriggered' so a crash-stranded row is reclaimed/re-driven. - ADR-0018 amended with the ownership invariants. Non-HA/memory: no behavior change (claim is a no-op, claim actor is HA-gated). The e2e re-trigger quarantine stays as defense-in-depth until 3-node runs confirm. Co-Authored-By: Claude Opus 4.8 (1M context) --- .../0018-high-availability-and-clustering.md | 12 +++++ internal/actors/workflow_handler.go | 50 +++++++++++++++++++ internal/actors/workflow_sup.go | 4 +- internal/repositories/claim.go | 6 +++ internal/repositories/claim_memory.go | 5 ++ internal/repositories/claim_memory_test.go | 9 ++++ internal/repositories/postgres/claim.go | 22 +++++++- tests/functional/claim_repository_test.go | 21 ++++++++ 8 files changed, 127 insertions(+), 2 deletions(-) diff --git a/docs/adr/0018-high-availability-and-clustering.md b/docs/adr/0018-high-availability-and-clustering.md index d8ed72b..ce665f1 100644 --- a/docs/adr/0018-high-availability-and-clustering.md +++ b/docs/adr/0018-high-availability-and-clustering.md @@ -38,6 +38,18 @@ infrastructure.** Three mechanisms: (`node_heartbeats`), and reassigns workflows from nodes whose heartbeat is older than `HA_LEASE_TIMEOUT` (default 30s). The memory backend is a no-op (HA off). + **Ownership invariants (fix #89).** Two invariants prevent a node's own sweep — or another + node's NOTIFY-driven immediate sweep — from stealing a workflow that is already being run: + (a) when a node spawns a workflow it **claims it for itself first** (`ClaimRepository.ClaimWorkflow` + in `WorkflowHandler.Init`, both the new-trigger and replay paths); if another live node already + owns it the handler bows out without running, so no duplicate executes; and (b) the + `untriggered → running/sleeping` state transition is **persisted immediately at `Init`** (not + only journaled), so the DB row never sits at a stale `untriggered` for a running workflow — which + keeps `GET /v1/workflows/{id}` accurate, lets `recoverWorkflows` (which now also includes + `untriggered`) re-drive a crash-stranded row, and lets the lease-expiry reclaim cover `untriggered` + too. Before this, a freshly triggered workflow stayed `claimed_by = NULL, state = untriggered` + for its whole run and was repeatedly stolen by other nodes' sweeps, intermittently stranding it. + 2. **Lower-latency claims.** A `PgListenerActor` subscribes to Postgres **LISTEN/NOTIFY** (`workflow_state_change`) on a dedicated connection; a newly available workflow nudges the claim actor (`ClaimSweepNowMsg`) to claim immediately instead of waiting for the next sweep. diff --git a/internal/actors/workflow_handler.go b/internal/actors/workflow_handler.go index c49ba06..17b93e9 100644 --- a/internal/actors/workflow_handler.go +++ b/internal/actors/workflow_handler.go @@ -4,6 +4,7 @@ import ( "context" "encoding/json" "fmt" + "os" "time" "github.com/open-source-cloud/fuse/internal/actors/actornames" @@ -45,6 +46,7 @@ func NewWorkflowHandlerFactory( fuseMetrics *metrics.FuseMetrics, tracingProvider *tracing.Provider, secretStore secrets.SecretStore, + claimRepo repositories.ClaimRepository, ) *WorkflowHandlerFactory { return &WorkflowHandlerFactory{ Factory: func() gen.ProcessBehavior { @@ -60,6 +62,7 @@ func NewWorkflowHandlerFactory( fuseMetrics: fuseMetrics, tracingProvider: tracingProvider, secretStore: secretStore, + claimRepo: claimRepo, } }, } @@ -81,6 +84,7 @@ type ( fuseMetrics *metrics.FuseMetrics tracingProvider *tracing.Provider secretStore secrets.SecretStore + claimRepo repositories.ClaimRepository workflow *internalworkflow.Workflow executionTimer *ExecutionTimer @@ -124,6 +128,11 @@ func (a *WorkflowHandler) Init(args ...any) error { if a.workflowRepository.Exists(initArgs.workflowID.String()) { a.workflow, _ = a.workflowRepository.Get(initArgs.workflowID.String()) + // Claim this workflow for this node before running it; if another live node owns it, bow + // out so we don't run a duplicate (ADR-0018). + if !a.claimForThisNode(initArgs.workflowID) { + return nil + } // Replay: the environment comes from the reconstructed workflow, not init args, so // resolution stays deterministic across restart/recovery (ADR-0031). a.workflow.SetSecretResolver(a.newSecretResolver(a.workflow.Environment())) @@ -149,6 +158,9 @@ func (a *WorkflowHandler) Init(args ...any) error { if action != nil { a.handleWorkflowAction(action) } + // Persist the (possibly advanced) state so the DB row reflects running/sleeping, not a + // stale untriggered (ADR-0018). + a.persistWorkflowState() return nil } @@ -168,6 +180,11 @@ func (a *WorkflowHandler) Init(args ...any) error { a.Log().Error("failed to save workflow for id %s: %s", initArgs.workflowID, err) return nil } + // Claim this freshly-created workflow for this node before running it, so the claim sweep on + // other nodes cannot steal it mid-run (ADR-0018, fixes #89). + if !a.claimForThisNode(initArgs.workflowID) { + return nil + } a.Log().Debug("created new workflow with id %s", initArgs.workflowID) a.startRootSpan() @@ -175,6 +192,9 @@ func (a *WorkflowHandler) Init(args ...any) error { a.persistJournal() a.startWorkflowTimeout() a.handleWorkflowAction(action) + // Persist the running/sleeping state immediately so the DB row no longer reads untriggered + // (status accuracy + recovery + lease reclaim) (ADR-0018). + a.persistWorkflowState() return nil } @@ -337,6 +357,36 @@ func (a *WorkflowHandler) handleMsgAsyncFunctionResult(msg messaging.Message) er return nil } +// claimForThisNode claims the workflow for this node before running it (ADR-0018, fixes #89). It +// returns true when this node owns the workflow and may run it. On loss (another live node already +// owns it) it stops this duplicate instance tree by telling the instance supervisor the workflow is +// done here — without touching the workflow's persisted state — and returns false. In non-HA / +// memory mode the claim is a no-op that always succeeds. On a claim-store error it fails open (runs) +// rather than dropping the workflow. +func (a *WorkflowHandler) claimForThisNode(wfID workflow.ID) bool { + owned, err := a.claimRepo.ClaimWorkflow(resolveNodeID(a.config), wfID.String()) + if err != nil { + a.Log().Error("failed to claim workflow %s (running anyway): %s", wfID, err) + return true + } + if !owned { + a.Log().Info("workflow %s is owned by another node; not running it here", wfID) + if sendErr := a.Send(a.Parent(), messaging.NewWorkflowCompletedMessage(wfID, "claimed-elsewhere")); sendErr != nil { + a.Log().Error("failed to stop duplicate instance for workflow %s: %s", wfID, sendErr) + } + } + return owned +} + +// resolveNodeID returns this node's HA identifier (HA_NODE_ID, else hostname). +func resolveNodeID(cfg *config.Config) string { + if cfg.HA.NodeID != "" { + return cfg.HA.NodeID + } + host, _ := os.Hostname() + return host +} + // persistWorkflowState persists both the journal and the workflow state to the repository. // Call this after any state transition that should be visible to external queries (e.g. running → sleeping). func (a *WorkflowHandler) persistWorkflowState() { diff --git a/internal/actors/workflow_sup.go b/internal/actors/workflow_sup.go index 811f592..2fe829c 100644 --- a/internal/actors/workflow_sup.go +++ b/internal/actors/workflow_sup.go @@ -210,7 +210,9 @@ func (a *WorkflowSupervisor) HandleEvent(event gen.MessageEvent) error { } func (a *WorkflowSupervisor) recoverWorkflows() { - ids, err := a.workflowRepository.FindByState(internalworkflow.StateRunning, internalworkflow.StateSleeping) + // Include untriggered: a workflow can be stranded in untriggered if its node crashed between + // the create-Save and the running-state persist, so recovery must re-drive those too (ADR-0018). + ids, err := a.workflowRepository.FindByState(internalworkflow.StateUntriggered, internalworkflow.StateRunning, internalworkflow.StateSleeping) if err != nil { a.Log().Error("failed to query workflows for recovery: %s", err) return diff --git a/internal/repositories/claim.go b/internal/repositories/claim.go index 1c9fe9a..ce5a391 100644 --- a/internal/repositories/claim.go +++ b/internal/repositories/claim.go @@ -15,6 +15,12 @@ type ClaimRepository interface { // Returns up to limit claimed workflows. ClaimWorkflows(nodeID string, limit int) ([]ClaimedWorkflow, error) + // ClaimWorkflow atomically claims a single workflow for nodeID. It returns true when this node + // now owns the workflow (it was unclaimed, already this node's, or the previous owner's lease + // expired) and false when another live node owns it. Used by a node to claim a workflow it is + // about to run, so the sweep on other nodes cannot steal it (ADR-0018). + ClaimWorkflow(nodeID, workflowID string) (bool, error) + // ReleaseWorkflows releases all workflows claimed by the given node. ReleaseWorkflows(nodeID string) error diff --git a/internal/repositories/claim_memory.go b/internal/repositories/claim_memory.go index 0f52b07..b87ce24 100644 --- a/internal/repositories/claim_memory.go +++ b/internal/repositories/claim_memory.go @@ -15,6 +15,11 @@ func (r *MemoryClaimRepository) ClaimWorkflows(_ string, _ int) ([]ClaimedWorkfl return nil, nil } +// ClaimWorkflow always succeeds in memory mode (single process; no other node can contend). +func (r *MemoryClaimRepository) ClaimWorkflow(_, _ string) (bool, error) { + return true, nil +} + // ReleaseWorkflows is a no-op in memory mode. func (r *MemoryClaimRepository) ReleaseWorkflows(_ string) error { return nil diff --git a/internal/repositories/claim_memory_test.go b/internal/repositories/claim_memory_test.go index 8f33dd2..a75be85 100644 --- a/internal/repositories/claim_memory_test.go +++ b/internal/repositories/claim_memory_test.go @@ -38,6 +38,15 @@ func TestMemoryClaimRepository_ClaimWorkflows_ConcurrentCallsNoPanic(_ *testing. // Assert — if we reach here without panic/race, test passes } +func TestMemoryClaimRepository_ClaimWorkflow_AlwaysOwned(t *testing.T) { + repo := repositories.NewMemoryClaimRepository() + + owned, err := repo.ClaimWorkflow("node-1", "wf-1") + + require.NoError(t, err) + assert.True(t, owned, "memory claim is a single-process no-op and always succeeds") +} + func TestMemoryClaimRepository_ReleaseWorkflows_NoError(t *testing.T) { // Arrange repo := repositories.NewMemoryClaimRepository() diff --git a/internal/repositories/postgres/claim.go b/internal/repositories/postgres/claim.go index 0fdb1f5..7834467 100644 --- a/internal/repositories/postgres/claim.go +++ b/internal/repositories/postgres/claim.go @@ -32,7 +32,7 @@ func (r *ClaimRepository) ClaimWorkflows(nodeID string, limit int) ([]repositori WHERE (claimed_by IS NULL AND state IN ('untriggered', 'running', 'sleeping')) OR (claimed_by IS NOT NULL AND claimed_by != $1 AND claimed_at < NOW() - INTERVAL '1 second' * $3 - AND state IN ('running', 'sleeping')) + AND state IN ('untriggered', 'running', 'sleeping')) ORDER BY id FOR UPDATE SKIP LOCKED LIMIT $2 @@ -55,6 +55,26 @@ func (r *ClaimRepository) ClaimWorkflows(nodeID string, limit int) ([]repositori return claimed, rows.Err() } +// ClaimWorkflow atomically claims a single workflow for nodeID. It succeeds when the workflow is +// unclaimed, already claimed by this node, or the previous owner's lease has expired; it returns +// false when another live node holds the claim. The owning node calls this on spawn so the sweep +// on other nodes cannot steal a workflow it is actively running (ADR-0018). +func (r *ClaimRepository) ClaimWorkflow(nodeID, workflowID string) (bool, error) { + ctx := context.Background() + tag, err := r.pool.Exec(ctx, ` + UPDATE workflows + SET claimed_by = $1, claimed_at = NOW(), updated_at = NOW() + WHERE workflow_id = $2 + AND (claimed_by IS NULL + OR claimed_by = $1 + OR claimed_at < NOW() - INTERVAL '1 second' * $3) + `, nodeID, workflowID, r.leaseTimeout.Seconds()) + if err != nil { + return false, err + } + return tag.RowsAffected() == 1, nil +} + // ReleaseWorkflows releases all workflows claimed by the given node. func (r *ClaimRepository) ReleaseWorkflows(nodeID string) error { ctx := context.Background() diff --git a/tests/functional/claim_repository_test.go b/tests/functional/claim_repository_test.go index e13ebc4..8be053e 100644 --- a/tests/functional/claim_repository_test.go +++ b/tests/functional/claim_repository_test.go @@ -91,6 +91,27 @@ func contractTestClaimRepository( assert.Empty(t, claimed) }) + t.Run("ClaimWorkflow claims a specific workflow, is idempotent for the owner, and blocks others", func(t *testing.T) { + reset() + repo := newClaimRepo() + ids := seedWorkflows(t, 1, internalworkflow.StateUntriggered) + + // The owning node claims it. + owned, err := repo.ClaimWorkflow("node-A", ids[0]) + require.NoError(t, err) + assert.True(t, owned) + + // Re-claiming by the same node is idempotent. + owned, err = repo.ClaimWorkflow("node-A", ids[0]) + require.NoError(t, err) + assert.True(t, owned) + + // Another live node cannot steal it (its lease has not expired). + owned, err = repo.ClaimWorkflow("node-B", ids[0]) + require.NoError(t, err) + assert.False(t, owned) + }) + t.Run("ReleaseWorkflows releases all claims for a node", func(t *testing.T) { reset() repo := newClaimRepo()