diff --git a/adr/20260608-pipeline-composition.md b/adr/20260608-pipeline-composition.md new file mode 100644 index 0000000000..112638b6ca --- /dev/null +++ b/adr/20260608-pipeline-composition.md @@ -0,0 +1,461 @@ +# Pipeline composition + +- Authors: Ben Sherman +- Status: accepted +- Date: 2026-06-08 +- Tags: pipelines, params, dsl +- Version: 1.4 + +## Updates + +### Version 1.4 (2026-09-25) + +- **Remove output block inclusion**: the `output` block of a pipeline can no longer be included. A pipeline call returns its outputs like a workflow call, and the meta-pipeline declares its outputs like any other pipeline. + +### Version 1.3 (2026-09-17) + +- **Included output block is only a record type**: including the `output` block of a pipeline provides a record type of its outputs and nothing else, mirroring the `params` block. Declaring an output with this type no longer redeclares each output of the included pipeline with its output directives. The meta-pipeline declares its outputs like any other pipeline. + +### Version 1.2 (2026-07-13) + +- **Reframe as pipeline composition**: the core feature is the ability to compose pipelines in a Nextflow-native manner. Meta-pipelines are the artifact. Remote pipeline inclusion is deferred to future work. + +### Version 1.1 (2026-06-22) + +- **Separate remote pipelines from remote workflows**: Workflows are treated separately by the [Workflow modules ADR](20260608-workflow-modules.md). +- **Replace core workflow distinction with pipeline inclusion**: Instead of isolating the *core workflow* of a pipeline, the include syntax is extended to support *pipeline inclusion*, in which the `params` / `workflow` / `output` trio is imported and used like a named workflow. + +## Summary + +Provide a way to compose pipelines using regular dataflow logic. + +## Problem Statement + +Processes and workflows can be composed into a larger workflow using dataflow logic. However, pipelines cannot be composed in the same way. The only way to call a pipeline is via `nextflow run`, which does not allow for dataflow composition. + +This ADR defines how a pipeline can be included like a named workflow and composed with other pipelines with dataflow logic. + +## Goals + +- **Preserve dataflow composition**: the included pipeline participates in the including pipeline's dataflow graph (same session, same DAG, same work dir), enabling incremental reaction to emitted outputs. + +- **Preserve reproducibility**: an included pipeline should produce the exact same results as it would when executed directly. Transitive dependencies should not be silently altered to reduce duplication. + +## Non-goals + +- **Nested pipeline execution**: avoid Nextflow-in-Nextflow execution, which forfeits dataflow composition. + +- **Remote pipeline inclusion**: deferred to future work. + +- **Remote pipeline execution**: out of scope. Registry-based execution (e.g. `nextflow pipeline run nf-core/rnaseq@3.0.0`) may be investigated in the future. + +## Decision + +Provide a way to include an entire pipeline (`params` block, entry workflow, `output` block) as a named workflow to facilitate workflow composition. + +## Core Capabilities + +### Pipeline composition + +A pipeline -- that is, a `params` / `workflow` / `output` trio -- can be included and called like a named workflow. This way, pipelines can be composed using regular dataflow logic. + +For example, given the following pipeline: + +```groovy +// pipelines/rnaseq.nf +params { + input: Path + aligner: String = 'star_salmon' + fasta: Path +} +workflow { + // ... +} +output { + bams: Channel { path 'bams' } + multiqc: Path { path 'multiqc' } +} +``` + +It can be included and called as follows: + +```groovy +// main.nf +include { workflow as RNASEQ } from './pipelines/rnaseq.nf' + +workflow { + rnaseq = RNASEQ(record( + input: file('input.csv'), + fasta: file('index.fasta') + )) + rnaseq.bams.view() // Channel + rnaseq.multiqc.view() // Value +} +``` + +Notes: + +- The pipeline must be included using the `workflow` keyword and aliased to a specific name (`RNASEQ`). +- Both the including script and the included pipeline must enable static typing. +- A pipeline with a `params` block is called with a single record of params, so that params with a default can be omitted. A pipeline without a `params` block is called with no arguments. +- The `output` block becomes the `emit:` section. As with a workflow call, a pipeline with a single output returns it directly, and a pipeline with multiple outputs returns a record. +- All outputs are either a `Channel` or wrapped as `Value`, allowing them to be used in regular dataflow logic. + +### Including the params block + +The `params` block of a pipeline can be included as a *record type*, so that a calling pipeline can declare the params of an included pipeline as a single param instead of redeclaring each one: + +```groovy +// main.nf +include { + params as RnaseqParams ; + workflow as RNASEQ +} from './pipelines/rnaseq.nf' + +params { + rnaseq: RnaseqParams +} + +workflow { + rnaseq = RNASEQ( params.rnaseq ) +} +``` + +Notes: + +- The included params record type is *partial*: every field is nullable, and defaults are not pre-filled. The included pipeline applies its own defaults and validates its required params when it is called. +- The user can provide each param of the included pipeline as `--rnaseq.`, and the calling pipeline can override specific params, e.g. `params.rnaseq + record(input: samples)`. + +### Best practices for including pipelines + +Pipeline inclusion only captures the pipeline's main script and included modules. It does not capture external context such as config or the `lib` directory. As a result, the pipeline should be written in a way that works when included in another pipeline: + +1. Pipeline parameters should be defined in the script `params` block and referenced only in the entry workflow and `output` block. The config should only declare *config params* (params that only affect config settings). + +2. Project-level assets (`projectDir`, `bin`, `lib`) should not be used since the meta-pipeline will have a different project root. Module-level assets can be safely used through the module `resources/` bundle and `moduleDir`. + +3. Default `ext` settings should be specified in the process definition or avoided in favor of process inputs. + +4. Software dependencies (`container`, `conda`) should be declared in the process definition, not in config. + +5. Workflow outputs should be published using the `output` block, not `publishDir`. + +None of these constraints are absolute. All of them can be circumvented by manually replicating the external context in the meta-pipeline. Following these constraints simply makes it easier to import a pipeline with minimal extra work. + +## Open Questions + +### Pipeline registry and CLI + +A pipeline registry would enable remote pipeline inclusion: + +```groovy +// module +include { BWA_MEM } from 'nf-core/bwa/mem' + +// pipeline +include { workflow as NFCORE_RNASEQ } from 'nf-core/rnaseq' +``` + +This would require a `nextflow pipeline` command group for publishing and installing pipelines, similar to modules. + +This effort is deferred, since it is not required for pipeline composition. Users can already compose pipelines by cloning or git submodules. + +### Using plugin functions in included pipeline + +If an included pipeline uses plugins, these plugins must be explicitly declared in the meta-pipeline config since they cannot be inferred from the pipeline inclusion. + +If we introduce a pipeline spec, these plugin dependencies could be specified there under `requires.plugins`. When installing a pipeline, Nextflow could copy these plugin declarations into the meta-pipeline config and/or spec. + +### Params overridden by the calling pipeline + +The calling pipeline can override a param of an included pipeline, e.g. `params.rnaseq + record(input: samples)`. If the user also provides `--rnaseq.input`, the value is resolved at launch and then discarded: + +- A `Channel` param is loaded from its samplesheet at launch, so an invalid samplesheet fails the run even though it would be discarded. +- A valid value is ignored without any notice to the user. + +Ideally, Nextflow would report an error when the user provides a value that is overridden. However, the override is an arbitrary expression in the entry workflow, so it can't be detected at launch. One option is to check after the entry workflow is executed (but before the dataflow network is started) for any `Channel` param that was never read. However, this approach only works for `Channel` params, and it can't distinguish an overridden param from one that is unused for another reason (e.g. a pipeline call that is conditionally skipped). + +## Alternatives + +### Pipeline chaining + +An alternative to pipeline composition is a *pipeline chain*, in which multiple Nextflow pipelines are called in sequence via `nextflow run`. + +For example, a fetchngs -> rnaseq pipeline chain can be implemented in a shell script: + +```bash +# fetch FASTQ samples from NCBI SRA +nextflow -q run nf-core/fetchngs \ + --input samplesheet.csv \ + -output-format json \ + > results/output-fetchngs.json + +# adapt fetchngs output to rnaseq input (add strandedness column) +nextflow -q run ./fetchngs-rnaseq.nf \ + -params-file results/output-fetchngs.json \ + --strandedness auto \ + -output-format json \ + > results/output-fetchngs-rnaseq.json + +# perform RNAseq analysis +nextflow -q run nf-core/rnaseq \ + -params-file results/output-fetchngs-rnaseq.json \ + -output-format json \ + > results/output-rnaseq.json +``` + +Or a Nextflow pipeline: + +```groovy +include { NEXTFLOW_RUN as NFCORE_FETCHNGS } from "./modules/local/nextflow/run" +include { NEXTFLOW_RUN as NFCORE_RNASEQ } from "./modules/local/nextflow/run" + +params { + // ... +} + +workflow { + // fetch FASTQ samples from NCBI SRA + fetchngs = NFCORE_FETCHNGS ( + 'nf-core/fetchngs', + // nextflow opts, pipeline inputs, etc ... + ) + // adapt fetchngs output to rnaseq input (add strandedness column) + ch_samples = fetchngs2rnaseq(fetchngs) + // perform RNAseq analysis + rnaseq = NFCORE_RNASEQ ( + 'nf-core/rnaseq', + // nextflow opts, pipeline inputs, etc ... + ) +} + +output { + // ... +} +``` + +The `NEXTFLOW_RUN` process simply calls `nextflow run` in a native process. See [nf-cascade](https://github.com/mahesh-panchal/nf-cascade) for more information about this approach. + +Pipeline chains can also be implemented in Seqera Platform using actions (e.g. when a fetchngs run completes -> launch rnaseq on the fetchngs output). + +Pipeline chaining works with any Nextflow pipeline out of the box, because it simply executes each pipeline directly. Language features such as [workflow outputs](20251020-workflow-outputs.md) and [record types](20260306-record-types.md) make pipeline chaining easier by allowing each pipeline to define structured inputs and outputs. + +However, there are a number of downsides: + +- It forfeits dataflow composition. The developer must serialize/deserialize samplesheet files instead of passing channels directly between pipelines. Each pipeline must complete before the next pipeline can start. + +- It requires an external workflow system instead of reusing the language that pipeline developers already know. Even the Nextflow-in-Nextflow approach shown above requires many tricks to orchestrate nested pipeline runs via the `NEXTFLOW_RUN` process. + +Pipeline chaining can be practical for certain use cases, such as simple chains (A -> B -> C) of off-the-shelf pipelines. But the general solution is to compose pipelines using dataflow logic, just like any other Nextflow pipeline. + +## Links + +- Community issues: [#6474](https://github.com/nextflow-io/nextflow/issues/6474) +- Related: [Workflow params](20250825-workflow-params.md) +- Related: [Workflow outputs](20251020-workflow-outputs.md) +- Related: [Module system](20251114-module-system.md) +- Related: [Workflow modules](20260608-workflow-modules.md) + +## Appendix + +### Example: fetchngs -> rnaseq + +This section walks through the aforementioned `fetchngs -> rnaseq` example as a meta-pipeline. See `examples/pipeline-composition/` for the full example. + +> NOTE: This example uses simplified and idealized versions of `nf-core/fetchngs` and `nf-core/rnaseq` and may not match the actual implementations. + +**Project layout** + +The meta-pipeline is an ordinary Nextflow project with `nf-core/fetchngs` and `nf-core/rnaseq` vendored under `pipelines/`: + +``` +examples/pipeline-composition/ +├── main.nf +├── nextflow.config +└── pipelines/ + └── nf-core/ + ├── fetchngs/ + │ ├── main.nf + │ ├── nextflow.config + │ ├── conf/modules.config + │ └── modules/ + └── rnaseq/ + ├── main.nf + ├── nextflow.config + ├── conf/modules.config + └── modules/ +``` + +Each pipeline has its own `modules/` directory, so the two pipelines can depend on different versions of the same module without conflict. Both pipelines are committed to the meta-pipeline repository. + +**Pipeline code** + +The included pipelines are defined as follows: + +```groovy +// nf-core/fetchngs — main.nf +params { + input: Path // file of SRA/ENA accessions +} +workflow { + main: + ch_ids = channel.fromPath(params.input).splitCsv() + ch_samples = // ... + publish: + samples = ch_samples +} +output { + samples: Channel { path 'fastq' } +} +``` + +```groovy +// nf-core/rnaseq — main.nf +params { + input: Channel // samplesheet + aligner: String = 'star_salmon' + fasta: Path +} +workflow { + main: + rnaseq = // ... + publish: + multiqc = rnaseq.multiqc + bams = rnaseq.bams + counts = rnaseq.counts +} +output { + multiqc: Path { path 'multiqc' } + bams: Channel { path 'bams' } + counts: Channel { path 'counts' } +} +``` + +The meta-pipeline includes each pipeline, along with its `params` block as a record type, and composes them into a new entry workflow: + +```groovy +include { + params as FetchngsParams ; + workflow as NFCORE_FETCHNGS +} from './pipelines/nf-core/fetchngs' + +include { + params as RnaseqParams ; + workflow as NFCORE_RNASEQ +} from './pipelines/nf-core/rnaseq' + +params { + fetchngs: FetchngsParams // input + strandedness: String = 'auto' // unique to meta-pipeline + rnaseq: RnaseqParams // input, aligner, fasta +} + +workflow { + main: + // fetch FASTQ samples from NCBI SRA + samples = NFCORE_FETCHNGS( params.fetchngs ) + + // adapt fetchngs output to rnaseq input (add strandedness) + ch_samples = samples.map { r -> + r + record(strandedness: params.strandedness) + } + + // perform RNAseq analysis (ch_samples overrides params.rnaseq.input) + rnaseq = NFCORE_RNASEQ( params.rnaseq + record(input: ch_samples) ) + + publish: + multiqc = rnaseq.multiqc + bams = rnaseq.bams + counts = rnaseq.counts +} + +output { + multiqc: Path { path 'multiqc' } + bams: Channel { path 'bams' } + counts: Channel { path 'counts' } +} +``` + +Notes: + +- **The handoff is a channel, not a file.** rnaseq declares its samplesheet input as `Channel` instead of `Path`, so that it can be executed directly from a CSV samplesheet or called by a meta-pipeline with a live channel. When rnaseq is launched directly, the `Channel` param is loaded from the samplesheet given on the command line. It allows rnaseq to begin aligning each sample as soon as it is emitted by fetchngs, whereas a pipeline chain would block until fetchngs finished completely. + +- **Included params are partial record types.** See [Including the params block](#including-the-params-block). `rnaseq.input` is supplied by the dataflow, which overrides any value given by the user. `rnaseq.fasta` must still be provided by the user, but the error surfaces at the `NFCORE_RNASEQ()` call rather than at launch. + +- **Params and outputs are not inherited.** The meta-pipeline declares its own `params` and `output` blocks and passes params explicitly to each included pipeline. The included pipelines do not contribute any of their own params or outputs. A meta-pipeline decides for itself which outputs to publish and where to publish them. + +**Configuration** + +Since each included pipeline is just part of the dataflow graph, configuration works like normal. However, the config of each included pipeline is not loaded, so the meta-pipeline must provide it. + +*Global config is redeclared.* Most config settings apply to the whole run, so the `nextflow.config` of each included pipeline can't be merged into the meta-pipeline. In practice, the meta-pipeline must recreate the configuration shell used by the included pipelines: + +- Config params (`outdir`, `publish_dir_mode`, `max_cpus`, etc) +- Resource settings (`cpus`, `memory`, `time`, etc) +- Environment profiles (executors, software dependencies, test profiles) +- Reports (execution, timeline, trace) +- Manifest (name, authors, description, etc) +- Plugins + +For example, each pipeline might declare config params to cap the resources of every task: + +```groovy +params { + max_cpus = 16 + max_memory = '128.GB' +} + +process { + resourceLimits = [ + cpus: params.max_cpus, + memory: params.max_memory + ] +} +``` + +The meta-pipeline declares the same config params in its own `nextflow.config`. + +*Process config is included.* Process config can be reused by the meta-pipeline as long as it is kept in its own config file. Each pipeline keeps its process config in `conf/modules.config`, which is included by the pipeline's own `nextflow.config` as well as the meta-pipeline: + +```groovy +// nextflow.config +includeConfig 'pipelines/nf-core/fetchngs/conf/modules.config' +includeConfig 'pipelines/nf-core/rnaseq/conf/modules.config' +``` + +Process selectors should be written in a way that is correct both when a pipeline is executed directly *and* when it is called by a meta-pipeline. + +Processes in an included pipeline are scoped by the include alias, so the meta-pipeline can override the included process config with qualified selectors: + +```groovy +process { + withName: 'NFCORE_FETCHNGS:.*:SRATOOLS_FASTERQDUMP' { + cpus = 6 + memory = 24.GB + } + withName: 'NFCORE_RNASEQ:.*:STAR_ALIGN' { + cpus = 12 + memory = 72.GB + } +} +``` + +Both the meta-pipeline developer and users can override whatever they want from config. + +**Remote inclusion** + +With a pipeline registry (see [Pipeline registry and CLI](#pipeline-registry-and-cli)), the included pipelines would no longer need to be vendored under `pipelines/`. The includes would refer to the remote pipeline instead of the local path: + +```groovy +include { + params as FetchngsParams ; + workflow as NFCORE_FETCHNGS +} from 'nf-core/fetchngs' + +include { + params as RnaseqParams ; + workflow as NFCORE_RNASEQ +} from 'nf-core/rnaseq' +``` + +Otherwise, the meta-pipeline would work the same way. diff --git a/docs/modules/using-modules.mdx b/docs/modules/using-modules.mdx index f2a02122b1..dccf5f356d 100644 --- a/docs/modules/using-modules.mdx +++ b/docs/modules/using-modules.mdx @@ -134,11 +134,7 @@ See [module run][cli-module-run] for the full command reference. Parameters are inferred from the module's declared inputs: the `input:` section for a process, or the `take:` section for a named workflow. -Type conversions are handled the same way as [typed parameters][typed-params]. Workflow modules use the following additional rules for dataflow types: - -- A `Channel` input accepts a samplesheet path, which Nextflow loads as a channel of records. The samplesheet file can be CSV, JSON, or YAML. The element type must be `Map`, `Record`, or a record type. Each row is validated against and converted to the declared type. - -- A `Value` input accepts a value of type `V`, which Nextflow wraps in a value channel. +Type conversions are handled the same way as [typed parameters][typed-params], including the rules for dataflow types (`Channel` and `Value`). Consider the following workflow: diff --git a/docs/reference/syntax.mdx b/docs/reference/syntax.mdx index ab3763fff0..107fc739f5 100644 --- a/docs/reference/syntax.mdx +++ b/docs/reference/syntax.mdx @@ -88,7 +88,7 @@ include { hello as sayHello } from './some/module' The include source should be a string literal. Each include clause should specify a name, and may also specify an *alias*. In the above example, `hello` is included under the alias `sayHello`. :::note -[Enum](#enum-type) and [record](#record-type) types cannot be aliased. They must be included under their original name. +[Enum](#enum-type) and [record](#record-type) types cannot be aliased. They must be included under their original name. The `params` block of an included pipeline is the exception -- it is a type created for the include, so it must be aliased. ::: Include clauses can be separated by semi-colons or newlines: @@ -111,6 +111,19 @@ include { hello } from './some/module' include { bye as goodbye } from './some/module' ``` + + +The `workflow` name refers to the entire pipeline defined by the included script, i.e., its params block, entry workflow, and output block. The `params` name refers to the params block as a record type. Each of these names must be given an alias: + +```nextflow +include { + params as RnaseqParams ; + workflow as RNASEQ +} from './pipelines/rnaseq.nf' +``` + +See [pipeline composition][pipeline-composition] for more information. + ### Params block The params block consists of one or more *parameter declarations*. A parameter declaration consists of a name, type, and an optional default value: @@ -1015,3 +1028,4 @@ See [strict syntax][strict-syntax-page] for more information. [strict-syntax-page]: ../strict-syntax [workflow-output-def]: ../workflow#outputs [workflow-typed-page]: ../workflow-typed +[pipeline-composition]: ../workflow-typed#pipeline-composition diff --git a/docs/typed-parameters.mdx b/docs/typed-parameters.mdx index 5eddf3efb3..2417fd93f0 100644 --- a/docs/typed-parameters.mdx +++ b/docs/typed-parameters.mdx @@ -30,7 +30,42 @@ The `params` block does not require the `nextflow.enable.types` feature flag. Fo ## Supported types -Parameters can use any [standard type][stdlib-types] except the dataflow types (`Channel` and `Value`). This includes primitive types such as `Path`, `String`, `Integer`, and `Boolean`, as well as collections and records. +Parameters can use any [standard type][stdlib-types]. This includes primitive types such as `Path`, `String`, `Integer`, and `Boolean`, as well as collections, records, and the dataflow types (`Channel` and `Value`). + + + +A parameter can be declared with a dataflow type: + +- A `Channel` parameter accepts a samplesheet path, which Nextflow loads as a channel of records. The samplesheet file can be CSV, JSON, or YAML. The element type must be `Map`, `Record`, or a record type. Each row is validated against and converted to the declared type. + +- A `Value` parameter accepts a value of type `V`, which Nextflow wraps in a value channel. + +For example: + +```nextflow +params { + samples: Channel + index: Value +} + +workflow { + RNASEQ(params.samples, params.index) +} + +record Sample { + id: String + fastq_1: Path + fastq_2: Path +} +``` + +The pipeline can be run with a samplesheet and an index file: + +```console +$ nextflow run main.nf --samples samples.csv --index genome.fa +``` + +Dataflow types are useful when a pipeline is [included by another pipeline][pipeline-composition], because they allow the calling pipeline to provide a parameter from dataflow logic. ## Default and required parameters @@ -66,6 +101,7 @@ The language server validates each parameter reference against its declared type - [Pipeline parameters][cli-params]: How parameter values are resolved across sources. [cli-params]: ./cli#pipeline-parameters +[pipeline-composition]: ./workflow-typed#pipeline-composition [static-typing-page]: ./static-typing [stdlib-types]: ./reference/stdlib-types [syntax-params-block]: ./reference/syntax#params-block diff --git a/docs/workflow-typed.mdx b/docs/workflow-typed.mdx index 54ca4cd00f..20e247dd8d 100644 --- a/docs/workflow-typed.mdx +++ b/docs/workflow-typed.mdx @@ -62,6 +62,130 @@ The operator library supports static typing and records. All operators work in b For more information about best practices when migrating existing code, see [Using operators with static typing][migrating-static-types-operators]. +## Pipeline composition + + + +:::warning{title="Experimental: may change in a future release."} +::: + +An entire pipeline -- the `params` block, entry workflow, and `output` block of a script -- can be included as a named workflow and called like any other workflow. The params block acts as the `take:` section and the output block acts as the `emit:` section. + +Given the following pipeline: + +```nextflow +// pipelines/rnaseq.nf +nextflow.enable.types = true + +params { + input: Channel + aligner: String = 'star_salmon' + fasta: Path +} + +workflow { + main: + // ... + + publish: + bams = ch_bams + multiqc = val_multiqc +} + +output { + bams: Channel { path 'bams' } + multiqc: Path { path 'multiqc' } +} +``` + +It can be included and called as follows: + +```nextflow +include { workflow as RNASEQ } from './pipelines/rnaseq.nf' + +workflow { + main: + rnaseq = RNASEQ(record( + input: samples, + fasta: file('index.fasta') + )) + rnaseq.bams.view() // Channel + rnaseq.multiqc.view() // Value +} +``` + +The included pipeline must be aliased to a specific name (`RNASEQ`), which is also used to scope its processes in the config: + +```nextflow +process { + withName: 'RNASEQ:STAR_ALIGN' { + cpus = 12 + memory = 72.GB + } +} +``` + +Because the included pipeline is part of the calling pipeline's dataflow graph, it can consume a channel produced by another pipeline, and begin working on each item as soon as it is emitted. + +Note the following: + +- Both the calling script and the included pipeline must enable static typing (`nextflow.enable.types = true`). + +- The pipeline is called with a single record, with one field for each param. Params with a default value can be omitted. + +- The outputs published by the included pipeline are emitted to the calling workflow instead of being published. Only the calling pipeline decides what is published, by declaring its own `output` block. + +- Params and outputs are not inherited. The calling pipeline declares its own params and outputs and passes them explicitly to the included pipeline. + +### Importing the params block + +Redeclaring every param of an included pipeline gets tedious. The `params` block can be imported as a record type instead: + +```nextflow +include { + params as RnaseqParams ; + workflow as NFCORE_RNASEQ +} from './pipelines/rnaseq.nf' + +params { + strandedness: String = 'auto' // unique to this pipeline + rnaseq: RnaseqParams // input, aligner, fasta +} + +workflow { + main: + rnaseq = NFCORE_RNASEQ( params.rnaseq + record(input: samples) ) +} +``` + +`RnaseqParams` is a *partial* record type: every field is nullable, so a user can provide any rnaseq param as `--rnaseq.`, the calling pipeline can override specific params with `params.rnaseq + record(input: samples)`, and the `NFCORE_RNASEQ()` call reports any param that is still missing. + +Note the following: + +- A param that the pipeline defaults (`aligner`) can be omitted from the record. The pipeline applies its own default when it is called. + +- `rnaseq.input` is supplied by the dataflow, which overrides any value given by the user. + +- `rnaseq.fasta` must still be provided, but the error surfaces at the `NFCORE_RNASEQ()` call rather than at launch. + +A pipeline receives its params when it is called, so `params` refers to a single execution of the pipeline. A pipeline can be included under any number of aliases, and each alias resolves its own params. As with a named workflow, a pipeline can be called only once per alias -- include it again under a different alias to call it again. + +### Best practices + +Pipeline inclusion only captures the pipeline script and the modules it includes. It does not capture external context such as the config or the `lib` directory. As a result, an included pipeline should be written so that it works when included by another pipeline: + +- Pipeline parameters should be declared in the `params` block. The config should only declare *config params*, i.e. params that only affect config settings. + +- Project-level assets (`projectDir`, `bin`, `lib`) should not be used, since the calling pipeline has a different project root. Module-level assets can be safely used through the module `resources/` bundle and `moduleDir`. + +- Params should be referred to only in the entry workflow and output block. A process or workflow that reads `params` directly should declare an explicit input instead. + +- Process configuration (`container`, `conda`, `ext`) should be specified in the process definition or avoided in favor of process inputs. + +- Workflow outputs should be published using the `output` block, not `publishDir`. + +None of these constraints are absolute. Each of them can be circumvented by replicating the external context in the calling pipeline. Following them simply makes it easier to include a pipeline with minimal extra work. + ## Validation The language server validates each workflow input and output against its declared type. Calling a workflow with an argument whose type does not match its `take:` declaration, or emitting a value that does not match its `emit:` declaration, is reported as an error before the pipeline runs. diff --git a/examples/pipeline-composition/.gitignore b/examples/pipeline-composition/.gitignore new file mode 100644 index 0000000000..cbfb8c695a --- /dev/null +++ b/examples/pipeline-composition/.gitignore @@ -0,0 +1,4 @@ +# Nextflow run artifacts +.nextflow* +work/ +results/ diff --git a/examples/pipeline-composition/README.md b/examples/pipeline-composition/README.md new file mode 100644 index 0000000000..1a3f27496b --- /dev/null +++ b/examples/pipeline-composition/README.md @@ -0,0 +1,174 @@ +# Pipeline composition + +A meta-pipeline that fetches FASTQ samples with `nf-core/fetchngs` and analyzes them with `nf-core/rnaseq`, composing the two pipelines with regular dataflow logic. + +The tools are fake. Every process writes a dummy file, so the example runs anywhere in a couple of seconds. + +## What it demonstrates + +A pipeline (a `params` block, an entry workflow, and an `output` block) can be included as a named workflow and called like any other workflow: + +```nextflow +include { workflow as NFCORE_FETCHNGS } from './pipelines/nf-core/fetchngs' + +workflow { + main: + samples = NFCORE_FETCHNGS( record(ids: file('data/ids.txt')) ) +} +``` + +The `params` block acts as the `take:` section, so the pipeline is called with a record of its params, and params with a default value can be omitted. The `output` block acts as the `emit:` section. `NFCORE_FETCHNGS` has a single output, so the call returns the channel of samples directly to calling workflow. + +The `params` block can also be imported as a record type. That way the meta-pipeline declares one param per included pipeline instead of replicating every param: + +```nextflow +include { + params as RnaseqParams ; + workflow as NFCORE_RNASEQ +} from './pipelines/nf-core/rnaseq' + +params { + rnaseq: RnaseqParams +} +``` + +## Layout + +``` +pipeline-composition/ +├── main.nf # the meta-pipeline +├── nextflow.config +├── data/ # accession list, samplesheet, and reference genome +└── pipelines/ + └── nf-core/ + ├── fetchngs/ + │ ├── main.nf + │ ├── nextflow.config + │ ├── conf/modules.config + │ └── modules/nf-core/ + │ └── sratools/fasterqdump/main.nf + └── rnaseq/ + ├── main.nf + ├── nextflow.config + ├── conf/modules.config + └── modules/nf-core/ + ├── multiqc/main.nf + ├── salmon/quant/main.nf + └── star/align/main.nf +``` + +Each included pipeline is an ordinary Nextflow pipeline, with its own `modules/` directory, vendored into the meta-pipeline repository. They can still be run on their own: + +```console +$ nextflow -C pipelines/nf-core/fetchngs/nextflow.config \ + run ./pipelines/nf-core/fetchngs --ids ./data/ids.txt + +$ nextflow -C pipelines/nf-core/rnaseq/nextflow.config \ + run ./pipelines/nf-core/rnaseq --input ./data/samplesheet.csv --fasta ./data/genome.fa +``` + +The `-C` option makes Nextflow use only the given config. Otherwise it would also load the meta-pipeline's `nextflow.config` from the launch directory. + +`rnaseq` declares `input` as a `Channel` param. On the command line it takes a samplesheet (CSV, JSON, or YAML), which Nextflow loads as a channel of `Sample` records, validating each row against the record type: + +```csv +id,fastq_1,fastq_2,strandedness +SRR1553606,data/fastq/SRR1553606_1.fastq,data/fastq/SRR1553606_2.fastq,auto +``` + +## Running it + +```console +$ nextflow run . +``` + +## What to look at + +**The handoff is a channel, not a file.** `rnaseq` declares its samplesheet input as `Channel`, so `main.nf` passes it a live channel instead of a CSV file: + +```nextflow +ch_samples = samples.map { sample -> + sample + record(strandedness: params.strandedness) +} + +rnaseq = NFCORE_RNASEQ( params.rnaseq + record(input: ch_samples) ) +``` + +Each sample starts aligning as soon as fetchngs emits it. A pipeline chain that shells out to `nextflow run` would wait for fetchngs to finish first. + +**Params flow through the imported record.** Every field of `RnaseqParams` is nullable, so the user fills in what they want: + +```console +$ nextflow run . --rnaseq.aligner hisat2 +``` + +`nextflow.config` sets `params.rnaseq.fasta`, and the two are merged. `aligner` is left unset by both, so rnaseq applies its own default. `input` is supplied by the dataflow, so `--rnaseq.input` is ignored. + +**Outputs are declared by the meta-pipeline.** The rnaseq outputs arrive as channels, and `main.nf` decides which ones to publish and where, in its own `output` block. The output directives of the included pipeline are not inherited. + +**Only the calling pipeline publishes.** `fetchngs` publishes its samples to `fastq/`, but nothing lands there when it is included, because `main.nf` doesn't declare that output. The outputs of an included pipeline are emitted to the calling workflow, which decides what to publish. + +**Global config is redeclared by the meta-pipeline.** Most config settings apply to the whole run (`manifest`, executors, reports, plugins), so the `nextflow.config` of each pipeline can't be merged into the meta-pipeline. The meta-pipeline declares its own instead, such as its `manifest`. + +The same goes for config params, i.e. params that only the config uses. For example, older nf-core pipelines used `max_cpus` and `max_memory` to cap the resources of every task. Each pipeline declares them in `nextflow.config`: + +```nextflow +params { + max_cpus = 4 + max_memory = '8.GB' +} + +process { + resourceLimits = [ + cpus: params.max_cpus, + memory: params.max_memory + ] +} +``` + +The meta-pipeline declares the same config params in its own `nextflow.config`. They are not namespaced like pipeline params (`--max_memory`, not `--rnaseq.max_memory`), because they belong to the config of the whole run: + +```console +$ nextflow run . --max_memory 2.GB +``` + +**Process config is included, not replicated.** Unlike global config, process config can be reused by the meta-pipeline, as long as it lives in its own file. Each pipeline keeps it in `conf/modules.config`: + +```nextflow +// pipelines/nf-core/rnaseq/conf/modules.config +process { + withName: 'STAR_ALIGN' { + cpus = 2 + memory = 2.GB + } +} +``` + +The pipeline's own `nextflow.config` includes it, and so does the meta-pipeline's: + +```nextflow +// nextflow.config +includeConfig 'pipelines/nf-core/fetchngs/conf/modules.config' +includeConfig 'pipelines/nf-core/rnaseq/conf/modules.config' +``` + +The selectors use the simple process name, so they match `STAR_ALIGN` when rnaseq runs on its own and `NFCORE_RNASEQ:STAR_ALIGN` when it is included. + +**Processes are scoped by the include alias.** The run log shows `NFCORE_RNASEQ:STAR_ALIGN`, and `nextflow.config` can target it the same way to override the included config: + +```nextflow +process { + withName: 'NFCORE_RNASEQ:STAR_ALIGN' { + memory = 3.GB + } +} +``` + +A selector on the qualified name is more specific than one on the simple name, so it wins regardless of the order of the config files. `STAR_ALIGN` still gets `cpus = 2` from rnaseq, and `memory` goes up to 3 GB. + +If two included pipelines both have a process called `STAR_ALIGN`, their unqualified selectors both match both processes, and the last one included wins. Qualified selectors in the meta-pipeline config resolve the conflict. + +## See also + +- [Pipeline composition](../../adr/20260608-pipeline-composition.md): the design decision behind this example. +- [Typed workflows](../../docs/workflow-typed.mdx): the syntax used here. diff --git a/examples/pipeline-composition/data/fastq/SRR1553606_1.fastq b/examples/pipeline-composition/data/fastq/SRR1553606_1.fastq new file mode 100644 index 0000000000..66f50bebd9 --- /dev/null +++ b/examples/pipeline-composition/data/fastq/SRR1553606_1.fastq @@ -0,0 +1 @@ +@SRR1553606/1 diff --git a/examples/pipeline-composition/data/fastq/SRR1553606_2.fastq b/examples/pipeline-composition/data/fastq/SRR1553606_2.fastq new file mode 100644 index 0000000000..6303372db8 --- /dev/null +++ b/examples/pipeline-composition/data/fastq/SRR1553606_2.fastq @@ -0,0 +1 @@ +@SRR1553606/2 diff --git a/examples/pipeline-composition/data/genome.fa b/examples/pipeline-composition/data/genome.fa new file mode 100644 index 0000000000..c240c562e0 --- /dev/null +++ b/examples/pipeline-composition/data/genome.fa @@ -0,0 +1,2 @@ +>chr1 +ACGTACGTACGTACGT diff --git a/examples/pipeline-composition/data/ids.txt b/examples/pipeline-composition/data/ids.txt new file mode 100644 index 0000000000..3647089d08 --- /dev/null +++ b/examples/pipeline-composition/data/ids.txt @@ -0,0 +1,3 @@ +SRR1553606 +SRR1553607 +SRR1553608 diff --git a/examples/pipeline-composition/data/samplesheet.csv b/examples/pipeline-composition/data/samplesheet.csv new file mode 100644 index 0000000000..2c0fbf785a --- /dev/null +++ b/examples/pipeline-composition/data/samplesheet.csv @@ -0,0 +1,2 @@ +id,fastq_1,fastq_2,strandedness +SRR1553606,data/fastq/SRR1553606_1.fastq,data/fastq/SRR1553606_2.fastq,auto diff --git a/examples/pipeline-composition/main.nf b/examples/pipeline-composition/main.nf new file mode 100644 index 0000000000..f2834b471e --- /dev/null +++ b/examples/pipeline-composition/main.nf @@ -0,0 +1,52 @@ +#!/usr/bin/env nextflow + +// A meta-pipeline: it fetches FASTQ samples with `nf-core/fetchngs` and +// analyzes them with `nf-core/rnaseq`, composing the two pipelines with +// regular dataflow logic. +// +// The params block of each included pipeline is imported as a record type, +// so that the meta-pipeline declares one param per pipeline instead of +// replicating every param. + +nextflow.enable.types = true + +include { + params as FetchngsParams ; + workflow as NFCORE_FETCHNGS +} from './pipelines/nf-core/fetchngs' + +include { + params as RnaseqParams ; + workflow as NFCORE_RNASEQ +} from './pipelines/nf-core/rnaseq' + +params { + fetchngs: FetchngsParams // ids + strandedness: String = 'auto' // unique to the meta-pipeline + rnaseq: RnaseqParams // input, aligner, fasta +} + +workflow { + main: + // fetch FASTQ samples from NCBI SRA + samples = NFCORE_FETCHNGS( params.fetchngs ) + + // adapt fetchngs output to rnaseq input (add strandedness) + ch_samples = samples.map { sample -> + sample + record(strandedness: params.strandedness) + } + + // perform RNAseq analysis (ch_samples overrides params.rnaseq.input) + rnaseq = NFCORE_RNASEQ( params.rnaseq + record(input: ch_samples) ) + + publish: + bams = rnaseq.bams + counts = rnaseq.counts + multiqc = rnaseq.multiqc +} + +output { + bams: Channel { path 'bams' } + counts: Channel { path 'counts' } + multiqc: Path { path 'multiqc' } +} diff --git a/examples/pipeline-composition/nextflow.config b/examples/pipeline-composition/nextflow.config new file mode 100644 index 0000000000..0981e644ff --- /dev/null +++ b/examples/pipeline-composition/nextflow.config @@ -0,0 +1,38 @@ + +// meta-pipeline config such as `manifest` replaces +// the main config of each included pipeline +manifest { + name = 'fetchngs-rnaseq' + description = 'Fetch FASTQ files from SRA and analyze them with rnaseq' +} + +params { + // supply test input data by default + fetchngs.ids = "${projectDir}/data/ids.txt" + rnaseq.fasta = "${projectDir}/data/genome.fa" + + // config params are not inherited from the included pipelines, + // so the meta-pipeline declares its own + max_cpus = 4 + max_memory = '8.GB' +} + +process { + resourceLimits = [ + cpus: params.max_cpus, + memory: params.max_memory + ] +} + +// meta-pipeline can include process config from included +// pipelines as long as it is isolated from the main config +includeConfig 'pipelines/nf-core/fetchngs/conf/modules.config' +includeConfig 'pipelines/nf-core/rnaseq/conf/modules.config' + +// meta-pipeline can override process config from included +// pipelines using normal process selectors +process { + withName: 'NFCORE_RNASEQ:STAR_ALIGN' { + memory = 3.GB + } +} diff --git a/examples/pipeline-composition/pipelines/nf-core/fetchngs/conf/modules.config b/examples/pipeline-composition/pipelines/nf-core/fetchngs/conf/modules.config new file mode 100644 index 0000000000..00536d70dc --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/fetchngs/conf/modules.config @@ -0,0 +1,8 @@ +// Process config for fetchngs. It is kept apart from `nextflow.config` +// so that a meta-pipeline can include it as-is. +process { + withName: 'SRATOOLS_FASTERQDUMP' { + cpus = 1 + memory = 1.GB + } +} diff --git a/examples/pipeline-composition/pipelines/nf-core/fetchngs/main.nf b/examples/pipeline-composition/pipelines/nf-core/fetchngs/main.nf new file mode 100644 index 0000000000..95a2d6e75c --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/fetchngs/main.nf @@ -0,0 +1,32 @@ +#!/usr/bin/env nextflow + +// A stand-in for `nf-core/fetchngs`: it "downloads" a FASTQ pair for each +// accession. The tools are replaced by dummy files so that the example +// runs anywhere. + +nextflow.enable.types = true + +include { SRATOOLS_FASTERQDUMP } from './modules/nf-core/sratools/fasterqdump' + +params { + ids: Path // file of SRA/ENA accessions, one per line +} + +workflow { + main: + ch_ids = channel.fromList( params.ids.readLines().findAll { line -> line } ) + ch_samples = SRATOOLS_FASTERQDUMP( ch_ids ) + + publish: + samples = ch_samples +} + +output { + samples: Channel { path 'fastq' } +} + +record Sample { + id: String + fastq_1: Path + fastq_2: Path +} diff --git a/examples/pipeline-composition/pipelines/nf-core/fetchngs/modules/nf-core/sratools/fasterqdump/main.nf b/examples/pipeline-composition/pipelines/nf-core/fetchngs/modules/nf-core/sratools/fasterqdump/main.nf new file mode 100644 index 0000000000..3417951ef9 --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/fetchngs/modules/nf-core/sratools/fasterqdump/main.nf @@ -0,0 +1,21 @@ +nextflow.enable.types = true + +process SRATOOLS_FASTERQDUMP { + tag "${id}" + + input: + id: String + + output: + record( + id: id, + fastq_1: file('*_1.fastq'), + fastq_2: file('*_2.fastq') + ) + + script: + """ + echo "@${id}/1" > ${id}_1.fastq + echo "@${id}/2" > ${id}_2.fastq + """ +} diff --git a/examples/pipeline-composition/pipelines/nf-core/fetchngs/nextflow.config b/examples/pipeline-composition/pipelines/nf-core/fetchngs/nextflow.config new file mode 100644 index 0000000000..210dcdcebf --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/fetchngs/nextflow.config @@ -0,0 +1,19 @@ +manifest { + name = 'nf-core/fetchngs' + description = 'Fetch FASTQ files from public databases' +} + +// config params: params that are only used by the config +params { + max_cpus = 4 + max_memory = '8.GB' +} + +process { + resourceLimits = [ + cpus: params.max_cpus, + memory: params.max_memory + ] +} + +includeConfig 'conf/modules.config' diff --git a/examples/pipeline-composition/pipelines/nf-core/rnaseq/conf/modules.config b/examples/pipeline-composition/pipelines/nf-core/rnaseq/conf/modules.config new file mode 100644 index 0000000000..f70d47a792 --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/rnaseq/conf/modules.config @@ -0,0 +1,12 @@ +// Process config for rnaseq. It is kept apart from `nextflow.config` +// so that a meta-pipeline can include it as-is. +process { + withName: 'STAR_ALIGN' { + cpus = 2 + memory = 2.GB + } + withName: 'SALMON_QUANT' { + cpus = 1 + memory = 1.GB + } +} diff --git a/examples/pipeline-composition/pipelines/nf-core/rnaseq/main.nf b/examples/pipeline-composition/pipelines/nf-core/rnaseq/main.nf new file mode 100644 index 0000000000..6b6b3bc670 --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/rnaseq/main.nf @@ -0,0 +1,42 @@ +#!/usr/bin/env nextflow + +// A stand-in for `nf-core/rnaseq`: it aligns each sample, quantifies it, and +// summarizes the run. The tools are replaced by dummy files so that the +// example runs anywhere. + +nextflow.enable.types = true + +include { STAR_ALIGN } from './modules/nf-core/star/align' +include { SALMON_QUANT } from './modules/nf-core/salmon/quant' +include { MULTIQC } from './modules/nf-core/multiqc' + +params { + input: Channel // samplesheet + aligner: String = 'star_salmon' + fasta: Path +} + +workflow { + main: + ch_bams = STAR_ALIGN( params.input, params.fasta, params.aligner ) + ch_counts = SALMON_QUANT( ch_bams ) + val_multiqc = MULTIQC( ch_counts.map { c -> c.counts }.collect() ) + + publish: + bams = ch_bams.map { a -> a.bam } + counts = ch_counts.map { c -> c.counts } + multiqc = val_multiqc +} + +output { + bams: Channel { path 'bams' } + counts: Channel { path 'counts' } + multiqc: Path { path 'multiqc' } +} + +record Sample { + id: String + fastq_1: Path + fastq_2: Path + strandedness: String +} diff --git a/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/multiqc/main.nf b/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/multiqc/main.nf new file mode 100644 index 0000000000..8db2eabd77 --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/multiqc/main.nf @@ -0,0 +1,14 @@ +nextflow.enable.types = true + +process MULTIQC { + input: + counts: Bag + + output: + file('multiqc_report.html') + + script: + """ + echo "summarized ${counts.size()} samples" > multiqc_report.html + """ +} diff --git a/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/salmon/quant/main.nf b/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/salmon/quant/main.nf new file mode 100644 index 0000000000..42bfee12be --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/salmon/quant/main.nf @@ -0,0 +1,22 @@ +nextflow.enable.types = true + +process SALMON_QUANT { + tag "${id}" + + input: + record( + id: String, + bam: Path + ) + + output: + record( + id: id, + counts: file('*.counts.tsv') + ) + + script: + """ + printf 'gene\\tcount\\nENSG0001\\t42\\n' > ${id}.counts.tsv + """ +} diff --git a/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/star/align/main.nf b/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/star/align/main.nf new file mode 100644 index 0000000000..fa3397d351 --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/rnaseq/modules/nf-core/star/align/main.nf @@ -0,0 +1,26 @@ +nextflow.enable.types = true + +process STAR_ALIGN { + tag "${id}" + + input: + record( + id: String, + fastq_1: Path, + fastq_2: Path, + strandedness: String + ) + fasta: Path + aligner: String + + output: + record( + id: id, + bam: file('*.bam') + ) + + script: + """ + echo "aligned ${id} against ${fasta.name} with ${aligner} (${strandedness})" > ${id}.bam + """ +} diff --git a/examples/pipeline-composition/pipelines/nf-core/rnaseq/nextflow.config b/examples/pipeline-composition/pipelines/nf-core/rnaseq/nextflow.config new file mode 100644 index 0000000000..31200780ea --- /dev/null +++ b/examples/pipeline-composition/pipelines/nf-core/rnaseq/nextflow.config @@ -0,0 +1,19 @@ +manifest { + name = 'nf-core/rnaseq' + description = 'RNA sequencing analysis pipeline' +} + +// config params: params that are only used by the config +params { + max_cpus = 4 + max_memory = '8.GB' +} + +process { + resourceLimits = [ + cpus: params.max_cpus, + memory: params.max_memory + ] +} + +includeConfig 'conf/modules.config' diff --git a/modules/nextflow/src/main/groovy/nextflow/Session.groovy b/modules/nextflow/src/main/groovy/nextflow/Session.groovy index 6d5141ca36..5815f1d567 100644 --- a/modules/nextflow/src/main/groovy/nextflow/Session.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/Session.groovy @@ -32,7 +32,6 @@ import groovy.transform.Memoized import groovy.transform.PackageScope import groovy.util.logging.Slf4j import groovyx.gpars.GParsConfig -import groovyx.gpars.dataflow.DataflowWriteChannel import groovyx.gpars.dataflow.operator.DataflowProcessor import nextflow.cache.CacheDB import nextflow.cache.CacheFactory @@ -121,8 +120,6 @@ class Session implements ISession { */ private volatile boolean dataflowNetworkFired - final Map outputs = [:] - /** * Creates process executors */ diff --git a/modules/nextflow/src/main/groovy/nextflow/config/parser/v2/ConfigDsl.groovy b/modules/nextflow/src/main/groovy/nextflow/config/parser/v2/ConfigDsl.groovy index d608e01e3c..722ad8fa0c 100644 --- a/modules/nextflow/src/main/groovy/nextflow/config/parser/v2/ConfigDsl.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/config/parser/v2/ConfigDsl.groovy @@ -86,7 +86,8 @@ class ConfigDsl extends Script { void setParams(Map params) { this.cliParams = params - (target.params as Map).putAll(params) + // deep copy so that nested config assignments don't modify the CLI params + (target.params as Map).putAll(Bolts.deepClone(params)) } void setConfigParams(Map params) { @@ -138,23 +139,40 @@ class ConfigDsl extends Script { /** * Assign a value to a config option. * - * When assigning a param, if the param was specified - * on the command line, then the command line value takes - * precedence. CLI params are applied here in order to ensure - * that if the param is referenced later in the config file, - * the command line value is used. + * When assigning a param (e.g. `params.input` or `params.rnaseq.fasta`), + * the corresponding command line value takes precedence. CLI params are + * applied here in order to ensure that if the param is referenced later + * in the config file, the command line value is used. * * @param names * @param value */ void assign(List names, Object value) { - if( names.size() == 2 && names.first() == 'params' ) { - final name = names.last() - declareParam(name, value) - if( cliParams.containsKey(name) ) - value = asDeclaredType(cliParams[name], value) - } + final isParam = names.size() > 1 && names.first() == 'params' + if( isParam ) + value = withCliOverride(names.tail(), value) navigate(names.init()).put(names.last(), value) + if( isParam ) + declareParam(names[1], (target.params as Map).get(names[1])) + } + + /** + * Apply the command line value of a param, if any, to its config value. + * A map value is merged with the corresponding command line values. + * + * @param path + * @param value + */ + private Object withCliOverride(List path, Object value) { + Object cliValue = cliParams + for( final name : path ) { + if( cliValue !instanceof Map || !((Map)cliValue).containsKey(name) ) + return value + cliValue = ((Map)cliValue).get(name) + } + return cliValue instanceof Map && value instanceof Map + ? Bolts.deepMerge((Map)value, (Map)cliValue) + : asDeclaredType(cliValue, value) } /** diff --git a/modules/nextflow/src/main/groovy/nextflow/script/BaseScript.groovy b/modules/nextflow/src/main/groovy/nextflow/script/BaseScript.groovy index 9a8839cba8..1650816f26 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/BaseScript.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/BaseScript.groovy @@ -20,6 +20,7 @@ import java.lang.reflect.InvocationTargetException import java.nio.file.Paths import groovy.transform.CompileStatic +import groovy.transform.PackageScope import groovy.util.logging.Slf4j import nextflow.NF import nextflow.NextflowMeta @@ -48,6 +49,8 @@ abstract class BaseScript extends Script implements ExecutionContext { private WorkflowDef entryFlow + private boolean moduleLoaded + private OutputDef outputDef BaseScript() { @@ -72,6 +75,32 @@ abstract class BaseScript extends Script implements ExecutionContext { return typingEnabled } + /** + * The declared params of this script, keyed by name. + */ + Map getParamDeclarations() { + return paramsDef.getDeclarations() + } + + /** + * The entry workflow of this script, or null if it doesn't have one. + */ + WorkflowDef getEntryFlow() { + return entryFlow + } + + /** + * Execute this script as an included module, only once + * no matter how many scripts include it. + */ + @PackageScope + void runModule() { + if( moduleLoaded ) + return + moduleLoaded = true + run() + } + /** * Holds the configuration object which will used to execution the user tasks */ @@ -284,11 +313,13 @@ abstract class BaseScript extends Script implements ExecutionContext { // Execute a single named workflow directly final handler = new WorkflowEntryHandler(this, session, meta) this.entryFlow = handler.createEntryWorkflow() + this.outputDef = handler.createOutputDef() } else if( moduleRun && meta.hasExecutableProcesses() ) { // Execute a single process directly final handler = new ProcessEntryHandler(this, session, meta) this.entryFlow = handler.createEntryWorkflow() + this.outputDef = handler.createOutputDef() } else if( meta.getLocalProcessNames() || meta.getLocalWorkflowNames() ) { throw new AbortOperationException("No entry workflow specified -- script must define an entry workflow, a single process or named workflow, or be a code snippet") @@ -303,9 +334,12 @@ abstract class BaseScript extends Script implements ExecutionContext { session.notifyBeforeWorkflowExecution() if( paramsDef ) paramsDef.apply(session) - final ret = entryFlow.invoke_a(BaseScriptConsts.EMPTY_ARGS) + final args = paramsDef + ? [ binding.getParams() ] as Object[] + : BaseScriptConsts.EMPTY_ARGS + final ret = entryFlow.run(args) if( outputDef ) - outputDef.apply(session) + outputDef.apply(session, entryFlow.getBinding().getPublished()) session.notifyAfterWorkflowExecution() return ret } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/BindableDef.groovy b/modules/nextflow/src/main/groovy/nextflow/script/BindableDef.groovy index 5860c03a17..5f8ccec5c5 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/BindableDef.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/BindableDef.groovy @@ -42,7 +42,7 @@ abstract class BindableDef extends ComponentDef { final fqName = prefix ? prefix+SCOPE_SEP+name : name if( this instanceof ProcessDef && !invocations.add(fqName) ) { log.debug "Bindable invocations=$invocations" - final msg = "Process '$name' has been already used -- If you need to reuse the same component, include it with a different name or include it in a different workflow context" + final msg = "Process '$fqName' was called twice -- include it with a different alias or call it in a different workflow" throw new DuplicateProcessInvocation(msg) } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/IncludeDef.groovy b/modules/nextflow/src/main/groovy/nextflow/script/IncludeDef.groovy index 6f4520393f..2500bd9967 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/IncludeDef.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/IncludeDef.groovy @@ -102,7 +102,7 @@ class IncludeDef { final moduleFile = realModulePath(path).normalize() // -- load the module final moduleScript = NF.isSyntaxParserV2() - ? loadModuleV2(moduleFile, ownerParams, session) + ? loadModuleV2(moduleFile, ownerParams) : loadModuleV1(moduleFile, resolveParams(ownerParams), session) // -- add it to the inclusions for( Module module : modules ) { @@ -131,16 +131,14 @@ class IncludeDef { * * @param path The included script path * @param params The params of the including script - * @param session The current workflow run */ @PackageScope - @Memoized - static BaseScript loadModuleV2(Path path, Map params, Session session) { + static BaseScript loadModuleV2(Path path, Map params) { final script = ScriptMeta.getScriptByPath(path) if( !script ) throw new IllegalStateException("Unable to find module script for path: $path") script.getBinding().setParams(params) - script.run() + script.runModule() return script } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/OutputDef.groovy b/modules/nextflow/src/main/groovy/nextflow/script/OutputDef.groovy index 7e0df472d5..219baaec94 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/OutputDef.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/OutputDef.groovy @@ -18,6 +18,7 @@ package nextflow.script import groovy.transform.CompileStatic import groovy.util.logging.Slf4j +import groovyx.gpars.dataflow.DataflowWriteChannel import nextflow.Session /** * Models the workflow output definition @@ -34,14 +35,14 @@ class OutputDef { this.closure = closure } - void apply(Session session) { + void apply(Session session, Map outputs) { final dsl = new OutputDsl() final cl = (Closure)closure.clone() cl.setDelegate(dsl) cl.setResolveStrategy(Closure.DELEGATE_FIRST) cl.call() - dsl.apply(session) + dsl.apply(session, outputs) } } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/OutputDsl.groovy b/modules/nextflow/src/main/groovy/nextflow/script/OutputDsl.groovy index d11228e28c..e9de914caa 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/OutputDsl.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/OutputDsl.groovy @@ -19,6 +19,7 @@ package nextflow.script import groovy.transform.CompileStatic import groovy.util.logging.Slf4j import groovyx.gpars.dataflow.DataflowVariable +import groovyx.gpars.dataflow.DataflowWriteChannel import nextflow.Session import nextflow.exception.ScriptRuntimeException import nextflow.extension.CH @@ -50,8 +51,7 @@ class OutputDsl { declarations[name] = dsl.getOptions() } - void apply(Session session) { - final outputs = session.outputs + void apply(Session session, Map outputs) { final defaults = session.config.navigate('workflow.output', Collections.emptyMap()) as Map // make sure every output was assigned diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ParamsDef.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ParamsDef.groovy index 18914cdd58..52f5ff5e2b 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ParamsDef.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ParamsDef.groovy @@ -35,14 +35,28 @@ class ParamsDef { this.closure = closure } + private ParamsDsl dsl + void apply(Session session) { - final dsl = new ParamsDsl(clazz) + dsl().apply(session) + } + + /** + * The declared params, keyed by name. + */ + Map getDeclarations() { + return dsl().getDeclarations() + } + + private ParamsDsl dsl() { + if( dsl != null ) + return dsl + dsl = new ParamsDsl(clazz) final cl = (Closure)closure.clone() cl.setDelegate(dsl) cl.setResolveStrategy(Closure.DELEGATE_FIRST) cl.call() - - dsl.apply(session) + return dsl } } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy index 68f2977c9f..f97a371a11 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy @@ -21,9 +21,6 @@ import java.lang.reflect.Type import groovy.transform.CompileStatic import groovy.util.logging.Slf4j import nextflow.Session -import nextflow.exception.ScriptRuntimeException -import nextflow.script.dsl.Types -import nextflow.util.TypeHelper /** * Implements the DSL for defining workflow params * @@ -48,41 +45,10 @@ class ParamsDsl { declarations[name] = new Param(name, type, optional, defaultValue) } - void apply(Session session) { - final cliParams = session.cliParams ?: [:] - final configParams = session.configParams ?: [:] - - for( final name : cliParams.keySet() ) { - if( !declarations.containsKey(name) && !configParams.containsKey(name) ) - throw new ScriptRuntimeException("Parameter `$name` was specified on the command line or params file but is not declared in the script or config") - } - - final params = new HashMap() - for( final name : declarations.keySet() ) { - final decl = declarations[name] - if( cliParams.containsKey(name) ) { - params[name] = ParamsHelper.resolveFromCli(decl, cliParams[name]) - } - else if( configParams.containsKey(name) ) { - params[name] = ParamsHelper.resolveFromCode(decl, configParams[name]) - } - else if( decl.defaultValue != null ) { - params[name] = ParamsHelper.resolveFromCode(decl, decl.defaultValue) - } - else { - params[name] = null - } + Map getDeclarations() { declarations } - if( params[name] == null && !decl.optional ) { - throw new ScriptRuntimeException("Parameter `$name` is required but no value was provided") - } - - final expectedType = TypeHelper.getRawType(decl.type) - final actualType = params[name]?.getClass() - if( actualType != null && !ParamsHelper.isAssignableFrom(expectedType, actualType) ) { - throw new ScriptRuntimeException("Parameter `$name` with type ${Types.getName(decl.type)} cannot be assigned to ${params[name]} [${Types.getName(actualType)}]") - } - } + void apply(Session session) { + final params = ParamsHelper.resolveParams(declarations.values(), session.cliParams ?: [:], session.configParams ?: [:]) // propagate resolved params to all scripts for legacy compatibility if( !session.binding.getScriptPath() ) diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy index a099b6070b..330c4791f5 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy @@ -16,12 +16,26 @@ package nextflow.script +import java.lang.reflect.ParameterizedType +import java.lang.reflect.Type import java.nio.file.Path +import java.util.function.BiFunction +import groovy.json.JsonSlurper import groovy.transform.CompileStatic +import groovy.yaml.YamlSlurper +import groovyx.gpars.dataflow.DataflowWriteChannel +import nextflow.dataflow.ChannelImpl +import nextflow.dataflow.ChannelNamespace +import nextflow.dataflow.ValueImpl import nextflow.exception.ScriptRuntimeException +import nextflow.script.dsl.Nullable +import nextflow.script.dsl.PipelineParams import nextflow.script.dsl.Types +import nextflow.script.types.Channel import nextflow.script.types.Record +import nextflow.script.types.Value +import nextflow.splitter.CsvSplitter import nextflow.util.Duration import nextflow.util.MemoryUnit import nextflow.util.RecordMap @@ -39,6 +53,224 @@ import org.codehaus.groovy.runtime.typehandling.GroovyCastException @CompileStatic class ParamsHelper { + /** + * Resolve declared params from the command line and config. + * + * The config params already include the command line overrides + * (see ConfigDsl), so they take precedence when present. + * + * @param declarations + * @param cliParams + * @param configParams + */ + static Map resolveParams(Collection declarations, Map cliParams, Map configParams) { + final names = declarations*.name as Set + for( final name : cliParams.keySet() ) { + if( name !in names && !configParams.containsKey(name) ) + throw new ScriptRuntimeException("Parameter `${name}` was specified on the command line or params file but is not declared in the script or config") + } + + final given = cliParams.subMap(names) + configParams.subMap(names) + return resolveParams(declarations, given, '') { Param decl, Object value -> + resolveParam(decl, value, cliParams.containsKey(decl.name)) + } + } + + /** + * Resolve declared params against the given values. A param + * with no given value is given its default value. + * + * @param declarations + * @param given + * @param context appended to the param name in error messages + * @param resolve resolves a given value against its declared param + */ + static Map resolveParams(Collection declarations, Map given, String context, BiFunction resolve) { + final result = new LinkedHashMap(declarations.size()) + for( final decl : declarations ) { + final name = decl.name + final value = given.containsKey(name) + ? resolve.apply(decl, given.get(name)) + : resolveDefault(decl) + + if( value == null && !decl.optional ) + throw new ScriptRuntimeException("Parameter `${name}`${context} is required but no value was provided") + + result.put(name, value) + } + return result + } + + /** + * Resolve the params given to the entry workflow of a pipeline against + * the params block of the pipeline. Called by the entry workflow (see + * WorkflowToGroovyVisitor). + * + * The session params of a top-level run are already resolved by the + * params block, so they are returned as-is. An included pipeline + * resolves its params even when given the session params (e.g. + * `RNASEQ(params)`), so that its defaults are applied. + * + * @param script the pipeline script + * @param value the params given to the entry workflow + */ + static Map resolveArguments(BaseScript script, Object value) { + if( !ScriptMeta.get(script).isModule() ) + return (Map)value + + final pipeline = ExecutionStack.workflow().name + if( value !instanceof RecordMap && value !instanceof ScriptBinding.ParamsMap ) + throw new ScriptRuntimeException("Pipeline `${pipeline}` should be called with a record") + + final given = (Map)value + final declarations = script.getParamDeclarations() + for( final name : given.keySet() ) { + if( !declarations.containsKey(name) ) + throw new ScriptRuntimeException("Pipeline `${pipeline}` does not declare a parameter named `${name}`") + } + + final params = resolveParams(declarations.values(), given, " of pipeline `${pipeline}`") { Param decl, Object v -> + isDataflow(v) ? DataflowTypeHelper.normalizeV2(v) : resolveParam(decl, v, false) + } + return new RecordMap(params) + } + + private static boolean isDataflow(Object value) { + return value instanceof ChannelImpl + || value instanceof ValueImpl + || value instanceof DataflowWriteChannel + || value instanceof ChannelOut + } + + /** + * Resolve a param value against its declared type. + * + * A {@code Channel} param is loaded from a samplesheet file, with each + * record converted to the element type. A {@code Value} param is + * converted to {@code V} and wrapped in a dataflow value. Any other param + * is converted directly to the declared type. + * + * @param decl + * @param value + * @param fromCli whether the value came from the command line (and is + * therefore a string that may need to be parsed) + */ + static Object resolveParam(Param decl, Object value, boolean fromCli) { + if( value == null ) + return null + + final rawType = TypeHelper.getRawType(decl.type) + + if( rawType == Channel ) + return ChannelNamespace.fromList(loadChannelInput(decl, value)) + + if( rawType == Value ) + return ChannelNamespace.value(resolveParam(elementDecl(decl), value, fromCli)) + + if( TypeHelper.isRecordType(decl.type) && value instanceof Map ) + return resolveRecord(decl, (Map)value, fromCli) + + final result = fromCli + ? resolveFromCli(decl, value) + : resolveFromCode(decl, value) + checkAssignable(decl, result) + return result + } + + private static RecordMap resolveRecord(Param decl, Map value, boolean fromCli) { + final type = (Class)decl.type + final result = new LinkedHashMap(value) + for( final field : type.getDeclaredFields() ) { + if( field.isSynthetic() ) + continue + final name = field.getName() + final optional = field.isAnnotationPresent(Nullable) + final fieldValue = value.get(name) + if( fieldValue == null ) { + if( !optional ) + throw new ScriptRuntimeException("Parameter `${decl.name}` with type ${type.getSimpleName()} is missing required field `${name}`") + continue + } + final fieldDecl = new Param("${decl.name}.${name}", field.getGenericType(), optional, null) + result.put(name, resolveParam(fieldDecl, fieldValue, fromCli)) + } + return new RecordMap(result) + } + + /** + * Load a channel param from a samplesheet file, converting each record + * to the declared element type. + * + * @param decl + * @param value + */ + private static List loadChannelInput(Param decl, Object value) { + if( value !instanceof CharSequence && value !instanceof Path ) + throw new ScriptRuntimeException("Parameter `${decl.name}` with type ${Types.getName(decl.type)} should be a samplesheet file, but received: ${value} [${Types.getName(value.getClass())}]") + + final path = value instanceof Path + ? (Path)value + : TypeHelper.asPathType(value.toString()) + final elementType = elementDecl(decl).type + final elementRawType = TypeHelper.getRawType(elementType) + + if( !Map.isAssignableFrom(elementRawType) && !Record.isAssignableFrom(elementRawType) ) + throw new ScriptRuntimeException("Parameter `${decl.name}` with type ${Types.getName(decl.type)} cannot be loaded from a samplesheet -- the element type should be Map, Record, or a record type") + + return loadFromFile(decl.name, path).collect { el -> + try { + TypeHelper.asType(el, elementType) + } + catch( Exception e ) { + throw new ScriptRuntimeException("Invalid record in samplesheet '${path}' for parameter `${decl.name}` -- ${e.message}") + } + } + } + + /** + * Get the declared param for the element type of a parameterized + * type, e.g. {@code Sample} for {@code Channel}. + * + * @param decl + */ + private static Param elementDecl(Param decl) { + final elementType = decl.type instanceof ParameterizedType + ? ((ParameterizedType)decl.type).getActualTypeArguments()[0] + : (Type)Object + return new Param(decl.name, elementType, decl.optional, null) + } + + /** + * Load the contents of a samplesheet file as a list of records. + * + * Supported formats: + * - CSV: header row required, comma-separated + * - JSON: must be a top-level array + * - YAML / YML: must be a top-level sequence + * + * @param name the param name (for error messages) + * @param file the samplesheet file to load + */ + static List loadFromFile(String name, Path file) { + final ext = file.getExtension() + final value = switch( ext ) { + case 'csv' -> loadFromCsv(file) + case 'json' -> new JsonSlurper().parse(file) + case 'yaml', 'yml' -> new YamlSlurper().parse(file) + default -> throw new ScriptRuntimeException("Unrecognized file format '${ext}' for input file '${file}' for parameter `${name}` -- should be CSV, JSON, or YAML") + } + if( value !instanceof List ) + throw new ScriptRuntimeException("Input file '${file}' for parameter `${name}` must contain a list of records, but got: ${value.class.simpleName}") + return (List)value + } + + private static List loadFromCsv(Path file) { + final rows = new CsvSplitter().options(header: true, sep: ',', quote: '"').target(file).list() + return rows.collect { row -> + ((Map)row).collectEntries { k, v -> [ k, v != '' ? v : null ] } + } + } + /** * Resolve a value given on the command line. Command-line values are * always strings, so they are parsed according to the declared type. @@ -208,6 +440,40 @@ class ParamsHelper { } } + /** + * The value of a param for which no value was provided: its + * default value, if any, otherwise an empty record or null. + * + * A param whose type is the params block of an included pipeline + * defaults to an empty record, so that the calling pipeline can + * supply the params by dataflow. Missing params are reported when + * the pipeline is called. + * + * @param decl + */ + static Object resolveDefault(Param decl) { + if( decl.defaultValue != null ) + return resolveParam(decl, decl.defaultValue, false) + final type = TypeHelper.getRawType(decl.type) + return type.isAnnotationPresent(PipelineParams) + ? new RecordMap([:]) + : null + } + + /** + * Check that a resolved value can be assigned to the declared type + * of a param. + * + * @param decl + * @param value + */ + private static void checkAssignable(Param decl, Object value) { + final expectedType = TypeHelper.getRawType(decl.type) + final actualType = value?.getClass() + if( actualType != null && !isAssignableFrom(expectedType, actualType) ) + throw new ScriptRuntimeException("Parameter `${decl.name}` with type ${Types.getName(decl.type)} cannot be assigned to ${value} [${Types.getName(actualType)}]") + } + static boolean isAssignableFrom(Class target, Class source) { // any numeric value can be assigned to Float if( target == Float.class ) diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ProcessEntryHandler.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ProcessEntryHandler.groovy index d165164a14..2b750c32e3 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ProcessEntryHandler.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ProcessEntryHandler.groovy @@ -85,8 +85,7 @@ class ProcessEntryHandler { // Execute the process final output = processDef.run(inputArgs as Object[]) as ChannelOut // Publish process outputs as workflow outputs - assignOutputs(output) - publishOutputs() + assignOutputs((WorkflowBinding)(Object)getDelegate(), output) return output } @@ -98,10 +97,9 @@ class ProcessEntryHandler { return new WorkflowDef(script, workflowBody) } - private void assignOutputs(ChannelOut output) { + private void assignOutputs(WorkflowBinding dsl, ChannelOut output) { final config = processDef.getProcessConfig() final outputNames = getProcessOutputs(config) - final dsl = script.getBinding() if( output.size() != outputNames.size() ) log.warn("Process ${processDef.name} is missing emit names for one or more outputs -- unnamed outputs will be omitted") if( output.size() == 1 && outputNames.size() == 1 ) { @@ -113,16 +111,19 @@ class ProcessEntryHandler { } } - private void publishOutputs() { - final config = processDef.getProcessConfig() - final outputNames = getProcessOutputs(config) - final dsl = new OutputDsl() - for( final name : outputNames ) - dsl.declare(name, { -> }) + /** + * Creates an output definition that declares each output of the + * entry workflow, without creating an output directory. + */ + OutputDef createOutputDef() { // disable the output directory -- report output files by // their work directory path instead of publishing them session.outputDir = null - dsl.apply(session) + return new OutputDef({ -> + final dsl = (OutputDsl)(Object)getDelegate() + for( final name : getProcessOutputs(processDef.getProcessConfig()) ) + dsl.declare(name, { -> }) + }) } private List getProcessOutputs(ProcessConfig config) { @@ -327,7 +328,7 @@ class ProcessEntryHandler { // non-file inputs: a missing value is a hard error (required) if( value == null ) - throw new IllegalArgumentException("Parameter `--${name}` is required but no value was provided") + throw new IllegalArgumentException("Parameter `${name}` is required but no value was provided") // handle env, stdin inputs switch( decl ) { @@ -435,7 +436,7 @@ class ProcessEntryHandler { if( result == null ) { if( decl.isOptional() ) return null - throw new IllegalArgumentException("Parameter `--${name}` is required but no value was provided") + throw new IllegalArgumentException("Parameter `${name}` is required but no value was provided") } // report a value that could not be converted diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ScriptMeta.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ScriptMeta.groovy index 5c488f8a3c..eff651963b 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ScriptMeta.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ScriptMeta.groovy @@ -377,13 +377,16 @@ class ScriptMeta { addModule(get(script), name, alias) } - void addModule(ScriptMeta script, String name, String alias) { - assert script + void addModule(ScriptMeta meta, String name, String alias) { + assert meta assert name - // include a specific - def item = script.getComponent(name) + // the pipeline of the module takes precedence over a definition + // with the same name, matching the compiler + final item = name == 'workflow' && NF.isSyntaxParserV2() && meta.script.getEntryFlow() + ? meta.script.getEntryFlow() + : meta.getComponent(name) if( !item ) - throw new MissingModuleComponentException(script, name) + throw new MissingModuleComponentException(meta, name) addModule0(item, alias) } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/WorkflowBinding.groovy b/modules/nextflow/src/main/groovy/nextflow/script/WorkflowBinding.groovy index e8baccbe07..db41df30c2 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/WorkflowBinding.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/WorkflowBinding.groovy @@ -49,6 +49,8 @@ class WorkflowBinding extends Binding { private ScriptMeta meta + private Map published = new LinkedHashMap<>() + WorkflowBinding() { } WorkflowBinding(Map vars) { @@ -184,6 +186,9 @@ class WorkflowBinding extends Binding { } } + @PackageScope + Map getPublished() { published } + void _publish_(String name, Object source) { if( source instanceof ChannelOut ) { if( source.size() > 1 ) @@ -191,7 +196,7 @@ class WorkflowBinding extends Binding { source = source[0] } - owner.session.outputs[name] = + published[name] = source instanceof ChannelImpl ? source.getSource() : source instanceof ValueImpl ? source.getSource() : source instanceof DataflowWriteChannel ? source : diff --git a/modules/nextflow/src/main/groovy/nextflow/script/WorkflowDef.groovy b/modules/nextflow/src/main/groovy/nextflow/script/WorkflowDef.groovy index f40b4b0fd7..8599ed75c4 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/WorkflowDef.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/WorkflowDef.groovy @@ -106,8 +106,6 @@ class WorkflowDef extends BindableDef implements ChainableDef, IterableDef, Exec @PackageScope List getDeclaredOutputs() { declaredOutputs } - @PackageScope Map getDeclaredPublish() { declaredPublish } - @PackageScope List getDeclaredVariables() { new ArrayList(variableNames) } String getType() { 'workflow' } @@ -128,7 +126,9 @@ class WorkflowDef extends BindableDef implements ChainableDef, IterableDef, Exec final params = ChannelOut.spread(args) if( params.size() != declaredInputs.size() ) { final prefix = name ? "Workflow `$name`" : "Main workflow" - throw new IllegalArgumentException("$prefix declares ${declaredInputs.size()} input channels but ${params.size()} were given") + final expected = declaredInputs.size() + final actual = params.size() + throw new IllegalArgumentException("$prefix declares ${expected} ${expected == 1 ? 'input' : 'inputs'} but was called with ${actual} ${actual == 1 ? 'argument' : 'arguments'}") } // attach declared inputs with the invocation arguments @@ -204,15 +204,12 @@ class WorkflowDef extends BindableDef implements ChainableDef, IterableDef, Exec closure.setDelegate(binding) closure.setResolveStrategy(Closure.DELEGATE_FIRST) final result = closure.call() - if( name == null ) { - // return the last statement if entry workflow (used for testing) - return result - } - else { - // otherwise collect the outputs from the workflow binding - output = collectOutputs(declaredOutputs) - return output - } + // the outputs of an entry workflow are its published outputs + this.output = declaredOutputs + ? collectOutputs(declaredOutputs) + : new ChannelOut(binding.getPublished()) + // return the last statement if entry workflow (used for testing) + return name == null ? result : output } } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/WorkflowEntryHandler.groovy b/modules/nextflow/src/main/groovy/nextflow/script/WorkflowEntryHandler.groovy index b6522a6329..6f05206b7b 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/WorkflowEntryHandler.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/WorkflowEntryHandler.groovy @@ -16,23 +16,10 @@ package nextflow.script -import java.lang.reflect.ParameterizedType -import java.lang.reflect.Type -import java.nio.file.Path - -import groovy.json.JsonSlurper import groovy.transform.CompileStatic import groovy.util.logging.Slf4j -import groovy.yaml.YamlSlurper import nextflow.Session -import nextflow.dataflow.ChannelNamespace import nextflow.exception.ScriptRuntimeException -import nextflow.script.dsl.Types -import nextflow.script.types.Channel -import nextflow.script.types.Record -import nextflow.script.types.Value -import nextflow.splitter.CsvSplitter -import nextflow.util.TypeHelper /** * Helper class for named workflow execution. @@ -97,8 +84,7 @@ class WorkflowEntryHandler { // Execute the named workflow final output = workflowDef.run(inputs as Object[]) as ChannelOut // Publish workflow emits as pipeline outputs - assignOutputs(output) - publishOutputs() + assignOutputs((WorkflowBinding)(Object)getDelegate(), output) return output } final sourceCode = " // Auto-generated workflow entry\n ${workflowName}(...)" @@ -107,9 +93,8 @@ class WorkflowEntryHandler { return new WorkflowDef(script, entryBody) } - private void assignOutputs(ChannelOut output) { + private void assignOutputs(WorkflowBinding dsl, ChannelOut output) { final outputNames = workflowDef.getDeclaredOutputs() - final dsl = script.getBinding() if( output.size() == 1 && outputNames.size() == 1 ) { dsl._publish_(outputNames.first(), output[0]) } @@ -119,15 +104,20 @@ class WorkflowEntryHandler { } } - private void publishOutputs() { + /** + * Creates an output definition that declares each output of the + * entry workflow, without creating an output directory. + */ + OutputDef createOutputDef() { final outputNames = workflowDef.getDeclaredOutputs() - final dsl = new OutputDsl() - for( final name : outputNames ) - dsl.declare(name, { -> }) // disable the output directory -- report output files by // their work directory path instead of publishing them session.outputDir = null - dsl.apply(session) + return new OutputDef({ -> + final dsl = (OutputDsl)(Object)getDelegate() + for( final name : outputNames ) + dsl.declare(name, { -> }) + }) } /** @@ -141,152 +131,8 @@ class WorkflowEntryHandler { * @param workflowDef */ protected List getWorkflowArguments(WorkflowDef workflowDef) { - final inputs = workflowDef.getDeclaredInputs() - final inputNames = inputs*.name - final cliParams = session.cliParams ?: [:] - final configParams = session.configParams ?: [:] - - for( final name : cliParams.keySet() ) { - if( name !in inputNames && !configParams.containsKey(name) ) - throw new ScriptRuntimeException("Parameter `${name}` was specified on the command line but is not an input of workflow `${workflowDef.name}`") - } - - final arguments = [] - for( final decl : inputs ) { - final name = decl.name - final value = - cliParams.containsKey(name) ? resolveInput(decl, cliParams.get(name), true) : - configParams.containsKey(name) ? resolveInput(decl, configParams.get(name), false) : - null - - if( value == null && !decl.optional ) { - throw new ScriptRuntimeException("Parameter `--${name}` is required but no value was provided") - } - - arguments.add(value) - } - return arguments - } - - /** - * Resolves a single workflow input value against its declared type. - * - * A {@code Channel} input is loaded from a samplesheet file, with each - * record converted to the element type. A {@code Value} input is - * converted to {@code V} and wrapped in a value channel. Any other input is - * converted directly to the declared type. - * - * @param decl the declared input - * @param value the raw param value - * @param fromCli whether the value came from the command line (and is - * therefore a string that may need to be parsed) - */ - protected Object resolveInput(Param decl, Object value, boolean fromCli) { - if( value == null ) - return null - - final rawType = TypeHelper.getRawType(decl.type) - - if( rawType == Channel ) - return ChannelNamespace.fromList(loadChannelInput(decl, value)) - - if( rawType == Value ) - return ChannelNamespace.value(resolveScalarInput(elementDecl(decl), value, fromCli)) - - return resolveScalarInput(decl, value, fromCli) - } - - private Object resolveScalarInput(Param decl, Object value, boolean fromCli) { - final result = fromCli - ? ParamsHelper.resolveFromCli(decl, value) - : ParamsHelper.resolveFromCode(decl, value) - - final expectedType = TypeHelper.getRawType(decl.type) - final actualType = result?.getClass() - if( actualType != null && !ParamsHelper.isAssignableFrom(expectedType, actualType) ) - throw new ScriptRuntimeException("Workflow input `${decl.name}` with type ${Types.getName(decl.type)} cannot be assigned to ${result} [${Types.getName(actualType)}]") - - return result - } - - /** - * Loads a channel input from a samplesheet file, converting each record - * to the declared element type. - * - * @param decl - * @param value - */ - private List loadChannelInput(Param decl, Object value) { - if( value !instanceof CharSequence && value !instanceof Path ) - throw new ScriptRuntimeException("Workflow input `${decl.name}` with type ${Types.getName(decl.type)} should be a samplesheet file, but received: ${value} [${Types.getName(value.getClass())}]") - - final path = value instanceof Path - ? (Path)value - : TypeHelper.asPathType(value.toString()) - final elementType = elementDecl(decl).type - final elementRawType = TypeHelper.getRawType(elementType) - - if( !Map.isAssignableFrom(elementRawType) && !Record.isAssignableFrom(elementRawType) ) - throw new ScriptRuntimeException("Workflow input `${decl.name}` with type ${Types.getName(decl.type)} cannot be loaded from a samplesheet -- the element type should be Map, Record, or a record type") - - return loadFromFile(decl.name, path).collect { el -> - try { - TypeHelper.asType(el, elementType) - } - catch( Exception e ) { - throw new ScriptRuntimeException("Invalid record in samplesheet '${path}' for workflow input `${decl.name}` -- ${e.message}") - } - } - } - - /** - * Returns the declared input for the element type of a parameterized - * input type, e.g. {@code Sample} for {@code Channel}. - * - * @param decl - */ - private static Param elementDecl(Param decl) { - final elementType = decl.type instanceof ParameterizedType - ? ((ParameterizedType)decl.type).getActualTypeArguments()[0] - : (Type)Object - return new Param(decl.name, elementType, decl.optional, null) - } - - /** - * Loads the contents of a samplesheet file as a list of records. - * - * Supported formats: - * - CSV: header row required, comma-separated - * - JSON: must be a top-level array - * - YAML / YML: must be a top-level sequence - * - * @param name the input name (for error messages) - * @param file the samplesheet file to load - * @return a list of raw records (maps) - */ - protected List loadFromFile(String name, Path file) { - final ext = file.getExtension() - final value = switch( ext ) { - case 'csv' -> loadFromCsv(file) - case 'json' -> new JsonSlurper().parse(file) - case 'yaml', 'yml' -> new YamlSlurper().parse(file) - default -> throw new ScriptRuntimeException("Unrecognized file format '${ext}' for input file '${file}' for workflow input `${name}` -- should be CSV, JSON, or YAML") - } - if( value !instanceof List ) - throw new ScriptRuntimeException("Input file '${file}' for workflow input `${name}` must contain a list of records, but got: ${value.class.simpleName}") - return (List)value - } - - /** - * Loads a CSV samplesheet as a list of records. - * - * @param file - */ - private static List loadFromCsv(Path file) { - final rows = new CsvSplitter().options(header: true, sep: ',').target(file).list() - return rows.collect { row -> - ((Map)row).collectEntries { k, v -> [ k, v != '' ? v : null ] } - } + final params = ParamsHelper.resolveParams(workflowDef.getDeclaredInputs(), session.cliParams ?: [:], session.configParams ?: [:]) + return new ArrayList(params.values()) } } diff --git a/modules/nextflow/src/test/groovy/nextflow/config/parser/v2/ConfigParserV2Test.groovy b/modules/nextflow/src/test/groovy/nextflow/config/parser/v2/ConfigParserV2Test.groovy index 5538725449..c506cd63bd 100644 --- a/modules/nextflow/src/test/groovy/nextflow/config/parser/v2/ConfigParserV2Test.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/config/parser/v2/ConfigParserV2Test.groovy @@ -445,6 +445,41 @@ class ConfigParserV2Test extends Specification { slurper.getDeclaredParams() == [a: 3, b: 2] } + def 'should apply nested CLI params to nested config params' () { + given: + def cliParams = [rnaseq: [aligner: 'hisat2'], qc: [limit: '5']] + def text = ''' + params { + rnaseq.aligner = 'star' + rnaseq.fasta = 'genome.fa' + qc = [limit: 10, skip: false] + } + ''' + + when: + def slurper = new ConfigParserV2().setParams(cliParams) + def config = slurper.parse(text) + then: + config.params.rnaseq == [aligner: 'hisat2', fasta: 'genome.fa'] + config.params.qc == [limit: '5', skip: false] + slurper.getDeclaredParams() == [rnaseq: [aligner: 'hisat2', fasta: 'genome.fa'], qc: [limit: '5', skip: false]] + } + + def 'should not modify nested CLI params when assigning nested config params' () { + given: + def cliParams = [rnaseq: [aligner: 'hisat2']] + def text = ''' + params.rnaseq.skip_qc = false + params.rnaseq.skip_qc = true + ''' + + when: + def config = new ConfigParserV2().setParams(cliParams).parse(text) + then: + config.params.rnaseq == [aligner: 'hisat2', skip_qc: true] + cliParams == [rnaseq: [aligner: 'hisat2']] + } + def 'should ignore config includes when specified' () { given: def text = ''' diff --git a/modules/nextflow/src/test/groovy/nextflow/script/OutputDslTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/OutputDslTest.groovy index 13e5c50743..55c9506bac 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/OutputDslTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/OutputDslTest.groovy @@ -66,9 +66,10 @@ class OutputDslTest extends Specification { when: def session = Spy(createSession(config)) + def outputs = [:] - session.outputs.put('foo', Channel.of(file1)) - session.outputs.put('bar', Channel.of(file2)) + outputs.put('foo', Channel.of(file1)) + outputs.put('bar', Channel.of(file2)) def dsl = new OutputDsl() dsl.declare('foo') { @@ -82,7 +83,7 @@ class OutputDslTest extends Specification { path 'index.csv' } } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() dsl.getOutput() @@ -121,13 +122,14 @@ class OutputDslTest extends Specification { when: def session = Spy(createSession(config)) + def outputs = [:] - session.outputs.put('foo', Channel.of(file1)) + outputs.put('foo', Channel.of(file1)) def dsl = new OutputDsl() dsl.declare('foo') { } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() dsl.getOutput() @@ -159,15 +161,16 @@ class OutputDslTest extends Specification { when: def session = Spy(createSession(config)) + def outputs = [:] // the output directory is disabled when a named workflow is executed directly session.outputDir = null - session.outputs.put('foo', Channel.of(file1)) + outputs.put('foo', Channel.of(file1)) def dsl = new OutputDsl() dsl.declare('foo') { } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() def output = dsl.getOutput() @@ -220,13 +223,14 @@ class OutputDslTest extends Specification { when: def session = Spy(createSession(config)) + def outputs = [:] - session.outputs.put('foo', Channel.of(record)) + outputs.put('foo', Channel.of(record)) def dsl = new OutputDsl() dsl.declare('foo') { } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() dsl.getOutput() @@ -290,14 +294,15 @@ class OutputDslTest extends Specification { def 'should report error for invalid path directive' () { when: def session = createSession(outputDir: Path.of('results')) + def outputs = [:] - session.outputs.put('foo', Channel.of(1, 2, 3)) + outputs.put('foo', Channel.of(1, 2, 3)) def dsl = new OutputDsl() dsl.declare('foo') { path { v -> 42 } } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() dsl.getOutput() @@ -310,15 +315,16 @@ class OutputDslTest extends Specification { def 'should report error for invalid publish target' () { when: def session = createSession(outputDir: Path.of('results')) + def outputs = [:] def file = Path.of('output.txt') - session.outputs.put('foo', Channel.of([file, file, file])) + outputs.put('foo', Channel.of([file, file, file])) def dsl = new OutputDsl() dsl.declare('foo') { path { files -> publish(files, 'foo.txt') } } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() dsl.getOutput() @@ -330,14 +336,15 @@ class OutputDslTest extends Specification { def 'should report error for invalid publish source' () { when: def session = createSession(outputDir: Path.of('results')) + def outputs = [:] - session.outputs.put('foo', Channel.of(42)) + outputs.put('foo', Channel.of(42)) def dsl = new OutputDsl() dsl.declare('foo') { path { v -> publish(v, 'foo') } } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() dsl.getOutput() @@ -350,8 +357,9 @@ class OutputDslTest extends Specification { def 'should report error for invalid index file extension' () { when: def session = createSession(outputDir: Path.of('results')) + def outputs = [:] - session.outputs.put('foo', Channel.empty()) + outputs.put('foo', Channel.empty()) def dsl = new OutputDsl() dsl.declare('foo') { @@ -359,7 +367,7 @@ class OutputDslTest extends Specification { path 'index.txt' } } - dsl.apply(session) + dsl.apply(session, outputs) session.fireDataflowNetwork() dsl.getOutput() diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy index bf83ef2c5b..905ad710f6 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy @@ -78,6 +78,27 @@ class ParamsDslTest extends Specification { noExceptionThrown() } + def 'should report error for missing required record param'() { + when: + runScript( + '''\ + params { + sample: Sample + } + + record Sample { + id: String? + } + + workflow { params } + ''', + params: [:] + ) + then: + def e = thrown(ScriptRuntimeException) + e.message == 'Parameter `sample` is required but no value was provided' + } + def 'should report error for missing required param'() { when: runScript( @@ -238,6 +259,45 @@ class ParamsDslTest extends Specification { result[2].id == 3 } + def 'should load dataflow params from the command line'() { + given: + def samplesheet = Files.createTempFile('test', '.csv') + samplesheet.text = 'id,count\na,1\nb,2\n' + def cliParams = [samples: samplesheet.toString(), limit: '5'] + + when: + def result = runScript( + '''\ + nextflow.enable.types = true + + params { + samples: Channel + limit: Value + } + + record Sample { + id: String + count: Integer + } + + workflow { + params.samples + .map { s -> s.count } + .collect() + .combine(params.limit) + } + ''', + params: cliParams + ) + then: + def (counts, limit) = result.val + counts.toSorted() == [1, 2] + limit == 5 + + cleanup: + samplesheet?.delete() + } + def 'should validate record param from nested map'() { when: 'a script is invoked as `nextflow run module.nf --sample.id a --sample.greeting hola`' def result = runScript( @@ -263,6 +323,52 @@ class ParamsDslTest extends Specification { result.greeting == 'hola' } + def 'should report the field of a record param that cannot be converted'() { + when: + runScript( + '''\ + params { + sample: Sample + } + + record Sample { + id: String + paired: Boolean + } + + workflow { + params.sample + } + ''', + params: [sample: [id: 'a', paired: 'yes']] + ) + then: + def e = thrown(ScriptRuntimeException) + e.message == 'Parameter `sample.paired` with type Boolean cannot be assigned to yes [String]' + + when: + runScript( + '''\ + params { + sample: Sample + } + + record Sample { + id: String + paired: Boolean + } + + workflow { + params.sample + } + ''', + params: [sample: [paired: 'true']] + ) + then: + e = thrown(ScriptRuntimeException) + e.message == 'Parameter `sample` with type Sample is missing required field `id`' + } + def 'should report error for invalid record type'() { when: runScript( diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy index 0959b3849c..927fa03088 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy @@ -19,6 +19,7 @@ package nextflow.script import java.nio.file.Files import java.nio.file.Path +import nextflow.exception.ScriptRuntimeException import nextflow.util.Duration import nextflow.util.MemoryUnit import nextflow.util.RecordMap @@ -177,6 +178,63 @@ class ParamsHelperTest extends Specification { SampleRec | Map | false } + def 'should load records from a samplesheet'() { + given: + def file = Files.createTempFile('test', ".${EXT}") + file.text = TEXT + + when: + def result = ParamsHelper.loadFromFile('samples', file.toAbsolutePath()) + + then: + result == EXPECTED + + cleanup: + file?.delete() + + where: + EXT | TEXT | EXPECTED + // CSV has no types, so every value is a string + 'csv' | 'id,name\n1,sample1\n2,sample2\n' | [[id: '1', name: 'sample1'], [id: '2', name: 'sample2']] + // quoted values, as written by an output index file + 'csv' | '"id","name"\n"1","sample1"\n"2",""\n' | [[id: '1', name: 'sample1'], [id: '2', name: null]] + 'csv' | 'id,name\n1,"sample 1, rep 1"\n' | [[id: '1', name: 'sample 1, rep 1']] + 'json' | '[{"id":1,"name":"s1"},{"id":2,"name":"s2"}]' | [[id: 1, name: 's1'], [id: 2, name: 's2']] + 'yml' | '- id: 1\n name: s1\n- id: 2\n name: s2\n' | [[id: 1, name: 's1'], [id: 2, name: 's2']] + } + + def 'should throw for unrecognized samplesheet format'() { + given: + def txtFile = Files.createTempFile('test', '.txt') + txtFile.text = 'some text' + + when: + ParamsHelper.loadFromFile('items', txtFile.toAbsolutePath()) + + then: + def e = thrown(ScriptRuntimeException) + e.message.contains("Unrecognized file format 'txt'") + + cleanup: + txtFile?.delete() + } + + def 'should throw for a JSON file whose top level is not a list'() { + given: + def jsonFile = Files.createTempFile('test', '.json') + jsonFile.text = '{"key":"value"}' // object, not array + + when: + ParamsHelper.loadFromFile('samples', jsonFile.toAbsolutePath()) + + then: + def e = thrown(ScriptRuntimeException) + e.message.contains('must contain a list of records') + + cleanup: + jsonFile?.delete() + } + static class SampleRec implements nextflow.script.types.Record { String id } diff --git a/modules/nextflow/src/test/groovy/nextflow/script/PipelineCompositionTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/PipelineCompositionTest.groovy new file mode 100644 index 0000000000..0871f368d6 --- /dev/null +++ b/modules/nextflow/src/test/groovy/nextflow/script/PipelineCompositionTest.groovy @@ -0,0 +1,642 @@ +/* + * Copyright 2013-2026, Seqera Labs + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package nextflow.script + +import java.nio.file.Files +import java.nio.file.Path + +import nextflow.exception.ScriptRuntimeException +import spock.lang.Timeout +import spock.lang.Unroll +import test.Dsl2Spec + +import static test.ScriptHelper.* + +/** + * Tests for pipeline composition -- a pipeline included as a named workflow. + * + * @author Ben Sherman + */ +@Timeout(30) +class PipelineCompositionTest extends Dsl2Spec { + + private Path folder + + def setup() { + folder = Files.createTempDirectory('test') + } + + def cleanup() { + folder?.deleteDir() + } + + private static final String GREET = ''' + params { + names: Channel + greeting: String = 'Hello' + } + + workflow { + main: + messages = params.names.map { name -> "${params.greeting}, ${name}!" } + + publish: + messages = messages + } + + output { + messages: Channel {} + } + ''' + + private static final String COUNT = ''' + params { + samples: Channel + factor: Value = 1 + } + + record Sample { + id: String + count: Integer + } + + workflow { + main: + totals = params.samples + .map { s -> s.count } + .collect() + .combine(params.factor) + .map { counts, factor -> counts.sum() * factor } + + publish: + total = totals + } + + output { + total: Channel {} + } + ''' + + /** + * Write typed scripts to the test folder and return the path of `main.nf`. + */ + private Path write(Map files) { + files.each { name, text -> + folder.resolve(name).text = 'nextflow.enable.types = true\n' + text.stripIndent() + } + return folder.resolve('main.nf') + } + + /** + * Write the greet pipeline and a main script. + */ + private Path pipeline(String text) { + return write([ + 'greet.nf': GREET, + 'main.nf': text + ]) + } + + @Unroll + def 'should call an included pipeline like a named workflow: #CALL' () { + given: + def script = pipeline(""" + include { workflow as GREET } from './greet.nf' + + workflow { + main: + ${CALL} + } + """) + + when: + def result = runScript(script) + then: + EXPECTED.collect { result.val } == EXPECTED + + where: + CALL | EXPECTED + "GREET( record(names: channel.of('World', 'Nextflow')) )" | ['Hello, World!', 'Hello, Nextflow!'] + "GREET( record(names: channel.of('World'), greeting: 'Hola') )" | ['Hola, World!'] + "GREET( record(greeting: 'Ciao') + record(names: channel.of('World')) )" | ['Ciao, World!'] + } + + def 'should resolve the params of a pipeline called with the session params' () { + given: + def script = write([ + 'hello.nf': ''' + params { + name: String + greeting: String = 'Hello' + } + + workflow { + main: + message = channel.value("${params.greeting}, ${params.name}!") + + publish: + message = message + } + + output { + message: Value {} + } + ''', + 'main.nf': ''' + include { workflow as HELLO } from './hello.nf' + + params { + name: String + } + + workflow { + main: + HELLO( params ) + } + ''' + ]) + + when: + def result = runScript(script, params: [name: 'World']) + then: + result.val == 'Hello, World!' + } + + @Unroll + def 'should fail on an invalid pipeline call: #CALL' () { + given: + def script = pipeline(""" + include { workflow as GREET } from './greet.nf' + + workflow { + main: + ${CALL} + } + """) + + when: + runScript(script) + then: + def e = thrown(ScriptRuntimeException) + e.message == ERROR + + where: + CALL | ERROR + "GREET( record(greeting: 'Hola') )" | 'Parameter `names` of pipeline `GREET` is required but no value was provided' + "GREET( record(names: channel.of('World'), greeting: null) )" | 'Parameter `greeting` of pipeline `GREET` is required but no value was provided' + "GREET( record(names: channel.of('World'), foo: 'bar') )" | 'Pipeline `GREET` does not declare a parameter named `foo`' + "GREET( record(names: channel.of('World'), foo: null) )" | 'Pipeline `GREET` does not declare a parameter named `foo`' + "GREET( names: channel.of('World') )" | 'Pipeline `GREET` should be called with a record' + "GREET( [names: channel.of('World')] )" | 'Pipeline `GREET` should be called with a record' + } + + @Unroll + def 'should resolve an included params record from the command line and config' () { + given: + def script = pipeline(''' + include { params as GreetParams ; workflow as GREET } from './greet.nf' + + params { + greet: GreetParams + } + + workflow { + main: + greet = GREET( params.greet + record(names: channel.of('World')) ) + greet + } + ''') + + when: + def result = runScript([params: CLI, configParams: CONFIG], script) + then: + result.val == RESULT + + where: + CLI | CONFIG | RESULT + [greet: [greeting: 'Hola']] | [:] | 'Hola, World!' + [:] | [greet: [greeting: 'Ciao']] | 'Ciao, World!' + // the pipeline applies its own default for an omitted param + [greet: [:]] | [:] | 'Hello, World!' + // a params record with no required fields is not required at launch + [:] | [:] | 'Hello, World!' + } + + @Unroll + def 'should accept dataflow values for Channel and Value params: #CALL' () { + given: + def script = write([ + 'count.nf': COUNT, + 'main.nf': """ + include { workflow as COUNT } from './count.nf' + + workflow { + main: + samples = channel.of( record(id: 'a', count: 1), record(id: 'b', count: 2) ) + ${CALL} + } + """ + ]) + + when: + def result = runScript(script) + then: + result.val == RESULT + + where: + CALL | RESULT + "COUNT( record(samples: samples, factor: channel.value(2)) )" | 6 + // the default of a Value param is wrapped in a dataflow value + "COUNT( record(samples: samples) )" | 3 + } + + def 'should load the dataflow params of an included pipeline from the command line' () { + given: + folder.resolve('samples.csv').text = 'id,count\na,1\nb,2\n' + def script = write([ + 'count.nf': COUNT, + 'main.nf': ''' + include { params as CountParams ; workflow as COUNT } from './count.nf' + + params { + count: CountParams + } + + workflow { + main: + COUNT( params.count ) + } + ''' + ]) + def cliParams = [count: [samples: folder.resolve('samples.csv').toString(), factor: '5']] + + when: + def result = runScript([params: cliParams], script) + then: + result.val == 15 + } + + def 'should scope the processes of an included pipeline by its name' () { + given: + def script = write([ + 'greet.nf': ''' + params { + names: Channel + } + + workflow { + main: + messages = FOO(params.names) + + publish: + messages = messages + } + + output { + messages: Channel {} + } + + process FOO { + input: + name: String + + output: + message: String + + exec: + message = "${task.ext.greeting}, ${name}!" + } + ''', + 'main.nf': ''' + include { workflow as GREET } from './greet.nf' + + workflow { + main: + greet = GREET( record(names: channel.of('World')) ) + greet + } + ''' + ]) + + when: + def result = runScript([config: [process: ['withName:GREET:FOO': [ext: [greeting: 'Bonjour']]]]], script) + then: + result.val == 'Bonjour, World!' + } + + def 'should publish the outputs of the calling pipeline only' () { + given: + def script = pipeline(''' + include { workflow as GREET } from './greet.nf' + + workflow { + main: + greet = GREET( record(names: channel.of('World')) ) + + publish: + out = greet + } + + output { + out: Channel {} + } + ''') + + when: + runScript(script) + then: + noExceptionThrown() + } + + def 'should not expose the params of the calling pipeline to an included pipeline' () { + given: + def script = write([ + 'greet.nf': ''' + params { + names: Channel + } + + workflow { + main: + messages = params.names.map { name -> "${params.secret}, ${name}!" } + + publish: + messages = messages + } + + output { + messages: Channel {} + } + ''', + 'main.nf': ''' + include { workflow as GREET } from './greet.nf' + + params { + secret: String = 'LEAKED' + } + + workflow { + main: + greet = GREET( record(names: channel.of('World')) ) + greet + } + ''' + ]) + + when: + def result = runScript(script) + then: + // a param that the included pipeline does not declare is not visible + // to its entry workflow, even though the calling pipeline declares it + result.val == 'null, World!' + } + + def 'should call an included pipeline without a params block' () { + given: + def script = write([ + 'hello.nf': ''' + workflow { + main: + messages = channel.of('Hello') + + publish: + messages = messages + } + + output { + messages: Channel {} + } + ''', + 'main.nf': ''' + include { workflow as HELLO } from './hello.nf' + + params { + greeting: String = 'Hola' + } + + workflow { + main: + HELLO() + } + ''' + ]) + + when: + def result = runScript(script) + then: + result.val == 'Hello' + } + + def 'should support multiple aliases of the same pipeline' () { + given: + def script = pipeline(''' + include { workflow as GREET } from './greet.nf' + include { workflow as GREET_AGAIN } from './greet.nf' + + workflow { + main: + a = GREET( record(names: channel.of('World'), greeting: 'Hola') ) + b = GREET_AGAIN( record(names: channel.of('Nextflow'), greeting: 'Ciao') ) + a.mix(b) + } + ''') + + when: + def result = runScript(script) + then: + // each alias resolves its own params, because the entry workflow + // receives them as an input instead of reading global state + [result.val, result.val].sort() == ['Ciao, Nextflow!', 'Hola, World!'] + } + + def 'should support calling the same alias more than once when the pipeline has no process' () { + given: + def script = pipeline(''' + include { workflow as GREET } from './greet.nf' + + workflow { + main: + a = GREET( record(names: channel.of('World'), greeting: 'Hola') ) + b = GREET( record(names: channel.of('Nextflow'), greeting: 'Ciao') ) + a.mix(b) + } + ''') + + when: + def result = runScript(script) + then: + // a pipeline that contains a process can be called only once per + // alias, like a named workflow -- see the spec below + [result.val, result.val].sort() == ['Ciao, Nextflow!', 'Hola, World!'] + } + + def 'should reject calling the same alias more than once when the pipeline has a process' () { + given: + def script = write([ + 'greet.nf': ''' + params { + greeting: String = 'Hello' + } + + process SAY { + input: + message: String + output: + stdout() + script: + "echo '${message}'" + } + + workflow { + main: + said = SAY( params.greeting ) + + publish: + said = said + } + + output { + said: Channel {} + } + ''', + 'main.nf': ''' + include { workflow as GREET } from './greet.nf' + + workflow { + main: + GREET( record(greeting: 'Hola') ) + GREET( record(greeting: 'Ciao') ) + } + ''' + ]) + + when: + runScript(script) + then: + def e = thrown(Exception) + e.message.contains("Process 'GREET:SAY' was called twice") + } + + def 'should support including the same params block from two scripts' () { + given: + def script = pipeline(''' + include { params as GreetParams ; workflow as GREET } from './greet.nf' + include { params as MidParams ; workflow as MID } from './mid.nf' + + workflow { + main: + MID( record(names: channel.of('World')) ) + } + ''') + write([ + 'mid.nf': ''' + include { params as GreetParams ; workflow as GREET } from './greet.nf' + + params { + names: Channel + } + + workflow { + main: + messages = GREET( record(names: params.names) ) + + publish: + messages = messages + } + + output { + messages: Channel {} + } + ''' + ]) + + when: + // the record type of an included params block is qualified by the + // including script, so two scripts can include the same one + def result = runScript(script) + then: + result.val == 'Hello, World!' + } + + def 'should execute a pipeline module only once when it is included by two scripts' () { + given: + def script = pipeline(''' + include { workflow as GREET } from './greet.nf' + include { workflow as MID } from './mid.nf' + + workflow { + main: + a = GREET( record(names: channel.of('World'), greeting: 'Hola') ) + b = MID( record(names: channel.of('Nextflow')) ) + a.mix(b) + } + ''') + write([ + 'mid.nf': ''' + include { workflow as GREET_INNER } from './greet.nf' + + params { + names: Channel + } + + workflow { + main: + messages = GREET_INNER( record(names: params.names, greeting: 'Ciao') ) + + publish: + messages = messages + } + + output { + messages: Channel {} + } + ''' + ]) + + when: + def result = runScript(script) + then: + [result.val, result.val].sort() == ['Ciao, Nextflow!', 'Hola, World!'] + } + + def 'should include a definition named output' () { + given: + def script = write([ + 'lib.nf': ''' + def output(value: String) -> String { + return value.toUpperCase() + } + ''', + 'main.nf': ''' + include { output } from './lib.nf' + + workflow { + main: + channel.of('World').map { name -> output(name) } + } + ''' + ]) + + when: + // a definition that happens to be named `output` is not the output + // block of a pipeline, so it needs no alias + def result = runScript(script) + then: + result.val == 'WORLD' + } + +} diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ProcessEntryHandlerTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ProcessEntryHandlerTest.groovy index 4456a16d0e..7371e603e5 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ProcessEntryHandlerTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ProcessEntryHandlerTest.groovy @@ -235,7 +235,7 @@ class ProcessEntryHandlerTest extends Specification { then: def e = thrown(IllegalArgumentException) - e.message == 'Parameter `--id` is required but no value was provided' + e.message == 'Parameter `id` is required but no value was provided' } def 'should resolve a typed input like a pipeline parameter (v2)' () { @@ -295,6 +295,6 @@ class ProcessEntryHandlerTest extends Specification { then: def e = thrown(IllegalArgumentException) - e.message == 'Parameter `--reads` is required but no value was provided' + e.message == 'Parameter `reads` is required but no value was provided' } } diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ScriptProcessRunTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ScriptProcessRunTest.groovy index 77bc950cd3..87b79224d3 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ScriptProcessRunTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ScriptProcessRunTest.groovy @@ -139,7 +139,7 @@ class ScriptProcessRunTest extends Dsl2Spec { then: def e = thrown(Exception) - e.message.contains('Parameter `--requiredParam` is required but no value was provided') + e.message.contains('Parameter `requiredParam` is required but no value was provided') } def 'should cast boolean parameter to boolean' () { diff --git a/modules/nextflow/src/test/groovy/nextflow/script/WorkflowEntryHandlerTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/WorkflowEntryHandlerTest.groovy index 6af62d74b5..f530ced2e2 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/WorkflowEntryHandlerTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/WorkflowEntryHandlerTest.groovy @@ -38,8 +38,6 @@ import static test.ScriptHelper.* @Timeout(10) class WorkflowEntryHandlerTest extends Dsl2Spec { - // ── unit: loadFromFile ──────────────────────────────────────────────────── - private WorkflowEntryHandler makeHandler(List inputs = [], List workflowNames = ['HELLO']) { def workflowDef = Mock(WorkflowDef) { getName() >> workflowNames.first() @@ -56,60 +54,6 @@ class WorkflowEntryHandlerTest extends Dsl2Spec { return new WorkflowEntryHandler(script, session, meta) } - def 'should load records from a samplesheet'() { - given: - def file = Files.createTempFile('test', ".${EXT}") - file.text = TEXT - - when: - def result = makeHandler().loadFromFile('samples', file.toAbsolutePath()) - - then: - result == EXPECTED - - cleanup: - file?.delete() - - where: - EXT | TEXT | EXPECTED - // CSV has no types, so every value is a string - 'csv' | 'id,name\n1,sample1\n2,sample2\n' | [[id: '1', name: 'sample1'], [id: '2', name: 'sample2']] - 'json' | '[{"id":1,"name":"s1"},{"id":2,"name":"s2"}]' | [[id: 1, name: 's1'], [id: 2, name: 's2']] - 'yml' | '- id: 1\n name: s1\n- id: 2\n name: s2\n' | [[id: 1, name: 's1'], [id: 2, name: 's2']] - } - - def 'should throw for unrecognized samplesheet format'() { - given: - def txtFile = Files.createTempFile('test', '.txt') - txtFile.text = 'some text' - - when: - makeHandler().loadFromFile('items', txtFile.toAbsolutePath()) - - then: - def e = thrown(ScriptRuntimeException) - e.message.contains("Unrecognized file format 'txt'") - - cleanup: - txtFile?.delete() - } - - def 'should throw for a JSON file whose top level is not a list'() { - given: - def jsonFile = Files.createTempFile('test', '.json') - jsonFile.text = '{"key":"value"}' // object, not array - - when: - makeHandler().loadFromFile('samples', jsonFile.toAbsolutePath()) - - then: - def e = thrown(ScriptRuntimeException) - e.message.contains('must contain a list of records') - - cleanup: - jsonFile?.delete() - } - def 'should throw error when multiple workflows are defined'() { when: makeHandler([], ['FIRST', 'SECOND']) @@ -237,7 +181,7 @@ class WorkflowEntryHandlerTest extends Dsl2Spec { then: def e = thrown(ScriptRuntimeException) - e.message.contains('Parameter `--name` is required but no value was provided') + e.message.contains('Parameter `name` is required but no value was provided') } def 'should convert samplesheet records to the declared element type'() { @@ -458,7 +402,7 @@ class WorkflowEntryHandlerTest extends Dsl2Spec { then: def e = thrown(ScriptRuntimeException) - e.message.contains('Parameter `bogus` was specified on the command line but is not an input of workflow `GREET`') + e.message.contains('Parameter `bogus` was specified on the command line or params file but is not declared in the script or config') } def 'should prefer explicit entry workflow over named workflow'() { @@ -593,7 +537,7 @@ class WorkflowEntryHandlerTest extends Dsl2Spec { then: def e = thrown(ScriptRuntimeException) - e.message.contains('Parameter `--name` is required but no value was provided') + e.message.contains('Parameter `name` is required but no value was provided') } def 'should pass a null param value to a nullable input'() { @@ -676,7 +620,7 @@ class WorkflowEntryHandlerTest extends Dsl2Spec { then: 'the param and file are named, rather than a bare NumberFormatException' def e = thrown(ScriptRuntimeException) e.message.contains('Invalid record in samplesheet') - e.message.contains('workflow input `samples`') + e.message.contains('parameter `samples`') cleanup: file?.delete() diff --git a/modules/nextflow/src/testFixtures/groovy/test/ScriptHelper.groovy b/modules/nextflow/src/testFixtures/groovy/test/ScriptHelper.groovy index e2540a2579..49d854cf6a 100644 --- a/modules/nextflow/src/testFixtures/groovy/test/ScriptHelper.groovy +++ b/modules/nextflow/src/testFixtures/groovy/test/ScriptHelper.groovy @@ -202,7 +202,7 @@ class ScriptHelper { def session = opts.config ? new MockSession(opts.config) : new MockSession() session.setBinding(new ScriptBinding()) - session.init( new ScriptFile(path), null, opts.params, null ) + session.init( new ScriptFile(path), null, opts.params, opts.configParams ) if( opts.moduleRun ) session.setModuleRun(true) session.start() diff --git a/modules/nf-lang/src/main/java/nextflow/script/ast/ASTNodeMarker.java b/modules/nf-lang/src/main/java/nextflow/script/ast/ASTNodeMarker.java index 965f98ab2d..239c835438 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/ast/ASTNodeMarker.java +++ b/modules/nf-lang/src/main/java/nextflow/script/ast/ASTNodeMarker.java @@ -54,9 +54,16 @@ public enum ASTNodeMarker { // the MethodNode targeted by a variable expression (PropertyNode) METHOD_VARIABLE_TARGET, + // the Parameter targeted by a named argument (MapEntryExpression) + NAMED_PARAM, + // denotes a nullable type annotation (ClassNode) NULLABLE, + // the ScriptNode that declares an entry workflow (WorkflowNode), so that + // the params and output blocks of a pipeline can be resolved from a call + PIPELINE_SCRIPT, + // the FieldNode targeted by a PropertyExpression PROPERTY_TARGET, diff --git a/modules/nf-lang/src/main/java/nextflow/script/ast/ScriptNode.java b/modules/nf-lang/src/main/java/nextflow/script/ast/ScriptNode.java index b46151b823..6af928fdcf 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/ast/ScriptNode.java +++ b/modules/nf-lang/src/main/java/nextflow/script/ast/ScriptNode.java @@ -18,7 +18,10 @@ import java.util.ArrayList; import java.util.List; +import nextflow.script.dsl.PipelineParams; import org.codehaus.groovy.ast.ASTNode; +import org.codehaus.groovy.ast.AnnotationNode; +import org.codehaus.groovy.ast.ClassHelper; import org.codehaus.groovy.ast.ClassNode; import org.codehaus.groovy.ast.ModuleNode; import org.codehaus.groovy.ast.expr.ConstantExpression; @@ -143,6 +146,33 @@ public void addParamV1(ParamNodeV1 paramNode) { public void setEntry(WorkflowNode entry) { this.entry = entry; + entry.putNodeMetaData(ASTNodeMarker.PIPELINE_SCRIPT, this); + } + + /** + * Get the script that declares an entry workflow, i.e. the pipeline that + * the workflow belongs to, or null if the workflow is not an entry. + * + * @param node + */ + public static ScriptNode getPipeline(WorkflowNode node) { + return (ScriptNode) node.getNodeMetaData(ASTNodeMarker.PIPELINE_SCRIPT); + } + + private static final ClassNode PIPELINE_PARAMS = ClassHelper.makeCached(PipelineParams.class); + + /** + * Returns true if a record type was synthesized from the params + * block of an included pipeline. + * + * @param node + */ + public static boolean isPipelineParams(ClassNode node) { + return !node.getAnnotations(PIPELINE_PARAMS).isEmpty(); + } + + public static void setPipelineParams(ClassNode node) { + node.addAnnotation(new AnnotationNode(PIPELINE_PARAMS)); } public void setOutputs(OutputBlockNode outputs) { diff --git a/modules/nf-lang/src/main/java/nextflow/script/control/CallArityVisitor.java b/modules/nf-lang/src/main/java/nextflow/script/control/CallArityVisitor.java index b69c671d81..5e65bfb143 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/control/CallArityVisitor.java +++ b/modules/nf-lang/src/main/java/nextflow/script/control/CallArityVisitor.java @@ -64,6 +64,9 @@ public void visitMethodCallExpression(MethodCallExpression node) { } private void checkMethodCallArguments(MethodCallExpression node, MethodNode defNode) { + // a pipeline call is checked by the type checker + if( defNode instanceof WorkflowNode wn && wn.isEntry() ) + return; var argsCount = asMethodCallArguments(node).size(); var paramsCount = defNode.getParameters().length; if( argsCount != paramsCount ) diff --git a/modules/nf-lang/src/main/java/nextflow/script/control/ResolveIncludeVisitor.java b/modules/nf-lang/src/main/java/nextflow/script/control/ResolveIncludeVisitor.java index 6f7a206552..b440d4c9b9 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/control/ResolveIncludeVisitor.java +++ b/modules/nf-lang/src/main/java/nextflow/script/control/ResolveIncludeVisitor.java @@ -15,6 +15,7 @@ */ package nextflow.script.control; +import java.lang.reflect.Modifier; import java.net.URI; import java.nio.file.Path; import java.util.ArrayList; @@ -22,12 +23,16 @@ import java.util.Set; import nextflow.script.ast.FunctionNode; +import nextflow.script.ast.IncludeEntryNode; import nextflow.script.ast.IncludeNode; import nextflow.script.ast.ScriptNode; +import nextflow.script.ast.RecordNode; import nextflow.script.ast.ScriptVisitorSupport; +import nextflow.script.ast.WorkflowNode; import org.codehaus.groovy.ast.ASTNode; import org.codehaus.groovy.ast.AnnotatedNode; import org.codehaus.groovy.ast.ClassNode; +import org.codehaus.groovy.ast.FieldNode; import org.codehaus.groovy.ast.MethodNode; import org.codehaus.groovy.control.SourceUnit; import org.codehaus.groovy.control.messages.SyntaxErrorMessage; @@ -111,18 +116,45 @@ public void visitInclude(IncludeNode node) { addError("Module could not be parsed: '" + includeUri.getPath() + "'", node); return; } + var scriptNode = (ScriptNode) includeUnit.getAST(); var definitions = getDefinitions(includeUri); + var hasPipeline = false; for( var entry : node.entries ) { var includedName = entry.name; - var includedNode = definitions.stream() + // a `params` entry that doesn't match a definition of the module + // refers to the params block of the pipeline + var target = definitions.stream() .filter(defNode -> includedName.equals(definitionName(defNode))) - .findFirst(); - if( !includedNode.isPresent() ) { + .findFirst() + .orElseGet(() -> paramsBlockType(scriptNode, entry)); + if( target == null ) { addError("Included name '" + includedName + "' is not defined in module '" + includeUri.getPath() + "'", node); continue; } - entry.setTarget(includedNode.get()); + hasPipeline |= target instanceof WorkflowNode wn && wn.isEntry() + || target instanceof ClassNode cn && ScriptNode.isPipelineParams(cn); + entry.setTarget(target); } + if( hasPipeline && !((ScriptNode) sourceUnit.getAST()).isTypingEnabled() ) + addError("Including a pipeline requires `nextflow.enable.types = true` in the including script", node); + if( hasPipeline && !scriptNode.isTypingEnabled() ) + addError("An included pipeline must enable static typing -- set `nextflow.enable.types = true` in '" + includeUri.getPath() + "'", node); + } + + /** + * Synthesize a record type from the params block of an included + * pipeline. The type checker makes it partial (all fields nullable). + */ + private static ClassNode paramsBlockType(ScriptNode sn, IncludeEntryNode entry) { + var block = sn.getParams(); + if( !"params".equals(entry.name) || block == null ) + return null; + var cn = new RecordNode(entry.getNameOrAlias()); + ScriptNode.setPipelineParams(cn); + for( var declaration : block.declarations ) { + cn.addField(new FieldNode(declaration.getName(), Modifier.PUBLIC, declaration.getType(), cn, null)); + } + return cn; } private static void setPlaceholderTargets(IncludeNode node) { @@ -155,7 +187,14 @@ private List getDefinitions(URI uri) { return result; } + /** + * An entire pipeline -- the `params` / `workflow` / `output` trio of a + * script -- can be included as a named workflow, using the `workflow` + * keyword to refer to the entry workflow of the included script. + */ private static String definitionName(AnnotatedNode node) { + if( node instanceof WorkflowNode wn && wn.isEntry() ) + return "workflow"; return node instanceof ClassNode cn ? cn.getNameWithoutPackage() : node instanceof MethodNode mn ? mn.getName() : diff --git a/modules/nf-lang/src/main/java/nextflow/script/control/ScriptToGroovyVisitor.java b/modules/nf-lang/src/main/java/nextflow/script/control/ScriptToGroovyVisitor.java index d8c9d5461c..7ef81627b9 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/control/ScriptToGroovyVisitor.java +++ b/modules/nf-lang/src/main/java/nextflow/script/control/ScriptToGroovyVisitor.java @@ -16,11 +16,10 @@ package nextflow.script.control; import java.lang.reflect.Modifier; +import java.util.ArrayList; import java.util.Arrays; import java.util.Comparator; -import java.util.List; import java.util.Set; -import java.util.stream.Collectors; import nextflow.script.ast.ASTNodeMarker; import nextflow.script.ast.AgentNode; @@ -145,15 +144,26 @@ public void visitFeatureFlag(FeatureFlagNode node) { @Override public void visitInclude(IncludeNode node) { - var entries = (List) node.entries.stream() - .map((entry) -> { - var name = constX(entry.name); - return entry.alias != null - ? createX("nextflow.script.IncludeDef.Module", args(name, constX(entry.alias))) - : createX("nextflow.script.IncludeDef.Module", args(name)); - }) - .collect(Collectors.toList()); + // an included params block is only a type -- add it to this + // script so that it is compiled, and don't include it at runtime + var entries = new ArrayList(); + for( var entry : node.entries ) { + if( entry.getTarget() instanceof ClassNode cn && ScriptNode.isPipelineParams(cn) ) { + // the type is qualified by the including script, so that two + // scripts can include the same block + cn.setName(sgh.packageName(moduleNode) + "." + cn.getName()); + addNullableAnnotations(cn); + moduleNode.addClass(cn); + continue; + } + var name = constX(entry.name); + entries.add(entry.alias != null + ? createX("nextflow.script.IncludeDef.Module", args(name, constX(entry.alias))) + : createX("nextflow.script.IncludeDef.Module", args(name))); + } + if( entries.isEmpty() ) + return; var include = callThisX("include", args(createX("nextflow.script.IncludeDef", args(listX(entries))))); var from = callX(include, "from", args(node.source)); var result = stmt(callX(from, "load0", args(varX("params")))); @@ -255,13 +265,16 @@ public void visitOutputs(OutputBlockNode node) { @Override public void visitRecord(RecordNode node) { + addNullableAnnotations(node); + var result = stmt(callThisX("declareType", args(classX(node)))); + moduleNode.addStatement(result); + } + + private static void addNullableAnnotations(ClassNode node) { for( var fn : node.getFields() ) { if( fn.getType().getNodeMetaData(ASTNodeMarker.NULLABLE) != null ) fn.addAnnotation(NULLABLE); } - - var result = stmt(callThisX("declareType", args(classX(node)))); - moduleNode.addStatement(result); } @Override diff --git a/modules/nf-lang/src/main/java/nextflow/script/control/TypeCheckingVisitor.java b/modules/nf-lang/src/main/java/nextflow/script/control/TypeCheckingVisitor.java index e5334c4060..b2472d0b73 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/control/TypeCheckingVisitor.java +++ b/modules/nf-lang/src/main/java/nextflow/script/control/TypeCheckingVisitor.java @@ -21,6 +21,7 @@ import java.util.Arrays; import java.util.List; import java.util.Map; +import java.util.stream.Collectors; import java.util.stream.IntStream; import java.util.stream.Stream; @@ -29,7 +30,10 @@ import nextflow.script.ast.AssignmentExpression; import nextflow.script.ast.FeatureFlagNode; import nextflow.script.ast.FunctionNode; +import nextflow.script.ast.IncludeNode; +import nextflow.script.ast.OutputBlockNode; import nextflow.script.ast.OutputNode; +import nextflow.script.ast.ParamBlockNode; import nextflow.script.ast.ProcessNode; import nextflow.script.ast.ProcessNodeV2; import nextflow.script.ast.RecordNode; @@ -134,6 +138,8 @@ public void visit() { return; for( var featureFlag : sn.getFeatureFlags() ) visitFeatureFlag(featureFlag); + for( var includeNode : sn.getIncludes() ) + visitInclude(includeNode); if( sn.getParams() != null ) visitParams(sn.getParams()); for( var functionNode : sn.getFunctions() ) @@ -150,6 +156,31 @@ public void visit() { // script declarations + /** + * The params record type of an included pipeline is partial, so each + * field is nullable. + */ + @Override + public void visitInclude(IncludeNode node) { + for( var entry : node.entries ) { + if( !(entry.getTarget() instanceof ClassNode cn) ) + continue; + if( !ScriptNode.isPipelineParams(cn) ) + continue; + for( var fn : cn.getFields() ) + fn.setType(nullableType(fn.getType())); + } + } + + private static ClassNode nullableType(ClassNode type) { + if( isNullable(type) ) + return type; + var result = type.getPlainNodeReference(); + result.setGenericsTypes(type.getGenericsTypes()); + result.putNodeMetaData(ASTNodeMarker.NULLABLE, Boolean.TRUE); + return result; + } + @Override public void visitFeatureFlag(FeatureFlagNode node) { var fn = node.target; @@ -533,6 +564,10 @@ public void visitMethodCallExpression(MethodCallExpression node) { } } + // resolve params and outputs for pipeline calls + if( checkPipelineCall(node) ) + return; + // resolve dataflow inputs and outputs for process calls if( checkProcessCall(node) ) return; @@ -779,7 +814,7 @@ private void checkNamedParams(Parameter param, MapExpression args) { var argType = getType(value); if( !Types.isAssignableFrom(namedParam.getType(), argType) ) addError("Named param `" + name + "` expects a " + Types.getName(namedParam.getType()) + " but received a " + Types.getName(argType), value); - entry.putNodeMetaData("_NAMED_PARAM", namedParam); + entry.putNodeMetaData(ASTNodeMarker.NAMED_PARAM, namedParam); } } @@ -805,6 +840,89 @@ private void visitClosureArguments(ClassNode receiverType, List argu } } + /** + * Check a call to an included pipeline against its params block, and + * resolve the return type from its output block. + * + * A pipeline declares its inputs with a params block rather than a `take:` + * section, so it is called with a single record (or no arguments). Each + * field must be assignable to the declared param type, like a workflow input. + * + * The return type follows the same rules as a workflow call, with the + * declared outputs as the emits. + * + * @param node + */ + private boolean checkPipelineCall(MethodCallExpression node) { + var mn = (MethodNode) node.getNodeMetaData(ASTNodeMarker.METHOD_TARGET); + var pipeline = mn instanceof WorkflowNode wn ? ScriptNode.getPipeline(wn) : null; + if( pipeline == null ) + return false; + node.putNodeMetaData(ASTNodeMarker.INFERRED_TYPE, pipelineOutputType(pipeline.getOutputs())); + + var name = node.getMethodAsString(); + var params = pipeline.getParams(); + var arguments = asMethodCallArguments(node); + if( params == null ) { + if( !arguments.isEmpty() ) + addError("Pipeline `" + name + "` does not declare any params, so it should be called with no arguments", node); + return true; + } + if( arguments.size() != 1 ) { + addError("Pipeline `" + name + "` should be called with a record", node); + return true; + } + var argument = arguments.get(0); + var argType = getType(argument); + if( !Types.isRecordType(argType) ) { + if( !ClassHelper.isDynamicTyped(argType) ) + addError("Pipeline `" + name + "` should be called with a record, but received a " + Types.getName(argType), argument); + return true; + } + if( argType.getFields().isEmpty() ) + return true; + checkPipelineParams(node, argument, params, argType); + return true; + } + + private void checkPipelineParams(MethodCallExpression node, ASTNode argument, ParamBlockNode params, ClassNode argType) { + var declarations = params.declarations; + var byName = Arrays.stream(declarations).collect(Collectors.toMap(Parameter::getName, p -> p, (a, b) -> a)); + + for( var fn : argType.getFields() ) { + var declaration = byName.get(fn.getName()); + if( declaration == null ) { + addError("Param `" + fn.getName() + "` is not defined by pipeline `" + node.getMethodAsString() + "`", argument); + continue; + } + var paramType = declaration.getType(); + if( !Types.isAssignableFrom(paramType, fn.getType()) ) + addError("Param `" + fn.getName() + "` expects a " + Types.getName(paramType) + " but received a " + Types.getName(fn.getType()), argument); + } + + var missing = Arrays.stream(declarations) + .filter(p -> argType.getField(p.getName()) == null) + .filter(p -> !p.hasInitialExpression() && !isNullable(p.getType()) && !ScriptNode.isPipelineParams(p.getType().redirect())) + .map(Parameter::getName) + .toList(); + if( !missing.isEmpty() ) + addError("Pipeline `" + node.getMethodAsString() + "` requires the following params: " + String.join(", ", missing), node); + } + + private static ClassNode pipelineOutputType(OutputBlockNode outputs) { + if( outputs == null ) + return ClassHelper.VOID_TYPE; + if( outputs.declarations.size() == 1 ) + return workflowEmitType(outputs.declarations.get(0).getType()); + var cn = new ClassNode(Record.class); + for( var declaration : outputs.declarations ) { + var fn = new FieldNode(declaration.getName(), Modifier.PUBLIC, workflowEmitType(declaration.getType()), cn, null); + fn.setDeclaringClass(cn); + cn.addField(fn); + } + return cn; + } + /** * Check process calls against the declared process inputs, and resolve * the return type based on the declared process outputs. diff --git a/modules/nf-lang/src/main/java/nextflow/script/control/VariableScopeVisitor.java b/modules/nf-lang/src/main/java/nextflow/script/control/VariableScopeVisitor.java index a665e14e0b..c3cbd08e98 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/control/VariableScopeVisitor.java +++ b/modules/nf-lang/src/main/java/nextflow/script/control/VariableScopeVisitor.java @@ -131,7 +131,16 @@ private void declareInclude(IncludeNode node) { for( var entry : node.entries ) { if( entry.getTarget() == null ) continue; - if( entry.getTarget() instanceof ClassNode && entry.alias != null ) { + // the parts of an included pipeline must be aliased, whereas other + // types cannot be aliased + var target = entry.getTarget(); + var isPipelinePart = target instanceof WorkflowNode wn && wn.isEntry() + || target instanceof ClassNode cn && ScriptNode.isPipelineParams(cn); + if( isPipelinePart && entry.alias == null ) { + vsc.addError("An included pipeline must be aliased, e.g. `" + entry.name + " as MY_PIPELINE`", entry); + continue; + } + if( !isPipelinePart && target instanceof ClassNode && entry.alias != null ) { vsc.addError("Included types cannot be aliased", entry); continue; } @@ -850,11 +859,8 @@ public void visitVariableExpression(VariableExpression node) { var name = node.getName(); Variable variable = vsc.findVariableDeclaration(name, node); if( variable == null ) { - if( "args".equals(name) ) { - vsc.addParanoidWarning("The use of `args` outside the entry workflow will not be supported in a future version", node); - } - else if( "params".equals(name) ) { - vsc.addParanoidWarning("The use of `params` outside the entry workflow will not be supported in a future version", node); + if( "args".equals(name) || "params".equals(name) ) { + vsc.addParanoidWarning("The use of `" + name + "` outside the entry workflow is discouraged", name, node); } else if( isStdinStdout(name) ) { // stdin, stdout can be declared without parentheses diff --git a/modules/nf-lang/src/main/java/nextflow/script/control/WorkflowToGroovyVisitor.java b/modules/nf-lang/src/main/java/nextflow/script/control/WorkflowToGroovyVisitor.java index 3d5c599d70..bc31cd6ad9 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/control/WorkflowToGroovyVisitor.java +++ b/modules/nf-lang/src/main/java/nextflow/script/control/WorkflowToGroovyVisitor.java @@ -25,6 +25,8 @@ import nextflow.script.ast.RecordNode; import nextflow.script.ast.ScriptNode; import nextflow.script.ast.WorkflowNode; +import org.codehaus.groovy.ast.ClassHelper; +import org.codehaus.groovy.ast.ClassNode; import org.codehaus.groovy.ast.FieldNode; import org.codehaus.groovy.ast.Parameter; import org.codehaus.groovy.ast.VariableScope; @@ -49,11 +51,20 @@ public class WorkflowToGroovyVisitor { private ScriptNode moduleNode; + private static final ClassNode PARAMS_HELPER = ClassHelper.makeWithoutCaching("nextflow.script.ParamsHelper"); + public WorkflowToGroovyVisitor(SourceUnit sourceUnit) { this.sourceUnit = sourceUnit; this.moduleNode = (ScriptNode) sourceUnit.getAST(); } + /** + * Transform a workflow definition. The entry workflow of a script + * with a params block takes the params as input, so that the + * pipeline can be called like a named workflow. + * + * @param node + */ public Statement transform(WorkflowNode node) { var main = node.main instanceof BlockStatement block ? block : new BlockStatement(); visitWorkflowEmits(node.emits, main); @@ -61,6 +72,13 @@ public Statement transform(WorkflowNode node) { visitWorkflowHandler(node.onComplete, "setOnComplete", main); visitWorkflowHandler(node.onError, "setOnError", main); + var takes = node.getParameters(); + if( node.isEntry() && moduleNode.getParams() != null ) { + takes = new Parameter[] { new Parameter(ClassHelper.dynamicType(), "params") }; + var params = varX("params"); + var stmt = assignS(params, callX(PARAMS_HELPER, "resolveArguments", args(varX("this"), params))); + main.getStatements().add(0, stmt); + } var bodyDef = stmt(createX( "nextflow.script.BodyDef", args( @@ -70,7 +88,7 @@ public Statement transform(WorkflowNode node) { ) )); var closure = closureX(null, block(new VariableScope(), List.of( - workflowTakes(node.getParameters(), node.isEntry() ? null : node.getName()), + workflowTakes(takes, node.isEntry() ? null : node.getName()), node.emits, bodyDef ))); diff --git a/modules/nf-lang/src/main/java/nextflow/script/dsl/PipelineParams.java b/modules/nf-lang/src/main/java/nextflow/script/dsl/PipelineParams.java new file mode 100644 index 0000000000..dc0edd4211 --- /dev/null +++ b/modules/nf-lang/src/main/java/nextflow/script/dsl/PipelineParams.java @@ -0,0 +1,30 @@ +/* + * Copyright 2013-2026, Seqera Labs + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package nextflow.script.dsl; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +/** + * Annotation for denoting a record type that was synthesized from + * the params block of an included pipeline. + */ +@Retention(RetentionPolicy.RUNTIME) +@Target(ElementType.TYPE) +public @interface PipelineParams { +} diff --git a/modules/nf-lang/src/main/java/nextflow/script/dsl/Types.java b/modules/nf-lang/src/main/java/nextflow/script/dsl/Types.java index e68df4d431..c695c974ea 100644 --- a/modules/nf-lang/src/main/java/nextflow/script/dsl/Types.java +++ b/modules/nf-lang/src/main/java/nextflow/script/dsl/Types.java @@ -456,6 +456,8 @@ public static Class normalize(Class type) { continue; if( STANDARD_TYPES.contains(c) ) return c; + if( TYPE_ALIASES.containsKey(c) ) + return TYPE_ALIASES.get(c); queue.add(c.getSuperclass()); for( var ic : c.getInterfaces() ) queue.add(ic); diff --git a/modules/nf-lang/src/test/groovy/nextflow/script/control/PipelineTypeCheckingTest.groovy b/modules/nf-lang/src/test/groovy/nextflow/script/control/PipelineTypeCheckingTest.groovy new file mode 100644 index 0000000000..18f1b475f1 --- /dev/null +++ b/modules/nf-lang/src/test/groovy/nextflow/script/control/PipelineTypeCheckingTest.groovy @@ -0,0 +1,257 @@ +/* + * Copyright 2013-2026, Seqera Labs + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package nextflow.script.control + +import spock.lang.Specification +import test.TestUtils + +import static test.TestUtils.deleteDir +import static test.TestUtils.tempDir +import static test.TestUtils.tempFile + +/** + * Type checking of calls to an included pipeline. + * + * @see nextflow.script.control.TypeCheckingVisitor + */ +class PipelineTypeCheckingTest extends Specification { + + static final String PIPELINE = '''\ + nextflow.enable.types = true + + params { + input: Channel + fasta: Path + aligner: String = 'star' + } + + workflow { + main: + ch_bams = params.input.map { s -> file(s) } + + publish: + bams = ch_bams + multiqc = params.fasta + } + + output { + bams: Channel {} + multiqc: Path {} + } + ''' + + List check(String main, Map modules = [:]) { + def root = tempDir() + try { + def mainFile = tempFile(root, 'main.nf', "nextflow.enable.types = true\n" + main.stripIndent()) + def pipelineFile = tempFile(root, 'rnaseq.nf', PIPELINE) + def moduleFiles = modules.collect { name, text -> tempFile(root, name, "nextflow.enable.types = true\n" + text.stripIndent()) } + def parser = new ScriptParser(root) + return TestUtils.check(parser, [mainFile, pipelineFile] + moduleFiles) + .findAll { e -> e.getSourceLocator().endsWith('main.nf') } + .collect { e -> e.getOriginalMessage() } + } + finally { + deleteDir(root) + } + } + + def 'should check an included params record with overrides' () { + expect: + check('''\ + include { params as RnaseqParams ; workflow as RNASEQ } from './rnaseq.nf' + + params { + rnaseq: RnaseqParams + } + + workflow { + RNASEQ( params.rnaseq + record(aligner: 42) ) + } + ''') == [ 'Param `aligner` expects a String but received a Integer' ] + } + + def 'should treat the fields of an included params record as nullable' () { + expect: + check('''\ + include { params as RnaseqParams } from './rnaseq.nf' + + workflow { + main: + SUMMARY( record(aligner: 'hisat2') ) + def p = record(aligner: 'hisat2') as RnaseqParams + } + + workflow SUMMARY { + take: + p: RnaseqParams + + main: + p.aligner + } + ''') == [] + } + + def 'should not require an included params record param' () { + expect: + check('''\ + include { workflow as META } from './meta.nf' + + workflow { + META( record(label: 'x') ) + } + ''', [ + 'meta.nf': '''\ + include { params as RnaseqParams ; workflow as RNASEQ } from './rnaseq.nf' + + params { + rnaseq: RnaseqParams + options: Options + sample: Sample + label: String + } + + record Options { + verbose: Boolean? + } + + record Sample { + id: String + } + + workflow { + RNASEQ( params.rnaseq ) + } + ''' + ]) == [ 'Pipeline `META` requires the following params: options, sample' ] + } + + def 'should check the arguments of a pipeline call' () { + expect: + check('''\ + include { workflow as RNASEQ } from './rnaseq.nf' + + workflow { + ''' + CALL + ''' + } + ''') == ERRORS + + where: + CALL | ERRORS + "RNASEQ( record(input: channel.of('a'), fasta: file('x')) )" | [] + "RNASEQ( record(input: channel.of('a'), fasta: channel.value(file('x'))) )" | ['Param `fasta` expects a Path but received a Value'] + "RNASEQ( record(input: channel.value('a'), fasta: file('x')) )" | ['Param `input` expects a Channel but received a Value'] + "RNASEQ( record(input: channel.of('a'), aligner: 42) )" | ['Pipeline `RNASEQ` requires the following params: fasta', 'Param `aligner` expects a String but received a Integer'] + "RNASEQ( record(input: channel.of('a'), fasta: file('x'), foo: 1) )" | ['Param `foo` is not defined by pipeline `RNASEQ`'] + "RNASEQ( record(input: channel.of('a')) )" | ['Pipeline `RNASEQ` requires the following params: fasta'] + "RNASEQ()" | ['Pipeline `RNASEQ` should be called with a record'] + "RNASEQ( input: channel.of('a'), fasta: file('x') )" | ['Pipeline `RNASEQ` should be called with a record, but received a Map'] + "RNASEQ( [input: channel.of('a'), fasta: file('x')] )" | ['Pipeline `RNASEQ` should be called with a record, but received a Map'] + "RNASEQ( channel.of('a'), file('x') )" | ['Pipeline `RNASEQ` should be called with a record'] + "RNASEQ( 'a' )" | ['Pipeline `RNASEQ` should be called with a record, but received a String'] + } + + def 'should report an included pipeline that is not aliased' () { + expect: + check('''\ + include { workflow ; params as RnaseqParams } from './rnaseq.nf' + ''') == [ 'An included pipeline must be aliased, e.g. `workflow as MY_PIPELINE`' ] + } + + def 'should call a pipeline without a params block with no arguments' () { + expect: + check('''\ + include { workflow as HELLO } from './hello.nf' + + workflow { + HELLO() + HELLO( record() ) + } + ''', [ + 'hello.nf': '''\ + workflow { + println 'Hello' + } + ''' + ]) == [ 'Pipeline `HELLO` does not declare any params, so it should be called with no arguments' ] + } + + def 'should return nothing from a pipeline without an output block' () { + expect: + check('''\ + include { workflow as HELLO } from './hello.nf' + + workflow { + HELLO().foo + } + ''', [ + 'hello.nf': '''\ + workflow { + println 'Hello' + } + ''' + ]) == [ 'Unrecognized property `foo` for type void' ] + } + + def 'should return the output of a pipeline with a single output' () { + expect: + check('''\ + include { workflow as HELLO } from './hello.nf' + + workflow { + HELLO().map { s -> s.toUpperCase() } + HELLO().messages + } + ''', [ + 'hello.nf': '''\ + workflow { + main: + messages = channel.of('Hello') + + publish: + messages = messages + } + + output { + messages: Channel {} + } + ''' + ]) == [ 'Unrecognized property `messages` for type Channel' ] + } + + def 'should type pipeline outputs as channels or values' () { + expect: + check('''\ + include { workflow as RNASEQ } from './rnaseq.nf' + + workflow { + r = RNASEQ( record(input: channel.of('a'), fasta: file('genome.fa')) ) + BAMS( r.bams ) + r.multiqc.map { p -> p.name }.view() + } + + workflow BAMS { + take: + bams: Channel + + main: + bams.view() + } + ''') == [] + } + +} diff --git a/modules/nf-lang/src/test/groovy/nextflow/script/control/ResolveIncludeTest.groovy b/modules/nf-lang/src/test/groovy/nextflow/script/control/ResolveIncludeTest.groovy index eb88409681..b3b681b76f 100644 --- a/modules/nf-lang/src/test/groovy/nextflow/script/control/ResolveIncludeTest.groovy +++ b/modules/nf-lang/src/test/groovy/nextflow/script/control/ResolveIncludeTest.groovy @@ -140,6 +140,71 @@ class ResolveIncludeTest extends Specification { deleteDir(root) } + def 'should require static typing in a script that includes a pipeline: #INCLUDE' () { + given: + def root = tempDir() + def main = tempFile(root, 'main.nf', + """\ + include { ${INCLUDE} } from './greet.nf' + """) + def module = tempFile(root, 'greet.nf', + '''\ + nextflow.enable.types = true + + params { + greeting: String = 'Hello' + } + + workflow { + println params.greeting + } + ''') + + when: + def errors = check(root, [main, module]) + then: + errors.size() == 1 + errors[0].getSourceLocator().endsWith('main.nf') + errors[0].getOriginalMessage() == 'Including a pipeline requires `nextflow.enable.types = true` in the including script' + + cleanup: + deleteDir(root) + + where: + INCLUDE << [ 'workflow as GREET', 'params as GreetParams' ] + } + + def 'should require static typing in an included pipeline' () { + given: + def root = tempDir() + def main = tempFile(root, 'main.nf', + '''\ + nextflow.enable.types = true + + include { workflow as GREET } from './greet.nf' + ''') + def module = tempFile(root, 'greet.nf', + '''\ + params { + greeting: String = 'Hello' + } + + workflow { + println params.greeting + } + ''') + + when: + def errors = check(root, [main, module]) + then: + errors.size() == 1 + errors[0].getSourceLocator().endsWith('main.nf') + errors[0].getOriginalMessage() == "An included pipeline must enable static typing -- set `nextflow.enable.types = true` in '${module}'" + + cleanup: + deleteDir(root) + } + def 'should resolve an include' () { given: def root = tempDir() diff --git a/modules/nf-lang/src/test/groovy/nextflow/script/types/TypesTest.groovy b/modules/nf-lang/src/test/groovy/nextflow/script/types/TypesTest.groovy index 2c6ba31081..57d2cfe3e1 100644 --- a/modules/nf-lang/src/test/groovy/nextflow/script/types/TypesTest.groovy +++ b/modules/nf-lang/src/test/groovy/nextflow/script/types/TypesTest.groovy @@ -48,6 +48,21 @@ class TypesTest extends Specification { Types.getName(cn) == 'Map' } + static class Sample implements Record {} + + def 'should render a runtime type' () { + expect: + Types.getName(TYPE) == NAME + + where: + TYPE | NAME + java.nio.file.Paths.get('/a').class | 'Path' + ArrayList | 'List' + Sample | 'Record' + Record | 'Record' + String | 'String' + } + def 'should render the return type of a method' () { when: def cn = ClassHelper.makeCached(TaskConfig) diff --git a/tests/checks/workflow-entry.nf/.checks b/tests/checks/workflow-entry.nf/.checks index d2bbff20d8..ed95cfd565 100644 --- a/tests/checks/workflow-entry.nf/.checks +++ b/tests/checks/workflow-entry.nf/.checks @@ -48,5 +48,4 @@ echo '' echo '=== Testing missing required input ===' $NXF_RUN --prefix Hello > stdout3 2>&1 || true -[[ $(grep -c 'Parameter `--samples` is required but no value was provided' stdout3) == 1 ]] || false -[[ $(grep -c 'samples' stdout3) -ge 1 ]] || false +[[ $(grep -c 'Parameter `samples` is required but no value was provided' stdout3) == 1 ]] || false