From 330cef9f4fbec14d1c1c0a5819c3d7a0a89b23a0 Mon Sep 17 00:00:00 2001 From: wanglu14 Date: Fri, 10 Apr 2026 16:20:05 +0800 Subject: [PATCH 1/5] fix --- .../traversal/runners/SuspenseTaskRunner.java | 24 ++- .../SuspenseTaskTraversalTest.groovy | 153 ++++++++++++++++++ .../interfaces/model/strategy/Timeline.java | 4 + 3 files changed, 179 insertions(+), 2 deletions(-) diff --git a/rill-flow-dag/olympicene-traversal/src/main/java/com/weibo/rill/flow/olympicene/traversal/runners/SuspenseTaskRunner.java b/rill-flow-dag/olympicene-traversal/src/main/java/com/weibo/rill/flow/olympicene/traversal/runners/SuspenseTaskRunner.java index 172ca2163..c0c6d9b16 100644 --- a/rill-flow-dag/olympicene-traversal/src/main/java/com/weibo/rill/flow/olympicene/traversal/runners/SuspenseTaskRunner.java +++ b/rill-flow-dag/olympicene-traversal/src/main/java/com/weibo/rill/flow/olympicene/traversal/runners/SuspenseTaskRunner.java @@ -19,6 +19,8 @@ import com.google.common.collect.ImmutableSet; import com.google.common.collect.Maps; import com.google.common.collect.Sets; +import com.weibo.rill.flow.interfaces.model.strategy.Timeline; +import com.weibo.rill.flow.interfaces.model.task.BaseTask; import com.weibo.rill.flow.interfaces.model.task.TaskInfo; import com.weibo.rill.flow.interfaces.model.task.TaskInvokeMsg; import com.weibo.rill.flow.interfaces.model.task.TaskStatus; @@ -31,11 +33,11 @@ import com.weibo.rill.flow.olympicene.core.runtime.DAGInfoStorage; import com.weibo.rill.flow.olympicene.core.runtime.DAGStorageProcedure; import com.weibo.rill.flow.olympicene.core.switcher.SwitcherManager; -import com.weibo.rill.flow.olympicene.traversal.utils.ConditionsUtil; import com.weibo.rill.flow.olympicene.traversal.constant.TraversalErrorCode; import com.weibo.rill.flow.olympicene.traversal.exception.DAGTraversalException; import com.weibo.rill.flow.olympicene.traversal.helper.ContextHelper; import com.weibo.rill.flow.olympicene.traversal.mappings.InputOutputMapping; +import com.weibo.rill.flow.olympicene.traversal.utils.ConditionsUtil; import lombok.extern.slf4j.Slf4j; import org.apache.commons.collections.CollectionUtils; import org.apache.commons.collections.MapUtils; @@ -117,7 +119,25 @@ public ExecutionResult finish(String executionId, NotifyInfo notifyInfo, Map event.eventCode == DAGEvent.TASK_FAILED.getCode() + }) + } + + def "suspense task should be FAILED on timeout when tolerance=false even if skipOnTimeout=true"() { + given: + // tolerance=false,即使 skipOnTimeout=true,超时后也应触发 TASK_FAILED 事件 + String text = "version: 0.0.1\n" + + "namespace: olympicene\n" + + "service: mca\n" + + "name: test\n" + + "type: flow\n" + + "tasks: \n" + + "- category: suspense\n" + + " name: A\n" + + " tolerance: false\n" + + " timeline:\n" + + " timeoutInSeconds: \"120\"\n" + + " skipOnTimeout: \"true\"\n" + + " conditions:\n" + + " - \$.input.[?(@.url == \"bbb\")]\n" + DAG dag = dagParser.parse(text) + + when: + olympicene.submit('timeout_tol_false_1', dag, [:]) + olympicene.wakeup('timeout_tol_false_1', [:], + NotifyInfo.builder() + .taskInfoName('A') + .taskStatus(TaskStatus.FAILED) + .taskInvokeMsg(TaskInvokeMsg.builder().msg("timeout").build()) + .build()) + + then: + // tolerance=false → 触发 TASK_FAILED 回调事件,而不是 TASK_SKIPPED + 1 * callback.onEvent({ + Event event -> event.eventCode == DAGEvent.TASK_FAILED.getCode() + }) + } + def "suspense task should remain FAILED on interruption even if skipOnTimeout=true and tolerance=true"() { + given: + // tolerance=true 且 skipOnTimeout=true,但触发的是打断而非超时,应保持 FAILED + String text = "version: 0.0.1\n" + + "namespace: olympicene\n" + + "service: mca\n" + + "name: test\n" + + "type: flow\n" + + "tasks: \n" + + "- category: suspense\n" + + " name: A\n" + + " tolerance: true\n" + + " timeline:\n" + + " timeoutInSeconds: \"120\"\n" + + " skipOnTimeout: \"true\"\n" + + " inputMappings:\n" + + " - target: \$.input.url\n" + + " source: \$.context.url\n" + + " - target: \$.input.text\n" + + " source: \$.context.text\n" + + " conditions:\n" + + " - \$.input.[?(@.url == \"bbb\")]\n" + + " interruptions:\n" + + " - \$.input.[?(@.text == \"aaa\")]\n" + DAG dag = dagParser.parse(text) + + when: + // 提交 DAG,初始上下文触发 interruption 条件(text=aaa) + olympicene.submit('interrupt_with_skip_flag_1', dag, ["text": "aaa"]) + + then: + // 虽然 tolerance=true && skipOnTimeout=true,但触发的是打断而非超时,应为 TASK_FAILED 而非 TASK_SKIPPED + 1 * callback.onEvent({ + Event event -> event.eventCode == DAGEvent.TASK_FAILED.getCode() + }) + } } \ No newline at end of file diff --git a/rill-flow-interfaces/src/main/java/com/weibo/rill/flow/interfaces/model/strategy/Timeline.java b/rill-flow-interfaces/src/main/java/com/weibo/rill/flow/interfaces/model/strategy/Timeline.java index 20fe95540..720ef795d 100644 --- a/rill-flow-interfaces/src/main/java/com/weibo/rill/flow/interfaces/model/strategy/Timeline.java +++ b/rill-flow-interfaces/src/main/java/com/weibo/rill/flow/interfaces/model/strategy/Timeline.java @@ -26,15 +26,19 @@ @Getter public class Timeline { private String timeoutInSeconds; + private String skipOnTimeout; private String suspenseIntervalSeconds; private String suspenseTimestamp; @JsonCreator public Timeline(@JsonProperty("timeoutInSeconds") String timeoutInSeconds, + @JsonProperty("skipOnTimeout") String skipOnTimeout, @JsonProperty("suspenseIntervalSeconds") String suspenseIntervalSeconds, @JsonProperty("suspenseTimestamp") String suspenseTimestamp) { this.timeoutInSeconds = timeoutInSeconds; + this.skipOnTimeout = skipOnTimeout; this.suspenseIntervalSeconds = suspenseIntervalSeconds; this.suspenseTimestamp = suspenseTimestamp; } } + From e29fc0dea6b7fd2d083ba32aabcbf49e2f677c3d Mon Sep 17 00:00:00 2001 From: wanglu14 Date: Fri, 10 Apr 2026 18:18:13 +0800 Subject: [PATCH 2/5] fix --- .github/workflows/build.yml | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 2bf0d0f86..8975a9e83 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -33,14 +33,14 @@ jobs: - name: Build and analyze env: SONAR_TOKEN: ${{ secrets.SONAR_TOKEN }} - run: mvn -B clean javadoc:javadoc clean verify -P coverage - if: env.SONAR_TOKEN == '' + run: mvn -B clean javadoc:javadoc verify -P coverage + if: env.SONAR_TOKEN == '' || env.SONAR_TOKEN == null - name: Build and analyze with sonar env: GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} SONAR_TOKEN: ${{ secrets.SONAR_TOKEN }} - run: mvn -B clean javadoc:javadoc verify org.sonarsource.scanner.maven:sonar-maven-plugin:sonar -P coverage -Dsonar.projectKey=weibocom_rill-flow - if: env.SONAR_TOKEN != '' + run: mvn -B clean javadoc:javadoc verify org.sonarsource.scanner.maven:sonar-maven-plugin:sonar -P coverage -Dsonar.projectKey=weibocom_rill-flow + if: env.SONAR_TOKEN != '' && env.SONAR_TOKEN != null - name: Upload coverage reports to Codecov uses: codecov/codecov-action@v3 env: From 06bb32bb4b36a89168b49d6cb21001a4650e6217 Mon Sep 17 00:00:00 2001 From: wanglu14 Date: Fri, 10 Apr 2026 18:46:50 +0800 Subject: [PATCH 3/5] fix --- .github/workflows/release.yml | 2 +- pom.xml | 19 +++++++++---------- 2 files changed, 10 insertions(+), 11 deletions(-) diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 54953a32d..2c3b719eb 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -32,7 +32,7 @@ jobs: with: distribution: 'temurin' java-version: 17 - server-id: ossrh + server-id: central cache: 'maven' server-username: MAVEN_USERNAME server-password: MAVEN_CENTRAL_TOKEN diff --git a/pom.xml b/pom.xml index c573d98ff..a6cda4df7 100644 --- a/pom.xml +++ b/pom.xml @@ -712,12 +712,12 @@ - ossrh - https://oss.sonatype.org/content/repositories/snapshots + central + https://central.sonatype.com/repository/maven-snapshots/ - ossrh - https://oss.sonatype.org/service/local/staging/deploy/maven2/ + central + https://central.sonatype.com/ @@ -736,14 +736,13 @@ - org.sonatype.plugins - nexus-staging-maven-plugin - 1.6.13 + org.sonatype.central + central-publishing-maven-plugin + 0.5.0 true - ossrh - https://oss.sonatype.org/ - true + central + true From 6ece0522f744dcd025f2dab4cce5149bd6d9b211 Mon Sep 17 00:00:00 2001 From: wanglu14 Date: Fri, 10 Apr 2026 18:47:27 +0800 Subject: [PATCH 4/5] fix --- .github/workflows/build.yml | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index 8975a9e83..fdb3186b1 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -33,15 +33,15 @@ jobs: - name: Build and analyze env: SONAR_TOKEN: ${{ secrets.SONAR_TOKEN }} - run: mvn -B clean javadoc:javadoc verify -P coverage - if: env.SONAR_TOKEN == '' || env.SONAR_TOKEN == null + run: mvn -B clean javadoc:javadoc clean verify -P coverage + if: env.SONAR_TOKEN == '' - name: Build and analyze with sonar env: GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} SONAR_TOKEN: ${{ secrets.SONAR_TOKEN }} - run: mvn -B clean javadoc:javadoc verify org.sonarsource.scanner.maven:sonar-maven-plugin:sonar -P coverage -Dsonar.projectKey=weibocom_rill-flow - if: env.SONAR_TOKEN != '' && env.SONAR_TOKEN != null + run: mvn -B clean javadoc:javadoc verify org.sonarsource.scanner.maven:sonar-maven-plugin:sonar -P coverage -Dsonar.projectKey=weibocom_rill-flow + if: env.SONAR_TOKEN != '' - name: Upload coverage reports to Codecov uses: codecov/codecov-action@v3 env: - CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }} + CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }} \ No newline at end of file From 501670b3d9f3dc65f16b79d77dbbce42e77660df Mon Sep 17 00:00:00 2001 From: wanglu14 Date: Fri, 10 Apr 2026 18:54:26 +0800 Subject: [PATCH 5/5] fix --- .github/workflows/build.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml index fdb3186b1..fe44b9cdf 100644 --- a/.github/workflows/build.yml +++ b/.github/workflows/build.yml @@ -44,4 +44,4 @@ jobs: - name: Upload coverage reports to Codecov uses: codecov/codecov-action@v3 env: - CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }} \ No newline at end of file + CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }}