Skip to content

kfpytorch: fail recoverably when the elastic agent is signalled instead of raising IgnoreOutputs - #60

Open
devin-ai-integration[bot] wants to merge 1 commit into
masterfrom
devin/1789071138-kfpytorch-sigterm-recoverable
Open

devin-ai-integration[bot] wants to merge 1 commit into
masterfrom
devin/1789071138-kfpytorch-sigterm-recoverable

Conversation

@devin-ai-integration

Copy link
Copy Markdown

Tracking issue

Related to upstream flyteorg#2064 (the change that introduced SignalException -> IgnoreOutputs). Exa incident: GS1 exec azvg5f9v9974m78fvlf7 (BIGMODEL, 2026-09-10).

Why are the changes needed?

PytorchElasticFunctionTask._execute catches torch elastic's SignalException (raised by the agent's SIGTERM/SIGINT handler) and turns it into IgnoreOutputs. flytekit/bin/entrypoint.py treats IgnoreOutputs as a clean exit: no error.pb, no outputs.pb, process exits 0, so propeller marks the attempt succeeded.

On a spot / low-priority preemption the kubelet SIGTERMs the pod while training is still running. Result: the execution ends SUCCEEDED with no outputs and no checkpoint, retries and auto_resume never fire, and the job dies silently green. Whether a preempted job retries or "succeeds" is a race on which pod gets the signal first.

Evidence (read-only, Loki + Overseer): azvg5f9v9974m78fvlf7 — 19:45:22Z Received 15 death signal, shutting down workers → 19:45:29Z Plugin [container] returned no outputReader ... PhaseSuccess → execution succeeded, s3://exa-models/checkpoints/gs1-distill/bigmodel-c42-nem1b-a2-lr5em5-v1/ empty.

The IgnoreOutputs path came from upstream (flyteorg#2064, to keep the operator's SIGTERM of sibling pods after another worker's failure from masking that failure). This is not what it does in practice: with EarliestErrorAggregationStrategy (hardcoded for the pytorch plugin in flyteplugins) the sibling's own, earlier error already wins on timestamp; swallowing the signal only hides genuine preemption.

What changes were proposed in this pull request?

 except SignalException as e:
     logger.exception(f"Elastic launch agent process terminating: {e}")
-    raise IgnoreOutputs()
+    raise FlyteRecoverableException(
+        f"Elastic launch agent process was terminated by signal {e.sigval.name} before the worker "
+        f"group finished: {e}",
+        timestamp=time.time(),
+    ) from e
  • FlyteRecoverableException (USER:Recoverable) is what the plugin already raises for recoverable worker errors, so propeller counts the attempt against retries and relaunches; auto_resume then finds the last checkpoint.
  • timestamp is carried into the ContainerError written by the entrypoint, so earliest-error aggregation still attributes a task failure to a sibling worker whose error predates the teardown signal.
  • Docstring Raises: updated. IgnoreOutputs is still raised for worker groups with index > 0 on success (unchanged).

How was this patch tested?

Two tests added to plugins/flytekit-kf-pytorch/tests/test_elastic_task.py:

  • test_agent_signal_is_recoverable_failure — deterministic: elastic_launch patched to raise SignalException(SIGTERM); asserts FlyteRecoverableException (not IgnoreOutputs), __cause__ is the SignalException, message names SIGTERM, timestamp set.
  • test_sigterm_to_agent_process_is_recoverable_failure[spawn|fork] — end to end: real Elastic(nnodes=1, nproc_per_node=2), rank 0 os.kill(os.getppid(), SIGTERM)s the agent after a gloo barrier (the barrier guarantees every worker is forked so the signal lands in the agent's monitor loop, not inside os.fork() where CPython's fork warning machinery drops the handler's exception). Both start methods raise FlyteRecoverableException. On the base branch the same test observes IgnoreOutputs.
cd plugins/flytekit-kf-pytorch && python -m pytest tests -q
# 39 passed; test_end_to_end[spawn|fork] fail identically on the base branch in this env
# (torch 2.14 weights_only unpickling of the nn.Module output), unrelated to this change.
ruff check plugins/flytekit-kf-pytorch && ruff format --check plugins/flytekit-kf-pytorch  # clean

Setup process

uv venv && uv pip install -e . -e plugins/flytekit-kf-pytorch[elastic] pytest

Screenshots

n/a

Check all the applicable boxes

  • I updated the documentation accordingly (docstring).
  • All new and existing tests passed (see note on test_end_to_end above).
  • All commits are signed-off.

Related PRs

Upstream origin of the swallowed signal: flyteorg#2064. Monorepo pin bump to follow once merged (flytekitplugins-kfpytorch[elastic] @ .../exa-labs/flytekit/<sha>.tar.gz in python/shared/exa_flyte/pyproject.toml and friends).

Docs link

n/a

Link to Devin session: https://app.devin.ai/sessions/84936c9760074a4793d713c917952f29
Open in Devin Desktop: https://app.devin.ai/desktop/session/84936c9760074a4793d713c917952f29?variant=devin
Requested by: @jld-adriano

…ad of raising IgnoreOutputs

A SIGTERM to the elastic agent (spot/preemption, eviction, operator
teardown) surfaced as IgnoreOutputs, which Flyte records as a succeeded
attempt without outputs: retries and auto-resume never fire and a
preempted job dies green. Re-raise it as FlyteRecoverableException
(cause + timestamp preserved) so the attempt fails and is retried.

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
@devin-ai-integration

Copy link
Copy Markdown
Author

🤖 Devin AI Engineer

I'll be helping with this pull request! Here's what you should know:

✅ I will automatically:

  • Address comments on this PR. Add '(aside)' to your comment to have me ignore it.
  • Look at CI failures and help fix them

Note: I can only respond to comments from users who have write access to this repository.

⚙️ Control Options:

  • Disable automatic comment, CI, and merge conflict monitoring

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.

1 participant