Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions docs/adr/0018-high-availability-and-clustering.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
50 changes: 50 additions & 0 deletions internal/actors/workflow_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"context"
"encoding/json"
"fmt"
"os"
"time"

"github.com/open-source-cloud/fuse/internal/actors/actornames"
Expand Down Expand Up @@ -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 {
Expand All @@ -60,6 +62,7 @@ func NewWorkflowHandlerFactory(
fuseMetrics: fuseMetrics,
tracingProvider: tracingProvider,
secretStore: secretStore,
claimRepo: claimRepo,
}
},
}
Expand All @@ -81,6 +84,7 @@ type (
fuseMetrics *metrics.FuseMetrics
tracingProvider *tracing.Provider
secretStore secrets.SecretStore
claimRepo repositories.ClaimRepository

workflow *internalworkflow.Workflow
executionTimer *ExecutionTimer
Expand Down Expand Up @@ -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()))
Expand All @@ -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
}

Expand All @@ -168,13 +180,21 @@ 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()

action := a.workflow.Trigger()
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
}

Expand Down Expand Up @@ -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() {
Expand Down
4 changes: 3 additions & 1 deletion internal/actors/workflow_sup.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions internal/repositories/claim.go
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
5 changes: 5 additions & 0 deletions internal/repositories/claim_memory.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
9 changes: 9 additions & 0 deletions internal/repositories/claim_memory_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
22 changes: 21 additions & 1 deletion internal/repositories/postgres/claim.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()
Expand Down
21 changes: 21 additions & 0 deletions tests/functional/claim_repository_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
Loading