From 9a8257423032e2069b2837576d7c9504eb80e61f Mon Sep 17 00:00:00 2001 From: Abhishek Date: Wed, 16 Sep 2026 20:10:17 +0530 Subject: [PATCH 1/2] feat/r-script-commands-triggers --- AGENTS.md | 2 + plugin-script-r/build.gradle | 3 + .../plugin/scripts/r/CommandsTrigger.java | 294 ++++++++++++++++++ .../plugin/scripts/r/ScriptTrigger.java | 290 +++++++++++++++++ .../r/CommandsTriggerConditionTest.java | 65 ++++ .../plugin/scripts/r/CommandsTriggerTest.java | 115 +++++++ .../scripts/r/ScriptTriggerConditionTest.java | 59 ++++ .../plugin/scripts/r/ScriptTriggerTest.java | 95 ++++++ 8 files changed, 923 insertions(+) create mode 100644 plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/CommandsTrigger.java create mode 100644 plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/ScriptTrigger.java create mode 100644 plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerConditionTest.java create mode 100644 plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerTest.java create mode 100644 plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerConditionTest.java create mode 100644 plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerTest.java diff --git a/AGENTS.md b/AGENTS.md index 84c7163d..f7458d3f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -108,7 +108,9 @@ This is a **multi-module** plugin with 19 submodules: **plugin-script-r:** - `io.kestra.plugin.scripts.r.Commands` +- `io.kestra.plugin.scripts.r.CommandsTrigger` - `io.kestra.plugin.scripts.r.Script` +- `io.kestra.plugin.scripts.r.ScriptTrigger` **plugin-script-ruby:** - `io.kestra.plugin.scripts.ruby.Commands` diff --git a/plugin-script-r/build.gradle b/plugin-script-r/build.gradle index cefb63f3..acccc504 100644 --- a/plugin-script-r/build.gradle +++ b/plugin-script-r/build.gradle @@ -16,4 +16,7 @@ dependencies { implementation project(':plugin-script') testImplementation project(path: ':plugin-script', configuration: 'testOutput') + + testImplementation group: "io.kestra", name: "scheduler", version: kestraVersion + testImplementation group: "io.kestra", name: "worker", version: kestraVersion } diff --git a/plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/CommandsTrigger.java b/plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/CommandsTrigger.java new file mode 100644 index 00000000..432d5f21 --- /dev/null +++ b/plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/CommandsTrigger.java @@ -0,0 +1,294 @@ +package io.kestra.plugin.scripts.r; + +import io.kestra.core.models.annotations.Example; +import io.kestra.core.models.annotations.Plugin; +import io.kestra.core.models.annotations.PluginProperty; +import io.kestra.core.models.conditions.ConditionContext; +import io.kestra.core.models.executions.Execution; +import io.kestra.core.models.property.Property; +import io.kestra.core.models.tasks.RunnableTaskException; +import io.kestra.core.models.tasks.runners.TaskException; +import io.kestra.core.models.triggers.*; +import io.kestra.core.runners.RunContext; +import io.kestra.plugin.scripts.exec.TriggerRunContext; +import io.kestra.plugin.scripts.exec.scripts.models.ScriptOutput; +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotNull; +import lombok.*; +import lombok.experimental.SuperBuilder; + +import java.time.Duration; +import java.time.Instant; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +@SuperBuilder +@ToString +@EqualsAndHashCode +@Getter +@NoArgsConstructor +@Schema( + title = "Trigger a flow when R commands match a condition", + description = "Polls by running R commands in a container (default image 'r-base') and starts the flow when their result matches the condition." +) +@Plugin( + examples = { + @Example( + title = "Trigger when an R command fails.", + full = true, + code = """ + id: r_commands_trigger + namespace: company.team + + triggers: + - id: on_fail + type: io.kestra.plugin.scripts.r.CommandsTrigger + interval: PT5S + exitCondition: "exit 1" + commands: + - Rscript -e 'stop("boom")' + + tasks: + - id: log + type: io.kestra.plugin.core.log.Log + message: "Triggered with exitCode={{ trigger.exitCode }} (condition={{ trigger.condition }})" + """ + ) + } +) +// TODO: extract shared trigger logic (evaluate, matchesCondition, extractFailure, Output) +// into an AbstractScriptTrigger in plugin-script to reduce duplication across Shell, Node, Ruby, R, etc. +public class CommandsTrigger extends AbstractTrigger + implements PollingTriggerInterface, TriggerOutput { + + private static final String DEFAULT_IMAGE = "r-base"; + + private static final Pattern EXIT_CONDITION_PATTERN = + Pattern.compile("^\\s*exit\\s+(\\d+)\\s*$", Pattern.CASE_INSENSITIVE); + + @Schema( + title = "Docker image used to execute the commands", + description = """ + Container image used by the underlying Commands task to run R commands. + Defaults to 'r-base'. + """ + ) + @Builder.Default + @PluginProperty(group = "execution") + protected Property containerImage = Property.ofValue(DEFAULT_IMAGE); + + @Schema( + title = "R commands to execute", + description = "Commands executed in order on each poll." + ) + @NotNull + @PluginProperty(group = "main") + protected Property> commands; + + @Schema( + title = "Condition to match", + description = """ + Condition evaluated after execution. + + Supported forms: + - 'exit N' + - regex / substring matched against vars + logs + """ + ) + @NotNull + @PluginProperty(group = "main") + protected Property exitCondition; + + @Schema( + title = "Check interval", + description = "Interval between polling evaluations." + ) + @Builder.Default + @PluginProperty(group = "execution") + private final Duration interval = Duration.ofSeconds(60); + + @Schema( + title = "Edge trigger mode", + description = """ + If true, the trigger emits only on a transition from 'not matching' to 'matching' (anti-spam). + If false, the trigger emits on every poll where the condition matches. + """ + ) + @Builder.Default + @PluginProperty(group = "advanced") + protected Property edge = Property.ofValue(true); + + // Known limitation: in-memory only — resets when the trigger is rehydrated (e.g. after restart), + // so edge mode may re-fire once after a scheduler restart. + @Builder.Default + @Getter(AccessLevel.NONE) + private final AtomicBoolean lastMatched = new AtomicBoolean(false); + + @Override + public Optional evaluate(ConditionContext conditionContext, TriggerContext context) throws Exception { + RunContext runContext = conditionContext.getRunContext(); + boolean edgeEnabled = runContext.render(this.edge).as(Boolean.class).orElse(true); + + Output out; + try { + out = runOnce(runContext); + } catch (Exception e) { + runContext.logger().warn("Trigger evaluation failed, returning empty result to avoid blocking the scheduler", e); + return Optional.empty(); + } + + boolean matched = matchesCondition(out); + + boolean emit = edgeEnabled + ? (!lastMatched.getAndSet(matched) && matched) + : matched; + + if (!emit) { + return Optional.empty(); + } + + return Optional.of( + TriggerService.generateExecution(this, conditionContext, context, out) + ); + } + + private Output runOnce(RunContext runContext) throws Exception { + Commands task = Commands.builder() + .id(this.getId()) + .type(Commands.class.getName()) + .containerImage(this.containerImage) + .commands(this.commands) + .build(); + + String renderedCondition = runContext.render(this.exitCondition) + .as(String.class) + .orElse(""); + + try { + ScriptOutput taskOutput = task.run(TriggerRunContext.forEmbeddedTask(runContext, task)); + + return new Output( + Instant.now(), + renderedCondition, + safeExitCode(taskOutput), + safeVars(taskOutput) + ); + } catch (RunnableTaskException e) { + ExtractedFailure failure = extractFailure(e); + return new Output( + Instant.now(), + renderedCondition, + failure.exitCode, + null + ); + } + } + + boolean matchesCondition(Output out) { + String cond = out.getCondition() == null ? "" : out.getCondition().trim(); + + Matcher exitMatcher = EXIT_CONDITION_PATTERN.matcher(cond); + + if (exitMatcher.matches()) { + int expected = Integer.parseInt(exitMatcher.group(1)); + return out.getExitCode() != null && out.getExitCode() == expected; + } + + String haystack = buildHaystack(out); + if (haystack.isEmpty() || cond.isEmpty()) { + return false; + } + + try { + // Guard against catastrophic backtracking (ReDoS) from user-supplied patterns + var pattern = Pattern.compile(cond); + var future = CompletableFuture.supplyAsync( + () -> pattern.matcher(haystack).find() + ); + return future.get(5, TimeUnit.SECONDS); + } catch (TimeoutException te) { + return haystack.contains(cond); + } catch (Exception e) { + return haystack.contains(cond); + } + } + + private String buildHaystack(Output out) { + if (out.getVars() == null || out.getVars().isEmpty()) { + return ""; + } + // Map.toString() produces {key=value, ...} — intentional for substring/regex matching. + return out.getVars().toString(); + } + + private Integer safeExitCode(ScriptOutput taskOutput) { + try { + return taskOutput.getExitCode(); + } catch (Exception ignored) { + return null; + } + } + + private Map safeVars(ScriptOutput taskOutput) { + try { + return taskOutput.getVars(); + } catch (Exception ignored) { + return null; + } + } + + private record ExtractedFailure(Integer exitCode) {} + + private ExtractedFailure extractFailure(RunnableTaskException e) { + Integer exitCode = null; + + Throwable cur = e.getCause(); + while (cur != null) { + if (cur instanceof TaskException te) { + exitCode = te.getExitCode(); + break; + } + cur = cur.getCause(); + } + + return new ExtractedFailure(exitCode); + } + + @Data + @AllArgsConstructor + public static class Output implements io.kestra.core.models.tasks.Output { + @Schema( + title = "Poll timestamp", + description = "Timestamp when this trigger evaluation occurred." + ) + private Instant timestamp; + + @Schema( + title = "Rendered condition", + description = "Rendered value of the exitCondition property for this poll." + ) + private String condition; + + @Schema( + title = "Commands exit code", + description = "Exit code returned by the R process (may be null if not available)." + ) + private Integer exitCode; + + @Schema( + title = "Commands vars", + description = """ + Vars produced by the task (e.g. via ::{"outputs":{...}}:: convention). This is the main structured + way to evaluate non-exit conditions on successful runs. + """ + ) + private Map vars; + } +} diff --git a/plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/ScriptTrigger.java b/plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/ScriptTrigger.java new file mode 100644 index 00000000..34c99e45 --- /dev/null +++ b/plugin-script-r/src/main/java/io/kestra/plugin/scripts/r/ScriptTrigger.java @@ -0,0 +1,290 @@ +package io.kestra.plugin.scripts.r; + +import io.kestra.core.models.annotations.Example; +import io.kestra.core.models.annotations.Plugin; +import io.kestra.core.models.annotations.PluginProperty; +import io.kestra.core.models.conditions.ConditionContext; +import io.kestra.core.models.enums.MonacoLanguages; +import io.kestra.core.models.executions.Execution; +import io.kestra.core.models.property.Property; +import io.kestra.core.models.tasks.RunnableTaskException; +import io.kestra.core.models.tasks.runners.TaskException; +import io.kestra.core.models.triggers.*; +import io.kestra.core.runners.RunContext; +import io.kestra.plugin.scripts.exec.TriggerRunContext; +import io.kestra.plugin.scripts.exec.scripts.models.ScriptOutput; +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotNull; +import lombok.*; +import lombok.experimental.SuperBuilder; + +import java.time.Duration; +import java.time.Instant; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +@SuperBuilder +@ToString +@EqualsAndHashCode +@Getter +@NoArgsConstructor +@Schema( + title = "Trigger a flow when an R script matches a condition", + description = "Polls by running an inline R script in a container (default image 'r-base') and starts the flow when its result matches the condition." +) +@Plugin( + examples = { + @Example( + title = "Trigger when the script fails with exit code 1.", + full = true, + code = """ + id: r_script_trigger + namespace: company.team + + triggers: + - id: script_failure + type: io.kestra.plugin.scripts.r.ScriptTrigger + interval: PT10S + exitCondition: "exit 1" + edge: true + script: | + stop("boom") + + tasks: + - id: log + type: io.kestra.plugin.core.log.Log + message: "Triggered with exitCode={{ trigger.exitCode }} (condition={{ trigger.condition }})" + """ + ) + } +) +// TODO: extract shared trigger logic (evaluate, matchesCondition, extractFailure, Output) +// into an AbstractScriptTrigger in plugin-script to reduce duplication across Shell, Node, Ruby, R, etc. +public class ScriptTrigger extends AbstractTrigger + implements PollingTriggerInterface, TriggerOutput { + + private static final String DEFAULT_IMAGE = "r-base"; + + private static final Pattern EXIT_CONDITION_PATTERN = + Pattern.compile("^\\s*exit\\s+(\\d+)\\s*$", Pattern.CASE_INSENSITIVE); + + @Schema( + title = "Container image for script execution", + description = "Image used by the Script task to run the inline R script; defaults to 'r-base'. Provide an image that includes the required CRAN packages, or install them in the script itself." + ) + @Builder.Default + @PluginProperty(group = "execution") + protected Property containerImage = Property.ofValue(DEFAULT_IMAGE); + + @Schema( + title = "Inline R script", + description = "Multi-line R script executed on each poll, with the same semantics as the R Script task." + ) + @NotNull + @PluginProperty(language = MonacoLanguages.R, group = "main") + protected Property script; + + @Schema( + title = "Condition to match", + description = """ + Condition evaluated after each execution. The trigger emits only when it matches. + 'exit N' compares the exit code, otherwise the string is used as a regex + (or substring fallback) against emitted vars and failure logs. + """ + ) + @NotNull + @PluginProperty(group = "main") + protected Property exitCondition; + + @Schema( + title = "Check interval", + description = "Interval between polling evaluations." + ) + @Builder.Default + @PluginProperty(group = "execution") + private final Duration interval = Duration.ofSeconds(60); + + @Schema( + title = "Edge trigger mode", + description = """ + If true, the trigger emits only on a transition from 'not matching' to 'matching' (anti-spam). + If false, the trigger emits on every poll where the condition matches. + """ + ) + @Builder.Default + @PluginProperty(group = "advanced") + protected Property edge = Property.ofValue(true); + + // Known limitation: in-memory only — resets when the trigger is rehydrated (e.g. after restart), + // so edge mode may re-fire once after a scheduler restart. + @Getter(AccessLevel.NONE) + @Builder.Default + private final AtomicBoolean lastMatched = new AtomicBoolean(false); + + @Override + public Optional evaluate(ConditionContext conditionContext, TriggerContext context) throws Exception { + RunContext runContext = conditionContext.getRunContext(); + boolean edgeEnabled = runContext.render(this.edge).as(Boolean.class).orElse(true); + + Output output; + try { + output = runOnce(runContext); + } catch (Exception e) { + runContext.logger().warn("Trigger evaluation failed, returning empty result to avoid blocking the scheduler", e); + return Optional.empty(); + } + + boolean matched = matchesCondition(output); + + boolean emit = edgeEnabled + ? (!lastMatched.getAndSet(matched) && matched) + : matched; + + if (!emit) { + return Optional.empty(); + } + + return Optional.of( + TriggerService.generateExecution(this, conditionContext, context, output) + ); + } + + private Output runOnce(RunContext runContext) throws Exception { + Script task = Script.builder() + .id(this.getId()) + .type(Script.class.getName()) + .containerImage(this.containerImage) + .script(this.script) + .build(); + + String renderedCondition = runContext.render(this.exitCondition) + .as(String.class) + .orElse(""); + + try { + ScriptOutput taskOutput = task.run(TriggerRunContext.forEmbeddedTask(runContext, task)); + + return new Output( + Instant.now(), + renderedCondition, + safeExitCode(taskOutput), + safeVars(taskOutput) + ); + } catch (RunnableTaskException e) { + ExtractedFailure failure = extractFailure(e); + return new Output( + Instant.now(), + renderedCondition, + failure.exitCode, + null + ); + } + } + + boolean matchesCondition(Output out) { + String cond = out.getCondition() == null ? "" : out.getCondition().trim(); + + Matcher exitMatcher = EXIT_CONDITION_PATTERN.matcher(cond); + + if (exitMatcher.matches()) { + int expected = Integer.parseInt(exitMatcher.group(1)); + return out.getExitCode() != null && out.getExitCode() == expected; + } + + String haystack = buildHaystack(out); + if (haystack.isEmpty() || cond.isEmpty()) { + return false; + } + + try { + // Guard against catastrophic backtracking (ReDoS) from user-supplied patterns + var pattern = Pattern.compile(cond); + var future = CompletableFuture.supplyAsync( + () -> pattern.matcher(haystack).find() + ); + return future.get(5, TimeUnit.SECONDS); + } catch (TimeoutException te) { + return haystack.contains(cond); + } catch (Exception e) { + return haystack.contains(cond); + } + } + + private String buildHaystack(Output out) { + if (out.getVars() == null || out.getVars().isEmpty()) { + return ""; + } + // Map.toString() produces {key=value, ...} — intentional for substring/regex matching. + return out.getVars().toString(); + } + + private Integer safeExitCode(ScriptOutput output) { + try { + return output.getExitCode(); + } catch (Exception ignored) { + return null; + } + } + + private Map safeVars(ScriptOutput output) { + try { + return output.getVars(); + } catch (Exception ignored) { + return null; + } + } + + private record ExtractedFailure(Integer exitCode) {} + + private ExtractedFailure extractFailure(RunnableTaskException e) { + Integer exitCode = null; + + Throwable cur = e.getCause(); + while (cur != null) { + if (cur instanceof TaskException te) { + exitCode = te.getExitCode(); + break; + } + cur = cur.getCause(); + } + + return new ExtractedFailure(exitCode); + } + + @Data + @AllArgsConstructor + public static class Output implements io.kestra.core.models.tasks.Output { + @Schema( + title = "Poll timestamp", + description = "Timestamp when this trigger evaluation occurred." + ) + private Instant timestamp; + + @Schema( + title = "Rendered condition", + description = "Rendered value of the exitCondition property for this poll." + ) + private String condition; + + @Schema( + title = "Script exit code", + description = "Exit code returned by the R process (may be null if not available)." + ) + private Integer exitCode; + + @Schema( + title = "Script vars", + description = """ + Vars produced by the task (e.g. via ::{"outputs":{...}}:: convention). This is the main structured + way to evaluate non-exit conditions on successful runs. + """ + ) + private Map vars; + } +} diff --git a/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerConditionTest.java b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerConditionTest.java new file mode 100644 index 00000000..2a120568 --- /dev/null +++ b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerConditionTest.java @@ -0,0 +1,65 @@ +package io.kestra.plugin.scripts.r; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import java.time.Instant; +import java.util.Map; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; + +class CommandsTriggerConditionTest { + + private final CommandsTrigger trigger = CommandsTrigger.builder().build(); + + private CommandsTrigger.Output output(String condition, Integer exitCode, Map vars) { + return new CommandsTrigger.Output(Instant.now(), condition, exitCode, vars); + } + + @ParameterizedTest + @CsvSource({ + "exit 0, 0, true", + "exit 1, 1, true", + "EXIT 1, 1, true", + "exit 0, 1, false", + "exit 1, 0, false", + "exit 42, 42, true", + }) + void exitCodeCondition(String condition, int exitCode, boolean expected) { + assertThat(trigger.matchesCondition(output(condition, exitCode, null)), is(expected)); + } + + @Test + void exitCondition_nullExitCode_doesNotMatch() { + assertThat(trigger.matchesCondition(output("exit 1", null, null)), is(false)); + } + + @Test + void substringMatch_inVars() { + assertThat(trigger.matchesCondition( + output("toto", 0, Map.of("key", "toto"))), is(true)); + } + + @Test + void regexMatch_inVars() { + assertThat(trigger.matchesCondition( + output("status=\\w+", 0, Map.of("status", "status=ready"))), is(true)); + } + + @Test + void noMatch_emptyHaystack() { + assertThat(trigger.matchesCondition(output("something", 0, null)), is(false)); + } + + @Test + void noMatch_emptyCondition() { + assertThat(trigger.matchesCondition(output("", 0, Map.of("k", "v"))), is(false)); + } + + @Test + void nullCondition_doesNotMatch() { + assertThat(trigger.matchesCondition(output(null, 0, Map.of("k", "v"))), is(false)); + } +} diff --git a/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerTest.java b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerTest.java new file mode 100644 index 00000000..a99cbca4 --- /dev/null +++ b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/CommandsTriggerTest.java @@ -0,0 +1,115 @@ +package io.kestra.plugin.scripts.r; + +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.jupiter.api.Test; + +import io.kestra.core.junit.annotations.KestraTest; +import io.kestra.core.models.executions.Execution; +import io.kestra.core.models.property.Property; +import io.kestra.core.runners.RunContextFactory; +import io.kestra.core.utils.TestsUtils; + +import jakarta.inject.Inject; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.Matchers.notNullValue; + +@KestraTest +class CommandsTriggerTest { + @Inject + private RunContextFactory runContextFactory; + + @Test + void commandsTrigger_shouldTriggerOnImplicitFailureExit1() throws Exception { + CommandsTrigger trigger = CommandsTrigger.builder() + .id("commands-trigger") + .type(CommandsTrigger.class.getName()) + .exitCondition(Property.ofValue("exit 1")) + .edge(Property.ofValue(true)) + .containerImage(Property.ofValue("r-base")) + .commands(Property.ofValue(List.of("Rscript -e 'quit(status = 1)'"))) + .build(); + + var context = TestsUtils.mockTrigger(runContextFactory, trigger); + Optional execution = trigger.evaluate(context.getKey(), context.getValue()); + + assertThat(execution.isPresent(), is(true)); + + Map triggerVars = execution.get().getTrigger().getVariables(); + assertThat("condition should be present", triggerVars.get("condition"), is("exit 1")); + assertThat("exitCode should be present", triggerVars.get("exitCode"), notNullValue()); + assertThat("exitCode should be 1", triggerVars.get("exitCode"), is(1)); + assertThat("timestamp should be present", triggerVars.get("timestamp"), notNullValue()); + } + + @Test + void commandsTrigger_shouldTriggerOnStdoutMatchUsingStructuredOutputs() throws Exception { + CommandsTrigger trigger = CommandsTrigger.builder() + .id("commands-stdout-match-trigger") + .type(CommandsTrigger.class.getName()) + .exitCondition(Property.ofValue("toto")) + .edge(Property.ofValue(true)) + .containerImage(Property.ofValue("r-base")) + .commands(Property.ofValue(List.of("echo '::{\"outputs\":{\"listing\":\"toto\"}}::'"))) + .build(); + + var context = TestsUtils.mockTrigger(runContextFactory, trigger); + Optional execution = trigger.evaluate(context.getKey(), context.getValue()); + + assertThat(execution.isPresent(), is(true)); + + Map triggerVars = execution.get().getTrigger().getVariables(); + assertThat("condition should be present", triggerVars.get("condition"), is("toto")); + assertThat("exitCode should be present", triggerVars.get("exitCode"), notNullValue()); + assertThat("exitCode should be 0", triggerVars.get("exitCode"), is(0)); + assertThat("timestamp should be present", triggerVars.get("timestamp"), notNullValue()); + assertThat("vars should be present", triggerVars.get("vars"), notNullValue()); + } + + @Test + void commandsTrigger_shouldNotEmitWhenConditionDoesNotMatch() throws Exception { + CommandsTrigger trigger = CommandsTrigger.builder() + .id("commands-no-match-trigger") + .type(CommandsTrigger.class.getName()) + .exitCondition(Property.ofValue("exit 1")) + .edge(Property.ofValue(true)) + .containerImage(Property.ofValue("r-base")) + .commands(Property.ofValue(List.of("Rscript -e 'quit(status = 0)'"))) + .build(); + + var context = TestsUtils.mockTrigger(runContextFactory, trigger); + Optional execution = trigger.evaluate(context.getKey(), context.getValue()); + + assertThat("successful run should not match 'exit 1'", execution.isPresent(), is(false)); + } + + @Test + void edgeMode_preventsConsecutiveEmit() { + AtomicBoolean lastMatched = new AtomicBoolean(false); + + // First match: transition false->true => should emit + boolean matched1 = true; + boolean emit1 = !lastMatched.getAndSet(matched1) && matched1; + assertThat("first match should emit", emit1, is(true)); + + // Second consecutive match: true->true => should NOT emit + boolean matched2 = true; + boolean emit2 = !lastMatched.getAndSet(matched2) && matched2; + assertThat("consecutive match should NOT emit in edge mode", emit2, is(false)); + + // Non-match: true->false => should not emit + boolean matched3 = false; + boolean emit3 = !lastMatched.getAndSet(matched3) && matched3; + assertThat("non-match should not emit", emit3, is(false)); + + // Match again after non-match: false->true => should emit + boolean matched4 = true; + boolean emit4 = !lastMatched.getAndSet(matched4) && matched4; + assertThat("match after non-match should emit", emit4, is(true)); + } +} diff --git a/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerConditionTest.java b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerConditionTest.java new file mode 100644 index 00000000..9779f859 --- /dev/null +++ b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerConditionTest.java @@ -0,0 +1,59 @@ +package io.kestra.plugin.scripts.r; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import java.time.Instant; +import java.util.Map; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; + +class ScriptTriggerConditionTest { + + private final ScriptTrigger trigger = ScriptTrigger.builder().build(); + + private ScriptTrigger.Output output(String condition, Integer exitCode, Map vars) { + return new ScriptTrigger.Output(Instant.now(), condition, exitCode, vars); + } + + @ParameterizedTest + @CsvSource({ + "exit 0, 0, true", + "exit 1, 1, true", + "EXIT 1, 1, true", + "exit 0, 1, false", + "exit 1, 0, false", + "exit 42, 42, true", + }) + void exitCodeCondition(String condition, int exitCode, boolean expected) { + assertThat(trigger.matchesCondition(output(condition, exitCode, null)), is(expected)); + } + + @Test + void exitCondition_nullExitCode_doesNotMatch() { + assertThat(trigger.matchesCondition(output("exit 1", null, null)), is(false)); + } + + @Test + void substringMatch_inVars() { + assertThat(trigger.matchesCondition( + output("toto", 0, Map.of("key", "toto"))), is(true)); + } + + @Test + void noMatch_emptyHaystack() { + assertThat(trigger.matchesCondition(output("something", 0, null)), is(false)); + } + + @Test + void noMatch_emptyCondition() { + assertThat(trigger.matchesCondition(output("", 0, Map.of("k", "v"))), is(false)); + } + + @Test + void nullCondition_doesNotMatch() { + assertThat(trigger.matchesCondition(output(null, 0, Map.of("k", "v"))), is(false)); + } +} diff --git a/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerTest.java b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerTest.java new file mode 100644 index 00000000..745e8bff --- /dev/null +++ b/plugin-script-r/src/test/java/io/kestra/plugin/scripts/r/ScriptTriggerTest.java @@ -0,0 +1,95 @@ +package io.kestra.plugin.scripts.r; + +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; + +/** + * Unit tests for ScriptTrigger's condition-matching logic and edge mode. + * + * These tests exercise matchesCondition via the Output model without requiring an R + * runtime, which may not be available on all CI machines. + * Integration coverage against an actual R runtime lives in CommandsTriggerTest. + */ +class ScriptTriggerTest { + + private final ScriptTrigger trigger = ScriptTrigger.builder().build(); + + private ScriptTrigger.Output output(String condition, Integer exitCode, Map vars) { + return new ScriptTrigger.Output(Instant.now(), condition, exitCode, vars); + } + + @Test + void exitCodeCondition_shouldMatchWhenExitCodeEquals() { + assertThat(trigger.matchesCondition(output("exit 1", 1, null)), is(true)); + } + + @Test + void exitCodeCondition_shouldNotMatchWhenExitCodeDiffers() { + assertThat(trigger.matchesCondition(output("exit 1", 127, null)), is(false)); + } + + @Test + void exitCodeCondition_shouldNotMatchWhenExitCodeIsNull() { + assertThat(trigger.matchesCondition(output("exit 1", null, null)), is(false)); + } + + @Test + void substringCondition_shouldMatchAgainstVars() { + assertThat(trigger.matchesCondition(output("toto", 0, Map.of("listing", "toto"))), is(true)); + } + + @Test + void substringCondition_shouldNotMatchWhenAbsent() { + assertThat(trigger.matchesCondition(output("toto", 0, Map.of("listing", "something_else"))), is(false)); + } + + @Test + void regexCondition_shouldMatchAgainstVars() { + assertThat(trigger.matchesCondition(output("status=\\w+", 0, Map.of("status", "status=ready"))), is(true)); + } + + @Test + void emptyCondition_shouldNotMatch() { + assertThat(trigger.matchesCondition(output("", 0, null)), is(false)); + } + + @Test + void nullCondition_shouldNotMatch() { + assertThat(trigger.matchesCondition(output(null, 0, null)), is(false)); + } + + @Test + void exitZeroCondition_shouldMatchSuccessfulExecution() { + assertThat(trigger.matchesCondition(output("exit 0", 0, null)), is(true)); + } + + @Test + void edgeMode_shouldEmitOnFirstMatch() { + var lastMatched = new AtomicBoolean(false); + boolean matched = true; + boolean emit = !lastMatched.getAndSet(matched) && matched; + assertThat("first match should emit", emit, is(true)); + } + + @Test + void edgeMode_shouldSuppressConsecutiveMatches() { + var lastMatched = new AtomicBoolean(true); + boolean matched = true; + boolean emit = !lastMatched.getAndSet(matched) && matched; + assertThat("consecutive match should not emit in edge mode", emit, is(false)); + } + + @Test + void edgeMode_shouldEmitAgainAfterNonMatch() { + var lastMatched = new AtomicBoolean(false); + boolean matched = true; + boolean emit = !lastMatched.getAndSet(matched) && matched; + assertThat("match after non-match should emit", emit, is(true)); + } +} From 61ad5d185e8003b3ab6385cc7b7a98fa633f19ef Mon Sep 17 00:00:00 2001 From: Abhishek Date: Sun, 20 Sep 2026 15:12:22 +0530 Subject: [PATCH 2/2] feat/perl-script-commands-triggers --- AGENTS.md | 2 + plugin-script-perl/build.gradle | 3 + .../plugin/scripts/perl/CommandsTrigger.java | 275 +++++++++++++++++ .../plugin/scripts/perl/ScriptTrigger.java | 276 ++++++++++++++++++ .../perl/CommandsTriggerConditionTest.java | 65 +++++ .../scripts/perl/CommandsTriggerTest.java | 136 +++++++++ .../perl/ScriptTriggerConditionTest.java | 65 +++++ .../scripts/perl/ScriptTriggerTest.java | 95 ++++++ 8 files changed, 917 insertions(+) create mode 100644 plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/CommandsTrigger.java create mode 100644 plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/ScriptTrigger.java create mode 100644 plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerConditionTest.java create mode 100644 plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerTest.java create mode 100644 plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerConditionTest.java create mode 100644 plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerTest.java diff --git a/AGENTS.md b/AGENTS.md index f7458d3f..e8e370a6 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -91,7 +91,9 @@ This is a **multi-module** plugin with 19 submodules: **plugin-script-perl:** - `io.kestra.plugin.scripts.perl.Commands` +- `io.kestra.plugin.scripts.perl.CommandsTrigger` - `io.kestra.plugin.scripts.perl.Script` +- `io.kestra.plugin.scripts.perl.ScriptTrigger` **plugin-script-php:** - `io.kestra.plugin.scripts.php.Commands` diff --git a/plugin-script-perl/build.gradle b/plugin-script-perl/build.gradle index 3bddfca5..2f8a80e6 100644 --- a/plugin-script-perl/build.gradle +++ b/plugin-script-perl/build.gradle @@ -16,4 +16,7 @@ dependencies { implementation project(':plugin-script') testImplementation project(path: ':plugin-script', configuration: 'testOutput') + + testImplementation group: "io.kestra", name: "scheduler", version: kestraVersion + testImplementation group: "io.kestra", name: "worker", version: kestraVersion } diff --git a/plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/CommandsTrigger.java b/plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/CommandsTrigger.java new file mode 100644 index 00000000..f4e2deb5 --- /dev/null +++ b/plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/CommandsTrigger.java @@ -0,0 +1,275 @@ +package io.kestra.plugin.scripts.perl; + +import java.time.Duration; +import java.time.Instant; +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import io.kestra.core.models.annotations.Example; +import io.kestra.core.models.annotations.Plugin; +import io.kestra.core.models.conditions.ConditionContext; +import io.kestra.core.models.executions.Execution; +import io.kestra.core.models.property.Property; +import io.kestra.core.models.tasks.RunnableTaskException; +import io.kestra.core.models.tasks.runners.TaskException; +import io.kestra.core.models.triggers.AbstractTrigger; +import io.kestra.core.models.triggers.PollingTriggerInterface; +import io.kestra.core.models.triggers.TriggerContext; +import io.kestra.core.models.triggers.TriggerOutput; +import io.kestra.core.models.triggers.TriggerService; +import io.kestra.core.runners.RunContext; +import io.kestra.plugin.scripts.exec.TriggerRunContext; +import io.kestra.plugin.scripts.exec.scripts.models.ScriptOutput; + +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotNull; +import lombok.AccessLevel; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.ToString; +import lombok.experimental.SuperBuilder; +import io.kestra.core.models.annotations.PluginProperty; + +@SuperBuilder +@ToString +@EqualsAndHashCode +@Getter +@NoArgsConstructor +@Schema( + title = "Trigger a flow when Perl commands match a condition", + description = "Polls and triggers a flow by running Perl commands within a script container." +) +@Plugin( + examples = { + @Example( + title = "Trigger when commands fail with an implicit error (exit 1).", + full = true, + code = """ + id: commands_trigger + namespace: company.team + + triggers: + - id: commands_failure + type: io.kestra.plugin.scripts.perl.CommandsTrigger + interval: PT10S + exitCondition: "exit 1" + edge: true + containerImage: perl + commands: + - perl missing.pl + + tasks: + - id: log + type: io.kestra.plugin.core.log.Log + message: "Triggered with exitCode={{ trigger.exitCode }} (condition={{ trigger.condition }})" + """ + ) + } +) +public class CommandsTrigger extends AbstractTrigger + implements PollingTriggerInterface, TriggerOutput { + + private static final String DEFAULT_IMAGE = "perl"; + private static final Pattern EXIT_CONDITION_PATTERN = Pattern.compile("^\\s*exit\\s+(\\d+)\\s*$", Pattern.CASE_INSENSITIVE); + + @Schema( + title = "Docker image used to execute the commands", + description = """ + Container image used by the underlying Commands task to run Perl commands. + Defaults to 'perl'. + """ + ) + @Builder.Default + @PluginProperty(group = "execution") + protected Property containerImage = Property.ofValue(DEFAULT_IMAGE); + + @Schema( + title = "Perl commands to execute", + description = "Commands executed on each poll (same semantics as the Perl Commands task)." + ) + @NotNull + @PluginProperty(group = "main") + protected Property> commands; + + @Schema( + title = "Condition to match", + description = """ + Condition evaluated after each commands execution. The trigger emits an event only when this condition matches. + + Supported forms: + - 'exit N' (example: 'exit 1'): matches when the process exit code equals N. + - Any other string: treated as a regex (or substring if regex is invalid) matched against: + - the task 'vars' (when commands emit ::{"outputs":...}::), + - and error logs when the task fails (TaskException). + """ + ) + @NotNull + @PluginProperty(group = "main") + protected Property exitCondition; + + @Schema( + title = "Check interval", + description = "Interval between polling evaluations." + ) + @Builder.Default + @PluginProperty(group = "execution") + private final Duration interval = Duration.ofSeconds(60); + + @Schema( + title = "Edge trigger mode", + description = """ + If true, the trigger emits only on a transition from 'not matching' to 'matching' (anti-spam). + If false, the trigger emits on every poll where the condition matches. + """ + ) + @Builder.Default + @PluginProperty(group = "advanced") + protected Property edge = Property.ofValue(true); + + @Builder.Default + @Getter(AccessLevel.NONE) + private final AtomicBoolean lastMatched = new AtomicBoolean(false); + + @Override + public Optional evaluate(ConditionContext conditionContext, TriggerContext context) throws Exception { + RunContext runContext = conditionContext.getRunContext(); + boolean renderedEdge = runContext.render(this.edge).as(Boolean.class).orElse(true); + + Output out; + try { + out = runOnce(runContext); + } catch (Exception e) { + runContext.logger().warn("Trigger evaluation failed, returning empty result to avoid blocking the scheduler", e); + return Optional.empty(); + } + + boolean matched = matchesCondition(out); + + boolean emit = renderedEdge + ? (!lastMatched.getAndSet(matched) && matched) + : matched; + + if (!emit) { + return Optional.empty(); + } + + return Optional.of(TriggerService.generateExecution(this, conditionContext, context, out)); + } + + private Output runOnce(RunContext runContext) throws Exception { + Commands task = Commands.builder() + .id(this.getId()) + .type(Commands.class.getName()) + .containerImage(this.containerImage) + .commands(this.commands) + .build(); + + String renderedCondition = runContext.render(this.exitCondition).as(String.class).orElse(""); + + try { + ScriptOutput taskOutput = task.run(TriggerRunContext.forEmbeddedTask(runContext, task)); + Integer exitCode = safeExitCode(taskOutput); + Map vars = safeVars(taskOutput); + + return new Output(Instant.now(), renderedCondition, exitCode, vars); + } catch (RunnableTaskException e) { + ExtractedFailure failure = extractFailure(e); + return new Output(Instant.now(), renderedCondition, failure.exitCode, null); + } + } + + boolean matchesCondition(Output out) { + String cond = out.getCondition() == null ? "" : out.getCondition().trim(); + + Matcher exitMatcher = EXIT_CONDITION_PATTERN.matcher(cond); + if (exitMatcher.matches()) { + int expected = Integer.parseInt(exitMatcher.group(1)); + return out.getExitCode() != null && out.getExitCode() == expected; + } + + String haystack = buildHaystack(out); + if (haystack.isEmpty() || cond.isEmpty()) { + return false; + } + + try { + return Pattern.compile(cond).matcher(haystack).find(); + } catch (Exception invalidRegex) { + return haystack.contains(cond); + } + } + + private String buildHaystack(Output out) { + if (out.getVars() == null || out.getVars().isEmpty()) { + return ""; + } + return out.getVars().toString(); + } + + private Integer safeExitCode(ScriptOutput taskOutput) { + try { + return taskOutput.getExitCode(); + } catch (Exception ignored) { + return null; + } + } + + private Map safeVars(ScriptOutput taskOutput) { + try { + return taskOutput.getVars(); + } catch (Exception ignored) { + return null; + } + } + + private record ExtractedFailure(Integer exitCode) { + } + + private ExtractedFailure extractFailure(RunnableTaskException e) { + Integer exitCode = null; + + Throwable cur = e.getCause(); + while (cur != null) { + if (cur instanceof TaskException te) { + exitCode = te.getExitCode(); + break; + } + cur = cur.getCause(); + } + + return new ExtractedFailure(exitCode); + } + + @Data + @AllArgsConstructor + public static class Output implements io.kestra.core.models.tasks.Output { + @Schema(title = "Timestamp of the event that fired the trigger") + private Instant timestamp; + + @Schema( + title = "Rendered condition", + description = "Rendered value of the exitCondition property for this poll." + ) + private String condition; + + @Schema( + title = "Commands exit code", + description = "Exit code returned by the Perl process (may be null if not available)." + ) + private Integer exitCode; + + @Schema( + title = "Commands vars", + description = "Vars produced by the task (e.g. via ::{\"outputs\":{...}}:: convention)." + ) + private Map vars; + } +} diff --git a/plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/ScriptTrigger.java b/plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/ScriptTrigger.java new file mode 100644 index 00000000..9594d02f --- /dev/null +++ b/plugin-script-perl/src/main/java/io/kestra/plugin/scripts/perl/ScriptTrigger.java @@ -0,0 +1,276 @@ +package io.kestra.plugin.scripts.perl; + +import java.time.Duration; +import java.time.Instant; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import io.kestra.core.models.annotations.Example; +import io.kestra.core.models.annotations.Plugin; +import io.kestra.core.models.annotations.PluginProperty; +import io.kestra.core.models.conditions.ConditionContext; +import io.kestra.core.models.executions.Execution; +import io.kestra.core.models.property.Property; +import io.kestra.core.models.tasks.RunnableTaskException; +import io.kestra.core.models.tasks.runners.TaskException; +import io.kestra.core.models.triggers.AbstractTrigger; +import io.kestra.core.models.triggers.PollingTriggerInterface; +import io.kestra.core.models.triggers.TriggerContext; +import io.kestra.core.models.triggers.TriggerOutput; +import io.kestra.core.models.triggers.TriggerService; +import io.kestra.core.runners.RunContext; +import io.kestra.plugin.scripts.exec.TriggerRunContext; +import io.kestra.plugin.scripts.exec.scripts.models.ScriptOutput; + +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.NotNull; +import lombok.AccessLevel; +import lombok.AllArgsConstructor; +import lombok.Builder; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.NoArgsConstructor; +import lombok.ToString; +import lombok.experimental.SuperBuilder; + +@SuperBuilder +@ToString +@EqualsAndHashCode +@Getter +@NoArgsConstructor +@Schema( + title = "Trigger on Perl script condition", + description = "Polls by running an inline Perl script in a container (default image perl) and emits when exitCondition matches. Supports edge mode to emit only on transitions and polls every 60s by default. Accepts 'exit N' or a regex (fallback substring) matched against emitted vars and failure logs." +) +@Plugin( + examples = { + @Example( + title = "Trigger when the script fails with an implicit error (exit 1).", + full = true, + code = """ + id: script_trigger + namespace: company.team + + triggers: + - id: script_failure + type: io.kestra.plugin.scripts.perl.ScriptTrigger + interval: PT10S + exitCondition: "exit 1" + edge: true + containerImage: perl + script: | + # This fails with a non-zero exit code. + exit 1; + + tasks: + - id: log + type: io.kestra.plugin.core.log.Log + message: "Triggered with exitCode={{ trigger.exitCode }} (condition={{ trigger.condition }})" + """ + ) + } +) +public class ScriptTrigger extends AbstractTrigger + implements PollingTriggerInterface, TriggerOutput { + + private static final String DEFAULT_IMAGE = "perl"; + private static final Pattern EXIT_CONDITION_PATTERN = Pattern.compile("^\\s*exit\\s+(\\d+)\\s*$", Pattern.CASE_INSENSITIVE); + + @Schema( + title = "Container image for script execution", + description = """ + Image used by the Script task to run the inline Perl script; defaults to 'perl'. + Provide an image that includes the Perl runtime and any required CPAN modules. + """ + ) + @Builder.Default + @PluginProperty(group = "execution") + protected Property containerImage = Property.ofValue(DEFAULT_IMAGE); + + @Schema( + title = "Inline Perl script", + description = """ + Multi-line Perl script executed on each poll, with the same semantics as the Perl Script task. + """ + ) + @NotNull + @PluginProperty(group = "main") + protected Property script; + + @Schema( + title = "Condition to match", + description = """ + Rendered condition evaluated after each execution; the trigger emits only when it matches. + 'exit N' compares the exit code, otherwise the string is used as a regex (or substring fallback) against emitted vars (from ::{"outputs":...}::) and failure logs. + """ + ) + @NotNull + @PluginProperty(group = "main") + protected Property exitCondition; + + @Schema( + title = "Check interval", + description = """ + Interval between polls; default PT60S. The scheduler uses this to schedule the next evaluation. + """ + ) + @Builder.Default + @PluginProperty(group = "execution") + private final Duration interval = Duration.ofSeconds(60); + + @Schema( + title = "Edge trigger mode", + description = """ + When true (default), emit only on a transition from not matching to matching. When false, emit on every poll that matches. + """ + ) + @Builder.Default + @PluginProperty(group = "advanced") + protected Property edge = Property.ofValue(true); + + @Builder.Default + @Getter(AccessLevel.NONE) + private final AtomicBoolean lastMatched = new AtomicBoolean(false); + + @Override + public Optional evaluate(ConditionContext conditionContext, TriggerContext context) throws Exception { + RunContext runContext = conditionContext.getRunContext(); + boolean renderedEdge = runContext.render(this.edge).as(Boolean.class).orElse(true); + + Output out; + try { + out = runOnce(runContext); + } catch (Exception e) { + runContext.logger().warn("Trigger evaluation failed, returning empty result to avoid blocking the scheduler", e); + return Optional.empty(); + } + + boolean matched = matchesCondition(out); + + boolean emit = renderedEdge + ? (!lastMatched.getAndSet(matched) && matched) + : matched; + + if (!emit) { + return Optional.empty(); + } + + return Optional.of(TriggerService.generateExecution(this, conditionContext, context, out)); + } + + private Output runOnce(RunContext runContext) throws Exception { + Script task = Script.builder() + .id(this.getId()) + .type(Script.class.getName()) + .containerImage(this.containerImage) + .script(this.script) + .build(); + + String renderedExitCondition = runContext.render(this.exitCondition).as(String.class).orElse(""); + + try { + ScriptOutput taskOutput = task.run(TriggerRunContext.forEmbeddedTask(runContext, task)); + Integer exitCode = safeExitCode(taskOutput); + Map vars = safeVars(taskOutput); + + return new Output(Instant.now(), renderedExitCondition, exitCode, vars); + } catch (RunnableTaskException e) { + ExtractedFailure failure = extractFailure(e); + return new Output(Instant.now(), renderedExitCondition, failure.exitCode, null); + } + } + + boolean matchesCondition(Output out) { + String cond = out.getCondition() == null ? "" : out.getCondition().trim(); + + Matcher exitMatcher = EXIT_CONDITION_PATTERN.matcher(cond); + if (exitMatcher.matches()) { + int expected = Integer.parseInt(exitMatcher.group(1)); + return out.getExitCode() != null && out.getExitCode() == expected; + } + + String haystack = buildHaystack(out); + if (haystack.isEmpty() || cond.isEmpty()) { + return false; + } + + try { + return Pattern.compile(cond).matcher(haystack).find(); + } catch (Exception invalidRegex) { + return haystack.contains(cond); + } + } + + private String buildHaystack(Output out) { + if (out.getVars() == null || out.getVars().isEmpty()) { + return ""; + } + return out.getVars().toString(); + } + + private Integer safeExitCode(ScriptOutput taskOutput) { + try { + return taskOutput.getExitCode(); + } catch (Exception ignored) { + return null; + } + } + + private Map safeVars(ScriptOutput taskOutput) { + try { + return taskOutput.getVars(); + } catch (Exception ignored) { + return null; + } + } + + private record ExtractedFailure(Integer exitCode) { + } + + private ExtractedFailure extractFailure(RunnableTaskException e) { + Integer exitCode = null; + + Throwable cur = e.getCause(); + while (cur != null) { + if (cur instanceof TaskException te) { + exitCode = te.getExitCode(); + break; + } + cur = cur.getCause(); + } + + return new ExtractedFailure(exitCode); + } + + @Data + @AllArgsConstructor + public static class Output implements io.kestra.core.models.tasks.Output { + @Schema(title = "Timestamp of the event that fired the trigger") + private Instant timestamp; + + @Schema( + title = "Rendered condition", + description = "Rendered value of the exitCondition property for this poll." + ) + private String condition; + + @Schema( + title = "Script exit code", + description = "Exit code returned by the Perl process (may be null if not available)." + ) + private Integer exitCode; + + @Schema( + title = "Script vars", + description = """ + Vars produced by the task (e.g. via ::{"outputs":{...}}:: convention). This is the main structured + way to evaluate non-exit conditions on successful runs. + """ + ) + private Map vars; + } +} diff --git a/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerConditionTest.java b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerConditionTest.java new file mode 100644 index 00000000..efde42a0 --- /dev/null +++ b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerConditionTest.java @@ -0,0 +1,65 @@ +package io.kestra.plugin.scripts.perl; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import java.time.Instant; +import java.util.Map; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; + +class CommandsTriggerConditionTest { + + private final CommandsTrigger trigger = CommandsTrigger.builder().build(); + + private CommandsTrigger.Output output(String condition, Integer exitCode, Map vars) { + return new CommandsTrigger.Output(Instant.now(), condition, exitCode, vars); + } + + @ParameterizedTest + @CsvSource({ + "exit 0, 0, true", + "exit 1, 1, true", + "EXIT 1, 1, true", + "exit 0, 1, false", + "exit 1, 0, false", + "exit 42, 42, true", + }) + void exitCodeCondition(String condition, int exitCode, boolean expected) { + assertThat(trigger.matchesCondition(output(condition, exitCode, null)), is(expected)); + } + + @Test + void exitCondition_nullExitCode_doesNotMatch() { + assertThat(trigger.matchesCondition(output("exit 1", null, null)), is(false)); + } + + @Test + void substringMatch_inVars() { + assertThat(trigger.matchesCondition( + output("toto", 0, Map.of("key", "toto"))), is(true)); + } + + @Test + void regexMatch_inVars() { + assertThat(trigger.matchesCondition( + output("status=\\w+", 0, Map.of("status", "status=ready"))), is(true)); + } + + @Test + void noMatch_emptyHaystack() { + assertThat(trigger.matchesCondition(output("something", 0, null)), is(false)); + } + + @Test + void noMatch_emptyCondition() { + assertThat(trigger.matchesCondition(output("", 0, Map.of("k", "v"))), is(false)); + } + + @Test + void nullCondition_doesNotMatch() { + assertThat(trigger.matchesCondition(output(null, 0, Map.of("k", "v"))), is(false)); + } +} diff --git a/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerTest.java b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerTest.java new file mode 100644 index 00000000..79ea87c1 --- /dev/null +++ b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/CommandsTriggerTest.java @@ -0,0 +1,136 @@ +package io.kestra.plugin.scripts.perl; + +import java.util.List; +import java.util.Map; +import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.junit.jupiter.api.Test; + +import io.kestra.core.junit.annotations.KestraTest; +import io.kestra.core.models.executions.Execution; +import io.kestra.core.models.property.Property; +import io.kestra.core.runners.RunContextFactory; +import io.kestra.core.utils.TestsUtils; + +import jakarta.inject.Inject; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.Matchers.notNullValue; + +@KestraTest +class CommandsTriggerTest { + @Inject + private RunContextFactory runContextFactory; + + @Test + void commandsTrigger_shouldTriggerOnImplicitFailureExit1() throws Exception { + CommandsTrigger trigger = CommandsTrigger.builder() + .id("commands-trigger") + .type(CommandsTrigger.class.getName()) + .exitCondition(Property.ofValue("exit 1")) + .edge(Property.ofValue(true)) + .containerImage(Property.ofValue("perl:latest")) + .commands(Property.ofValue(List.of("perl -e 'exit 1;'"))) + .build(); + + var context = TestsUtils.mockTrigger(runContextFactory, trigger); + Optional execution = trigger.evaluate(context.getKey(), context.getValue()); + + assertThat(execution.isPresent(), is(true)); + + Map triggerVars = execution.get().getTrigger().getVariables(); + assertThat("condition should be present", triggerVars.get("condition"), is("exit 1")); + assertThat("exitCode should be present", triggerVars.get("exitCode"), notNullValue()); + assertThat("exitCode should be 1", triggerVars.get("exitCode"), is(1)); + assertThat("timestamp should be present", triggerVars.get("timestamp"), notNullValue()); + } + + @Test + void commandsTrigger_shouldTriggerOnStdoutMatchUsingStructuredOutputs() throws Exception { + CommandsTrigger trigger = CommandsTrigger.builder() + .id("commands-stdout-match-trigger") + .type(CommandsTrigger.class.getName()) + .exitCondition(Property.ofValue("toto")) + .edge(Property.ofValue(true)) + .containerImage(Property.ofValue("perl:latest")) + .commands(Property.ofValue(List.of("echo '::{\"outputs\":{\"listing\":\"toto\"}}::'"))) + .build(); + + var context = TestsUtils.mockTrigger(runContextFactory, trigger); + Optional execution = trigger.evaluate(context.getKey(), context.getValue()); + + assertThat(execution.isPresent(), is(true)); + + Map triggerVars = execution.get().getTrigger().getVariables(); + assertThat("condition should be present", triggerVars.get("condition"), is("toto")); + assertThat("exitCode should be present", triggerVars.get("exitCode"), notNullValue()); + assertThat("exitCode should be 0", triggerVars.get("exitCode"), is(0)); + assertThat("timestamp should be present", triggerVars.get("timestamp"), notNullValue()); + assertThat("vars should be present", triggerVars.get("vars"), notNullValue()); + } + + @Test + void commandsTrigger_shouldNotEmitWhenConditionDoesNotMatch() throws Exception { + CommandsTrigger trigger = CommandsTrigger.builder() + .id("commands-no-match-trigger") + .type(CommandsTrigger.class.getName()) + .exitCondition(Property.ofValue("exit 1")) + .edge(Property.ofValue(true)) + .containerImage(Property.ofValue("perl:latest")) + .commands(Property.ofValue(List.of("perl -e 'exit 0;'"))) + .build(); + + var context = TestsUtils.mockTrigger(runContextFactory, trigger); + Optional execution = trigger.evaluate(context.getKey(), context.getValue()); + + assertThat("successful run should not match 'exit 1'", execution.isPresent(), is(false)); + } + + @Test + void commandsTrigger_shouldMatchRegexAgainstStructuredOutputs() throws Exception { + CommandsTrigger trigger = CommandsTrigger.builder() + .id("commands-regex-trigger") + .type(CommandsTrigger.class.getName()) + .exitCondition(Property.ofValue("status=\\w+")) + .edge(Property.ofValue(true)) + .containerImage(Property.ofValue("perl:latest")) + .commands(Property.ofValue(List.of("echo '::{\"outputs\":{\"status\":\"status=ready\"}}::'"))) + .build(); + + var context = TestsUtils.mockTrigger(runContextFactory, trigger); + Optional execution = trigger.evaluate(context.getKey(), context.getValue()); + + assertThat("Regex condition should match", execution.isPresent(), is(true)); + + Map triggerVars = execution.get().getTrigger().getVariables(); + assertThat("exitCode should be 0", triggerVars.get("exitCode"), is(0)); + assertThat("vars should be present", triggerVars.get("vars"), notNullValue()); + } + + @Test + void edgeMode_preventsConsecutiveEmit() { + AtomicBoolean lastMatched = new AtomicBoolean(false); + + // First match: transition false->true => should emit + boolean matched1 = true; + boolean emit1 = !lastMatched.getAndSet(matched1) && matched1; + assertThat("first match should emit", emit1, is(true)); + + // Second consecutive match: true->true => should NOT emit + boolean matched2 = true; + boolean emit2 = !lastMatched.getAndSet(matched2) && matched2; + assertThat("consecutive match should NOT emit in edge mode", emit2, is(false)); + + // Non-match: true->false => should not emit + boolean matched3 = false; + boolean emit3 = !lastMatched.getAndSet(matched3) && matched3; + assertThat("non-match should not emit", emit3, is(false)); + + // Match again after non-match: false->true => should emit + boolean matched4 = true; + boolean emit4 = !lastMatched.getAndSet(matched4) && matched4; + assertThat("match after non-match should emit", emit4, is(true)); + } +} diff --git a/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerConditionTest.java b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerConditionTest.java new file mode 100644 index 00000000..ecbc5aa8 --- /dev/null +++ b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerConditionTest.java @@ -0,0 +1,65 @@ +package io.kestra.plugin.scripts.perl; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.CsvSource; + +import java.time.Instant; +import java.util.Map; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; + +class ScriptTriggerConditionTest { + + private final ScriptTrigger trigger = ScriptTrigger.builder().build(); + + private ScriptTrigger.Output output(String condition, Integer exitCode, Map vars) { + return new ScriptTrigger.Output(Instant.now(), condition, exitCode, vars); + } + + @ParameterizedTest + @CsvSource({ + "exit 0, 0, true", + "exit 1, 1, true", + "EXIT 1, 1, true", + "exit 0, 1, false", + "exit 1, 0, false", + "exit 42, 42, true", + }) + void exitCodeCondition(String condition, int exitCode, boolean expected) { + assertThat(trigger.matchesCondition(output(condition, exitCode, null)), is(expected)); + } + + @Test + void exitCondition_nullExitCode_doesNotMatch() { + assertThat(trigger.matchesCondition(output("exit 1", null, null)), is(false)); + } + + @Test + void substringMatch_inVars() { + assertThat(trigger.matchesCondition( + output("toto", 0, Map.of("key", "toto"))), is(true)); + } + + @Test + void regexMatch_inVars() { + assertThat(trigger.matchesCondition( + output("status=\\w+", 0, Map.of("status", "status=ready"))), is(true)); + } + + @Test + void noMatch_emptyHaystack() { + assertThat(trigger.matchesCondition(output("something", 0, null)), is(false)); + } + + @Test + void noMatch_emptyCondition() { + assertThat(trigger.matchesCondition(output("", 0, Map.of("k", "v"))), is(false)); + } + + @Test + void nullCondition_doesNotMatch() { + assertThat(trigger.matchesCondition(output(null, 0, Map.of("k", "v"))), is(false)); + } +} diff --git a/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerTest.java b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerTest.java new file mode 100644 index 00000000..4ae863ca --- /dev/null +++ b/plugin-script-perl/src/test/java/io/kestra/plugin/scripts/perl/ScriptTriggerTest.java @@ -0,0 +1,95 @@ +package io.kestra.plugin.scripts.perl; + +import org.junit.jupiter.api.Test; + +import java.time.Instant; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; + +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; + +/** + * Unit tests for ScriptTrigger's condition-matching logic and edge mode. + * + * These tests exercise matchesCondition via the Output model without requiring a Perl + * runtime, which may not be available on all CI machines. + * Integration coverage against an actual Perl runtime lives in CommandsTriggerTest. + */ +class ScriptTriggerTest { + + private final ScriptTrigger trigger = ScriptTrigger.builder().build(); + + private ScriptTrigger.Output output(String condition, Integer exitCode, Map vars) { + return new ScriptTrigger.Output(Instant.now(), condition, exitCode, vars); + } + + @Test + void exitCodeCondition_shouldMatchWhenExitCodeEquals() { + assertThat(trigger.matchesCondition(output("exit 1", 1, null)), is(true)); + } + + @Test + void exitCodeCondition_shouldNotMatchWhenExitCodeDiffers() { + assertThat(trigger.matchesCondition(output("exit 1", 127, null)), is(false)); + } + + @Test + void exitCodeCondition_shouldNotMatchWhenExitCodeIsNull() { + assertThat(trigger.matchesCondition(output("exit 1", null, null)), is(false)); + } + + @Test + void substringCondition_shouldMatchAgainstVars() { + assertThat(trigger.matchesCondition(output("toto", 0, Map.of("listing", "toto"))), is(true)); + } + + @Test + void substringCondition_shouldNotMatchWhenAbsent() { + assertThat(trigger.matchesCondition(output("toto", 0, Map.of("listing", "something_else"))), is(false)); + } + + @Test + void regexCondition_shouldMatchAgainstVars() { + assertThat(trigger.matchesCondition(output("status=\\w+", 0, Map.of("status", "status=ready"))), is(true)); + } + + @Test + void emptyCondition_shouldNotMatch() { + assertThat(trigger.matchesCondition(output("", 0, null)), is(false)); + } + + @Test + void nullCondition_shouldNotMatch() { + assertThat(trigger.matchesCondition(output(null, 0, null)), is(false)); + } + + @Test + void exitZeroCondition_shouldMatchSuccessfulExecution() { + assertThat(trigger.matchesCondition(output("exit 0", 0, null)), is(true)); + } + + @Test + void edgeMode_shouldEmitOnFirstMatch() { + var lastMatched = new AtomicBoolean(false); + boolean matched = true; + boolean emit = !lastMatched.getAndSet(matched) && matched; + assertThat("first match should emit", emit, is(true)); + } + + @Test + void edgeMode_shouldSuppressConsecutiveMatches() { + var lastMatched = new AtomicBoolean(true); + boolean matched = true; + boolean emit = !lastMatched.getAndSet(matched) && matched; + assertThat("consecutive match should not emit in edge mode", emit, is(false)); + } + + @Test + void edgeMode_shouldEmitAgainAfterNonMatch() { + var lastMatched = new AtomicBoolean(false); + boolean matched = true; + boolean emit = !lastMatched.getAndSet(matched) && matched; + assertThat("match after non-match should emit", emit, is(true)); + } +}