Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 41 additions & 3 deletions docs/plugins/developing-plugins.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -398,9 +398,46 @@ 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.
:::
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.

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

class MyExtension extends PluginExtensionPoint {

@Override
void init(Session session) {}

@Function
DataflowWriteChannel fromQuery(Map opts, String query) {
// ...
}

@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')
}
```

### Process directives

Expand Down Expand Up @@ -537,3 +574,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
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -213,6 +216,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
Expand All @@ -238,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")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,11 @@ 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.Session
import nextflow.plugin.extension.Function
import nextflow.plugin.extension.PluginExtensionPoint
import nextflow.plugin.extension.PluginExtensionProvider
import spock.lang.Specification

Expand Down Expand Up @@ -62,4 +66,17 @@ class PluginExtensionProviderTest extends Specification {
result.val == 4
result.val == 9
}

static class ChannelFunctions extends PluginExtensionPoint {
@Override protected void init(Session session) {}
@Function DataflowWriteChannel fromFoo(String value) { null }
@Function DataflowWriteChannel mapFoo(DataflowWriteChannel 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
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -492,4 +492,104 @@ 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 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()
}

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']
}

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']
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,27 @@ 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(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(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(){
throw new IllegalAccessException("This function can't be imported")
}
Expand Down
Loading