Repository navigation
Conversation
43571ab to
c1a7257
Compare
There was a problem hiding this comment.
Nice work getting nf-core/demo running end to end. A few notes inline, plus two things here.
Can we hold the merge until nf-flyte is published? CI discovers everything under plugins/ and publishes it, so merging now would release flyteplugins-nextflow with a default image that doesn't build.
The bigger question for me is what Flyte adds while retries, caching and staging all live in Nextflow. Right now it's mostly a pod launcher. I opened unionai/nf-flyte#1 about letting Flyte stage task data and cache on Nextflow's task hash, which would also drop the S3 requirement.
Smaller follow-ups that could fit here: render -with-report/-with-timeline/-with-dag into the task's Flyte report, and in local mode run Nextflow with its own local or docker executor instead of raising.
…w fixes - Keep resume state in the work dir: cloud cache (or NXF_CACHE_DIR for local work dirs), a session id derived from the work dir, history lookup off and an explicit -name. A bare -resume was ignored, since .nextflow/ lived in the pod's temp launch dir. - Default work dir and relative outdir to the task's raw data prefix, minus the per-attempt segment, and pass -resume on retries so a retried head picks up where the failed attempt stopped. - Parent task actions under tctx.task_action. - Forward SIGTERM to Nextflow while it runs (the runtime installs no handler, so the CancelledError path never ran on abort); keep the shutdown wait inside the pod grace period and keep relaying output while it stops. - Read Nextflow output in chunks, so lines over 64 KiB don't fail the task. - report=True renders Nextflow's report, timeline and DAG into the task's Flyte report. - Local runs use Nextflow's own executors instead of raising. - Pre-install nf-cloudcache in nextflow_image(). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Signed-off-by: Kyle Hazen <kyle@union.ai>
|
Thanks for the careful review. Most of the inline points are fixed in the latest push, and the resume changes are verified on a cluster. Details are in the threads. On the broader points:
|
nf-flyte now stages task data through Flyte's copilot (unionai/nf-flyte#3), so task pods no longer access the work dir and it no longer has to be on S3. On a cluster the work dir must be an object store location Flyte can read; the default, under the task's raw data prefix, always is. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Signed-off-by: Kyle Hazen <kyle@union.ai>
The repo's docstring check rejects Sphinx field lists (:param:, :return:); convert nextflow_image() and run_nextflow() to Args: / Returns:. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Signed-off-by: Kyle Hazen <kyle@union.ai>
|
@k1sauce one more case: the head pod dying while tasks are running. Your retry test failed the head after all 8 tasks had finished, so nothing was in flight. The 409 attach in FlyteClient.enqueue is meant for this, but I don't think it fires. On resume, Nextflow gives an unfinished task dir a new hash (it bumps tries in checkCachedOrLaunchTask), and the action name comes from task.hash. So the in-flight task gets a new action name and runs a second time next to the old one. Could the name come from cacheKey() plus failCount instead? cacheKey() doesn't change across sessions, so a retried head would hit the 409 and attach. Can we add the pod-kill case as an example: start nf-core/demo, delete the head pod during FASTQC, and check that attempt 2 finishes without resubmitting the in-flight tasks? With that in, I think the PR is done. |
Summary
Adds
flyteplugins-nextflow, which runs a Nextflow pipeline inside a Flyte task. Each Nextflow task runs as a child action of that task, so the whole pipeline is one Flyte run with one action per Nextflow task.The executor side lives in a separate Nextflow plugin, nf-flyte. That plugin submits each task as a plain container action through the Actions/Run services, and uses the credentials Flyte injects into the task pod. This package is the Flyte side. It builds the image the Nextflow head runs in, and launches Nextflow with everything nf-flyte needs.
What's included
nextflow_image(): an image containing a JRE, a pinned Nextflow version (default25.10.6) and the nf-flyte plugin. The plugin comes from the Nextflow registry, or from a local zip for development. The Nextflow runtime and nf-amazon are downloaded when the image is built, not on every run.run_nextflow(...), which is called from any task:process.executor = 'flyte'work_dirplusresume=Trueresumes across runsrevision,profile,params(passed as a params file),config(a file path or config text) andextra_argsoutdiras aDirNextflowErrorwith the end of the Nextflow output and of.nextflow.logon failureexamples/nf_core_demo.pyand a READMEDesign notes
errorStrategy), caching (-resume) and file staging. nf-flyte disables Flyte retries and caching on task actions.s5cmdinto each task pod for S3 staging. The only image requirement isbash.nextflow_image()doesn't check that a local plugin zip exists. The task module is imported again inside the pod, where the zip isn't present; a missing zip fails the image build instead.Testing
plugins/nextflow/tests) runrun_nextflowagainst a fakenextflowexecutable. They check the command line, the generated config, params, env vars, the defaultoutdir/work_dir, and error reporting.nf-core/demo1.2.0 with-profile testsucceeded. It ran 8 task actions (COWPY, FASTQC ×3, SEQTK_TRIM ×3, MULTIQC) on the pipeline's unmodified biocontainers images, and the results came back as aDir. A separate test pipeline confirmed S3 staging in and out, and that Nextflow'serrorStrategyretries work.Before merging
github.com/unionai/nf-flyteisn't public yet. Until then, the defaultnextflow_image()fails to build. Passnf_flyte=Path(".../nf-flyte-<version>.zip")instead.