From 7672e01cc8ab585d81b36ba07dda0c06c06dcb02 Mon Sep 17 00:00:00 2001 From: Ben Sherman Date: Tue, 29 Sep 2026 10:29:03 -0500 Subject: [PATCH 1/3] Fix plugin functions that return a channel A plugin method annotated with @Function that returns a DataflowWriteChannel was picked up by the legacy factory detection and registered as a channel factory, so it could not be called as a plain function. Skip @Function methods in the factory fallback so they are registered as functions. Signed-off-by: Ben Sherman --- .../extension/PluginExtensionProvider.groovy | 2 + .../PluginExtensionMethodsTest.groovy | 71 +++++++++++++++++++ .../plugin/hello/HelloExtension.groovy | 8 +++ 3 files changed, 81 insertions(+) diff --git a/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy b/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy index 109f68884f..aaacf2b366 100644 --- a/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy @@ -213,6 +213,8 @@ class PluginExtensionProvider implements ExtensionProvider { result.add(handle.name) continue } + // skip functions that return a channel + if( handle.isAnnotationPresent(Function) ) continue // skip non-public methods if( !Modifier.isPublic(handle.getModifiers()) ) continue // skip static methods diff --git a/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy b/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy index 4d03fb7c7a..9b709d646a 100644 --- a/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy +++ b/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy @@ -492,4 +492,75 @@ class PluginExtensionMethodsTest extends Dsl2Spec { e.cause.message.contains('`sayHello` is already included') } + def 'should execute custom function returning a channel in typed script'() { + given: + def SCRIPT_TEXT = ''' + nextflow.enable.types = true + + include { reverseFn } from 'plugin/nf-test-plugin-hello' + + workflow { + reverseFn('a string') + } + ''' + + when: + def result = runScript(SCRIPT_TEXT) + + then: + result.val == 'a string'.reverse() + result.val == Channel.STOP + } + + def 'should wrap custom function channel as typed channel'() { + given: + def SCRIPT_TEXT = ''' + nextflow.enable.types = true + + include { reverseFn } from 'plugin/nf-test-plugin-hello' + + workflow { + def ch = reverseFn('a string') + channel.value(ch.getClass().getName()) + } + ''' + + when: + def result = runScript(SCRIPT_TEXT) + + then: + result.val == 'nextflow.dataflow.ChannelImpl' + } + + def 'should apply typed operators to custom function channel'() { + given: + def SCRIPT_TEXT = ''' + nextflow.enable.types = true + + include { reverseFn } from 'plugin/nf-test-plugin-hello' + + workflow { + reverseFn('a string').map { s -> s.toUpperCase() }.collect() + } + ''' + + when: + def result = runScript(SCRIPT_TEXT) + + then: + result.val as List == ['a string'.reverse().toUpperCase()] + } + + def 'should execute custom function returning a channel'() { + when: + def result = runScript(''' + include { reverseFn } from 'plugin/nf-test-plugin-hello' + workflow { + reverseFn('a string') + } + ''') + then: + result.val == 'a string'.reverse() + } + } diff --git a/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy b/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy index 82b66ca248..1b6094338f 100644 --- a/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy +++ b/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy @@ -136,6 +136,14 @@ class HelloExtension extends PluginExtensionPoint { return functions.sayHello(lang) } + /** + * A @Function that returns a channel, i.e. a factory as a plain function + */ + @Function + DataflowWriteChannel reverseFn(String message) { + return reverse(message) + } + String aNonImportedFunction(){ throw new IllegalAccessException("This function can't be imported") } From b9b978239b754e7d2146deb7ff6638ad9c0e4f69 Mon Sep 17 00:00:00 2001 From: Ben Sherman Date: Wed, 30 Sep 2026 11:10:30 -0500 Subject: [PATCH 2/3] Support plugin operators as functions, reject conflicting annotations Skip @Function methods in the legacy operator detection, so that a function which takes a channel is registered as a function instead of an operator. Fail when a method is annotated with @Function and @Factory or @Operator. Document how to define channel factories and operators as functions. Signed-off-by: Ben Sherman --- docs/plugins/developing-plugins.mdx | 48 +++++++++++++++++-- .../extension/PluginExtensionProvider.groovy | 5 ++ .../plugin/PluginExtensionProviderTest.groovy | 39 +++++++++++++++ .../PluginExtensionMethodsTest.groovy | 29 +++++++++++ .../plugin/hello/HelloExtension.groovy | 8 ++++ 5 files changed, 126 insertions(+), 3 deletions(-) diff --git a/docs/plugins/developing-plugins.mdx b/docs/plugins/developing-plugins.mdx index 74955e3ffe..89db1fb5f7 100644 --- a/docs/plugins/developing-plugins.mdx +++ b/docs/plugins/developing-plugins.mdx @@ -398,9 +398,50 @@ channel The above snippet is based on the [nf-sqldb](https://github.com/nextflow-io/nf-sqldb) plugin. The `fromQuery` factory is included under the alias `fromTable`. ::: -:::tip -Before creating a custom operator, consider whether the operator can be defined as a [function](#functions) that can be composed with existing operators such as `map` or `subscribe`. Functions are easier to implement and can be used anywhere in your pipeline, not just channel logic. -::: + + +Channel factories and operators can also be defined as [functions](#functions). Functions are easier to implement and easier to use, since they are called directly rather than through the `channel` namespace or as a method on a channel. Plugin factories and operators are not supported when [static typing][static-typing-page] is enabled, so plugins that are used in typed pipelines should provide them as functions. + +A factory function returns a channel: + +```groovy +import groovyx.gpars.dataflow.DataflowWriteChannel +import nextflow.plugin.extension.Function + +@Function +DataflowWriteChannel fromQuery(Map opts, String query) { + // ... +} +``` + +An operator function takes a channel as an argument. The channel is received as a `DataflowWriteChannel`, which can be converted into a channel for reading with `CH.getReadChannel()`. Named arguments are received as a `Map` in the first parameter: + +```groovy +import groovyx.gpars.dataflow.DataflowWriteChannel +import nextflow.extension.CH +import nextflow.plugin.extension.Function + +@Function +DataflowWriteChannel sqlInsert(Map opts, DataflowWriteChannel source) { + final input = CH.getReadChannel(source) + // ... +} +``` + +These functions can be used in your pipeline like any other function: + +```nextflow +include { fromQuery; sqlInsert } from 'plugin/nf-sqldb' + +workflow { + def rows = fromQuery('select * from FOO', db: 'test') + sqlInsert(rows, into: 'BAR', columns: 'id', db: 'test') +} +``` + +Before creating a custom operator, consider whether the operator can be composed from existing operators such as `map` or `subscribe`, using a function to implement the custom logic. This approach works in any pipeline, with or without static typing. + +A method cannot be annotated with both `Function` and `Factory` or `Operator`. ### Process directives @@ -537,3 +578,4 @@ This variable is useful for testing a plugin release before publishing it to the [process-directives]: ../process#directives [process-ext]: ../reference/process/directives/ext [remote-files]: ../working-with-files#remote-files +[static-typing-page]: ../static-typing diff --git a/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy b/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy index aaacf2b366..e3d03428e0 100644 --- a/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy +++ b/modules/nextflow/src/main/groovy/nextflow/plugin/extension/PluginExtensionProvider.groovy @@ -178,6 +178,9 @@ class PluginExtensionProvider implements ExtensionProvider { continue } + // skip functions that take a channel + if( handle.isAnnotationPresent(Function) ) + continue // skip non-public methods if( !Modifier.isPublic(handle.getModifiers()) ) continue @@ -240,6 +243,8 @@ class PluginExtensionProvider implements ExtensionProvider { // custom functions must to be annotated with @Function if( !handle.isAnnotationPresent(Function)) continue + if( handle.isAnnotationPresent(Factory) || handle.isAnnotationPresent(Operator) ) + throw new IllegalStateException("Function extension '$handle.name' in `$clazz.name` cannot also be declared as a factory or operator") // skip non-public methods if( !Modifier.isPublic(handle.getModifiers()) ) throw new IllegalStateException("Function extension '$handle.name' in `$clazz.name` should be declared public") diff --git a/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy b/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy index 0f1dd37a44..2771cbe1e5 100644 --- a/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy @@ -20,7 +20,12 @@ import groovyx.gpars.dataflow.DataflowBroadcast import groovyx.gpars.dataflow.DataflowQueue import groovyx.gpars.dataflow.DataflowReadChannel import groovyx.gpars.dataflow.DataflowVariable +import groovyx.gpars.dataflow.DataflowWriteChannel import nextflow.Channel +import nextflow.plugin.extension.Factory +import nextflow.plugin.extension.Function +import nextflow.plugin.extension.Operator +import nextflow.plugin.extension.PluginExtensionPoint import nextflow.plugin.extension.PluginExtensionProvider import spock.lang.Specification @@ -62,4 +67,38 @@ class PluginExtensionProviderTest extends Specification { result.val == 4 result.val == 9 } + + static class ChannelFunctions extends PluginExtensionPoint { + @Override protected void init(nextflow.Session session) {} + @Function DataflowWriteChannel fromFoo(String value) { null } + @Function DataflowWriteChannel mapFoo(DataflowReadChannel source) { null } + } + + static class FunctionAndFactory extends PluginExtensionPoint { + @Override protected void init(nextflow.Session session) {} + @Function @Factory DataflowWriteChannel fromFoo(String value) { null } + } + + static class FunctionAndOperator extends PluginExtensionPoint { + @Override protected void init(nextflow.Session session) {} + @Function @Operator DataflowWriteChannel mapFoo(DataflowReadChannel source) { null } + } + + def 'should not detect channel functions as factories or operators' () { + expect: + PluginExtensionProvider.getDeclaredFactoryExtensionMethods0(ChannelFunctions).isEmpty() + PluginExtensionProvider.getDeclaredOperatorExtensionMethods0(ChannelFunctions).isEmpty() + PluginExtensionProvider.getDeclaredFunctionsExtensionMethods0(ChannelFunctions) == ['fromFoo', 'mapFoo'] as Set + } + + def 'should reject function that is also a factory or operator' () { + when: + PluginExtensionProvider.getDeclaredFunctionsExtensionMethods0(CLAZZ) + then: + def e = thrown(IllegalStateException) + e.message.contains('cannot also be declared as a factory or operator') + + where: + CLAZZ << [FunctionAndFactory, FunctionAndOperator] + } } diff --git a/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy b/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy index 9b709d646a..1c148cb297 100644 --- a/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy +++ b/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy @@ -563,4 +563,33 @@ class PluginExtensionMethodsTest extends Dsl2Spec { result.val == 'a string'.reverse() } + def 'should execute custom function taking a channel'() { + when: + def result = runScript(''' + include { goodbyeFn } from 'plugin/nf-test-plugin-hello' + workflow { + goodbyeFn(channel.of('Bye bye folks')) + } + ''') + then: + result.val == 'Bye bye folks' + result.val == Channel.STOP + } + + def 'should execute custom function taking a channel in typed script'() { + when: + def result = runScript(''' + nextflow.enable.types = true + + include { goodbyeFn } from 'plugin/nf-test-plugin-hello' + + workflow { + def ch = channel.of('Bye bye folks') + goodbyeFn(ch).map { s -> s.toUpperCase() }.collect() + } + ''') + then: + result.val as List == ['BYE BYE FOLKS'] + } + } diff --git a/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy b/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy index 1b6094338f..4e90c4ebb6 100644 --- a/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy +++ b/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy @@ -144,6 +144,14 @@ class HelloExtension extends PluginExtensionPoint { return reverse(message) } + /** + * A @Function that takes a channel, i.e. an operator as a plain function + */ + @Function + DataflowWriteChannel goodbyeFn(DataflowWriteChannel source) { + return goodbye(CH.getReadChannel(source)) + } + String aNonImportedFunction(){ throw new IllegalAccessException("This function can't be imported") } From 39af8677c322bc0a36acc82038415e9be668ab67 Mon Sep 17 00:00:00 2001 From: Ben Sherman Date: Wed, 30 Sep 2026 11:47:54 -0500 Subject: [PATCH 3/3] Verify documented plugin function patterns Align the test plugin functions with the documented signatures, test the documented pipeline in typed and untyped scripts, and test that plugin factories and operators are unsupported in typed scripts. Simplify the docs for factories and operators as functions. Signed-off-by: Ben Sherman --- docs/plugins/developing-plugins.mdx | 40 +++++++++---------- .../plugin/PluginExtensionProviderTest.groovy | 28 ++----------- .../PluginExtensionMethodsTest.groovy | 40 +++++++++---------- .../plugin/hello/HelloExtension.groovy | 13 ++++-- 4 files changed, 50 insertions(+), 71 deletions(-) diff --git a/docs/plugins/developing-plugins.mdx b/docs/plugins/developing-plugins.mdx index 89db1fb5f7..55ee46e74c 100644 --- a/docs/plugins/developing-plugins.mdx +++ b/docs/plugins/developing-plugins.mdx @@ -398,33 +398,33 @@ channel The above snippet is based on the [nf-sqldb](https://github.com/nextflow-io/nf-sqldb) plugin. The `fromQuery` factory is included under the alias `fromTable`. ::: - +Before creating a custom operator, consider whether the operator can be composed from existing operators such as `map` or `subscribe`, using a function to implement the custom logic. -Channel factories and operators can also be defined as [functions](#functions). Functions are easier to implement and easier to use, since they are called directly rather than through the `channel` namespace or as a method on a channel. Plugin factories and operators are not supported when [static typing][static-typing-page] is enabled, so plugins that are used in typed pipelines should provide them as functions. - -A factory function returns a channel: +Plugin factories and operators are not supported when [static typing][static-typing-page] is enabled. Provide them as [functions](#functions) instead: ```groovy import groovyx.gpars.dataflow.DataflowWriteChannel +import nextflow.Session +import nextflow.extension.CH import nextflow.plugin.extension.Function +import nextflow.plugin.extension.PluginExtensionPoint -@Function -DataflowWriteChannel fromQuery(Map opts, String query) { - // ... -} -``` +class MyExtension extends PluginExtensionPoint { -An operator function takes a channel as an argument. The channel is received as a `DataflowWriteChannel`, which can be converted into a channel for reading with `CH.getReadChannel()`. Named arguments are received as a `Map` in the first parameter: + @Override + void init(Session session) {} -```groovy -import groovyx.gpars.dataflow.DataflowWriteChannel -import nextflow.extension.CH -import nextflow.plugin.extension.Function + @Function + DataflowWriteChannel fromQuery(Map opts, String query) { + // ... + } + + @Function + DataflowWriteChannel sqlInsert(Map opts, DataflowWriteChannel source) { + final input = CH.getReadChannel(source) + // ... + } -@Function -DataflowWriteChannel sqlInsert(Map opts, DataflowWriteChannel source) { - final input = CH.getReadChannel(source) - // ... } ``` @@ -439,10 +439,6 @@ workflow { } ``` -Before creating a custom operator, consider whether the operator can be composed from existing operators such as `map` or `subscribe`, using a function to implement the custom logic. This approach works in any pipeline, with or without static typing. - -A method cannot be annotated with both `Function` and `Factory` or `Operator`. - ### Process directives Plugins that implement a custom executor will likely need to access [process directives][process-directives] that affect the task execution. When an executor receives a task, the process directives can be accessed through that task’s configuration. Custom executors should try to support all process directives that have executor-specific behavior and are relevant to the executor. diff --git a/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy b/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy index 2771cbe1e5..486ef514f4 100644 --- a/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy +++ b/modules/nextflow/src/test/groovy/nextflow/extension/plugin/PluginExtensionProviderTest.groovy @@ -22,9 +22,8 @@ import groovyx.gpars.dataflow.DataflowReadChannel import groovyx.gpars.dataflow.DataflowVariable import groovyx.gpars.dataflow.DataflowWriteChannel import nextflow.Channel -import nextflow.plugin.extension.Factory +import nextflow.Session import nextflow.plugin.extension.Function -import nextflow.plugin.extension.Operator import nextflow.plugin.extension.PluginExtensionPoint import nextflow.plugin.extension.PluginExtensionProvider import spock.lang.Specification @@ -69,19 +68,9 @@ class PluginExtensionProviderTest extends Specification { } static class ChannelFunctions extends PluginExtensionPoint { - @Override protected void init(nextflow.Session session) {} + @Override protected void init(Session session) {} @Function DataflowWriteChannel fromFoo(String value) { null } - @Function DataflowWriteChannel mapFoo(DataflowReadChannel source) { null } - } - - static class FunctionAndFactory extends PluginExtensionPoint { - @Override protected void init(nextflow.Session session) {} - @Function @Factory DataflowWriteChannel fromFoo(String value) { null } - } - - static class FunctionAndOperator extends PluginExtensionPoint { - @Override protected void init(nextflow.Session session) {} - @Function @Operator DataflowWriteChannel mapFoo(DataflowReadChannel source) { null } + @Function DataflowWriteChannel mapFoo(DataflowWriteChannel source) { null } } def 'should not detect channel functions as factories or operators' () { @@ -90,15 +79,4 @@ class PluginExtensionProviderTest extends Specification { PluginExtensionProvider.getDeclaredOperatorExtensionMethods0(ChannelFunctions).isEmpty() PluginExtensionProvider.getDeclaredFunctionsExtensionMethods0(ChannelFunctions) == ['fromFoo', 'mapFoo'] as Set } - - def 'should reject function that is also a factory or operator' () { - when: - PluginExtensionProvider.getDeclaredFunctionsExtensionMethods0(CLAZZ) - then: - def e = thrown(IllegalStateException) - e.message.contains('cannot also be declared as a factory or operator') - - where: - CLAZZ << [FunctionAndFactory, FunctionAndOperator] - } } diff --git a/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy b/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy index 1c148cb297..c4b3c1f803 100644 --- a/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy +++ b/modules/nf-commons/src/test/nextflow/plugin/extension/PluginExtensionMethodsTest.groovy @@ -512,26 +512,6 @@ class PluginExtensionMethodsTest extends Dsl2Spec { result.val == Channel.STOP } - def 'should wrap custom function channel as typed channel'() { - given: - def SCRIPT_TEXT = ''' - nextflow.enable.types = true - - include { reverseFn } from 'plugin/nf-test-plugin-hello' - - workflow { - def ch = reverseFn('a string') - channel.value(ch.getClass().getName()) - } - ''' - - when: - def result = runScript(SCRIPT_TEXT) - - then: - result.val == 'nextflow.dataflow.ChannelImpl' - } - def 'should apply typed operators to custom function channel'() { given: def SCRIPT_TEXT = ''' @@ -592,4 +572,24 @@ class PluginExtensionMethodsTest extends Dsl2Spec { result.val as List == ['BYE BYE FOLKS'] } + def 'should compose channel functions with named args'() { + when: + def result = runScript(""" + ${HEADER} + + include { reverseFn; goodbyeFn } from 'plugin/nf-test-plugin-hello' + + workflow { + def rows = reverseFn('a string', upper: true) + goodbyeFn(rows, prefix: 'x-') + } + """) + then: + result.val == 'x-GNIRTS A' + result.val == Channel.STOP + + where: + HEADER << ['', 'nextflow.enable.types = true'] + } + } diff --git a/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy b/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy index 4e90c4ebb6..22f3bb6406 100644 --- a/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy +++ b/modules/nf-commons/src/testFixtures/groovy/nextflow/plugin/hello/HelloExtension.groovy @@ -140,16 +140,21 @@ class HelloExtension extends PluginExtensionPoint { * A @Function that returns a channel, i.e. a factory as a plain function */ @Function - DataflowWriteChannel reverseFn(String message) { - return reverse(message) + DataflowWriteChannel reverseFn(Map opts = [:], String message) { + return reverse(opts.upper ? message.toUpperCase() : message) } /** * A @Function that takes a channel, i.e. an operator as a plain function */ @Function - DataflowWriteChannel goodbyeFn(DataflowWriteChannel source) { - return goodbye(CH.getReadChannel(source)) + DataflowWriteChannel goodbyeFn(Map opts = [:], DataflowWriteChannel source) { + final prefix = opts.prefix ?: '' + final target = CH.create() + final next = { target.bind(prefix + it) } + final done = { target.bind(Channel.STOP) } + DataflowHelper.subscribeImpl(CH.getReadChannel(source), [onNext: next, onComplete: done]) + return target } String aNonImportedFunction(){