Skip to content

Fix task executor worker identity reconciliation races - #874

Open
Andyz26 wants to merge 3 commits into
masterfrom
andyz/fix-task-executor-worker-identity-v2
Open

Andyz26 wants to merge 3 commits into
masterfrom
andyz/fix-task-executor-worker-identity-v2

Conversation

@Andyz26

@Andyz26 Andyz26 commented Sep 1, 2026 •

Copy link
Copy Markdown
Collaborator

Problem

Task ownership spans submission, preparation, execution, cancellation, and cleanup. Previously, the agent did not expose ownership until preparation completed, allowing the control plane and agent to disagree about which worker occupied an executor.

Ambiguous assignment failures could then discard ownership, cancel the wrong worker, or release an executor prematurely. Two additional lifecycle races could leave executors permanently unschedulable:

  • Disconnecting while Assigned(A) preserved that state across reconnect, but an Available heartbeat could not transition it back to Pending.
  • Stale disconnect, timeout, and status messages could revive an archived executor as active but unregistered. A late assignment failure could then mutate it, throw, and be dropped by the actor.

Summary

  • Reserve accepted worker identity before acknowledging submission.
  • Keep ambiguous executors fenced until ownership is reconciled.
  • Reconcile Assigned disconnects before returning an executor to scheduling.
  • Prevent stale lifecycle messages from reviving archived executors.
  • Fence only executors that can reserve their accepted task, so the control plane is safe to deploy ahead of the agents.
  • Make cancellation idempotent, prompt to acknowledge, and bounded during shutdown.
  • Dispose prepared tasks and classloader leases when cancellation races with preparation.

State Transition Analysis

Agent task ownership

flowchart LR
    I[Idle] -->|submit A: reserve before Ack| P[Preparing A]
    P -->|prepared| R[Running A]
    P -->|cancel accepted| C[Cancelling A]
    R -->|cancel accepted| C
    C -->|cleanup settles| I
    P -->|submit B| X[Reject with worker A]
    R -->|submit B| X
Loading

The reserved slot is the source of worker identity for submission checks, reports, cancellation, and cleanup. Cancellation acknowledgement confirms acceptance rather than completion of teardown.

Assignment failure direction

A submitGateClaimed CAS races gateway completion against the assignment timeout to decide NotSent versus MayHaveRun, so a gateway that arrives after the timeout can never start a submit. Conflict is not produced at those sites; it is derived downstream from TaskAlreadyRunningException.

flowchart TB
    F[Assignment attempt] --> G{Did submit begin?}
    G -->|No| N[NotSent]
    N -->|retry remains| R[Retry through fresh gateway]
    N -->|expired| B[Unassign and disconnect]
    G -->|Yes or unknown| M[MayHaveRun]
    M --> Q[Fence executor and cancel expected worker]
    F -->|throwable is AlreadyRunning B| C[Quarantine and cancel B]
    F -->|stale assignment epoch| S[Ignore failure]
Loading

Issue 1: reconnect while assigned

flowchart LR
    subgraph Before
        A1[Assigned A] --> B1[Disconnect and archive]
        B1 --> C1[Reconnect as Assigned A]
        C1 -->|Available heartbeat| D1[Assigned A forever]
    end

    subgraph After
        A2[Assigned A] --> B2[Disconnect]
        B2 --> C2[Cancelling A and mark A Lost]
        C2 --> D2[Archive and reconnect]
        D2 -->|Occupied A| E2[Quarantine and cancel A]
        D2 -->|Available heartbeat| F2[Release, see fencing rules]
    end
Loading

Assigned.onTaskExecutorStatusChange(Available) returns the same instance, so without the disconnect-time transition into Cancelling no Available report could ever release the executor.

Issue 3: stale lifecycle messages

flowchart LR
    subgraph Before
        A1[Archived executor] --> B1[Stale disconnect or status]
        B1 --> C1[Revived active but unregistered]
        C1 --> D1[Assignment gate passes]
        D1 --> E1[Mutation throws and message is dropped]
    end

    subgraph After
        A2[Archived executor] -->|registration or heartbeat| B2[Revive and reconcile]
        A2 -->|duplicate disconnect or timeout| C2[Ignore]
        A2 -->|late status| D2[Return NotFound]
        A2 -->|late assignment failure| E2[Ignore]
    end
Loading

Archived state is now revived only by messages that establish renewed executor activity. isCurrentAssignment additionally requires a registered, current, unreconciled assignment, and isRunningOrAssigned requires registration, which closes the same unhandled throw in onMarkExecutorTaskCancelledRequest and onTerminateWorkerRequest.

Worker identity mismatch

Occupied(B) on an executor expected to hold A is newly detected. Previously Running(A) ignored the report and Assigned(A) recorded Running(A), silently discarding B. The reported worker is now quarantined and cancelled, and the expected worker is marked lost.

Fencing and rollout

Fencing depends on the agent being able to expose ownership, which it advertises with the mantis_task_executor_accepted_task_reservation registration attribute. An executor that cannot reserve its accepted task cannot produce the signal the fence waits for, so fencing it would strand it. Those executors keep the previous behaviour instead.

Situation Reserves accepted task Does not
Running(A) + Available heartbeat mark A lost, Verifying; next heartbeat releases (2 heartbeats) mark A lost, release (1 heartbeat, unchanged)
Reconciling + Available heartbeat Verifying, then release (2 heartbeats) release (1 heartbeat, unchanged)
Reconciling + Available status change never releases never releases
Occupied(B) mismatch quarantine B, mark expected lost same

Two heartbeats are sufficient because ResourceManagerGatewayCxn.runIteration blocks on the in-flight heartbeat, so reports are strictly serialised and two consecutive Available reports guarantee one was sampled after the control plane acted. At the default 10s interval the fence costs roughly 20s.

Status changes never release a fenced executor in either column: setStatus is fire and forget with no sequencing, so a delayed Available cannot prove it follows the cancellation.

This makes the deployment order safe. The control plane can ship first with no behavioural change for agents that have not been upgraded; protection engages per executor as agents roll out and begin declaring the attribute. Agents force re-registration after upgrade by switching to a new durable registration-state file, so the control plane learns the attribute. The attribute key deliberately does not match MANTIS_SCHEDULING_ATTRIBUTE_*, so getSchedulingAttributes and TaskExecutorGroupKey are unchanged and the scheduling pool does not split during the rollout.

Not addressed here

  • MayHaveRun sends cancelTask(A) while a submitTask(A) may still be in flight; if the cancel lands first the agent can still accept A as an orphan. The CAS fences the control plane's call site, not an RPC already on the wire.
  • A stale Available status change still releases a Running executor that is not fenced. This is pre-existing and unchanged.
  • Submit failures are no longer retried. Any failure after submission begins is terminal, costing a worker launch failure plus the fence interval where the previous code retried immediately.
  • A reservation-capable executor left in Assigned with no reconciliation cannot self-release if its assignment-failure event is lost entirely. This matches the previous behaviour; a periodic sweep over long-lived Assigned states would close it.

Tests

./gradlew :mantis-control-plane:mantis-control-plane-server:test \
  :mantis-server:mantis-server-agent:test

Coverage includes preparation ownership, cancellation cleanup, stop failures, bounded shutdown, conflicting assignments, delayed failures, reconnect reconciliation, duplicate disconnects, late status messages, and the release path for executors in both columns of the fencing table.

@github-actions

github-actions Bot commented Sep 1, 2026 •

Copy link
Copy Markdown

Test Results

853 tests  +30   842 ✅ +30   10m 16s ⏱️ -2s
167 suites ± 0    11 💤 ± 0 
167 files   ± 0     0 ❌ ± 0 

Results for commit febe75c. ± Comparison against base commit 9ff3102.

♻️ This comment has been updated with latest results.

registration = null;
// Store the current WorkerId as previousWorkerId for potential reconnection notification
previousWorkerId = getWorkerId();
setAvailabilityState(null);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P0 — must fix before merge

  1. TaskExecutorState.java:167 — reconnecting executor becomes permanently unschedulable
    onDisconnection now only clears availabilityState when workerId == null. Since trackIfAbsent (ExecutorStateManagerImpl.java:186) revives the archived state object, an agent that drops off while Assigned(w) comes back with availabilityState == Assigned(w). isAvailable() requires Pending; Assigned.onTaskExecutorStatusChange(Available) returns this; the Running+Available branch at :328 doesn't cover Assigned. Nothing ever moves it back.
    Telling detail: getPreviousWorkerId()/clearPreviousWorkerId() (:464, :468) have zero callers tree-wide — the reconcile-on-reconnect half of this design isn't wired up.
    Fix: on re-registration, reconcile previousWorkerId against the first report and reset availability to Pending when the executor comes back idle.

  2. TaskExecutorState.java:280 (with :333) — rolling upgrade drains cluster capacity
    reservesAcceptedTask() defaults to false, i.e. every agent not yet redeployed with this PR. For those, an Available heartbeat while reconciliationState != None hits return false at :287 — no clearReconciliation(), no pending(); and :333 puts a legacy agent straight into Quarantined on its first Running+Available heartbeat. One MayHaveRun assignment failure then removes it from the pool forever. Master self-healed twice over (fell through to Running.onTaskExecutorStatusChange(Available) → pending(), and cleared a stale cancelledWorkerOnTask on mismatch).
    Fix: give legacy agents a recovery path — either clear reconciliation after N consecutive Available heartbeats, or treat !reservesAcceptedTask() as "trust the report" and go straight to Pending.

  3. ExecutorStateManagerActor.java:668 — ESM actor restarts on assignment failure
    Master wrapped this block in try { … } catch (IllegalStateException e) { log.error(…) } (master :660-671); the PR deletes it. Meanwhile setCancelledWorkerOnTask/quarantineWorkerOnTask (TaskExecutorState.java:129,137) now throw when unregistered, and the new isCurrentAssignment gate checks isAssigned() && epoch && workerId — not isRegistered(). A TE that disconnects while Assigned(w) and is later revived unregistered passes the gate, and :689/:691 throws out of receive. Akka drops the message and cancels every AbstractActorWithTimers heartbeat timer.
    Fix: add isRegistered() to isCurrentAssignment, and restore the catch.

P1 — should fix

  1. ExecutorStateManagerActor.java:964 and :982 — MarkExecutorTaskCancelled/TerminateWorker throw instead of replying
    Both select via isRunningOrAssigned(workerId) (TaskExecutorState.java:417), which compares only getWorkerId() and ignores registration, so they match the same unregistered phantom.setCancelledWorkerOnTask throws before sender().tell(Ack), with no try/catch — the job actor or leader-exclusive REST route gets no reply and blocks to its ask timeout, plus the actor restarts. This was a harmless field assignment before.

  2. ExecutorStateManagerActor.java:993 — cancelTaskOnExecutor swallows every failure
    The whenComplete only logs: no retry, no state correction, no disconnect. But the caller has already forced Cancelling/Quarantined + tryMarkUnavailable. So on the most likely failure — TaskNotFoundException because the agent never got the task, exactly the MayHaveRun-after-timeout case — the executor stays unavailable, and per Temporarily disable subprojects. #2 a legacy agent never recovers. The wholereconciliation design depends on this call succeeding, yet its failure is a no-op.

  3. ExecutorStateManagerActor.java:686 — foreign worker killed silently
    On TaskAlreadyRunningException, cancellationTarget becomes getCurrentlyRunningWorkerTask() — a worker belonging to a different job/attempt — which gets quarantined and cancelTask-ed. Theonly jobMessageRouter event emitted is WorkerLaunchFailed for expectedWorker (:674); nothing is routed for the worker actually being killed. Its job actor believes it's alive until its own heartbeat timeout expires, delaying replacement by the full window.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Also, what triggers us to make this change?

@hellolittlej hellolittlej Sep 4, 2026 •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I gave up reading the whole summary and reviewing whole things by myself lol coz it's too agentic, therefore I delegate the review process to my agent. :)

@hellolittlej hellolittlej left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do we have any metrics that can validate this change after the merge?

The reconciliation state machine added earlier in this branch left three
paths on which an executor could be stranded permanently.

Disconnecting while Assigned preserved availabilityState, and reviving the
archived state on reconnect brought it back. No Available report can leave
that state, because Assigned.onTaskExecutorStatusChange(Available) returns
the same instance, so the executor stayed registered and heartbeating while
never rejoining the schedulable pool. Verified end to end: five Available
heartbeats after a reconnect left it registered with assignedTask=true and
absent from getAvailableTaskExecutors, where master recovered on the first
heartbeat. Disconnecting from Assigned now enters Cancelling and marks the
worker lost, so reconciliation can release it.

isCurrentAssignment only consulted the availability state, so with Assigned
surviving a disconnect it accepted a late assignment failure for an
unregistered executor, and the setCancelledWorkerOnTask that followed threw.
MantisActorSupervisorStrategy resumes rather than restarts, so the message
was dropped after WorkerLaunchFailed had already been routed to the job and
before any cancellation was sent. It now also requires a registered executor
with no reconciliation in progress, and the handler no longer lets
IllegalStateException escape. isRunningOrAssigned gained the same
registration guard, which closes the identical unhandled throw in
onMarkExecutorTaskCancelledRequest and onTerminateWorkerRequest.

The two-heartbeat release was gated on reservesAcceptedTask, so an executor
without that attribute could never leave reconciliation at all: the branch
returned false unconditionally, and the only escape was an Available status
change, which an idle agent never sends because setStatus(available) is
reached only from clearTaskSlot. Fencing now applies only to executors that
reserve their accepted task; the rest keep the previous behaviour, in which
Available is authoritative and releases in a single heartbeat. An executor
that cannot expose ownership cannot produce the signal the fence waits for,
so holding it is strictly worse than leaving it as it was. This also keeps
the control plane safe to deploy ahead of the agents: the new protection
engages per executor as agents are upgraded and begin declaring the
attribute.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants