diff --git a/k8s-dvfs/.gitignore b/k8s-dvfs/.gitignore new file mode 100644 index 0000000..dbef60b --- /dev/null +++ b/k8s-dvfs/.gitignore @@ -0,0 +1,8 @@ +# Ignore Gradle project-specific cache directory +.gradle +.idea +.nextflow* + +# Ignore Gradle build output directory +build +work diff --git a/k8s-dvfs/COPYING b/k8s-dvfs/COPYING new file mode 100644 index 0000000..68c771a --- /dev/null +++ b/k8s-dvfs/COPYING @@ -0,0 +1,176 @@ + + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + diff --git a/k8s-dvfs/Makefile b/k8s-dvfs/Makefile new file mode 100644 index 0000000..1ad90a4 --- /dev/null +++ b/k8s-dvfs/Makefile @@ -0,0 +1,25 @@ +# Use the Gradle wrapper by default; override with e.g. `make GRADLE=gradle ...` +# or `export GRADLE=gradle` (useful in pixi/conda environments). +GRADLE ?= ./gradlew + +# Build the plugin +assemble: + $(GRADLE) assemble + +clean: + rm -rf .nextflow* + rm -rf work + rm -rf build + $(GRADLE) clean + +# Run plugin unit tests +test: + $(GRADLE) test + +# Install the plugin into local nextflow plugins dir +install: + $(GRADLE) install + +# Publish the plugin +release: + $(GRADLE) releasePlugin diff --git a/k8s-dvfs/README.md b/k8s-dvfs/README.md new file mode 100644 index 0000000..a83d19a --- /dev/null +++ b/k8s-dvfs/README.md @@ -0,0 +1,86 @@ +# k8s-dvfs + +## Summary + +`k8s-dvfs` is a Nextflow plugin scaffolded from the official plugin +template. Out of the box it provides: + +- A custom function `sayHello` that can be imported into Nextflow scripts. +- A workflow observer that reacts to pipeline lifecycle events (start and + completion). + +Replace this section with a description of what your plugin actually does. + +Note: The **Summary**, **Get Started**, **Examples**, and **License** sections are +mandatory: they are required by the Nextflow Registry, which uses this +file as the plugin description. The **Plugin development** section below is +guidance for working on the plugin and can be removed before publishing. + +## Get Started + +Enable the plugin in your pipeline `nextflow.config`: + +```groovy +plugins { + id 'k8s-dvfs@0.1.0' +} +``` + +Nextflow downloads the plugin from the Nextflow Registry the first time +the pipeline runs. + +## Examples + +Import and call the `sayHello` function from a Nextflow script: + +```nextflow +include { sayHello } from 'plugin/k8s-dvfs' + +workflow { + channel.of('Mundo', 'World').map { target -> sayHello(target) } +} +``` + +The bundled observer prints a message when the pipeline starts and completes, +so running any pipeline with the plugin enabled produces: + +``` +Pipeline is starting! 🚀 +Pipeline complete! 👋 +``` + +## Plugin development + +This project was created from the [Nextflow plugin template](https://www.nextflow.io/docs/latest/guides/gradle-plugin.html#gradle-plugin-create). + +### Building + +To build the plugin: + +```bash +make assemble +``` + +### Testing with Nextflow + +The plugin can be tested without a local Nextflow installation: + +1. Build and install the plugin to your local Nextflow installation: `make install` +2. Run a pipeline with the plugin: `nextflow run hello -plugins k8s-dvfs@0.1.0` + +### Publishing + +Plugins can be published to a central Nextflow registry to make them accessible to the Nextflow community. + +Follow these steps to publish the plugin to the Nextflow Registry: + +1. Create a file named `$HOME/.gradle/gradle.properties`, where `$HOME` is your home directory. Add the following properties: + * `npr.apiKey`: Your Nextflow Registry access token. +2. Package your plugin and publish it to the registry: `make release`. + +## License + +Apache License 2.0. See the [`COPYING`](COPYING) file for details. + +Note: The above license is given for guidance only; however the Nextflow Registry +requires the plugin to include an OSS (open source software) license. diff --git a/k8s-dvfs/build.gradle b/k8s-dvfs/build.gradle new file mode 100644 index 0000000..1f507d8 --- /dev/null +++ b/k8s-dvfs/build.gradle @@ -0,0 +1,17 @@ +plugins { + id 'io.nextflow.nextflow-plugin' version '1.0.0-beta.15' +} + +version = '0.1.0' + +nextflowPlugin { + nextflowVersion = '25.10.0' + + provider = 'recreational.tech' + className = 'recreationaltech.plugin.K8sDvfsPlugin' + extensionPoints = [ + 'recreationaltech.plugin.K8sDvfsExtension', + 'recreationaltech.plugin.K8sDvfsFactory' + ] + +} diff --git a/k8s-dvfs/gradle/wrapper/gradle-wrapper.jar b/k8s-dvfs/gradle/wrapper/gradle-wrapper.jar new file mode 100644 index 0000000..7f93135 Binary files /dev/null and b/k8s-dvfs/gradle/wrapper/gradle-wrapper.jar differ diff --git a/k8s-dvfs/gradle/wrapper/gradle-wrapper.properties b/k8s-dvfs/gradle/wrapper/gradle-wrapper.properties new file mode 100644 index 0000000..ca025c8 --- /dev/null +++ b/k8s-dvfs/gradle/wrapper/gradle-wrapper.properties @@ -0,0 +1,7 @@ +distributionBase=GRADLE_USER_HOME +distributionPath=wrapper/dists +distributionUrl=https\://services.gradle.org/distributions/gradle-8.14-bin.zip +networkTimeout=10000 +validateDistributionUrl=true +zipStoreBase=GRADLE_USER_HOME +zipStorePath=wrapper/dists diff --git a/k8s-dvfs/gradlew b/k8s-dvfs/gradlew new file mode 100755 index 0000000..1aa94a4 --- /dev/null +++ b/k8s-dvfs/gradlew @@ -0,0 +1,249 @@ +#!/bin/sh + +# +# Copyright © 2015-2021 the original authors. +# +# 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 +# +# https://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. +# + +############################################################################## +# +# Gradle start up script for POSIX generated by Gradle. +# +# Important for running: +# +# (1) You need a POSIX-compliant shell to run this script. If your /bin/sh is +# noncompliant, but you have some other compliant shell such as ksh or +# bash, then to run this script, type that shell name before the whole +# command line, like: +# +# ksh Gradle +# +# Busybox and similar reduced shells will NOT work, because this script +# requires all of these POSIX shell features: +# * functions; +# * expansions «$var», «${var}», «${var:-default}», «${var+SET}», +# «${var#prefix}», «${var%suffix}», and «$( cmd )»; +# * compound commands having a testable exit status, especially «case»; +# * various built-in commands including «command», «set», and «ulimit». +# +# Important for patching: +# +# (2) This script targets any POSIX shell, so it avoids extensions provided +# by Bash, Ksh, etc; in particular arrays are avoided. +# +# The "traditional" practice of packing multiple parameters into a +# space-separated string is a well documented source of bugs and security +# problems, so this is (mostly) avoided, by progressively accumulating +# options in "$@", and eventually passing that to Java. +# +# Where the inherited environment variables (DEFAULT_JVM_OPTS, JAVA_OPTS, +# and GRADLE_OPTS) rely on word-splitting, this is performed explicitly; +# see the in-line comments for details. +# +# There are tweaks for specific operating systems such as AIX, CygWin, +# Darwin, MinGW, and NonStop. +# +# (3) This script is generated from the Groovy template +# https://github.com/gradle/gradle/blob/HEAD/subprojects/plugins/src/main/resources/org/gradle/api/internal/plugins/unixStartScript.txt +# within the Gradle project. +# +# You can find Gradle at https://github.com/gradle/gradle/. +# +############################################################################## + +# Attempt to set APP_HOME + +# Resolve links: $0 may be a link +app_path=$0 + +# Need this for daisy-chained symlinks. +while + APP_HOME=${app_path%"${app_path##*/}"} # leaves a trailing /; empty if no leading path + [ -h "$app_path" ] +do + ls=$( ls -ld "$app_path" ) + link=${ls#*' -> '} + case $link in #( + /*) app_path=$link ;; #( + *) app_path=$APP_HOME$link ;; + esac +done + +# This is normally unused +# shellcheck disable=SC2034 +APP_BASE_NAME=${0##*/} +# Discard cd standard output in case $CDPATH is set (https://github.com/gradle/gradle/issues/25036) +APP_HOME=$( cd "${APP_HOME:-./}" > /dev/null && pwd -P ) || exit + +# Use the maximum available, or set MAX_FD != -1 to use that value. +MAX_FD=maximum + +warn () { + echo "$*" +} >&2 + +die () { + echo + echo "$*" + echo + exit 1 +} >&2 + +# OS specific support (must be 'true' or 'false'). +cygwin=false +msys=false +darwin=false +nonstop=false +case "$( uname )" in #( + CYGWIN* ) cygwin=true ;; #( + Darwin* ) darwin=true ;; #( + MSYS* | MINGW* ) msys=true ;; #( + NONSTOP* ) nonstop=true ;; +esac + +CLASSPATH=$APP_HOME/gradle/wrapper/gradle-wrapper.jar + + +# Determine the Java command to use to start the JVM. +if [ -n "$JAVA_HOME" ] ; then + if [ -x "$JAVA_HOME/jre/sh/java" ] ; then + # IBM's JDK on AIX uses strange locations for the executables + JAVACMD=$JAVA_HOME/jre/sh/java + else + JAVACMD=$JAVA_HOME/bin/java + fi + if [ ! -x "$JAVACMD" ] ; then + die "ERROR: JAVA_HOME is set to an invalid directory: $JAVA_HOME + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +else + JAVACMD=java + if ! command -v java >/dev/null 2>&1 + then + die "ERROR: JAVA_HOME is not set and no 'java' command could be found in your PATH. + +Please set the JAVA_HOME variable in your environment to match the +location of your Java installation." + fi +fi + +# Increase the maximum file descriptors if we can. +if ! "$cygwin" && ! "$darwin" && ! "$nonstop" ; then + case $MAX_FD in #( + max*) + # In POSIX sh, ulimit -H is undefined. That's why the result is checked to see if it worked. + # shellcheck disable=SC2039,SC3045 + MAX_FD=$( ulimit -H -n ) || + warn "Could not query maximum file descriptor limit" + esac + case $MAX_FD in #( + '' | soft) :;; #( + *) + # In POSIX sh, ulimit -n is undefined. That's why the result is checked to see if it worked. + # shellcheck disable=SC2039,SC3045 + ulimit -n "$MAX_FD" || + warn "Could not set maximum file descriptor limit to $MAX_FD" + esac +fi + +# Collect all arguments for the java command, stacking in reverse order: +# * args from the command line +# * the main class name +# * -classpath +# * -D...appname settings +# * --module-path (only if needed) +# * DEFAULT_JVM_OPTS, JAVA_OPTS, and GRADLE_OPTS environment variables. + +# For Cygwin or MSYS, switch paths to Windows format before running java +if "$cygwin" || "$msys" ; then + APP_HOME=$( cygpath --path --mixed "$APP_HOME" ) + CLASSPATH=$( cygpath --path --mixed "$CLASSPATH" ) + + JAVACMD=$( cygpath --unix "$JAVACMD" ) + + # Now convert the arguments - kludge to limit ourselves to /bin/sh + for arg do + if + case $arg in #( + -*) false ;; # don't mess with options #( + /?*) t=${arg#/} t=/${t%%/*} # looks like a POSIX filepath + [ -e "$t" ] ;; #( + *) false ;; + esac + then + arg=$( cygpath --path --ignore --mixed "$arg" ) + fi + # Roll the args list around exactly as many times as the number of + # args, so each arg winds up back in the position where it started, but + # possibly modified. + # + # NB: a `for` loop captures its iteration list before it begins, so + # changing the positional parameters here affects neither the number of + # iterations, nor the values presented in `arg`. + shift # remove old arg + set -- "$@" "$arg" # push replacement arg + done +fi + + +# Add default JVM options here. You can also use JAVA_OPTS and GRADLE_OPTS to pass JVM options to this script. +DEFAULT_JVM_OPTS='"-Xmx64m" "-Xms64m"' + +# Collect all arguments for the java command: +# * DEFAULT_JVM_OPTS, JAVA_OPTS, JAVA_OPTS, and optsEnvironmentVar are not allowed to contain shell fragments, +# and any embedded shellness will be escaped. +# * For example: A user cannot expect ${Hostname} to be expanded, as it is an environment variable and will be +# treated as '${Hostname}' itself on the command line. + +set -- \ + "-Dorg.gradle.appname=$APP_BASE_NAME" \ + -classpath "$CLASSPATH" \ + org.gradle.wrapper.GradleWrapperMain \ + "$@" + +# Stop when "xargs" is not available. +if ! command -v xargs >/dev/null 2>&1 +then + die "xargs is not available" +fi + +# Use "xargs" to parse quoted args. +# +# With -n1 it outputs one arg per line, with the quotes and backslashes removed. +# +# In Bash we could simply go: +# +# readarray ARGS < <( xargs -n1 <<<"$var" ) && +# set -- "${ARGS[@]}" "$@" +# +# but POSIX shell has neither arrays nor command substitution, so instead we +# post-process each arg (as a line of input to sed) to backslash-escape any +# character that might be a shell metacharacter, then use eval to reverse +# that process (while maintaining the separation between arguments), and wrap +# the whole thing up as a single "set" statement. +# +# This will of course break if any of these variables contains a newline or +# an unmatched quote. +# + +eval "set -- $( + printf '%s\n' "$DEFAULT_JVM_OPTS $JAVA_OPTS $GRADLE_OPTS" | + xargs -n1 | + sed ' s~[^-[:alnum:]+,./:=@_]~\\&~g; ' | + tr '\n' ' ' + )" '"$@"' + +exec "$JAVACMD" "$@" diff --git a/k8s-dvfs/settings.gradle b/k8s-dvfs/settings.gradle new file mode 100644 index 0000000..5401841 --- /dev/null +++ b/k8s-dvfs/settings.gradle @@ -0,0 +1 @@ +rootProject.name = 'k8s-dvfs' diff --git a/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsExtension.groovy b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsExtension.groovy new file mode 100644 index 0000000..e4c2a86 --- /dev/null +++ b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsExtension.groovy @@ -0,0 +1,45 @@ +/* + * Copyright 2025, 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 recreationaltech.plugin + +import groovy.transform.CompileStatic +import nextflow.Session +import nextflow.plugin.extension.Function +import nextflow.plugin.extension.PluginExtensionPoint + +/** + * Implements a custom function which can be imported by + * Nextflow scripts. + */ +@CompileStatic +class K8sDvfsExtension extends PluginExtensionPoint { + + @Override + protected void init(Session session) { + } + + /** + * Say hello to the given target. + * + * @param target + */ + @Function + void sayHello(String target) { + println "Hello, ${target}!" + } + +} diff --git a/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsFactory.groovy b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsFactory.groovy new file mode 100644 index 0000000..4779973 --- /dev/null +++ b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsFactory.groovy @@ -0,0 +1,36 @@ +/* + * Copyright 2025, 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 recreationaltech.plugin + +import groovy.transform.CompileStatic +import nextflow.Session +import nextflow.trace.TraceObserver +import nextflow.trace.TraceObserverFactory + +/** + * Implements a factory object required to create + * the {@link K8sDvfsObserver} instance. + */ +@CompileStatic +class K8sDvfsFactory implements TraceObserverFactory { + + @Override + Collection create(Session session) { + return List.of(new K8sDvfsObserver()) + } + +} diff --git a/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsObserver.groovy b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsObserver.groovy new file mode 100644 index 0000000..748116b --- /dev/null +++ b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsObserver.groovy @@ -0,0 +1,41 @@ +/* + * Copyright 2025, 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 recreationaltech.plugin + +import groovy.transform.CompileStatic +import groovy.util.logging.Slf4j +import nextflow.Session +import nextflow.trace.TraceObserver + +/** + * Implements an observer that allows implementing custom + * logic on nextflow execution events. + */ +@Slf4j +@CompileStatic +class K8sDvfsObserver implements TraceObserver { + + @Override + void onFlowCreate(Session session) { + println "Pipeline is starting! 🚀" + } + + @Override + void onFlowComplete() { + println "Pipeline complete! 👋" + } +} diff --git a/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsPlugin.groovy b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsPlugin.groovy new file mode 100644 index 0000000..0fb4fbb --- /dev/null +++ b/k8s-dvfs/src/main/groovy/recreationaltech/plugin/K8sDvfsPlugin.groovy @@ -0,0 +1,32 @@ +/* + * Copyright 2025, 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 recreationaltech.plugin + +import groovy.transform.CompileStatic +import nextflow.plugin.BasePlugin +import org.pf4j.PluginWrapper + +/** + * The plugin entry point + */ +@CompileStatic +class K8sDvfsPlugin extends BasePlugin { + + K8sDvfsPlugin(PluginWrapper wrapper) { + super(wrapper) + } +} diff --git a/k8s-dvfs/src/test/groovy/recreationaltech/plugin/K8sDvfsObserverTest.groovy b/k8s-dvfs/src/test/groovy/recreationaltech/plugin/K8sDvfsObserverTest.groovy new file mode 100644 index 0000000..b76b4f8 --- /dev/null +++ b/k8s-dvfs/src/test/groovy/recreationaltech/plugin/K8sDvfsObserverTest.groovy @@ -0,0 +1,22 @@ +package recreationaltech.plugin + +import nextflow.Session +import spock.lang.Specification + +/** + * Implements a basic factory test + * + */ +class K8sDvfsObserverTest extends Specification { + + def 'should create the observer instance' () { + given: + def factory = new K8sDvfsFactory() + when: + def result = factory.create(Mock(Session)) + then: + result.size() == 1 + result.first() instanceof K8sDvfsObserver + } + +} diff --git a/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sConfig.groovy b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sConfig.groovy index 5261ab9..2cdaaa6 100644 --- a/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sConfig.groovy +++ b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sConfig.groovy @@ -260,6 +260,18 @@ class K8sConfig implements ConfigScope { """) final Duration noiseRuntimeEstimatorNoiseMagnitude + @ConfigOption + @Description(""" + Max. amount of time two runtime estimates can differ to be considered equal. + """) + final Duration runtimeComparisonEpsilon + + @ConfigOption + @Description(""" + Number of saved top runtimes used to classify a task as critical. + """) + final int dvfsSchedulingNumTopRuntimes + /* required by extension point -- do not remove */ K8sConfig() { this(Collections.emptyMap()) @@ -300,7 +312,9 @@ class K8sConfig implements ConfigScope { schedulingStrategy = opts.schedulingStrategy as String ?: "Hash" runtimeEstimator = opts.runtimeEstimator as String ?: "LinearFit" - noiseRuntimeEstimatorNoiseMagnitude = opts.noiseRuntimeEstimatorNoiseMagnitude as Duration ?: new Duration(10, TimeUnit.SECONDS) + noiseRuntimeEstimatorNoiseMagnitude = opts.noiseRuntimeEstimatorNoiseMagnitude as Duration ?: new Duration(30, TimeUnit.SECONDS) + runtimeComparisonEpsilon = opts.runtimeComparisonEpsilon as Duration ?: new Duration(10, TimeUnit.SECONDS) + dvfsSchedulingNumTopRuntimes = opts.dvfsSchedulingTopRuntimes as int ?: 3 // -- shortcut to pod image pull-policy if( imagePullPolicy ) diff --git a/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sExecutor.groovy b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sExecutor.groovy index ceff9b9..72b9dca 100644 --- a/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sExecutor.groovy +++ b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/K8sExecutor.groovy @@ -112,7 +112,7 @@ class K8sExecutor extends Executor implements ExtensionPoint { K8sSchedulingStrategy strategy = null if (k8sConfig.schedulingStrategy == "Hash") { strategy = new K8sHashSchedulingStrategy() - } else if (k8sConfig.schedulingStrategy == "DVFS") { + } else if (k8sConfig.schedulingStrategy == "DVFS" || k8sConfig.schedulingStrategy == "DVFS-SPEED") { String[] ips = new String[nodes.length] for (int i = 0; i < nodes.length; i++) { ips[i] = client.getPodIpAddress(K8sNodeInitDeployer.buildPodName(nodes[i])) @@ -120,18 +120,10 @@ class K8sExecutor extends Executor implements ExtensionPoint { } strategy = new K8sDVFSSchedulingStrategy(this.runtimeEstimator, new K8sDVFSClient(nodes, ips), - () -> getClient()) - strategy.fullSpeedMode = false - } else if (k8sConfig.schedulingStrategy == "DVFS-SPEED") { - String[] ips = new String[nodes.length] - for (int i = 0; i < nodes.length; i++) { - ips[i] = client.getPodIpAddress(K8sNodeInitDeployer.buildPodName(nodes[i])) - log.info "[K8s] node ${nodes[i]} -> ${ips[i]}" - } - strategy = new K8sDVFSSchedulingStrategy(this.runtimeEstimator, - new K8sDVFSClient(nodes, ips), - () -> getClient()) - strategy.fullSpeedMode = true + () -> getClient(), + k8sConfig.runtimeComparisonEpsilon, + k8sConfig.dvfsSchedulingNumTopRuntimes) + strategy.fullSpeedMode = k8sConfig.schedulingStrategy == "DVFS-SPEED" } else { log.error "[K8s] invalid scheduling strategy $k8sConfig.schedulingStrategy, falling back on \"Hash\"" strategy = new K8sHashSchedulingStrategy() diff --git a/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/strategies/K8sDVFSSchedulingStrategy.groovy b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/strategies/K8sDVFSSchedulingStrategy.groovy index 2b14d44..8373d53 100644 --- a/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/strategies/K8sDVFSSchedulingStrategy.groovy +++ b/nextflow/plugins/nf-k8s/src/main/nextflow/k8s/strategies/K8sDVFSSchedulingStrategy.groovy @@ -11,7 +11,7 @@ import nextflow.k8s.K8sTaskHandler import nextflow.k8s.K8sTaskScheduler import nextflow.k8s.client.K8sClient import nextflow.processor.TaskRun -import nextflow.util.ArrayTuple +import nextflow.util.Duration /** * Implements a scheduling strategy utilizing dvfs to reduce the energy consumption @@ -142,16 +142,17 @@ class K8sDVFSSchedulingStrategy implements K8sSchedulingStrategy { class SchedulingRequestComparator implements Comparator { K8sRuntimeEstimator runtimeEstimator long currentTime - double avgRuntime + double epsilon @Override int compare(K8sSchedulingRequest o1, K8sSchedulingRequest o2) { // First, check if one of the tasks is (estimated to be) on the critical path double t1 = runtimeEstimator.estimate(o1.handler) double t2 = runtimeEstimator.estimate(o2.handler) - if (t1 > avgRuntime && t2 <= avgRuntime) + + if (t1 > t2 + epsilon) return -1 - else if (t1 < avgRuntime && t2 > avgRuntime) + else if (t2 > t1 + epsilon) return 1 // Both are not on the critical path. Sort based on the time they spent in the queue @@ -170,23 +171,57 @@ class K8sDVFSSchedulingStrategy implements K8sSchedulingStrategy { private ArrayList nodes private HashMap taskToNode - - private double averageRuntime - private long finishedTaskCount - private long globalMaxFrequency private long globalMinFrequency private K8sClientGetter clientGetter + private double comparisonEpsilonMillis + + private double[] topRuntimes + private double averageRuntime + private double finishedTaskCount + boolean fullSpeedMode - K8sDVFSSchedulingStrategy(K8sRuntimeEstimator runtimeEstimator, K8sDVFSClient dvfsClient, K8sClientGetter clientGetter) { + K8sDVFSSchedulingStrategy(K8sRuntimeEstimator runtimeEstimator, + K8sDVFSClient dvfsClient, + K8sClientGetter clientGetter, + Duration runtimeComparisonEpsilon, + int topRuntimeCount) { this.runtimeEstimator = runtimeEstimator this.dvfsClient = dvfsClient this.nodes = new ArrayList<>() this.taskToNode = new HashMap<>(); this.clientGetter = clientGetter + this.comparisonEpsilonMillis = (double)runtimeComparisonEpsilon.toMillis() + this.topRuntimes = new double[topRuntimeCount] + for (int i = 0; i < topRuntimeCount; i++) { + this.topRuntimes[i] = 0.0 + } + this.averageRuntime = 0.0 + this.finishedTaskCount = 0.0 + } + + private boolean isInTopRuntimes(double rt) { + for (int i = 0; i < topRuntimes.size(); i++) { + if (rt >= topRuntimes[i]) + return true + } + return false + } + + private void updateTopRuntimes(double rt) { + for (int i = 0; i < topRuntimes.size(); i++) { + if (rt > topRuntimes[i]) { + /* Move all one down */ + for (int j = topRuntimes.size() - 1; j > i; j--) { + topRuntimes[j] = topRuntimes[j - 1]; + } + topRuntimes[i] = rt + break + } + } } @Override @@ -199,11 +234,11 @@ class K8sDVFSSchedulingStrategy implements K8sSchedulingStrategy { /* Step 1: Sort by task priority. We will attempt to schedule tasks "in order", so that the * highest priority tasks are assigned to nodes as soon as possible. * - * Priority is based on a) an estimation if the task is on the critical path and b) the wait time of the task. + * Priority is based on a) the tasks estimated runtime and b) the wait time of the task. */ SchedulingRequestComparator comparator = new SchedulingRequestComparator() comparator.runtimeEstimator = runtimeEstimator - comparator.avgRuntime = averageRuntime + comparator.epsilon = comparisonEpsilonMillis comparator.currentTime = System.currentTimeMillis() queue.sort(comparator) @@ -214,7 +249,7 @@ class K8sDVFSSchedulingStrategy implements K8sSchedulingStrategy { * frequency. If not, we determine a frequency (see below). */ final double taskEstimation = runtimeEstimator.estimate(req.handler) - final boolean isCriticalPath = taskEstimation > averageRuntime + final boolean isCriticalPath = isInTopRuntimes(taskEstimation) long frequency = globalMaxFrequency if (!isCriticalPath && !fullSpeedMode) { /* Set frequency so that we expect the runtime to be close to the mean runtime. */ @@ -280,8 +315,9 @@ class K8sDVFSSchedulingStrategy implements K8sSchedulingStrategy { * elapsed time based on that. */ double runtime = (double)(task.getCompleteTimeMillis() - task.getStartTimeMillis()) - averageRuntime = (runtime + finishedTaskCount * averageRuntime) / (finishedTaskCount + 1) - finishedTaskCount += 1 + averageRuntime = (runtime + finishedTaskCount * averageRuntime) / (finishedTaskCount + 1.0) + finishedTaskCount += 1.0 + updateTopRuntimes(runtime) /* Free resources allocated by this task */ WorkerNode node = taskToNode.get(task.task.hash.toString()) diff --git a/test/nextflow.config b/test/nextflow.config index 5f28eb1..88caf5e 100644 --- a/test/nextflow.config +++ b/test/nextflow.config @@ -22,11 +22,13 @@ k8s { projectDir = '/workspace/projects' cleanup = true - nextflowImage = 'gitea.kleine.eulenhexe.de/kevin/ma/nextflow-dvfs:0.8.18' + nextflowImage = 'gitea.kleine.eulenhexe.de/kevin/ma/nextflow-dvfs:0.8.19' imagePullPolicy = 'IfNotPresent' schedulerInterval = '1m' recordTaskRuntimes = true - schedulingStrategy = 'DVFS-SPEED' + schedulingStrategy = 'DVFS' + runtimeComparisonEpsilon = '10s' + dvfsSchedulingTopRuntimes = 3 nodeInit { enabled = true