Add Slurm connector plugin - #1550
davidmirror-ops wants to merge 38 commits into
Conversation
Run Flyte 2 tasks on an existing Slurm cluster, including Soperator-managed clusters, without changing the cluster. Two task types are served by one connector: - `slurm`: a Python task submitted as an sbatch job that runs the task's own container image and Flyte entrypoint under Pyxis/Enroot, so typed I/O, caching, retries and error reporting work as they do for a Kubernetes pod. Removing `plugin_config` runs the identical task on Kubernetes. - `slurm_script`: an existing sbatch script submitted unchanged; scalar inputs are exposed as FLYTE_INPUT_<NAME>. Phase, exit code and logs only. Transport is SSH to a login node (asyncssh) behind a small SlurmTransport protocol so a slurmrestd transport can be added later. One connection per cluster is reused and status queries are batched (squeue, then sacct for finished jobs). Slurm states map to Flyte phases with PENDING reported as QUEUED rather than RUNNING, PREEMPTED as retryable, and CANCELLED as ABORTED; failures carry Slurm's reason and the tail of stderr. Connection details may be set per task or cluster-wide on the connector via FLYTE_SLURM_HOST / FLYTE_SLURM_USERNAME / FLYTE_SLURM_SSH_PRIVATE_KEY. The plugin is added to the default connector image build. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01X4XZW8jVev8ueGeDPfEiWx
Two defects found bringing the connector up against a live Soperator cluster on Nebius. Enroot resets PATH from its own environ.d rather than applying the image's, so the bare `a0` entrypoint the task template emits was never found and every native `slurm` task died with `a0: command not found` — while the identical task ran fine as a Kubernetes pod, where the kubelet does apply image ENV. VIRTUAL_ENV *is* propagated, so the generated srun line now restores the venv's bin directory inside the container before exec'ing the entrypoint. `exec "$@"` hands argv over untouched, so no task argument is re-quoted, and the shim is inert for images that do not use a virtualenv. `get()` also returned the job's stdout and stderr as TaskLogs built from POSIX paths on the login node. A TaskLog uri renders as a hyperlink, so those became dead links in the UI. The paths are named in the phase message instead, ahead of the stderr tail so that block stays last; streaming stdout remains `get_logs`' job. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
A native task exited 2 with `Missing option '--run-base-dir'`. That option is required by `a0` and its only other source is the `_U_RUN_BASE` environment variable, which the backend injects directly into the pod spec on Kubernetes (leaseworker/lifecycle/context.go) rather than putting it in the TaskTemplate. A connector therefore only ever sees it through `TaskExecutionMetadata.environment_variables` -- which this plugin was discarding, along with ACTION_NAME, RUN_NAME, _U_ORG_NAME and the endpoint-discovery pair _U_EP_OVERRIDE / _U_INSECURE. Merge them into the job environment, ordered so task config still wins over platform vars, which win over whatever the template carries. Only the deployed path was affected: for local runs the executor's `_render_task_template` passes `--run-base-dir` on the command line instead, and passes no TaskExecutionMetadata at all, so the merge is a no-op there. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The entrypoint aborted with `Project is required`. `_run_action` asserts on org, project, domain, run_name and action name. The connector metadata's environment supplies org, run and action, but project and domain are added on the Kubernetes path by `flytek8s.GetExecutionEnvVars`, which builds a pod spec and so never runs for a connector. Derive them from `TaskExecutionMetadata.task_execution_id` instead, along with the execution name and retry attempt, matching what that function writes into a task pod. This completes the runtime's required set: with `_U_RUN_BASE` from the previous commit, nothing else it asserts on is missing. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Three examples under examples/plugins/slurm, each covering something the others do not. slurm_script_example.py submits an existing sbatch script unchanged and drives srun itself, which is currently the only way to run gang-scheduled multi-node work through this plugin. slurm_example.py is a Python task with the same typed I/O it would have as a pod; deleting plugin_config runs it on Kubernetes, which is also the quickest way to tell a Slurm problem from a task problem. slurm_pipeline_example.py crosses backends: prepare on Kubernetes hands a Dir to train on Slurm, which hands a File to evaluate back on Kubernetes. It exists mainly to show that the handoff is object storage in both directions, and that returning a cluster filesystem path instead of a File would typecheck and then fail in the pod. The README covers the three unrelated credentials involved -- SSH for the connector, Enroot for the compute nodes, object storage for the job -- and how to read a failure from the generated sbatch script. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The example declared `depends_on=[k8s_env]` on the Slurm environment, which is backwards. `depends_on` lists environments to deploy alongside this one, so it points from the environment holding the invoked task to the environments its tasks call into. `pipeline` lives in the pod environment and calls `train` in the Slurm one, so the pod environment is what must declare the dependency. As written only one image was built and the run failed with `Environment 'slurm-pipeline-train' not found in image cache`. Reorders the definitions and explains the rule, since the failure surfaces at run time and points at an image cache rather than at the dependency. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
`examples/` is linted (only `examples/reports` is excluded, and the per-file-ignore there covers E402 only), so the unsorted import blocks in these three files would have failed CI with I001. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The examples asked for `gres="gpu:1"`, which a cluster without GRES
configured rejects outright at submission:
sbatch: error: Invalid generic resource (gres) specification.
GRES is opt-in per cluster (`GresTypes` and per-node `Gres`), so an example
that hard-codes it does not run on a stock install. Use `cpus_per_task` and
`mem` instead, and say how to check before adding a GPU request.
Also swaps the `sbatch_options` demonstration from `exclusive` to `requeue`.
Both are valid, but `exclusive` reserves the entire node, which is not
something an example should teach people to copy onto a shared cluster.
Directives verified with `sbatch --test-only` against a real cluster.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
`evaluate` called `.decode()` straight on the result of `fh.read()`, which
fails against object storage:
AttributeError: 'builtins.Bytes' object has no attribute 'decode'
The reader returns a Rust-backed `Bytes` that supports the buffer protocol
but has no `decode`. `bytes(await fh.read()).decode("utf-8")` is the idiom
in `File.open`'s own docstring and works for both backends.
Worth knowing: a local raw-data path returns real Python `bytes`, so the
broken form passes locally and only fails once the raw-data path is remote.
Verified both ways against GCS -- the unwrapped form reproduces the error,
the wrapped one returns the value.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
`uv lock --check` failed for this plugin in CI. The plugin resolves `flyte` as a path dependency, so the root's exact pins flow into its lockfile, and those moved on main: flyteidl2 to 2.0.48 and pydantic-monty to 0.0.22. The lock still recorded the versions from when the branch was cut. Part of this is the change I mistakenly reverted while committing the log_links fix, having read it as unrelated churn from an editable install. It was neither unrelated nor optional. Verified with the same loop CI runs: `uv lock --check` passes for every plugin, and the plugin's tests still pass. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
|
Do we want to use SSH or the API-based thing? Not sure if the v1 supports both: https://github.com/flyteorg/flytekit/tree/d69b3fbd5ffe2a94d47e4fdba98200ce8b13e49c/plugins/flytekit-slurm/flytekitplugins/slurm |
`make check-docstrings` reported 43 problems in this plugin: 39 RST double-backtick inline literals and 4 bare command-line flags. CONTRIBUTING asks for Markdown with Google-style sections, because the double-backtick form renders correctly only by accident and reads as markup in an IDE tooltip and in `help()`. Converts every ``value`` to `value` across the plugin. The bare-flag reports went with them: those flags sat inside double-backtick spans the checker could not see into, and are now in real code spans. Fenced ```python blocks are untouched. check-docstrings now reports 543 files clean; tests and ruff unchanged. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Seems like the V1 connector was also SSH only. |
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Critical script-injection and terminal-state handling issues, plus additional connector correctness gaps, remain unresolved.
Get a fresh assessment by requesting another Copilot review.
Review effort: Lite
Findings: 2
Open (3)
What changed in this PR
Adds a Slurm connector for containerized Flyte tasks and existing sbatch scripts over SSH, with state handling, documentation, tests, and default-image integration.
Changes:
- Added Slurm task types, connector lifecycle, transport, and script rendering.
- Added tests, documentation, and runnable examples.
- Integrated the plugin into the default connector image.
| File | Description |
|---|---|
plugins/slurm/tests/test_slurm.py |
Slurm connector and rendering tests |
plugins/slurm/tests/__init__.py |
Test package marker |
plugins/slurm/src/flyteplugins/slurm/transport.py |
SSH transport and Slurm state handling |
plugins/slurm/src/flyteplugins/slurm/task.py |
Slurm task definitions |
plugins/slurm/src/flyteplugins/slurm/script.py |
Batch script rendering |
plugins/slurm/src/flyteplugins/slurm/connector.py |
Connector submission and lifecycle |
plugins/slurm/src/flyteplugins/slurm/__init__.py |
Public plugin exports |
plugins/slurm/README.md |
Plugin documentation |
plugins/slurm/pyproject.toml |
Package metadata and dependencies |
maint_tools/build_default_image.py |
Default image integration |
examples/plugins/slurm/slurm_script_example.py |
Script task example |
examples/plugins/slurm/slurm_pipeline_example.py |
Mixed-backend pipeline example |
examples/plugins/slurm/slurm_example.py |
Native Slurm task example |
examples/plugins/slurm/README.md |
Example documentation |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Three findings from the PR review, each verified against a live cluster. Reject CR/LF in sbatch directive values. A newline ends the `#SBATCH` comment and everything after it becomes script body, which Slurm runs as the SSH user. Values reach the renderer from task config -- `sbatch_options` and `working_dir` -- so they are not necessarily trusted. Shell arguments were already quoted; directive values were not. Normalize decorated job states. `sacct` reports `CANCELLED by 1234` and suffixes a state with `+` when the field carries extra information. Only the whitespace form was handled, so `CANCELLED+` fell through to the unrecognized-state path, which reports RUNNING -- leaving a finished task polling until it timed out. Hoist a script task's own leading directives above the generated exports. `sbatch` stops reading `#SBATCH` at the first executable line, so emitting the exports before the user's script silently dropped every directive it carried, contradicting the documented behaviour. Confirmed on a cluster: a `--cpus-per-task` after an `echo` has no effect, and of two identical options the last wins. The plugin's directives are therefore emitted after the user's, so non-conflicting options survive and ours still win a duplicate. A directive below an executable line was already inert and stays in the body rather than being promoted, so wrapping cannot start applying an option the cluster was ignoring. Verified end to end: a wrapped script carrying `--cpus-per-task=3` now allocates 3 CPUs where it previously got the default, while the plugin's job name still overrides the script's. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Self-review after the Copilot findings, which covered only what it flagged. A non-scalar input to a script task was dropped without a word, leaving the author to debug a variable the script never receives. Name it instead, and say what to do about it. The transport docstring and the README both claimed the connector batches status queries. It does not: `AsyncConnector.get` is called per resource, so `status` receives one id at a time and issues one `squeue` per job per poll. Connection reuse is real and is what keeps this to one login-node session rather than one per job; the batching parameter stays because a coalescing cache belongs in the connector, but the docs should not describe a capability that is not exercised. Also records that job files accumulate: nothing removes the `.sbatch`, `.out` and `.err` left in `working_dir`. Keeping them is deliberate -- they are the first thing to read when a job fails -- but a busy cluster needs to prune them, and that was previously unsaid. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Fourteen findings from @samhita-alla. Verified against a live cluster where the behaviour is observable. Privilege boundaries. The connector's environment now wins over task config for host, username, port and known_hosts: the SSH key belongs to the deployment and is shared by every task, so letting a task redirect it to a host of its choosing hands that key over. `skip_host_key_verification` is no longer task-settable at all -- it reads FLYTE_SLURM_SKIP_HOST_KEY_VERIFICATION and warns when task config tries. `sbatch_options` can no longer set job-name, output, error, chdir, wrap, uid or gid; redirecting `--output` at the submitting user's authorized_keys was a shell as that shared account. Directive values now reject all whitespace, not just newlines, since a space smuggles in a second option. The rendered script is no longer logged: it carries every exported variable, including what the deployment injects. Liveness. Remote commands take a 60s timeout, so a wedged login node no longer blocks the poll for every job the connector tracks. squeue and sacct failures are read rather than discarded: "unable to contact slurm controller" used to be indistinguishable from "job finished" and failed a running task. When a job is in neither -- a cluster without accounting -- `scontrol show job` is consulted before giving up. An unrecognized state now fails with the raw string instead of being reported as RUNNING, which polled to the task's own timeout with nothing explaining why. Failing loudly. `resources` on a Slurm environment raised nothing and did nothing; it now explains which Slurm fields to use instead. An input with no environment representation fails at submission rather than warning into a log the task author never sees, and File/Dir inputs export their URI. Correctness. The native task pins `srun --nodes=1 --ntasks=1`. With `nodes=2` and no `ntasks`, srun's default of one task per node started the entrypoint twice against the same output prefix. Confirmed on a two-node allocation: one invocation now, where the unpinned form would have run on both. A RUNNING task's message now carries an `srun --jobid=... --overlap --pty bash` line, so the UI offers a shell in the live allocation -- Flyte's own debug SSH cannot reach a Slurm node. Adds tests/test_transport.py, which drives SSHTransport with a fake `conn.run`: the squeue -> sacct -> scontrol chain, id batching, cancel idempotency, the timeout, and the unreachable-controller path all had no coverage because every existing test replaced the transport wholesale. README states the Pyxis/Enroot requirement up front with the command to check it, notes that PREEMPTED only re-submits when retries > 0, and replaces the multi-node claim with what actually happens. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Native tasks assumed Pyxis, which ships with NVIDIA-shaped GPU clusters but not with the traditional HPC sites that run Apptainer. That made the task type unusable on a good share of clusters, and the only workaround was to fall back to `slurm_script`. `container_runtime` selects between them, defaulting to `pyxis` so nothing changes for existing users. Apptainer is a command the job runs rather than flags on srun, and addresses registries with a URI scheme instead of Enroot's `#` form, so the image reference and the mount and workdir flags are rendered per runtime. Everything else -- directives, exports, the single-task pin, the entrypoint -- is identical, so a task moves between clusters by changing this one field. The runtime-specific part is isolated in `_container_invocation`, which is where a third runtime would go. An unknown value is refused in `Slurm`'s `__post_init__` rather than at submission, so a typo reaches the task author instead of the connector's logs. The Pyxis path is verified on a live cluster after the refactor; the Apptainer path is covered by unit tests only, since that cluster has neither apptainer nor singularity installed. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
|
Thank you @samhita-alla this is great feedback. I'm incorporating it and testing again. I'll add also a known gaps section to the docs. Goal for now is to get early feedback from customers |
The old list named three things and had gone stale against the review changes. This states what the plugin does not do and what to do instead, grouped by execution, data and I/O, operations, and security. New entries are the ones the review and bring-up surfaced: only Pyxis and Apptainer as container runtimes, `resources` being refused rather than ignored, script inputs limited to scalars and URIs, per-job polling rather than batched, job files accumulating in the working directory, the `MinJobAge` blind spot on clusters without accounting, `env` values landing on the cluster in plain text, and the Apptainer path being unit-tested only. The multi-node entry is rewritten: it now describes the `--nodes=1 --ntasks=1` pin and why removing it would run the entrypoint once per node, rather than implying the allocation simply goes unused. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Two gaps in the Apptainer path from review. Apptainer does not expose the host's GPU driver and libraries unless asked, so a GPU job under it started, saw no device, and failed in a way that reads as a broken CUDA install. `--nv` is now added when the job requests GPUs through `gres` or `gpus_per_node`. An explicit GPU flag in the new `container_args` suppresses it, so an AMD site passing `--rocm` does not also get `--nv`. Enroot binds the NVIDIA stack itself, so Pyxis needs no equivalent. Many HPC sites keep tooling off the default PATH and expose it through Lmod or environment-modules, where `apptainer` simply is not found. `modules` emits `module load` lines in the sbatch body, before srun, so the step inherits what they put on PATH. Module names reject whitespace for the same reason directive values do. Both apply to script tasks as well, which may need the same tooling. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
A `slurm_script` task was terminal: it reported phase, exit code and logs,
so nothing downstream could consume its results. Not a policy choice -- the
interface was empty because an arbitrary sbatch script has no way to write
`outputs.pb` in Flyte's literal format, which is what the entrypoint does for
a native task.
Declared outputs close that without asking the script to know anything about
Flyte:
SlurmScriptTask(..., outputs={"model": File, "shards": Dir})
The plugin exports a destination URI per output as `FLYTE_OUTPUT_<NAME>`,
under the same output prefix `outputs.pb` would use, alongside the
`FLYTE_INPUT_*` variables it already exports. The script writes there with
whatever tooling the site has -- `aws s3 cp`, `gcloud storage cp`, rclone --
and the connector records each destination as the declared File or Dir once
the job succeeds, through `Resource.outputs`. The bytes never pass through
the connector: the job already holds credentials to read its inputs, so a
large checkpoint goes straight to object storage rather than through a
connector pod.
Two constraints, both deliberate. Only `File` and `Dir` may be declared,
because a scalar would have to be parsed out of stdout and that is silently
wrong for any script that logs. And a declared output the script never wrote
fails the task even on exit 0, since the alternative is handing a downstream
task a URI to nothing, which surfaces much later as an unexplained read
error; that costs one existence check per output, once, on the poll that
finds the job finished.
Native tasks are untouched -- no destinations, no existence checks, and
`Resource.outputs` stays None, because their entrypoint writes its own.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Matching what Spark and the other backend plugins do: everything the connector needs can come from the config object, leaving the data plane values with nothing to carry but the connector image. Most of this already worked. `ssh_private_key` names a Flyte secret, which `to_custom_config` puts in `connection.secrets`; the backend resolves it through the secret manager and hands the value to the connector (flyte2 webapi/connector/plugin.go:177-191). Host, username and port are plain config fields. `known_hosts` was the exception, because it is a path on the connector and so needed the deployment to mount a file. `known_hosts_secret` names a secret holding the entries themselves instead: asyncssh reads bytes as known_hosts content rather than a filename, so they can be passed straight through with nothing written to disk. That removes the last reason for a volume. The connector environment still exists and still wins where it is set. It is the right model when one cluster serves every task -- a platform team configures it once and tasks carry only scheduling options -- and it is what stops a task redirecting the deployment's shared key at a host of its choosing. The two compose: a connector with nothing set leaves everything to the task config, which is also how local execution works. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Caching was listed as unsupported for script tasks, and enabling it as-is
would have been worse than leaving it off. `cache="auto"` uses
FunctionBodyPolicy, which hashes the task function's source. A script task
has no function, so serialization passes VersionParameters(func=None,
image=None) and the policy returns sha256("") --
e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855, the same
constant for every script task. An edited script would have gone on hitting
its old cache entry, and two unrelated script tasks could have shared one.
Native tasks were never affected: their function body is hashed as usual.
The plugin now substitutes an explicit version over the script body, the
configuration that shapes execution, and the declared outputs. Connection
details and secret names are excluded -- moving the cluster to a new login
node does not change what the job computes, and a rotated secret is not a
change in the work. An author who pinned their own version_override keeps it,
and a disabled cache stays disabled.
Caching only became meaningful now that script tasks can declare outputs;
before there was nothing to restore on a hit.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The script example ended at the job: it submitted an sbatch script and
nothing downstream could use the result, which is no longer the shape of the
task type.
It now declares `outputs={"summary": File}`, and the script's last line
uploads to the destination the plugin gave it -- `aws s3 cp ./summary.json
"$FLYTE_OUTPUT_SUMMARY"`, chosen because `aws` and `rclone` are what a bare
compute node tends to have, while gcloud and the Python storage clients
generally are not. A second environment with no plugin_config then reads it
back as an ordinary File from a Kubernetes pod, so the example shows the
whole path: declare, write, consume.
Also demonstrates `retries` (PREEMPTED only re-submits when it is set) and
`cache="auto"`, whose version for a script task comes from the script body.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Declared outputs made the script do its own upload -- `aws s3 cp` to a URI it was handed -- which asks the node for tooling and credentials it may not have, and made the awkward path the only path. `FLYTE_OUTPUT_<NAME>` is now a local path by default. The script writes an ordinary file and the connector streams it to object storage afterwards, over the SSH connection it already holds, chunked through `storage.put_stream` so nothing is buffered whole: the connector pod has 100Mi of ephemeral storage and is polling every other job at the same time. A declared Dir is walked and each file streamed to the matching key. Above roughly 100MB it warns and uploads anyway. Refusing would throw away a finished job for a transfer that does work, just less well -- two hops instead of one, through a shared pod, on the connector's bandwidth rather than the cluster's. The warning names the fix: `output_upload="job"` hands the script the URI so it uploads directly, which is what a large artifact wants. Verified against a live cluster: `stat` on a 3MB file and on a missing path, a nested directory walk, and a chunked read whose sha256 matches the source. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The plugin has no cloud in it -- SSH to a login node, sbatch, and whatever object store Flyte's storage layer supports -- but the examples and the README read as though GCS on a Soperator cluster were the supported configuration. That is just the cluster this was developed against. The credentials mount is now neutral: a `/etc/cloud` path with the variable named per store rather than GOOGLE_APPLICATION_CREDENTIALS hardcoded, and a comment giving the S3, GCS and Azure equivalents. Raw-data paths use s3:// with a note that gs:// and abfs:// behave the same. The pipeline example no longer calls it "the Nebius filesystem" when it means the Slurm cluster's own filesystem, which is the point being made. The README now says outright that the plugin assumes nothing about the cloud the cluster runs in, and Soperator is named as one operator among on-premise HPC and cloud GPU clusters rather than as the frame. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Warning and streaming anyway meant the bad path stayed available and silent to anyone not reading connector logs. Above 100 MB the connector now fails before moving a byte, naming the fix. The reasoning it refuses on: streaming works, but every byte takes two hops instead of one, through a pod concurrently polling every other job this connector tracks, on its bandwidth rather than the cluster's. The efficient path is the job uploading from the compute node, and that has to be chosen before the job runs, because it decides whether the script is handed a local path or a URI. Nothing can switch after the fact. The cost is real and worth stating: the job's work is lost when this fires. An output that might be large should be declared output_upload='job' up front. The limit is on what the connector will carry, not on how large an output may be -- with 'job' there is no limit at all. Documents both modes with the upload one-liners a script can actually use -- aws s3 cp with --endpoint-url for S3-compatible stores, rclone copyto, gcloud storage cp, azcopy copy, and the --recursive form for a directory output -- plus the reminder that a script task runs on the bare node, so the tooling is the site's rather than the image's. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
100 MB was a reasonable default and a bad constant: a team that gave the connector real bandwidth had no way to use it, and a team running it lean had no way to say so. FLYTE_SLURM_CONNECTOR_UPLOAD_MAX_BYTES now sets it -- a byte count, a suffixed size (500MB, 2GB, 512MiB), or 0/unlimited for no ceiling. It sits on the connector deployment rather than on the task, alongside the other FLYTE_SLURM_* defaults, because the resource being spent is the connector pod's bandwidth and scratch space, shared with every job it polls. A task raising the limit would be spending headroom it does not own, and lowering it per task is not something anyone wants -- so the only useful direction for a task-level knob is the harmful one. Whoever sized the pod owns the number. Two details that matter more than the knob itself: A malformed value raises instead of falling back to 100 MB. A ceiling that silently reverts because of a typo is the kind of thing that costs a day, and an empty string -- what an unset Helm value renders as -- means "default", not "refuse nothing". create() parses it before submitting. Otherwise a bad value would first surface at output collection, after the job had already succeeded, throwing away the compute that produced the output it then refused to move. That is the same mistake the refusal itself exists to avoid. The refusal now names the configured limit and the variable, so an operator reading a task failure learns that a ceiling exists and where it lives. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The section had grown to eighty unbroken lines that answered questions in the wrong order. "Only File and Dir may be declared" sat at the very bottom, after everything about uploads, though it is the first thing that decides whether a script task can do what a reader wants. output_upload was introduced before the reader knew a choice existed. And the 100 MB limit, its rationale, and the variable that configures it ran together in one paragraph. Now three sub-headings in the order the work happens: declare what the script writes (with the File/Dir constraint and the never-wrote failure right there), write to the destination it is given, then choose who uploads. The two modes get a comparison table -- what the variable holds, who uploads, what the node needs, how many hops, whether there is a limit -- so the decision can be made from five rows instead of from prose. Also promotes the #SBATCH-hoisting paragraph to its own heading next to the inputs it belongs with. It had been stranded at the end of the outputs section, and the new sub-headings would have buried it three levels deep under an upload discussion it has nothing to do with. No behaviour change. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
`aws s3 cp ./summary.json "$FLYTE_OUTPUT_SUMMARY"` was left over from when a declared output was always handed to the script as an object-storage URI. The default is now output_upload="connector", where the destination is a local path under ~/.flyte/jobs/<job>.outputs -- so that line was a local-to-local `aws s3 cp`, which the CLI rejects outright. It is a `cp` now, and the example states which mode it is using and what the other one changes. Verified by rendering the example's own sbatch script through the connector rather than by reading it: FLYTE_OUTPUT_SUMMARY comes out as /home/flyte/.flyte/jobs/flyte-train-<hash>.outputs/summary, with the plugin's mkdir -p ahead of the body. The module docstring had the same drift, describing bytes going "straight from the job to object storage, never through the connector" above a task that does the opposite. Also in this pass: slurm_example.py had two copies of the credentials comment, one of them truncated mid-thought, and neither example mentioned container_runtime, container_args or modules -- so an Apptainer cluster had nothing to copy. The multi-node comment now says how to actually get more than one node (raise `nodes`) rather than implying srun alone does it. The README's outputs section gets the same mode table as the plugin README, replacing text that described "job" behaviour as if it were the only behaviour. Checked every Slurm(...) and SlurmScriptTask(...) keyword in the examples against the real signatures: no unknown fields. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Its opening paragraph still described the pre-output_upload behaviour -- a destination URI, bytes going straight from the job, nothing large passing through the connector -- three paragraphs above the text explaining that the default is the opposite. Anyone reading the class docstring top to bottom got the old contract first. Also repairs the indentation of the output_upload paragraph, which an earlier edit of mine left at four spaces inside an eight-space docstring, and swaps one em-dash for the `--` this file uses everywhere else. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
… access
Found by running the script example on a real GKE data plane: the upload failed
with 403 storage.objects.create denied, after the job had already succeeded.
The cause is structural, not a misconfiguration. output_upload="connector" makes
the connector pod write to the run's output prefix, which a connector otherwise
never does -- it talks to its remote system and nothing else. So the dataplane
chart gives the backend cloud identity to union-system, webhook, dataproxy and
fluent-bit, but not to flyteconnector, whose serviceAccount.annotations default
to {}. Every existing GCP and AWS data plane therefore fails this default until
someone annotates that service account.
Worth being plain about: this makes the easy path the one that needs a
deployment change, while output_upload="job" needs none. The default still earns
its place -- a node with no cloud client and no credentials can still produce
outputs -- but the prerequisite has to be unmissable, and it was not documented
at all.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
Found by running the example with `outputs` commented out: the job failed on the
cluster with `FLYTE_OUTPUT_SUMMARY: unbound variable`, because the variable is
only exported for a declared output. Correct, but late, and only loud because
that script happens to `set -u`.
Without `set -u` the same mistake expands to the empty string: `cp ./summary.json
""` fails with `cp: cannot create regular file ''`, and if that is not the last
command in a script without `set -e`, the job exits 0 having written nothing.
Nothing downstream catches it, because an output the task never declared is not
something the connector knows to look for. So the failure mode ranges from a
confusing message to silence.
SlurmScriptTask now scans the script at definition time, alongside the existing
File/Dir type check, and refuses a `$FLYTE_OUTPUT_*` reference it will not
export. The error names the variable, what the task does declare, and both fixes.
It also catches a typo against a real output -- $FLYTE_OUTPUT_SUMARY against a
declared `summary`.
Only `$NAME` and `${NAME}` expansions count. Prose in a comment is not a
reference, and a name assembled at run time (`${FLYTE_OUTPUT_${kind^^}}`) does
not match at all, so a script doing something dynamic is left alone rather than
wrongly refused -- a heuristic that fails a valid script would be worse than the
runtime error it replaces.
The variable-name transform moves to script.py, which task.py and connector.py
both already import, so the names the connector exports and the names the check
accepts cannot drift apart. A test asserts they agree.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The README covered one direction (a declared output the script never wrote) and not the new one, which is the direction someone hits first: writing the cp before remembering to declare the output. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The message was 494 characters and spent most of them explaining what would have
happened without the check -- unbound variable, empty path, silent exit 0. That
belongs in the docs, which is where someone goes to understand the convention.
An error's job is to say what is wrong and what to do.
Now 143, and more useful for being shorter: the no-outputs branch derives the
declaration to add from the variable itself, so it says
The script for 'train' writes to $FLYTE_OUTPUT_SUMMARY, but the task declares
no outputs. Add outputs={"summary": File}, or drop the reference.
and the mismatch branch lists what is declared, which is what identifies a typo.
The reasoning stays in the function's docstring, and the README now leads with
the requirement itself -- both sides must match -- as a three-row table of script
against outputs, replacing two prose paragraphs that stated each half separately.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8
The README's table claimed caching for both task types and then never said what that means for a script, which is the half that isn't obvious: there is no function to hash, so the version is computed from the script body, the declared outputs and the scheduling config, while connection details and output_upload are deliberately excluded. A table of what invalidates against what is reused, since "it caches" is not actionable when the interesting question is which edit costs you a queue wait. Every row checked against the code rather than written from memory. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NMB7XGvvC1qZ2NpTDWymP8


Run Flyte 2 tasks on an existing Slurm cluster, including Soperator-managed clusters, without changing the cluster.
Two task types are served by one connector:
slurm: a Python task submitted as an sbatch job that runs the task's own container image and Flyte entrypoint under Pyxis/Enroot, so typed I/O, caching, retries and error reporting work as they do for a Kubernetes pod. Removingplugin_configruns the identical task on Kubernetes.slurm_script: an existing sbatch script submitted unchanged; scalar inputs are exposed as FLYTE_INPUT_. Phase, exit code and logs only.Transport is SSH to a login node (asyncssh) behind a small SlurmTransport protocol so a slurmrestd transport can be added later. One connection per cluster is reused and status queries are batched (squeue, then sacct for finished jobs). Slurm states map to Flyte phases with PENDING reported as QUEUED rather than RUNNING, PREEMPTED as retryable, and CANCELLED as ABORTED; failures carry Slurm's reason and the tail of stderr.
Connection details may be set per task or cluster-wide on the connector via FLYTE_SLURM_HOST / FLYTE_SLURM_USERNAME / FLYTE_SLURM_SSH_PRIVATE_KEY. The plugin is added to the default connector image build.
Companion docs PR: unionai/unionai-docs#1615
Known gaps
nodes > 1allocates the nodes, but the Flyteentrypoint runs as a single process. No distributed launcher is wired up.
slurmrestdwould slot in behind theSlurmTransportprotocol(
transport.py), but is not written.attribution on the cluster.
slurm_scriptreports phase, exit code and logsonly — no typed outputs, no caching.
Tested scenarios
Fileand the downstream Python task takes it as inputView from the login node:
Execution: https://davide-sm-test.cloud-staging.union.ai/v2/domain/development/project/mlops-union-test/runs/ud9z6t272j6b7j6hdd9l?i=2af5kp2y499ba5nlx1dfhydd0