From 3f8caeaa038df4d1eb54bd9cbc4dcb0e9db9bc9b Mon Sep 17 00:00:00 2001 From: Jonathan Manning Date: Mon, 5 Oct 2026 13:28:36 +0100 Subject: [PATCH 1/3] Fix Seqera Platform and lineage records for Channel and Value params A Channel or Value param is a dataflow value that is bound only when the dataflow network starts, so serializing the session params on flow begin blocked the Seqera Platform begin request forever, and the lineage workflow run record failed to encode and was not saved. The params block also resolves each declared param to a plain value (the samplesheet of a Channel param, the converted value of a Value param). ParamsMap.toPlainMap() replaces each dataflow value in the params, including the fields of a record param, with its plain value, and keeps every other value as is. The Seqera Platform and lineage observers use it in place of the session params. Assisted-by: Claude Code (Opus 5.5) Signed-off-by: Jonathan Manning --- .../groovy/nextflow/script/ParamsDsl.groovy | 6 +- .../nextflow/script/ParamsHelper.groovy | 95 +++++++++++- .../nextflow/script/ScriptBinding.groovy | 30 ++++ .../nextflow/script/ParamsDslTest.groovy | 137 ++++++++++++++++++ .../nextflow/script/ParamsHelperTest.groovy | 57 ++++++++ .../script/PipelineCompositionTest.groovy | 33 +++++ .../main/nextflow/lineage/LinObserver.groovy | 2 +- .../nextflow/lineage/LinObserverTest.groovy | 101 +++++++++++++ .../seqera/tower/plugin/TowerObserver.groovy | 4 +- .../tower/plugin/TowerObserverTest.groovy | 76 ++++++++++ 10 files changed, 529 insertions(+), 12 deletions(-) diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy index f97a371a11..2f59106f74 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ParamsDsl.groovy @@ -48,7 +48,9 @@ class ParamsDsl { Map getDeclarations() { declarations } void apply(Session session) { - final params = ParamsHelper.resolveParams(declarations.values(), session.cliParams ?: [:], session.configParams ?: [:]) + final cliParams = session.cliParams ?: [:] + final configParams = session.configParams ?: [:] + final params = ParamsHelper.resolveParams(declarations.values(), cliParams, configParams) // propagate resolved params to all scripts for legacy compatibility if( !session.binding.getScriptPath() ) @@ -59,6 +61,8 @@ class ParamsDsl { final script = ScriptMeta.getScriptByPath(scriptPath) script.binding.setParams(params, true) } + + session.getParams().setPlainValues(ParamsHelper.resolvePlainParams(declarations.values(), cliParams, configParams)) } } diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy index 330c4791f5..a3144ba7c9 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy @@ -70,12 +70,44 @@ class ParamsHelper { 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) + final given = givenParams(names, cliParams, configParams) return resolveParams(declarations, given, '') { Param decl, Object value -> resolveParam(decl, value, cliParams.containsKey(decl.name)) } } + /** + * Resolve declared params from the command line and config to plain + * values, which can be serialized before the dataflow network has + * started (e.g. for lineage or Seqera Platform). + * + * Each param is resolved as in {@link #resolveParams(Collection,Map,Map)}, + * except that a {@code Channel} param is resolved to its samplesheet + * and a {@code Value} param to its value of type {@code V}, instead + * of a dataflow value (see {@link #toPlainValue}). The params are assumed + * to be valid, i.e. already resolved by {@link #resolveParams(Collection,Map,Map)}. + * + * @param declarations + * @param cliParams + * @param configParams + */ + static Map resolvePlainParams(Collection declarations, Map cliParams, Map configParams) { + final given = givenParams(declarations*.name as Set, cliParams, configParams) + final result = new LinkedHashMap(declarations.size()) + for( final decl : declarations ) { + final name = decl.name + final value = given.containsKey(name) + ? resolveParam0(decl, given.get(name), cliParams.containsKey(name), false) + : resolveDefault0(decl, false) + result.put(name, value) + } + return result + } + + private static Map givenParams(Set names, Map cliParams, Map configParams) { + return cliParams.subMap(names) + configParams.subMap(names) + } + /** * Resolve declared params against the given values. A param * with no given value is given its default value. @@ -142,6 +174,33 @@ class ParamsHelper { || value instanceof ChannelOut } + /** + * Replace each dataflow value in a resolved param with the + * corresponding plain value (see {@link #resolvePlainParams}), + * including the fields of a record. A value that is not and does + * not contain a dataflow value is returned as is. + * + * @param value the resolved value + * @param plainValue the plain value + */ + static Object toPlainValue(Object value, Object plainValue) { + if( isDataflow(value) ) + return plainValue + if( value !instanceof RecordMap || plainValue !instanceof Map ) + return value + final record = (RecordMap)value + Map result = null + for( final entry : record.entrySet() ) { + final plainField = toPlainValue(entry.value, ((Map)plainValue).get(entry.key)) + if( plainField.is(entry.value) ) + continue + if( result == null ) + result = new LinkedHashMap(record) + result.put(entry.key, plainField) + } + return result != null ? new RecordMap(result) : value + } + /** * Resolve a param value against its declared type. * @@ -156,19 +215,35 @@ class ParamsHelper { * therefore a string that may need to be parsed) */ static Object resolveParam(Param decl, Object value, boolean fromCli) { + return resolveParam0(decl, value, fromCli, true) + } + + /** + * Resolve a param value against its declared type, either as a + * dataflow value (see {@link #resolveParam(Param,Object,boolean)}) or + * as a plain value (see {@link #resolvePlainParams}). + * + * @param decl + * @param value + * @param fromCli + * @param dataflow + */ + private static Object resolveParam0(Param decl, Object value, boolean fromCli, boolean dataflow) { if( value == null ) return null final rawType = TypeHelper.getRawType(decl.type) if( rawType == Channel ) - return ChannelNamespace.fromList(loadChannelInput(decl, value)) + return dataflow ? ChannelNamespace.fromList(loadChannelInput(decl, value)) : value - if( rawType == Value ) - return ChannelNamespace.value(resolveParam(elementDecl(decl), value, fromCli)) + if( rawType == Value ) { + final result = resolveParam0(elementDecl(decl), value, fromCli, dataflow) + return dataflow ? ChannelNamespace.value(result) : result + } if( TypeHelper.isRecordType(decl.type) && value instanceof Map ) - return resolveRecord(decl, (Map)value, fromCli) + return resolveRecord(decl, (Map)value, fromCli, dataflow) final result = fromCli ? resolveFromCli(decl, value) @@ -177,7 +252,7 @@ class ParamsHelper { return result } - private static RecordMap resolveRecord(Param decl, Map value, boolean fromCli) { + private static RecordMap resolveRecord(Param decl, Map value, boolean fromCli, boolean dataflow) { final type = (Class)decl.type final result = new LinkedHashMap(value) for( final field : type.getDeclaredFields() ) { @@ -192,7 +267,7 @@ class ParamsHelper { continue } final fieldDecl = new Param("${decl.name}.${name}", field.getGenericType(), optional, null) - result.put(name, resolveParam(fieldDecl, fieldValue, fromCli)) + result.put(name, resolveParam0(fieldDecl, fieldValue, fromCli, dataflow)) } return new RecordMap(result) } @@ -452,8 +527,12 @@ class ParamsHelper { * @param decl */ static Object resolveDefault(Param decl) { + return resolveDefault0(decl, true) + } + + private static Object resolveDefault0(Param decl, boolean dataflow) { if( decl.defaultValue != null ) - return resolveParam(decl, decl.defaultValue, false) + return resolveParam0(decl, decl.defaultValue, false, dataflow) final type = TypeHelper.getRawType(decl.type) return type.isAnnotationPresent(PipelineParams) ? new RecordMap([:]) diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy index 5aaf9b34e7..19d6297888 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy @@ -216,6 +216,8 @@ class ScriptBinding extends WorkflowBinding { private List scriptAssignment = [] + private Map plainValues + @Delegate private Map target @@ -239,6 +241,34 @@ class ScriptBinding extends WorkflowBinding { return new ParamsMap(this, overrides) } + /** + * Set the plain values of the params declared in the params + * block (see {@link ParamsHelper#resolvePlainParams}). + * + * @param values + */ + void setPlainValues(Map values) { + plainValues = values + } + + /** + * Get the params with each dataflow value of a declared param + * (e.g. a {@code Channel} param, or a field of a record param) + * replaced by the plain value it was resolved from, so that the + * params can be serialized before the dataflow network has started. + * Any other param value is the original object. + */ + Map toPlainMap() { + if( !plainValues ) + return this + final result = new LinkedHashMap(target) + for( final entry : plainValues.entrySet() ) { + if( result.containsKey(entry.key) ) + result.put(entry.key, ParamsHelper.toPlainValue(result.get(entry.key), entry.value)) + } + return result + } + @Override Object get(Object key) { if( !target.containsKey(key) ) { diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy index 905ad710f6..69f7dbde66 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy @@ -19,7 +19,11 @@ package nextflow.script import java.nio.file.Files import java.nio.file.Path +import groovyx.gpars.dataflow.DataflowBroadcast +import groovyx.gpars.dataflow.DataflowVariable import nextflow.Session +import nextflow.dataflow.ChannelImpl +import nextflow.dataflow.ValueImpl import nextflow.exception.AbortOperationException import nextflow.exception.ScriptRuntimeException import nextflow.file.FileHelper @@ -298,6 +302,139 @@ class ParamsDslTest extends Specification { samplesheet?.delete() } + def 'should give dataflow params as plain values'() { + given: + def samplesheet = Files.createTempFile('test', '.csv') + samplesheet.text = 'id,count\na,1\nb,2\n' + def cliParams = [samples: samplesheet.toString(), limit: '5'] + def configParams = [outdir: 'results'] + cliParams + + when: + def params = runScript( + '''\ + nextflow.enable.types = true + + params { + samples: Channel + limit: Value + factor: Value = 3 + label: String = 'demo' + } + + record Sample { + id: String + count: Integer + } + + workflow { params } + ''', + config: [params: configParams], + params: cliParams, + configParams: configParams + ) + then: + params.samples instanceof ChannelImpl + ((ChannelImpl)params.samples).getSource() instanceof DataflowBroadcast + params.limit instanceof ValueImpl + ((ValueImpl)params.limit).getSource() instanceof DataflowVariable + and: + params.toPlainMap() == [outdir: 'results', samples: samplesheet.toString(), limit: 5, factor: 3, label: 'demo'] + + cleanup: + samplesheet?.delete() + } + + def 'should keep the non-dataflow fields of a record param'() { + given: + def samplesheet = Files.createTempFile('test', '.csv') + samplesheet.text = 'id\na\n' + def reference = Files.createTempFile('test', '.fa') + def cliParams = [inputs: [samples: samplesheet.toString(), reference: reference.toString()]] + + when: + def params = runScript( + '''\ + nextflow.enable.types = true + + params { + inputs: Inputs + } + + record Inputs { + samples: Channel + reference: Path + } + + record Sample { + id: String + } + + workflow { params } + ''', + config: [params: cliParams], + params: cliParams, + configParams: cliParams + ) + then: + def plain = params.toPlainMap() + params.inputs.samples instanceof ChannelImpl + plain.inputs.samples == samplesheet.toString() + plain.inputs.reference.is(params.inputs.reference) + + cleanup: + samplesheet?.delete() + reference?.delete() + } + + def 'should give non-dataflow params unchanged as plain values'() { + given: + def inputFile = Files.createTempFile('test', '.csv') + def cliParams = [input: inputFile.toString(), chunk_size: '3', sample: [id: 'a', greeting: 'hola']] + def configParams = [outdir: 'results'] + cliParams + + when: + def params = runScript( + '''\ + params { + input: Path + chunk_size: Integer = 1 + save_intermeds: Boolean + sample: Sample + } + + record Sample { + id: String + greeting: String + } + + workflow { params } + ''', + config: [params: configParams], + params: cliParams, + configParams: configParams + ) + then: + def plain = params.toPlainMap() + plain == params + plain.every { k, v -> v.is(params[k]) } + + cleanup: + inputFile?.delete() + } + + def 'should give the params as plain values without a params block'() { + when: + def params = runScript( + '''\ + params.input = 'samples.csv' + + workflow { params } + ''' + ) + then: + params.toPlainMap().is(params) + } + 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( diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy index 927fa03088..36c234c796 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ParamsHelperTest.groovy @@ -20,11 +20,14 @@ import java.nio.file.Files import java.nio.file.Path import nextflow.exception.ScriptRuntimeException +import nextflow.script.types.Channel +import nextflow.script.types.Value import nextflow.util.Duration import nextflow.util.MemoryUnit import nextflow.util.RecordMap import nextflow.util.VersionNumber import spock.lang.Specification +import spock.lang.Unroll /** * Tests for {@link ParamsHelper}. @@ -235,6 +238,60 @@ class ParamsHelperTest extends Specification { jsonFile?.delete() } + @Unroll + def 'should resolve dataflow params to plain values: #CLI #CONFIG'() { + given: + def declarations = [ + declaredParam('samples', 'default.csv'), + declaredParam('limit', 3), + declaredParam('label', 'demo') + ] + + when: + def result = ParamsHelper.resolvePlainParams(declarations, CLI, CONFIG) + then: + result == EXPECTED + result.limit.getClass() == Integer + + where: + CLI | CONFIG | EXPECTED + [:] | [:] | [samples: 'default.csv', limit: 3, label: 'demo'] + [:] | [samples: 'config.csv', limit: 7] | [samples: 'config.csv', limit: 7, label: 'demo'] + // the config params include the command line values (see ConfigDsl) + [samples: 'cli.csv', limit: '5'] | [samples: 'cli.csv', limit: '5'] | [samples: 'cli.csv', limit: 5, label: 'demo'] + [limit: '5'] | [samples: 'config.csv', limit: 5] | [samples: 'config.csv', limit: 5, label: 'demo'] + } + + def 'should resolve the dataflow fields of a record param to plain values'() { + given: + def declarations = [ declaredParam('pipeline') ] + def cliParams = [pipeline: [samples: 'cli.csv', limit: '5']] + def configParams = [pipeline: [samples: 'cli.csv', limit: '5', label: 'config']] + + when: + def result = ParamsHelper.resolvePlainParams(declarations, cliParams, configParams) + then: + result == [pipeline: [samples: 'cli.csv', limit: 5, label: 'config']] + result.pipeline instanceof RecordMap + } + + private static Param declaredParam(String name, Object defaultValue = null) { + new Param(name, TypedParams.getField(name).getGenericType(), false, defaultValue) + } + + static class TypedParams { + public Channel samples + public Value limit + public String label + public PipelineRec pipeline + } + + static class PipelineRec implements nextflow.script.types.Record { + Channel samples + Value limit + String label + } + 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 index 0871f368d6..4b189dc515 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/PipelineCompositionTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/PipelineCompositionTest.groovy @@ -19,6 +19,8 @@ package nextflow.script import java.nio.file.Files import java.nio.file.Path +import nextflow.dataflow.ChannelImpl +import nextflow.dataflow.ValueImpl import nextflow.exception.ScriptRuntimeException import spock.lang.Timeout import spock.lang.Unroll @@ -290,6 +292,37 @@ class PipelineCompositionTest extends Dsl2Spec { result.val == 15 } + def 'should give the dataflow params of an included pipeline as plain values' () { + 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 ) + params + } + ''' + ]) + def samples = folder.resolve('samples.csv').toString() + def cliParams = [count: [samples: samples, factor: '5']] + + when: + def params = runScript([params: cliParams], script) + then: + params.count.samples instanceof ChannelImpl + params.count.factor instanceof ValueImpl + and: + params.toPlainMap() == [count: [samples: samples, factor: 5]] + } + def 'should scope the processes of an included pipeline by its name' () { given: def script = write([ diff --git a/modules/nf-lineage/src/main/nextflow/lineage/LinObserver.groovy b/modules/nf-lineage/src/main/nextflow/lineage/LinObserver.groovy index ed790ae7d3..52046780c5 100644 --- a/modules/nf-lineage/src/main/nextflow/lineage/LinObserver.groovy +++ b/modules/nf-lineage/src/main/nextflow/lineage/LinObserver.groovy @@ -179,7 +179,7 @@ class LinObserver implements TraceObserverV2 { workflow, session.uniqueId.toString(), session.runName, - getNormalizedParams(session.params, normalizer), + getNormalizedParams(session.params.toPlainMap(), normalizer), SecretHelper.hideSecrets(session.config.deepClone()) as Map, collectWorkflowMetadata(normalizer) ) diff --git a/modules/nf-lineage/src/test/nextflow/lineage/LinObserverTest.groovy b/modules/nf-lineage/src/test/nextflow/lineage/LinObserverTest.groovy index dcd44ba7c5..925b9f901c 100644 --- a/modules/nf-lineage/src/test/nextflow/lineage/LinObserverTest.groovy +++ b/modules/nf-lineage/src/test/nextflow/lineage/LinObserverTest.groovy @@ -33,6 +33,7 @@ import java.nio.file.Path import java.nio.file.attribute.BasicFileAttributes import com.google.common.hash.HashCode +import nextflow.Global import nextflow.NextflowMeta import nextflow.Session import nextflow.file.FileHolder @@ -50,6 +51,8 @@ import nextflow.processor.TaskConfig import nextflow.processor.TaskHandler import nextflow.processor.TaskId import nextflow.processor.TaskRun +import nextflow.script.Param +import nextflow.script.ParamsHelper import nextflow.script.ScriptBinding import nextflow.script.PlatformMetadata import nextflow.script.ScriptMeta @@ -66,7 +69,9 @@ import nextflow.script.params.ValueInParam import nextflow.script.params.ValueOutParam import nextflow.script.params.v2.ProcessInput import nextflow.script.params.v2.ProcessOutput +import nextflow.script.types.Channel import nextflow.script.types.Record +import nextflow.script.types.Value import nextflow.trace.event.FilePublishEvent import nextflow.trace.event.TaskEvent import nextflow.trace.event.WorkflowOutputEvent @@ -253,6 +258,102 @@ class LinObserverTest extends Specification { folder?.deleteDir() } + def 'should save workflow with the plain values of dataflow params' (){ + given: + // the dataflow network is never started, so the params are never bound + Global.session = Mock(Session) + def folder = Files.createTempDirectory('test') + def config = [lineage:[enabled: true, store:[location:folder.toString()]]] + def store = new DefaultLinStore(); + def uniqueId = UUID.randomUUID() + def scriptFile = folder.resolve("main.nf") + def samplesheet = folder.resolve("samples.csv"); samplesheet.text = 'id\na\n' + def cliParams = [input: samplesheet.toString(), factor: '5'] + def params = resolveParams([param('input'), param('factor'), param('label', 'demo')], cliParams) + def map = [ + repository: "https://nextflow.io/nf-test/", + commitId: "123456", + scriptId: "78910", + scriptFile: scriptFile, + projectDir: folder.resolve("projectDir"), + revision: "main", + projectName: "nextflow.io/nf-test", + workDir: folder.resolve("workDir") + ] + def metadata = Mock(WorkflowMetadata){ + getRepository() >> map.repository + getCommitId() >> map.commitId + getScriptId() >> map.scriptId + getScriptFile() >> map.scriptFile + getProjectDir() >> map.projectDir + getRevision() >> map.revision + getProjectName() >> map.projectName + getWorkDir() >> map.workDir + toMap() >> map + } + def session = Mock(Session) { + getConfig() >> config + getUniqueId() >> uniqueId + getRunName() >> "test_run" + getWorkflowMetadata() >> metadata + getParams() >> params + } + store.open(LineageConfig.create(session)) + def observer = new LinObserver(session, store) + def mainScript = new DataPath("file://${scriptFile.toString()}", new Checksum("78910", "nextflow", "standard")) + def workflow = new Workflow([mainScript], map.repository, map.commitId) + def expectedParams = LinObserver.getNormalizedParams([input: samplesheet.toString(), factor: 5, label: 'demo'], new PathNormalizer(metadata)) + def workflowRun = new WorkflowRun(workflow, uniqueId.toString(), "test_run", expectedParams, config, map) + when: + observer.onFlowCreate(session) + observer.onFlowBegin() + then: + folder.resolve("${observer.executionHash}/.data.json").text == new LinEncoder().encode(workflowRun) + + cleanup: + Global.session = null + folder?.deleteDir() + } + + def 'should normalize non-dataflow params identically to the session params' () { + given: + def folder = Files.createTempDirectory('test') + def metadata = Mock(WorkflowMetadata){ + getProjectDir() >> folder.resolve("projectDir") + getWorkDir() >> folder.resolve("workDir") + } + def normalizer = new PathNormalizer(metadata) + def cliParams = [outdir: folder.toString(), chunks: '3'] + def params = resolveParams([param('outdir'), param('chunks'), param('label', 'demo')], cliParams) + def encode = { Map value -> + new LinEncoder().encode(new WorkflowRun(null, 'uuid', 'test_run', LinObserver.getNormalizedParams(value, normalizer), [:], [:])) + } + + expect: + encode(params.toPlainMap()) == encode(params) + + cleanup: + folder?.deleteDir() + } + + private static Param param(String name, Object defaultValue = null) { + new Param(name, TypedParams.getField(name).getGenericType(), false, defaultValue) + } + + private static ScriptBinding.ParamsMap resolveParams(List declarations, Map cliParams) { + final result = new ScriptBinding.ParamsMap(ParamsHelper.resolveParams(declarations, cliParams, cliParams)) + result.setPlainValues(ParamsHelper.resolvePlainParams(declarations, cliParams, cliParams)) + return result + } + + static class TypedParams { + public Channel input + public Value factor + public String label + public Path outdir + public Integer chunks + } + def 'should strip sensitive user data from platform metadata in lineage' () { given: def folder = Files.createTempDirectory('test') diff --git a/plugins/nf-tower/src/main/io/seqera/tower/plugin/TowerObserver.groovy b/plugins/nf-tower/src/main/io/seqera/tower/plugin/TowerObserver.groovy index fa405be08e..79edf6aad0 100644 --- a/plugins/nf-tower/src/main/io/seqera/tower/plugin/TowerObserver.groovy +++ b/plugins/nf-tower/src/main/io/seqera/tower/plugin/TowerObserver.groovy @@ -380,7 +380,7 @@ class TowerObserver implements TraceObserverV2 { protected Map makeBeginReq(Session session) { def workflow = session.getWorkflowMetadata().toMap() - workflow.params = session.getParams() + workflow.params = session.getParams().toPlainMap() workflow.id = getWorkflowId() workflow.remove('stats') @@ -437,7 +437,7 @@ class TowerObserver implements TraceObserverV2 { if( workflow.platform ) workflow.remove('platform') - workflow.params = session.getParams() + workflow.params = session.getParams().toPlainMap() workflow.id = getWorkflowId() // render as a string workflow.container = mapToString(workflow.container) diff --git a/plugins/nf-tower/src/test/io/seqera/tower/plugin/TowerObserverTest.groovy b/plugins/nf-tower/src/test/io/seqera/tower/plugin/TowerObserverTest.groovy index c7bbdd6f2e..bd2c82c3a0 100644 --- a/plugins/nf-tower/src/test/io/seqera/tower/plugin/TowerObserverTest.groovy +++ b/plugins/nf-tower/src/test/io/seqera/tower/plugin/TowerObserverTest.groovy @@ -17,11 +17,14 @@ package io.seqera.tower.plugin import java.nio.file.Files +import java.nio.file.Path import java.time.Instant import java.time.OffsetDateTime import java.time.ZoneId +import groovy.json.JsonSlurper import groovyx.gpars.dataflow.DataflowQueue +import nextflow.Global import nextflow.Session import nextflow.SysEnv import nextflow.cloud.types.CloudMachineInfo @@ -30,14 +33,19 @@ import nextflow.container.DockerConfig import nextflow.container.resolver.ContainerMeta import nextflow.dag.DAG import nextflow.exception.AbortRunException +import nextflow.script.Param +import nextflow.script.ParamsHelper import nextflow.script.PlatformMetadata import nextflow.script.ScriptBinding import nextflow.script.WorkflowMetadata +import nextflow.script.types.Channel +import nextflow.script.types.Value import nextflow.trace.TraceRecord import nextflow.trace.WorkflowStats import nextflow.trace.WorkflowStatsObserver import nextflow.util.ProcessHelper import spock.lang.Specification +import spock.lang.Timeout /** * * @author Paolo Di Tommaso @@ -749,5 +757,73 @@ class TowerObserverTest extends Specification { thrown(AbortRunException) } + @Timeout(10) + def 'should send the plain values of dataflow params in the begin and complete requests' () { + given: + // the dataflow network is never started, so the params are never bound + Global.session = Mock(Session) + def samplesheet = Files.createTempFile('test', '.csv') + samplesheet.text = 'id\na\n' + def cliParams = [input: samplesheet.toString(), factor: '5'] + def params = resolveParams([param('input'), param('factor'), param('label', 'demo')], cliParams) + and: + def session = Mock(Session) + session.getParams() >> params + session.getWorkflowMetadata() >> Mock(WorkflowMetadata) { toMap() >> [:] } + def observer = Spy(newObserver(session)) + observer.getMetricsList() >> [] + observer.getWorkflowProgress(false) >> new WorkflowProgress() + def generator = TowerJsonGenerator.create([:]) + + when: + def begin = new JsonSlurper().parseText(generator.toJson(observer.makeBeginReq(session).workflow)) + def complete = new JsonSlurper().parseText(generator.toJson(observer.makeCompleteReq(session).workflow)) + then: + begin.params == [input: samplesheet.toString(), factor: 5, label: 'demo'] + complete.params == [input: samplesheet.toString(), factor: 5, label: 'demo'] + + cleanup: + Global.session = null + samplesheet?.delete() + } + + def 'should serialize non-dataflow params identically to the session params' () { + given: + def outdir = Files.createTempDirectory('test') + def cliParams = [outdir: outdir.toString(), chunks: '3'] + def params = resolveParams([param('outdir'), param('chunks'), param('label', 'demo')], cliParams) + and: + def session = Mock(Session) + session.getParams() >> params + session.getWorkflowMetadata() >> Mock(WorkflowMetadata) { toMap() >> [:] } + def observer = Spy(newObserver(session)) + def generator = TowerJsonGenerator.create([:]) + + when: + def req = observer.makeBeginReq(session) + then: + generator.toJson(req.workflow.params) == generator.toJson(params) + + cleanup: + outdir?.deleteDir() + } + + private static Param param(String name, Object defaultValue = null) { + new Param(name, TypedParams.getField(name).getGenericType(), false, defaultValue) + } + + private static ScriptBinding.ParamsMap resolveParams(List declarations, Map cliParams) { + final result = new ScriptBinding.ParamsMap(ParamsHelper.resolveParams(declarations, cliParams, cliParams)) + result.setPlainValues(ParamsHelper.resolvePlainParams(declarations, cliParams, cliParams)) + return result + } + + static class TypedParams { + public Channel input + public Value factor + public String label + public Path outdir + public Integer chunks + } } From bcf59edcbee6d69d5d736a381e0e2e2c2c226c57 Mon Sep 17 00:00:00 2001 From: Jonathan Manning Date: Mon, 5 Oct 2026 14:08:58 +0100 Subject: [PATCH 2/3] Test a Channel param given as a Path in the config Assisted-by: Claude Code (Opus 5.5) Signed-off-by: Jonathan Manning --- .../nextflow/script/ParamsDslTest.groovy | 32 +++++++++++++++++++ 1 file changed, 32 insertions(+) diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy index 69f7dbde66..e495f9dccb 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy @@ -344,6 +344,38 @@ class ParamsDslTest extends Specification { samplesheet?.delete() } + def 'should give a channel param set to a path in the config as that path'() { + given: + def samplesheet = Files.createTempFile('test', '.csv') + samplesheet.text = 'id\na\n' + def configParams = [samples: samplesheet] + + when: + def params = runScript( + '''\ + nextflow.enable.types = true + + params { + samples: Channel + } + + record Sample { + id: String + } + + workflow { params } + ''', + config: [params: configParams], + configParams: configParams + ) + then: + params.samples instanceof ChannelImpl + params.toPlainMap().samples.is(samplesheet) + + cleanup: + samplesheet?.delete() + } + def 'should keep the non-dataflow fields of a record param'() { given: def samplesheet = Files.createTempFile('test', '.csv') From 2ace6c7468613718392fc635a5bf6ccb0baf34c2 Mon Sep 17 00:00:00 2001 From: Ben Sherman Date: Mon, 5 Oct 2026 10:56:19 -0500 Subject: [PATCH 3/3] Simplify plain param resolution Replace ParamsHelper.toPlainValue() with a direct overlay of the plain values in ParamsMap.toPlainMap(), and replace the resolveParam0() and resolveDefault0() overloads with a `plain` default argument. Signed-off-by: Ben Sherman --- .../nextflow/script/ParamsHelper.groovy | 75 +++++-------------- .../nextflow/script/ScriptBinding.groovy | 10 +-- .../nextflow/script/ParamsDslTest.groovy | 3 +- 3 files changed, 19 insertions(+), 69 deletions(-) diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy index a3144ba7c9..e37007445d 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ParamsHelper.groovy @@ -84,8 +84,8 @@ class ParamsHelper { * Each param is resolved as in {@link #resolveParams(Collection,Map,Map)}, * except that a {@code Channel} param is resolved to its samplesheet * and a {@code Value} param to its value of type {@code V}, instead - * of a dataflow value (see {@link #toPlainValue}). The params are assumed - * to be valid, i.e. already resolved by {@link #resolveParams(Collection,Map,Map)}. + * of a dataflow value. The params are assumed to be valid, i.e. already + * resolved by {@link #resolveParams(Collection,Map,Map)}. * * @param declarations * @param cliParams @@ -97,8 +97,8 @@ class ParamsHelper { for( final decl : declarations ) { final name = decl.name final value = given.containsKey(name) - ? resolveParam0(decl, given.get(name), cliParams.containsKey(name), false) - : resolveDefault0(decl, false) + ? resolveParam(decl, given.get(name), cliParams.containsKey(name), true) + : resolveDefault(decl, true) result.put(name, value) } return result @@ -174,33 +174,6 @@ class ParamsHelper { || value instanceof ChannelOut } - /** - * Replace each dataflow value in a resolved param with the - * corresponding plain value (see {@link #resolvePlainParams}), - * including the fields of a record. A value that is not and does - * not contain a dataflow value is returned as is. - * - * @param value the resolved value - * @param plainValue the plain value - */ - static Object toPlainValue(Object value, Object plainValue) { - if( isDataflow(value) ) - return plainValue - if( value !instanceof RecordMap || plainValue !instanceof Map ) - return value - final record = (RecordMap)value - Map result = null - for( final entry : record.entrySet() ) { - final plainField = toPlainValue(entry.value, ((Map)plainValue).get(entry.key)) - if( plainField.is(entry.value) ) - continue - if( result == null ) - result = new LinkedHashMap(record) - result.put(entry.key, plainField) - } - return result != null ? new RecordMap(result) : value - } - /** * Resolve a param value against its declared type. * @@ -213,37 +186,26 @@ class ParamsHelper { * @param value * @param fromCli whether the value came from the command line (and is * therefore a string that may need to be parsed) + * @param plain whether to give the plain value of a {@code Channel} + * or {@code Value} param instead of a dataflow value + * (see {@link #resolvePlainParams}) */ - static Object resolveParam(Param decl, Object value, boolean fromCli) { - return resolveParam0(decl, value, fromCli, true) - } - - /** - * Resolve a param value against its declared type, either as a - * dataflow value (see {@link #resolveParam(Param,Object,boolean)}) or - * as a plain value (see {@link #resolvePlainParams}). - * - * @param decl - * @param value - * @param fromCli - * @param dataflow - */ - private static Object resolveParam0(Param decl, Object value, boolean fromCli, boolean dataflow) { + static Object resolveParam(Param decl, Object value, boolean fromCli, boolean plain=false) { if( value == null ) return null final rawType = TypeHelper.getRawType(decl.type) if( rawType == Channel ) - return dataflow ? ChannelNamespace.fromList(loadChannelInput(decl, value)) : value + return plain ? value : ChannelNamespace.fromList(loadChannelInput(decl, value)) if( rawType == Value ) { - final result = resolveParam0(elementDecl(decl), value, fromCli, dataflow) - return dataflow ? ChannelNamespace.value(result) : result + final result = resolveParam(elementDecl(decl), value, fromCli, plain) + return plain ? result : ChannelNamespace.value(result) } if( TypeHelper.isRecordType(decl.type) && value instanceof Map ) - return resolveRecord(decl, (Map)value, fromCli, dataflow) + return resolveRecord(decl, (Map)value, fromCli, plain) final result = fromCli ? resolveFromCli(decl, value) @@ -252,7 +214,7 @@ class ParamsHelper { return result } - private static RecordMap resolveRecord(Param decl, Map value, boolean fromCli, boolean dataflow) { + private static RecordMap resolveRecord(Param decl, Map value, boolean fromCli, boolean plain) { final type = (Class)decl.type final result = new LinkedHashMap(value) for( final field : type.getDeclaredFields() ) { @@ -267,7 +229,7 @@ class ParamsHelper { continue } final fieldDecl = new Param("${decl.name}.${name}", field.getGenericType(), optional, null) - result.put(name, resolveParam0(fieldDecl, fieldValue, fromCli, dataflow)) + result.put(name, resolveParam(fieldDecl, fieldValue, fromCli, plain)) } return new RecordMap(result) } @@ -525,14 +487,11 @@ class ParamsHelper { * the pipeline is called. * * @param decl + * @param plain see {@link #resolveParam} */ - static Object resolveDefault(Param decl) { - return resolveDefault0(decl, true) - } - - private static Object resolveDefault0(Param decl, boolean dataflow) { + static Object resolveDefault(Param decl, boolean plain=false) { if( decl.defaultValue != null ) - return resolveParam0(decl, decl.defaultValue, false, dataflow) + return resolveParam(decl, decl.defaultValue, false, plain) final type = TypeHelper.getRawType(decl.type) return type.isAnnotationPresent(PipelineParams) ? new RecordMap([:]) diff --git a/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy b/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy index 19d6297888..c43258c0aa 100644 --- a/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/script/ScriptBinding.groovy @@ -256,17 +256,9 @@ class ScriptBinding extends WorkflowBinding { * (e.g. a {@code Channel} param, or a field of a record param) * replaced by the plain value it was resolved from, so that the * params can be serialized before the dataflow network has started. - * Any other param value is the original object. */ Map toPlainMap() { - if( !plainValues ) - return this - final result = new LinkedHashMap(target) - for( final entry : plainValues.entrySet() ) { - if( result.containsKey(entry.key) ) - result.put(entry.key, ParamsHelper.toPlainValue(result.get(entry.key), entry.value)) - } - return result + return plainValues ? target + plainValues : this } @Override diff --git a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy index e495f9dccb..e0526c8008 100644 --- a/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/script/ParamsDslTest.groovy @@ -411,7 +411,7 @@ class ParamsDslTest extends Specification { def plain = params.toPlainMap() params.inputs.samples instanceof ChannelImpl plain.inputs.samples == samplesheet.toString() - plain.inputs.reference.is(params.inputs.reference) + plain.inputs.reference == params.inputs.reference cleanup: samplesheet?.delete() @@ -448,7 +448,6 @@ class ParamsDslTest extends Specification { then: def plain = params.toPlainMap() plain == params - plain.every { k, v -> v.is(params[k]) } cleanup: inputFile?.delete()