diff --git a/.github/workflows/e2e.yml b/.github/workflows/e2e.yml index 7f9feef50f34..3d7f2accb2a2 100644 --- a/.github/workflows/e2e.yml +++ b/.github/workflows/e2e.yml @@ -122,6 +122,8 @@ jobs: class: org.apache.dolphinscheduler.e2e.cases.WorkflowHttpTaskE2ETest - name: WorkflowJavaTaskE2ETest class: org.apache.dolphinscheduler.e2e.cases.WorkflowJavaTaskE2ETest + - name: WorkflowSeaTunnelE2ETest + class: org.apache.dolphinscheduler.e2e.cases.WorkflowSeaTunnelE2ETest # - name: WorkflowForSwitch # class: org.apache.dolphinscheduler.e2e.cases.WorkflowSwitchE2ETest - name: FileManageE2ETest diff --git a/.github/workflows/frontend.yml b/.github/workflows/frontend.yml index 24f9b902b0be..130b5987dabe 100644 --- a/.github/workflows/frontend.yml +++ b/.github/workflows/frontend.yml @@ -89,6 +89,10 @@ jobs: echo "Code format check failed! Please run \`pnpm run lint\` to format the code." exit -1 fi + - name: Task execution type regressions + run: | + node tests/task-execution-type.cjs + node tests/task-execution-graph.cjs - name: Compile and Build on ${{ matrix.os }} run: | pnpm run build:prod diff --git a/docs/docs/en/guide/task/seatunnel.md b/docs/docs/en/guide/task/seatunnel.md index 509ed19a297b..f1c0081b20cf 100644 --- a/docs/docs/en/guide/task/seatunnel.md +++ b/docs/docs/en/guide/task/seatunnel.md @@ -13,6 +13,8 @@ Click [here](https://seatunnel.apache.org/) for more information about `Apache S ## Task Parameter - Please refer to [DolphinScheduler Task Parameters Appendix](appendix.md) `Default Task Parameters` section for default parameters. +- Task execution type: Defaults to `Batch Task`. For a streaming SeaTunnel job, select `Stream Task` to list its workflow task instances under `Task Instance -> Stream Task`. This selection is saved separately from the SeaTunnel configuration and applies to both custom scripts and resource files. Configure SeaTunnel's `env.job.mode = "STREAMING"` in the selected configuration as well. + Streaming tasks are started through their workflow and must use an attached SeaTunnel submission; detached submission (such as `--async`) does not let the task track the running job. SeaTunnel savepoint is not supported by this task plugin. - Startup script: Select script name to start the task (it may vary across SeaTunnel distributions, please check `${SEATUNNEL_HOME}/bin/`), including `seatunnel.sh`, `start-seatunnel-flink-13-connector-v2.sh`, `start-seatunnel-flink-15-connector-v2.sh`, `start-seatunnel-flink-connector-v2.sh`, `start-seatunnel-flink.sh`, `start-seatunnel-spark-2-connector-v2.sh`, `start-seatunnel-spark-3-connector-v2.sh`, `start-seatunnel-spark-connector-v2.sh`, `start-seatunnel-spark.sh` - FLINK - Run model: supports `run` and `run-application` modes diff --git a/docs/docs/zh/guide/task/seatunnel.md b/docs/docs/zh/guide/task/seatunnel.md index ec79c39d8a2c..bbb5999826a0 100644 --- a/docs/docs/zh/guide/task/seatunnel.md +++ b/docs/docs/zh/guide/task/seatunnel.md @@ -13,6 +13,8 @@ ## 任务参数 - 默认参数说明请参考[DolphinScheduler任务参数附录](appendix.md)`默认任务参数`一栏。 +- 任务执行类型:默认为“批量任务”。运行 SeaTunnel 流作业时,选择“实时任务”,工作流中的任务实例将显示在“任务实例 -> 实时任务”列表。该选项独立于 SeaTunnel 配置保存,同时适用于自定义脚本和资源文件;所选配置中仍需设置 SeaTunnel 的 `env.job.mode = "STREAMING"`。 + 实时任务通过工作流启动,必须使用前台等待作业结束的提交方式;使用 `--async` 等后台提交选项时,任务无法跟踪持续运行的作业。当前任务插件不支持 SeaTunnel 保存点操作。 - 启动脚本:选择你想要运行任务的启动脚本(不同 SeaTunnel 发行包可能存在差异,以实际 `${SEATUNNEL_HOME}/bin/` 为准),包括 `seatunnel.sh`, `start-seatunnel-flink-13-connector-v2.sh`, `start-seatunnel-flink-15-connector-v2.sh`, `start-seatunnel-flink-connector-v2.sh`, `start-seatunnel-flink.sh`, `start-seatunnel-spark-2-connector-v2.sh`, `start-seatunnel-spark-3-connector-v2.sh`, `start-seatunnel-spark-connector-v2.sh`, `start-seatunnel-spark.sh` - FLINK - 运行模型:支持 `run` 和 `run-application` 两种模式 diff --git a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/cases/WorkflowSeaTunnelE2ETest.java b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/cases/WorkflowSeaTunnelE2ETest.java new file mode 100644 index 000000000000..462ea04ea075 --- /dev/null +++ b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/cases/WorkflowSeaTunnelE2ETest.java @@ -0,0 +1,190 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 org.apache.dolphinscheduler.e2e.cases; + +import static org.assertj.core.api.Assertions.assertThat; + +import org.apache.dolphinscheduler.e2e.core.DolphinScheduler; +import org.apache.dolphinscheduler.e2e.core.WebDriverWaitFactory; +import org.apache.dolphinscheduler.e2e.models.users.AdminUser; +import org.apache.dolphinscheduler.e2e.pages.LoginPage; +import org.apache.dolphinscheduler.e2e.pages.project.ProjectPage; +import org.apache.dolphinscheduler.e2e.pages.project.workflow.TaskInstanceTab; +import org.apache.dolphinscheduler.e2e.pages.project.workflow.WorkflowDefinitionTab; +import org.apache.dolphinscheduler.e2e.pages.project.workflow.WorkflowForm; +import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.SeaTunnelTaskForm; +import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.ShellTaskForm; +import org.apache.dolphinscheduler.e2e.pages.security.SecurityPage; +import org.apache.dolphinscheduler.e2e.pages.security.TenantPage; +import org.apache.dolphinscheduler.e2e.pages.security.UserPage; + +import java.time.Duration; +import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; + +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.openqa.selenium.By; +import org.openqa.selenium.StaleElementReferenceException; +import org.openqa.selenium.remote.RemoteWebDriver; +import org.openqa.selenium.support.ui.ExpectedConditions; + +@DolphinScheduler(composeFiles = "docker/basic/docker-compose.yaml") +class WorkflowSeaTunnelE2ETest { + + private static RemoteWebDriver browser; + + @Test + @Order(1) + void testExecutionTypeCopyEdgesAndInstanceClassification() { + String project = "seatunnel-execution-type"; + String workflow = "seatunnel-modes"; + AdminUser admin = new AdminUser(); + TenantPage tenants = new LoginPage(browser).login(admin) + .goToNav(SecurityPage.class).goToTab(TenantPage.class); + if (tenants.tenants().stream().noneMatch(tenant -> tenant.tenantCode().equals(admin.getTenant()))) { + tenants.create(admin.getTenant()).goToNav(SecurityPage.class).goToTab(UserPage.class).update(admin); + } + WorkflowDefinitionTab definitions = tenants + .goToNav(ProjectPage.class) + .create(project) + .goTo(project) + .goToTab(WorkflowDefinitionTab.class); + WorkflowForm form = definitions.createWorkflow(); + form.addTask(WorkflowForm.TaskType.SHELL, 100, 100) + .script("echo batch-control\n").name("before").submit(); + browser.findElement(By.cssSelector(".task-cate-di .n-collapse-item__header")).click(); + + SeaTunnelTaskForm task = form.addTask(WorkflowForm.TaskType.SEATUNNEL, 350, 100); + assertThat(task.isExecutionType("BATCH")).isTrue(); + task.name("seatunnel").preTask("before").submit(); + browser.findElement(By.cssSelector(".task-cate-universal .n-collapse-item__header")).click(); + form.addTask(WorkflowForm.TaskType.SHELL, 600, 100) + .script("echo after\n").name("after").preTask("seatunnel").submit(); + form.waitForEdgeStyle(2, false); + form.submit().name(workflow).submit(); + + form = reopenWorkflow(workflow); + String sourceCode = form.taskCode("seatunnel"); + form.waitForEdgeStyle(2, false); + SeaTunnelTaskForm savedBatch = openTask(form, sourceCode); + assertThat(savedBatch.isExecutionType("BATCH")).isTrue(); + savedBatch.executionType("STREAM").submit(); + form.waitForEdgeStyle(2, true); + form.submit().submit(); + + form = reopenWorkflow(workflow); + form.waitForEdgeStyle(2, true); + SeaTunnelTaskForm savedStream = openTask(form, sourceCode); + assertThat(savedStream.isExecutionType("STREAM")).isTrue(); + savedStream.executionType("BATCH").submit(); + form.waitForEdgeStyle(2, false); + form.submit().submit(); + + form = reopenWorkflow(workflow); + form.waitForEdgeStyle(2, false); + SeaTunnelTaskForm edited = openTask(form, sourceCode); + assertThat(edited.isExecutionType("BATCH")).isTrue(); + edited.executionType("STREAM").submit(); + form.waitForEdgeStyle(2, true); + String copyCode = form.copyTask(sourceCode); + assertThat(copyCode).isNotEqualTo(sourceCode); + SeaTunnelTaskForm copy = openTask(form, copyCode); + assertThat(copy.isExecutionType("STREAM")).isTrue(); + String copyName = copy.inputNodeName().getAttribute("value"); + assertThat(copyName).startsWith("seatunnel_"); + copy.submit(); + form.submit().submit(); + + form = reopenWorkflow(workflow); + form.waitForEdgeStyle(2, true); + SeaTunnelTaskForm reloadedCopy = openTask(form, copyCode); + assertThat(reloadedCopy.isExecutionType("STREAM")).isTrue(); + reloadedCopy.submit(); + SeaTunnelTaskForm reloadedSource = openTask(form, sourceCode); + assertThat(reloadedSource.isExecutionType("STREAM")).isTrue(); + reloadedSource.submit(); + + definitions = new ProjectPage(browser).goToNav(ProjectPage.class).goTo(project) + .goToTab(WorkflowDefinitionTab.class); + definitions.publish(workflow).run(workflow).submit(); + + // The basic image has no SeaTunnel engine. Run the real scheduler (not dry-run) + // and verify persisted instance classification, which precedes worker execution. + // The sample SeaTunnel script and eventual engine result are not tested here. + TaskInstanceTab instances = new ProjectPage(browser).goToNav(ProjectPage.class).goTo(project) + .goToTab(TaskInstanceTab.class).selectType("Stream"); + waitForStreamInstance(instances, "seatunnel"); + waitForStreamInstance(instances, copyName); + + instances.selectType("Batch"); + // Positive loaded-table witness: the BATCH predecessor must execute before + // the source SeaTunnel task can be submitted. Avoid an empty-table false pass. + WebDriverWaitFactory.createWebDriverWait(browser, Duration.ofSeconds(120)) + .ignoring(StaleElementReferenceException.class) + .until(unused -> { + List names = instances.instances().stream().map(TaskInstanceTab.Row::taskInstanceName) + .collect(Collectors.toList()); + if (!names.contains("before")) { + return false; + } + assertThat(names).doesNotContain("seatunnel", copyName); + return true; + }); + // Testcontainers disposes this isolated project's workflow and instances. + } + + private void waitForStreamInstance(TaskInstanceTab instances, String name) { + // Refreshes can replace rows between locating them and reading their cells. + WebDriverWaitFactory.createWebDriverWait(browser, Duration.ofSeconds(120)) + .ignoring(StaleElementReferenceException.class) + .until(unused -> { + TaskInstanceTab.Row instance = instances.streamInstances().stream() + .filter(row -> row.taskInstanceName().equals(name)).findFirst().orElse(null); + if (instance == null) { + return false; + } + assertThat(instance.taskType()).isEqualTo("SEATUNNEL"); + assertThat(instance.dryRun()).isEqualTo("NO"); + return true; + }); + } + + private SeaTunnelTaskForm openTask(WorkflowForm form, String code) { + form.openTask(code); + return new SeaTunnelTaskForm(form); + } + + private WorkflowForm reopenWorkflow(String workflow) { + WebDriverWaitFactory.createWebDriverWait(browser).until(ExpectedConditions.visibilityOfElementLocated( + By.className("btn-create-workflow"))); + Set windows = browser.getWindowHandles(); + WebDriverWaitFactory.createWebDriverWait(browser).until(ExpectedConditions.elementToBeClickable( + By.xpath("//*[contains(@class, 'workflow-name')]//button[normalize-space(.)='" + workflow + "']"))) + .click(); + WebDriverWaitFactory.createWebDriverWait(browser) + .until(driver -> driver.getWindowHandles().size() > windows.size()); + String detailWindow = browser.getWindowHandles().stream().filter(window -> !windows.contains(window)) + .findFirst().orElseThrow(IllegalStateException::new); + browser.close(); + browser.switchTo().window(detailWindow); + browser.navigate().refresh(); + return new WorkflowForm(browser); + } +} diff --git a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/TaskInstanceTab.java b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/TaskInstanceTab.java index eedd9eccb531..ae0e5989d602 100644 --- a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/TaskInstanceTab.java +++ b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/TaskInstanceTab.java @@ -17,6 +17,7 @@ package org.apache.dolphinscheduler.e2e.pages.project.workflow; +import org.apache.dolphinscheduler.e2e.core.WebDriverWaitFactory; import org.apache.dolphinscheduler.e2e.pages.common.NavBarPage; import org.apache.dolphinscheduler.e2e.pages.project.ProjectDetailPage; @@ -30,6 +31,7 @@ import org.openqa.selenium.WebElement; import org.openqa.selenium.remote.RemoteWebDriver; import org.openqa.selenium.support.FindBy; +import org.openqa.selenium.support.ui.ExpectedConditions; @Getter public final class TaskInstanceTab extends NavBarPage implements ProjectDetailPage.Tab { @@ -49,6 +51,21 @@ public List instances() { .collect(Collectors.toList()); } + public TaskInstanceTab selectType(String type) { + By tab = By.cssSelector(".n-tabs-tab[data-name='" + type + "']"); + WebDriverWaitFactory.createWebDriverWait(driver).until(ExpectedConditions.elementToBeClickable(tab)).click(); + WebDriverWaitFactory.createWebDriverWait(driver).until(ExpectedConditions.attributeContains( + tab, "class", "n-tabs-tab--active")); + return this; + } + + public List streamInstances() { + return driver.findElements(By.cssSelector(".n-tab-pane .n-data-table-tbody tr")).stream() + .filter(WebElement::isDisplayed) + .filter(row -> !row.findElements(By.cssSelector("td[data-col-key=workflowDefinitionName]")).isEmpty()) + .map(Row::new).collect(Collectors.toList()); + } + @RequiredArgsConstructor public static class Row { @@ -62,6 +79,14 @@ public String workflowInstanceName() { return row.findElement(By.cssSelector("td[data-col-key=workflowInstanceName]")).getText(); } + public String taskType() { + return row.findElement(By.cssSelector("td[data-col-key=taskType]")).getText(); + } + + public String dryRun() { + return row.findElement(By.cssSelector("td[data-col-key=dryRun]")).getText(); + } + public int retryTimes() { return Integer.parseInt(row.findElement(By.cssSelector("td[data-col-key=retryTimes]")).getText()); } diff --git a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowForm.java b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowForm.java index 7c47eb36a1dd..e5cf94f802bb 100644 --- a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowForm.java +++ b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowForm.java @@ -21,12 +21,15 @@ import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.HttpTaskForm; import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.JavaTaskForm; import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.PythonTaskForm; +import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.SeaTunnelTaskForm; import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.ShellTaskForm; import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.SubWorkflowTaskForm; import org.apache.dolphinscheduler.e2e.pages.project.workflow.task.SwitchTaskForm; import java.nio.charset.StandardCharsets; import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; import lombok.Getter; import lombok.SneakyThrows; @@ -67,13 +70,19 @@ public WorkflowForm(WebDriver driver) { @SneakyThrows @SuppressWarnings("unchecked") public T addTask(TaskType type) { + return addTask(type, null, null); + } + + @SneakyThrows + @SuppressWarnings("unchecked") + public T addTask(TaskType type, Integer x, Integer y) { final WebElement task = driver.findElement(By.className("task-item-" + type.name())); final WebElement canvas = driver.findElement(By.className("dag-container")); final JavascriptExecutor js = (JavascriptExecutor) driver; final String dragAndDrop = String.join("\n", Resources.readLines(Resources.getResource("dragAndDrop.js"), StandardCharsets.UTF_8)); - js.executeScript(dragAndDrop, task, canvas); + js.executeScript(dragAndDrop, task, canvas, x, y); WebDriverWaitFactory.createWebDriverWait(driver).until(ExpectedConditions .visibilityOfElementLocated(By.xpath("//*[contains(text(), 'Current node settings')]"))); @@ -90,6 +99,8 @@ public T addTask(TaskType type) { return (T) new JavaTaskForm(this); case PYTHON: return (T) new PythonTaskForm(this); + case SEATUNNEL: + return (T) new SeaTunnelTaskForm(this); } throw new UnsupportedOperationException("Unknown task type"); } @@ -110,6 +121,84 @@ public WebElement getTask(String taskName) { return task; } + public String taskCode(String taskName) { + return WebDriverWaitFactory.createWebDriverWait(driver).until(unused -> driver + .findElements(By.cssSelector(".dag-container .x6-graph-scroller g[data-shape='dag-task']")).stream() + .filter(node -> node.findElement(By.xpath("./*[local-name()='text']")).getText() + .equals(taskName)) + .map(node -> node.getAttribute("data-cell-id")) + .findFirst().orElse(null)); + } + + private WebElement taskBodyInView(String code) { + By body = By.cssSelector(".dag-container .x6-graph-scroller g[data-shape='dag-task'][data-cell-id='" + + code + "'] > .dag-task-body"); + return WebDriverWaitFactory.createWebDriverWait(driver) + .withMessage("Task body is hidden or its center is covered: " + code).until(unused -> { + WebElement node = ExpectedConditions.visibilityOfElementLocated(body).apply(driver); + if (node == null) { + return null; + } + // Require an unclipped center so the hit test matches the native pointer location. + boolean target = Boolean.TRUE.equals(((JavascriptExecutor) driver).executeScript( + "const node = arguments[0]; node.scrollIntoView({block: 'center', inline: 'center'});" + + "const firstRect = el => Array.from(el.getClientRects())" + + ".find(r => r.width > 0 && r.height > 0); const r = firstRect(node);" + + "if (!r || r.left < 0 || r.top < 0 || r.right > innerWidth" + + " || r.bottom > innerHeight) return false;" + + "const clips = value => ['auto', 'scroll', 'hidden'].includes(value);" + + "for (let p = node.parentElement; p && p !== document.body; p = p.parentElement) {" + + "const rects = p.getClientRects(); if (!rects.length) break;" + + "const q = firstRect(p) || rects[0], s = getComputedStyle(p);" + + "if ((clips(s.overflowX) && (r.left < q.left || r.right > q.right))" + + " || (clips(s.overflowY) && (r.top < q.top || r.bottom > q.bottom))) return false; }" + + "const x = Math.floor((r.left + r.right) / 2);" + + "const y = Math.floor((r.top + r.bottom) / 2);" + + "const hit = document.elementFromPoint(x, y);" + + "return hit && hit.closest('g[data-shape=\"dag-task\"]') === node.parentElement;", + node)); + return target ? node : null; + }); + } + + public void openTask(String code) { + new Actions(driver).doubleClick(taskBodyInView(code)).perform(); + // Visible controls still move while their modal is entering or leaving. + WebDriverWaitFactory.createWebDriverWait(driver).until(ExpectedConditions.visibilityOfElementLocated( + By.cssSelector(".n-modal:not(.fade-in-scale-up-transition-enter-active)" + + ":not(.fade-in-scale-up-transition-leave-active) .input-node-name input"))); + } + + public String copyTask(String code) { + Set existing = driver + .findElements(By.cssSelector(".dag-container .x6-graph-scroller g[data-shape='dag-task']")).stream() + .map(node -> node.getAttribute("data-cell-id")).collect(Collectors.toSet()); + new Actions(driver).contextClick(taskBodyInView(code)).perform(); + WebDriverWaitFactory.createWebDriverWait(driver).until(ExpectedConditions.elementToBeClickable(By.xpath( + "//div[contains(@class, 'dag-context-menu')]//button[normalize-space(.)='Copy']"))).click(); + return WebDriverWaitFactory.createWebDriverWait(driver).until(unused -> { + List added = driver + .findElements(By.cssSelector(".dag-container .x6-graph-scroller g[data-shape='dag-task']")).stream() + .map(node -> node.getAttribute("data-cell-id")) + .filter(id -> !existing.contains(id)).collect(Collectors.toList()); + return added.size() == 1 ? added.get(0) : null; + }); + } + + public void waitForEdgeStyle(int count, boolean stream) { + WebDriverWaitFactory.createWebDriverWait(driver).until(unused -> { + List edges = + driver.findElements(By.cssSelector(".dag-container .x6-graph-scroller g[data-shape='dag-edge']")); + return edges.size() == count && edges.stream().allMatch(edge -> { + List lines = edge.findElements(By.cssSelector("path[stroke-dasharray]")); + return !lines.isEmpty() && lines.stream() + .allMatch(line -> line.getAttribute("stroke-dasharray").replace(',', ' ').trim() + .replaceAll("\\s+", " ") + .equals(stream ? "5 5" : "none")); + }); + }); + } + public WorkflowSaveDialog submit() { buttonSave().click(); WebDriverWaitFactory.createWebDriverWait(driver) @@ -129,6 +218,7 @@ public enum TaskType { SWITCH, HTTP, JAVA, - PYTHON + PYTHON, + SEATUNNEL } } diff --git a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowSaveDialog.java b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowSaveDialog.java index 7e98306b1a12..004302a1878b 100644 --- a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowSaveDialog.java +++ b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/WorkflowSaveDialog.java @@ -41,7 +41,11 @@ public final class WorkflowSaveDialog { }) private WebElement inputName; - @FindBy(xpath = "//div[contains(text(), 'Basic Information')]/../following-sibling::div[contains(@class, 'n-card__footer')]//button[contains(@class, 'btn-submit')]") + // The modal transition moves this button even while it is visible and enabled. + @FindBy(xpath = "//div[contains(concat(' ', normalize-space(@class), ' '), ' n-modal ')" + + " and not(contains(@class, 'fade-in-scale-up-transition-enter-active'))" + + " and not(contains(@class, 'fade-in-scale-up-transition-leave-active'))]" + + "//div[contains(text(), 'Basic Information')]/../following-sibling::div[contains(@class, 'n-card__footer')]//button[contains(@class, 'btn-submit')]") private WebElement buttonSubmit; @FindBys({ diff --git a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/task/SeaTunnelTaskForm.java b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/task/SeaTunnelTaskForm.java new file mode 100644 index 000000000000..aef526cb5081 --- /dev/null +++ b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/task/SeaTunnelTaskForm.java @@ -0,0 +1,54 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 org.apache.dolphinscheduler.e2e.pages.project.workflow.task; + +import org.apache.dolphinscheduler.e2e.core.WebDriverWaitFactory; +import org.apache.dolphinscheduler.e2e.pages.project.workflow.WorkflowForm; + +import org.openqa.selenium.By; +import org.openqa.selenium.JavascriptExecutor; +import org.openqa.selenium.WebElement; +import org.openqa.selenium.support.ui.ExpectedConditions; + +public class SeaTunnelTaskForm extends TaskNodeForm { + + public SeaTunnelTaskForm(WorkflowForm parent) { + super(parent); + } + + public SeaTunnelTaskForm executionType(String executionType) { + WebElement input = executionTypeInput(executionType); + WebElement label = input.findElement(By.xpath("..")); + ((JavascriptExecutor) parent().driver()).executeScript("arguments[0].scrollIntoView({block: 'center'});", + label); + label.click(); + WebDriverWaitFactory.createWebDriverWait(parent().driver()) + .until(ExpectedConditions.elementToBeSelected(input)); + return this; + } + + public boolean isExecutionType(String executionType) { + return executionTypeInput(executionType).isSelected(); + } + + private WebElement executionTypeInput(String executionType) { + return WebDriverWaitFactory.createWebDriverWait(parent().driver()) + .until(ExpectedConditions.presenceOfElementLocated(By.cssSelector( + ".seatunnel-execute-type input[value='" + executionType + "']"))); + } +} diff --git a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/task/TaskNodeForm.java b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/task/TaskNodeForm.java index bec099385d00..090d665ed7af 100644 --- a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/task/TaskNodeForm.java +++ b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/java/org/apache/dolphinscheduler/e2e/pages/project/workflow/task/TaskNodeForm.java @@ -145,14 +145,15 @@ public TaskNodeForm selectEnv(String envName) { public TaskNodeForm preTask(String preTaskName) { ((JavascriptExecutor) parent().driver()).executeScript("arguments[0].click();", selectPreTasks); - final By optionsLocator = By.className("option-pre-tasks"); + final By optionsLocator = By.cssSelector(".n-base-select-menu .n-base-select-option__content"); WebDriverWaitFactory.createWebDriverWait(parent.driver()) .until(ExpectedConditions.visibilityOfElementLocated(optionsLocator)); List webElements = parent.driver().findElements(optionsLocator); webElements.stream() - .filter(it -> it.getText().contains(preTaskName)) + .filter(WebElement::isDisplayed) + .filter(it -> it.getText().equals(preTaskName)) .findFirst() .orElseThrow(() -> new RuntimeException("No such task: " + preTaskName)) .click(); @@ -184,6 +185,8 @@ public TaskNodeForm selectResource(String resourceName) { public WorkflowForm submit() { buttonSubmit.click(); + WebDriverWaitFactory.createWebDriverWait(parent.driver()) + .until(ExpectedConditions.invisibilityOf(inputNodeName)); return parent(); } diff --git a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/resources/dragAndDrop.js b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/resources/dragAndDrop.js index 96011d9c327b..91252fc09e07 100644 --- a/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/resources/dragAndDrop.js +++ b/dolphinscheduler-e2e/dolphinscheduler-e2e-case/src/test/resources/dragAndDrop.js @@ -41,10 +41,19 @@ function dispatchEvent(element, event, transferData) { } } -function simulateHTML5DragAndDrop(element, destination) { +function simulateHTML5DragAndDrop(element, destination, x, y) { const dragStartEvent = createEvent('dragstart'); + if (x != null && y != null) { + dragStartEvent.offsetX = 0; + dragStartEvent.offsetY = 0; + } dispatchEvent(element, dragStartEvent); const dropEvent = createEvent('drop'); + if (x != null && y != null) { + const bounds = destination.getBoundingClientRect(); + dropEvent.clientX = bounds.left + x; + dropEvent.clientY = bounds.top + y; + } dispatchEvent(destination, dropEvent, dragStartEvent.dataTransfer); const dragEndEvent = createEvent('dragend'); dispatchEvent(element, dragEndEvent, dropEvent.dataTransfer); @@ -52,4 +61,4 @@ function simulateHTML5DragAndDrop(element, destination) { const source = arguments[0]; const destination = arguments[1]; -simulateHTML5DragAndDrop(source, destination); +simulateHTML5DragAndDrop(source, destination, arguments[2], arguments[3]); diff --git a/dolphinscheduler-ui/src/locales/en_US/project.ts b/dolphinscheduler-ui/src/locales/en_US/project.ts index 626410631650..4f25153ec7f0 100644 --- a/dolphinscheduler-ui/src/locales/en_US/project.ts +++ b/dolphinscheduler-ui/src/locales/en_US/project.ts @@ -383,6 +383,7 @@ export default { name_tips: 'Please enter name (required)', task_type: 'Task Type', task_type_tips: 'Please select a task type (required)', + task_execute_type: 'Task execution type', workflow_name: 'Workflow Name', child_node: 'Child Node', child_node_tips: 'Please select a child node (required)', diff --git a/dolphinscheduler-ui/src/locales/zh_CN/project.ts b/dolphinscheduler-ui/src/locales/zh_CN/project.ts index 90d526b175b7..246e9ba1f2fa 100644 --- a/dolphinscheduler-ui/src/locales/zh_CN/project.ts +++ b/dolphinscheduler-ui/src/locales/zh_CN/project.ts @@ -377,6 +377,7 @@ export default { name_tips: '请输入名称(必填)', task_type: '任务类型', task_type_tips: '请选择任务类型(必选)', + task_execute_type: '任务执行类型', child_node: '子节点', child_node_tips: '请选择子节点(必选)', run_flag: '运行标志', diff --git a/dolphinscheduler-ui/src/views/projects/task/components/node/detail-modal.tsx b/dolphinscheduler-ui/src/views/projects/task/components/node/detail-modal.tsx index e784947d5247..13c644fd1c90 100644 --- a/dolphinscheduler-ui/src/views/projects/task/components/node/detail-modal.tsx +++ b/dolphinscheduler-ui/src/views/projects/task/components/node/detail-modal.tsx @@ -224,6 +224,11 @@ const NodeDetailModal = defineComponent({ } const onTaskTypeChange = (taskType: ITaskType) => { + if (props.data.taskType !== taskType) { + // Let the new task type initialize its execution mode. + // eslint-disable-next-line vue/no-mutating-props + delete props.data.taskExecuteType + } // eslint-disable-next-line vue/no-mutating-props props.data.taskType = taskType initHeaderLinks(props.workflowInstance, props.data.taskType) diff --git a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sea-tunnel.ts b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sea-tunnel.ts index 7931e8304922..5e153af4b613 100644 --- a/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sea-tunnel.ts +++ b/dolphinscheduler-ui/src/views/projects/task/components/node/fields/use-sea-tunnel.ts @@ -55,6 +55,16 @@ export function useSeaTunnel(model: { [field: string]: any }): IJsonItem[] { ) return [ + { + type: 'radio', + field: 'taskExecuteType', + name: t('project.node.task_execute_type'), + class: 'seatunnel-execute-type', + options: [ + { label: t('project.task.batch_task'), value: 'BATCH' }, + { label: t('project.task.stream_task'), value: 'STREAM' } + ] + }, { type: 'select', field: 'startupScript', diff --git a/dolphinscheduler-ui/src/views/projects/task/components/node/use-task.ts b/dolphinscheduler-ui/src/views/projects/task/components/node/use-task.ts index c9c804188c93..977ff866c5c6 100644 --- a/dolphinscheduler-ui/src/views/projects/task/components/node/use-task.ts +++ b/dolphinscheduler-ui/src/views/projects/task/components/node/use-task.ts @@ -68,7 +68,9 @@ export function useTask({ model.preTasks = taskStore.getPreTasks model.name = taskStore.getName model.taskExecuteType = - TASK_TYPES_MAP[data.taskType || 'SHELL'].taskExecuteType || 'BATCH' + data.taskExecuteType || + TASK_TYPES_MAP[data.taskType || 'SHELL'].taskExecuteType || + 'BATCH' const getElements = () => { const { rules, elements } = getElementByJson(jsonRef.value, model) diff --git a/dolphinscheduler-ui/src/views/projects/task/instance/use-stream-table.ts b/dolphinscheduler-ui/src/views/projects/task/instance/use-stream-table.ts index 1aff5c4d6712..bf6957148072 100644 --- a/dolphinscheduler-ui/src/views/projects/task/instance/use-stream-table.ts +++ b/dolphinscheduler-ui/src/views/projects/task/instance/use-stream-table.ts @@ -160,6 +160,7 @@ export function useTable() { circle: true, type: 'info', size: 'small', + disabled: row.taskType === 'SEATUNNEL', onClick: () => onSavePoint(row.id) }, { diff --git a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts index 5e7497503eb7..a092de3a01dd 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts +++ b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/types.ts @@ -15,7 +15,7 @@ * limitations under the License. */ -import type { TaskType } from '@/store/project/types' +import type { TaskType, TaskExecuteType } from '@/store/project/types' export type { ITaskState } from '@/common/types' export interface WorkflowDefinition { @@ -70,6 +70,7 @@ export interface TaskDefinition { projectCode: any userId: number taskType: TaskType + taskExecuteType?: TaskExecuteType taskParams: any taskParamList: any[] taskParamMap: any diff --git a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-cell-update.ts b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-cell-update.ts index b4f87a76e6bc..771bfa8857d9 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-cell-update.ts +++ b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-cell-update.ts @@ -17,7 +17,7 @@ import type { Ref } from 'vue' import type { Graph } from '@antv/x6' -import type { TaskType } from '@/store/project/types' +import type { TaskType, TaskExecuteType } from '@/store/project/types' import type { Coordinate } from './types' import { TASK_TYPES_MAP } from '@/store/project/task-type' import { useCustomCellBuilder } from './dag-hooks' @@ -64,6 +64,16 @@ export function useCellUpdate(options: Options) { node.attr('rect/fill', color) } + function setNodeExecuteType(id: string, taskExecuteType: TaskExecuteType) { + graph.value?.getCellById(id)?.setData({ taskExecuteType }) + getNodeEdge(id).forEach((edge) => { + const isStream = + edge.getSourceNode()?.getData().taskExecuteType === 'STREAM' || + edge.getTargetNode()?.getData().taskExecuteType === 'STREAM' + edge.attr('line/strokeDasharray', isStream ? '5 5' : 'none') + }) + } + /** * Add a node to the graph * @param {string} id @@ -75,12 +85,13 @@ export function useCellUpdate(options: Options) { type: TaskType, name: string, flag: string, - coordinate: Coordinate = { x: 100, y: 100 } + coordinate: Coordinate = { x: 100, y: 100 }, + taskExecuteType?: TaskExecuteType ) { if (!TASK_TYPES_MAP[type as TaskType]) { return } - const node = buildNode(id, type, name, flag, coordinate) + const node = buildNode(id, type, name, flag, coordinate, taskExecuteType) graph.value?.addNode(node) } @@ -138,6 +149,7 @@ export function useCellUpdate(options: Options) { return { setNodeName, setNodeFillColor, + setNodeExecuteType, setNodeEdge, addNode, removeNode, diff --git a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-custom-cell-builder.ts b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-custom-cell-builder.ts index c9ca1a1277a7..09acbd087023 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-custom-cell-builder.ts +++ b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-custom-cell-builder.ts @@ -18,7 +18,7 @@ import type { Node, Edge } from '@antv/x6' import { X6_NODE_NAME, X6_EDGE_NAME } from './dag-config' import utils from '@/utils' -import { TaskType } from '@/store/project/types' +import { TaskType, TaskExecuteType } from '@/store/project/types' import { TASK_TYPES_MAP } from '@/store/project/task-type' import { WorkflowDefinition, Coordinate } from './types' @@ -75,7 +75,8 @@ export function useCustomCellBuilder() { type: TaskType, taskName: string, flag: string, - coordinate: Coordinate = { x: 100, y: 100 } + coordinate: Coordinate = { x: 100, y: 100 }, + taskExecuteType?: TaskExecuteType ): Node.Metadata { const truncation = taskName ? utils.truncateText(taskName, 18) : id return { @@ -87,7 +88,8 @@ export function useCustomCellBuilder() { taskType: type, taskName: taskName || id, flag: flag, - taskExecuteType: TASK_TYPES_MAP[type].taskExecuteType + taskExecuteType: + taskExecuteType || TASK_TYPES_MAP[type].taskExecuteType || 'BATCH' }, attrs: { image: { @@ -122,7 +124,7 @@ export function useCustomCellBuilder() { parseLocationStr(definition.workflowDefinition.locations) || [] const tasks = definition.taskDefinitionList const connects = definition.workflowTaskRelationList - const taskTypeMap = {} as { [key in string]: TaskType } + const taskExecuteTypeMap = {} as { [key in string]: TaskExecuteType } tasks.forEach((task) => { const location = locations.find((l) => l.taskCode === task.code) || {} @@ -134,20 +136,19 @@ export function useCustomCellBuilder() { { x: location.x, y: location.y - } + }, + task.taskExecuteType ) nodes.push(node) - taskTypeMap[String(task.code)] = task.taskType + taskExecuteTypeMap[String(task.code)] = node.data.taskExecuteType }) connects .filter((r) => !!r.preTaskCode) .forEach((c) => { const isStream = - TASK_TYPES_MAP[taskTypeMap[c.preTaskCode]].taskExecuteType === - 'STREAM' || - TASK_TYPES_MAP[taskTypeMap[c.postTaskCode]].taskExecuteType === - 'STREAM' + taskExecuteTypeMap[c.preTaskCode] === 'STREAM' || + taskExecuteTypeMap[c.postTaskCode] === 'STREAM' const edge = buildEdge( c.preTaskCode + '', c.postTaskCode + '', diff --git a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-task-edit.ts b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-task-edit.ts index 1c23885c61c1..253d4e451c74 100644 --- a/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-task-edit.ts +++ b/dolphinscheduler-ui/src/views/projects/workflow/components/dag/use-task-edit.ts @@ -18,6 +18,7 @@ import { ref, onMounted, watch } from 'vue' import { remove, cloneDeep } from 'lodash' import { TaskType } from '@/store/project/types' +import { TASK_TYPES_MAP } from '@/store/project/task-type' import { formatParams } from '@/views/projects/task/components/node/format-data' import { useCellUpdate } from './dag-hooks' import type { Ref } from 'vue' @@ -48,6 +49,7 @@ export function useTaskEdit(options: Options) { getTargets, setNodeName, setNodeFillColor, + setNodeExecuteType, setNodeEdge } = useCellUpdate({ graph @@ -91,10 +93,17 @@ export function useTaskEdit(options: Options) { flag: string, coordinate: Coordinate ) { - addNode(code + '', type, name, flag, coordinate) const definition = workflowDefinition.value.taskDefinitionList.find( (t) => t.code === targetCode ) + addNode( + code + '', + type, + name, + flag, + coordinate, + definition?.taskExecuteType + ) const newDefinition = { ...cloneDeep(definition), @@ -149,7 +158,7 @@ export function useTaskEdit(options: Options) { (t) => t.code === code ) if (definition) { - currTask.value = definition + currTask.value = cloneDeep(definition) } updatePreTasks(getSources(String(code)), code) updatePostTasks(code) @@ -177,6 +186,12 @@ export function useTaskEdit(options: Options) { setNodeFillColor(task.code + '', fillColor) setNodeEdge(String(task.code), data.preTasks) + setNodeExecuteType( + String(task.code), + taskDef.taskExecuteType || + TASK_TYPES_MAP[currTask.value.taskType].taskExecuteType || + 'BATCH' + ) updatePreTasks(data.preTasks, task.code) return { ...taskDef, diff --git a/dolphinscheduler-ui/tests/task-execution-graph.cjs b/dolphinscheduler-ui/tests/task-execution-graph.cjs new file mode 100644 index 000000000000..e4743e80715e --- /dev/null +++ b/dolphinscheduler-ui/tests/task-execution-graph.cjs @@ -0,0 +1,499 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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. + */ + +// Run from dolphinscheduler-ui: node tests/task-execution-graph.cjs [commit-SHA] +// Executes production editor, graph update/builder, serializer, business mapper, +// and stream-table renderer/service code. X6 storage and HTTP are adapters; +// this does not render a browser, validate forms, or execute task engines. +const assert = require('assert/strict') +const fs = require('fs') +const path = require('path') +const vm = require('vm') +const { execFileSync } = require('child_process') +const ts = require('typescript') +const vue = require('vue') +const { cloneDeep, get, set } = require('lodash') +const { NButton } = require('naive-ui') +const ui = path.resolve(__dirname, '..') +const revision = process.argv[2] +if (revision && !/^[a-f0-9]{40}$/.test(revision)) { + throw new Error('Provide an exact 40-character commit SHA') +} +const dagRoot = 'views/projects/workflow/components/dag/' +const cache = new Map() +const requests = [] +const messages = [] +const axios = (request) => { + requests.push(cloneDeep(request)) + return Promise.resolve({ total: 0, totalList: [], totalPage: 0 }) +} +const source = (file) => + revision + ? execFileSync( + 'git', + ['show', `${revision}:dolphinscheduler-ui/src/${file}`], + { cwd: path.dirname(ui), encoding: 'utf8' } + ) + : fs.readFileSync(path.join(ui, 'src', file), 'utf8') + +function load(file) { + if (cache.has(file)) return cache.get(file) + const module = { exports: {} } + // Vite's asset URL is irrelevant to graph data and unavailable in CommonJS. + const code = ts.transpileModule( + source(file).replace(/import\.meta\.env\.BASE_URL/g, JSON.stringify('/')), + { + compilerOptions: { + module: ts.ModuleKind.CommonJS, + target: ts.ScriptTarget.ES2020, + esModuleInterop: true + } + } + ).outputText + function localRequire(id) { + if (id === './dag-hooks') { + return { + useCellUpdate: (...args) => + load(dagRoot + 'use-cell-update.ts').useCellUpdate(...args), + useCustomCellBuilder: (...args) => + load(dagRoot + 'use-custom-cell-builder.ts').useCustomCellBuilder( + ...args + ) + } + } + if (id === '@/utils') { + return { + __esModule: true, + default: { truncateText: (text, length) => text.slice(0, length) } + } + } + if (id === '@/service/service') return { axios } + if (id === '@/service/modules/task-instances') { + return load('service/modules/task-instances/index.ts') + } + if (id === 'vue-i18n') return { useI18n: () => ({ t: (key) => key }) } + if (id === 'vue-router') { + return { useRoute: () => ({ params: { projectCode: '123' } }) } + } + if (id === '@/common/common') { + return { parseTime: (value) => value, renderTableTime: () => '' } + } + if (!id.startsWith('.') && !id.startsWith('@/')) return require(id) + const target = id.startsWith('@/') + ? id.slice(2) + : path.posix.join(path.posix.dirname(file), id) + return load(target + (path.posix.extname(target) ? '' : '.ts')) + } + vm.runInThisContext(`(function(require,module,exports,window){${code}\n})`, { + filename: file + })(localRequire, module, module.exports, { + $message: { success: (message) => messages.push(message) } + }) + cache.set(file, module.exports) + return module.exports +} + +// Only storage/endpoints/attributes are modeled here. In particular this +// adapter never derives execution modes or edge styling from task types. +class Graph { + constructor(json) { + this.nodes = new Map() + this.edges = [] + json.nodes.forEach((node) => this.addNode(node)) + json.edges.forEach((edge) => this.addEdge(edge)) + } + on() {} + addNode(metadata) { + const node = cloneDeep(metadata) + node.getData = () => node.data + node.setData = (data) => Object.assign(node.data, data) + node.getPosition = () => ({ x: node.x, y: node.y }) + node.attr = (key, value) => attribute(node, key, value) + this.nodes.set(node.id, node) + return node + } + addEdge(metadata) { + const edge = cloneDeep(metadata) + edge.getSourceCellId = () => edge.source.cell + edge.getTargetCellId = () => edge.target.cell + edge.getSourceNode = () => this.getCellById(edge.source.cell) + edge.getTargetNode = () => this.getCellById(edge.target.cell) + edge.getLabels = () => edge.labels || [] + edge.attr = (key, value) => attribute(edge, key, value) + this.edges.push(edge) + return edge + } + removeEdge(edge) { + this.edges = this.edges.filter((item) => item !== edge) + } + getCellById(id) { + return this.nodes.get(id) + } + getConnectedEdges(node) { + return this.edges.filter( + (edge) => edge.source.cell === node.id || edge.target.cell === node.id + ) + } + getNodes() { + return [...this.nodes.values()] + } + getEdges() { + return this.edges + } +} +function attribute(cell, key, value) { + const parts = key.split('/') + if (value === undefined) return get(cell.attrs, parts) + set(cell.attrs, parts, value) +} + +const { buildGraph } = load( + dagRoot + 'use-custom-cell-builder.ts' +).useCustomCellBuilder() +const { getConnects, getLocations } = load( + dagRoot + 'use-business-mapper.ts' +).useBusinessMapper() +const { formatModel } = load( + 'views/projects/task/components/node/format-data.ts' +) + +function task(code, taskType = 'SHELL', taskExecuteType = 'BATCH') { + return { + code, + id: code + 100, + version: 7, + name: `task-${code}`, + taskType, + taskExecuteType, + flag: 'YES', + taskParams: { + useCustom: true, + startupScript: 'seatunnel.sh', + rawScript: 'env { job.mode = "BATCH" }', + localParams: [{ prop: 'key', value: 'source' }], + resourceList: [] + } + } +} +function makeEditor(tasks, connections = []) { + const definition = vue.ref({ + workflowDefinition: { locations: '[]' }, + taskDefinitionList: tasks, + workflowTaskRelationList: connections.map( + ([preTaskCode, postTaskCode]) => ({ + preTaskCode, + postTaskCode + }) + ) + }) + const graph = new Graph(buildGraph(definition.value)) + let editor + const app = vue + .createRenderer({ + createComment: () => ({}), + insert() {}, + remove() {}, + parentNode() {}, + nextSibling() {} + }) + .createApp({ + setup() { + editor = load(dagRoot + 'use-task-edit.ts').useTaskEdit({ + graph: vue.shallowRef(graph), + definition + }) + return () => null + } + }) + app.mount({}) + return { editor, graph, definition, close: () => app.unmount() } +} +function confirm(fixture, mode, preTasks = [10], code = 20) { + fixture.editor.editTask(code) + const model = formatModel(fixture.editor.currTask.value) + model.taskExecuteType = mode + model.preTasks = preTasks + fixture.editor.taskConfirm({ data: model }) + assert.equal(fixture.editor.taskModalVisible.value, false) +} +function reload(fixture) { + const { graph, definition } = fixture + // The workflow save payload uses these production mappers and task list. + // JSON round-tripping represents persistence without calling an API. + const saved = JSON.parse( + JSON.stringify({ + workflowDefinition: { + locations: JSON.stringify(getLocations(graph.getNodes())) + }, + taskDefinitionList: definition.value.taskDefinitionList, + workflowTaskRelationList: getConnects( + graph.getNodes(), + graph.getEdges(), + definition.value.taskDefinitionList + ) + }) + ) + return { saved, graph: new Graph(buildGraph(saved)) } +} +function styles(graph) { + return graph + .getEdges() + .map((edge) => [ + `${edge.getSourceCellId()}->${edge.getTargetCellId()}`, + edge.attr('line/strokeDasharray') + ]) + .sort(([a], [b]) => a.localeCompare(b)) +} + +const checks = [] +async function check(name, run) { + try { + await run() + checks.push({ name, result: 'PASS' }) + } catch (error) { + checks.push({ name, result: 'FAIL', error: error.stack }) + } +} +async function main() { + await check( + 'STREAM SeaTunnel copy preserves graph mode and independent parameters', + () => { + const fixture = makeEditor([task(20, 'SEATUNNEL', 'STREAM')]) + try { + fixture.editor.copyTask('copied', 30, 20, 'SEATUNNEL', 'YES', { + x: 100, + y: 100 + }) + const [original, copied] = fixture.definition.value.taskDefinitionList + assert.equal(copied.code, 30) + assert.equal(copied.name, 'copied') + assert.equal(copied.taskType, 'SEATUNNEL') + assert.equal(copied.taskExecuteType, 'STREAM') + assert.equal(original.taskExecuteType, 'STREAM') + assert.equal( + fixture.graph.getCellById('30').data.taskExecuteType, + 'STREAM' + ) + assert.notEqual(copied.taskParams, original.taskParams) + copied.taskParams.localParams[0].value = 'copy' + assert.equal(original.taskParams.localParams[0].value, 'source') + const persisted = reload(fixture) + assert.equal( + persisted.graph.getCellById('30').data.taskExecuteType, + 'STREAM' + ) + assert.equal( + persisted.saved.taskDefinitionList[1].taskExecuteType, + 'STREAM' + ) + confirm(fixture, 'BATCH', [], 30) + assert.equal( + fixture.graph.getCellById('30').data.taskExecuteType, + 'BATCH' + ) + assert.equal( + fixture.graph.getCellById('20').data.taskExecuteType, + 'STREAM' + ) + assert.equal(original.taskExecuteType, 'STREAM') + } finally { + fixture.close() + } + } + ) + + await check( + 'Confirm updates incoming and outgoing edges BATCH -> STREAM -> BATCH', + () => { + const fixture = makeEditor( + [task(10), task(20, 'SEATUNNEL'), task(30), task(40), task(50)], + [ + [10, 20], + [20, 30], + [40, 50] + ] + ) + try { + const unrelated = fixture.graph.getEdges()[2] + assert.deepEqual(styles(fixture.graph), [ + ['10->20', 'none'], + ['20->30', 'none'], + ['40->50', 'none'] + ]) + for (const [mode, dash] of [ + ['STREAM', '5 5'], + ['BATCH', 'none'] + ]) { + const incoming = fixture.graph + .getEdges() + .find((edge) => edge.getTargetCellId() === '20') + confirm(fixture, mode) + assert.equal(fixture.graph.getEdges().includes(incoming), false) + assert.equal(fixture.graph.getEdges().includes(unrelated), true) + assert.equal( + fixture.graph.getCellById('20').data.taskExecuteType, + mode + ) + const expected = [ + ['10->20', dash], + ['20->30', dash], + ['40->50', 'none'] + ] + assert.deepEqual(styles(fixture.graph), expected) + const persisted = reload(fixture) + assert.deepEqual(styles(persisted.graph), expected) + assert.equal( + persisted.saved.taskDefinitionList.find((item) => item.code === 20) + .taskExecuteType, + mode + ) + } + } finally { + fixture.close() + } + } + ) + + for (const direction of ['incoming', 'outgoing']) { + await check( + `BATCH SeaTunnel retains dashed ${direction} edge to a STREAM neighbor`, + () => { + const fixture = makeEditor( + [ + direction === 'incoming' + ? task(10, 'FLINK_STREAM', 'STREAM') + : task(10), + task(20, 'SEATUNNEL', 'STREAM'), + direction === 'outgoing' + ? task(30, 'FLINK_STREAM', 'STREAM') + : task(30) + ], + [ + [10, 20], + [20, 30] + ] + ) + try { + confirm(fixture, 'BATCH') + const expected = [ + ['10->20', direction === 'incoming' ? '5 5' : 'none'], + ['20->30', direction === 'outgoing' ? '5 5' : 'none'] + ] + assert.deepEqual(styles(fixture.graph), expected) + assert.deepEqual(styles(reload(fixture).graph), expected) + const neighbor = direction === 'incoming' ? '10' : '30' + assert.equal( + fixture.graph.getCellById(neighbor).data.taskExecuteType, + 'STREAM' + ) + } finally { + fixture.close() + } + } + ) + } + + await check( + 'Legacy FLINK_STREAM null mode remains dashed after Confirm and reload', + () => { + const fixture = makeEditor( + [task(10), task(20, 'FLINK_STREAM', null), task(30)], + [ + [10, 20], + [20, 30] + ] + ) + try { + confirm(fixture, null) + assert.equal( + fixture.graph.getCellById('20').data.taskExecuteType, + 'STREAM' + ) + const expected = [ + ['10->20', '5 5'], + ['20->30', '5 5'] + ] + assert.deepEqual(styles(fixture.graph), expected) + const persisted = reload(fixture) + assert.deepEqual(styles(persisted.graph), expected) + assert.equal( + persisted.saved.taskDefinitionList[1].taskExecuteType, + null + ) + } finally { + fixture.close() + } + } + ) + + const table = load( + 'views/projects/task/instance/use-stream-table.ts' + ).useTable() + table.createColumns(table.variables) + const operation = table.variables.columns.find( + (column) => column.key === 'operation' + ) + function savepoint(taskType) { + const tooltips = operation.render({ id: 456, taskType }).children.default() + const matches = tooltips.filter( + (tooltip) => tooltip.children.default() === 'project.task.savepoint' + ) + assert.equal(matches.length, 1) + const button = matches[0].children.trigger() + assert.equal(button.type, NButton) + return button + } + await check('SeaTunnel savepoint renderer disables its button', () => { + assert.equal(savepoint('SEATUNNEL').props.disabled, true) + assert.equal(requests.length, 0) + }) + await check('Flink savepoint renderer keeps its button enabled', () => { + assert.equal(savepoint('FLINK_STREAM').props.disabled, false) + assert.equal(requests.length, 0) + }) + await check( + 'Enabled Flink savepoint calls the correct endpoint and refreshes Stream tasks', + async () => { + // Calling a disabled VNode callback would bypass Naive UI's click guard. + // Only invoke the enabled callback; disabled DOM interaction is an E2E concern. + const button = savepoint('FLINK_STREAM') + assert.equal(button.props.disabled, false) + button.props.onClick() + await new Promise((resolve) => setImmediate(resolve)) + assert.equal(requests.length, 2) + assert.deepEqual(requests[0], { + url: 'projects/123/task-instances/456/savepoint', + method: 'post' + }) + assert.equal(requests[1].url, '/projects/123/task-instances') + assert.equal(requests[1].method, 'get') + assert.equal(requests[1].params.taskExecuteType, 'STREAM') + assert.deepEqual(messages, ['project.task.success']) + } + ) + + for (const item of checks) { + process.stdout.write( + `${item.result} ${item.name}${item.error ? ': ' + item.error : ''}\n` + ) + } + process.exitCode = checks.some((item) => item.result === 'FAIL') ? 1 : 0 +} +main().catch((error) => { + process.stderr.write(`${error.stack}\n`) + process.exitCode = 1 +}) diff --git a/dolphinscheduler-ui/tests/task-execution-type.cjs b/dolphinscheduler-ui/tests/task-execution-type.cjs new file mode 100644 index 000000000000..72310f8d7ec7 --- /dev/null +++ b/dolphinscheduler-ui/tests/task-execution-type.cjs @@ -0,0 +1,394 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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. + */ + +// Run from dolphinscheduler-ui: node tests/task-execution-type.cjs [commit-SHA] +// Uses the installed UI dependencies; no additional test framework is required. +// Executes the production callback, models, defaults, and serializers. Unrelated +// field rendering and store services are stubbed. This is not a browser/API E2E. +const assert = require('assert/strict') +const fs = require('fs') +const path = require('path') +const vm = require('vm') +const { execFileSync } = require('child_process') +const ts = require('typescript') +const vue = require('vue') +const ui = path.resolve(__dirname, '..') +const revision = process.argv[2] +if (revision && !/^[a-f0-9]{40}$/.test(revision)) { + throw new Error('Provide an exact 40-character commit SHA') +} +const taskRoot = 'views/projects/task/components/node/' +const source = (file) => + revision + ? execFileSync( + 'git', + ['show', `${revision}:dolphinscheduler-ui/src/${file}`], + { + cwd: path.dirname(ui), + encoding: 'utf8' + } + ) + : fs.readFileSync(path.join(ui, 'src', file), 'utf8') +const compilerOptions = { + module: ts.ModuleKind.CommonJS, + target: ts.ScriptTarget.ES2020, + esModuleInterop: true +} +const cache = new Map() +const graphModes = new Map() +const fields = new Proxy( + { __esModule: true }, + { get: (target, key) => (key === '__esModule' ? true : () => []) } +) +const store = { + updateDefinition() {}, + getName: 'execution-mode-test', + getPreTasks: [] +} +function load(file) { + if (cache.has(file)) return cache.get(file) + const module = { exports: {} } + const code = ts.transpileModule(source(file), { compilerOptions }).outputText + function localRequire(id) { + if (id === './dag-hooks') { + return { + useCellUpdate: () => ({ + addNode() {}, + removeNode() {}, + getSources: () => [], + getTargets: () => [], + setNodeName() {}, + setNodeFillColor() {}, + setNodeExecuteType(code, mode) { + graphModes.set(code, mode) + }, + setNodeEdge() {} + }) + } + } + if (id === '../fields/index') return fields + if (id === './tasks') { + return { + __esModule: true, + default: { + SHELL: load(taskRoot + 'tasks/use-shell.ts').useShell, + FLINK: load(taskRoot + 'tasks/use-flink.ts').useFlink, + FLINK_STREAM: load(taskRoot + 'tasks/use-flink-stream.ts') + .useFlinkStream, + SEATUNNEL: load(taskRoot + 'tasks/use-sea-tunnel.ts').useSeaTunnel + } + } + } + if (id === '@/components/form/get-elements-by-json') { + return { __esModule: true, default: () => ({ rules: {}, elements: [] }) } + } + if (id === '@/store/project/task-node') { + return { useTaskNodeStore: () => store } + } + if (!id.startsWith('.') && !id.startsWith('@/')) return require(id) + const target = id.startsWith('@/') + ? id.slice(2) + : path.posix.join(path.posix.dirname(file), id) + return load(target + '.ts') + } + vm.runInThisContext(`(function(require,module,exports){${code}\n})`, { + filename: file + })(localRequire, module, module.exports) + cache.set(file, module.exports) + return module.exports +} + +// Extract the actual TSX callback so the test cannot diverge from its guard. +const modal = ts.createSourceFile( + 'detail-modal.tsx', + source(taskRoot + 'detail-modal.tsx'), + ts.ScriptTarget.Latest, + true, + ts.ScriptKind.TSX +) +let callback +function visit(node) { + if ( + ts.isVariableDeclaration(node) && + node.name.getText(modal) === 'onTaskTypeChange' + ) { + callback = node.initializer.getText(modal) + } + ts.forEachChild(node, visit) +} +visit(modal) +assert(callback, 'onTaskTypeChange callback must exist') +const changeCode = ts.transpileModule(`(${callback})(nextType)`, { + compilerOptions +}).outputText +function changeTaskType(data, nextType) { + vm.runInNewContext(changeCode, { + props: { data }, + nextType, + initHeaderLinks() {} + }) +} +const { useTask } = load(taskRoot + 'use-task.ts') +const { formatParams, formatModel } = load(taskRoot + 'format-data.ts') +const { useForm } = load('components/form/use-form.ts') +const makeForm = (data, from = 0) => + useTask({ data, projectCode: 1, from }).model +const save = (model) => formatParams(model).taskDefinitionJsonObj +function reopen(data) { + const model = makeForm(data) + const form = useForm() + form.state.formRef = { model } + form.setValues(formatModel(data)) + return form.getValues() +} +const checks = [] +function check(name, run) { + try { + run() + checks.push({ name, result: 'PASS' }) + } catch (error) { + checks.push({ name, result: 'FAIL', error: error.message }) + } +} + +// Mount the real editor hook so its lifecycle and reactive ownership are active. +// Graph rendering is irrelevant to the task-definition Cancel/Save contract. +function makeEditor(taskType = 'SEATUNNEL', taskExecuteType = 'STREAM') { + graphModes.clear() + const definition = vue.ref({ + workflowDefinition: {}, + workflowTaskRelationList: [], + taskDefinitionList: [ + { + id: 5, + code: 20, + version: 3, + name: 'saved', + taskType, + taskExecuteType, + taskParams: { localParams: [{ prop: 'key', value: 'original' }] } + } + ] + }) + let editor + const app = vue + .createRenderer({ + createComment: () => ({}), + insert() {}, + remove() {}, + parentNode() {}, + nextSibling() {} + }) + .createApp({ + setup() { + editor = load( + 'views/projects/workflow/components/dag/use-task-edit.ts' + ).useTaskEdit({ graph: vue.ref(), definition }) + return () => null + } + }) + app.mount({}) + return { definition, editor, close: () => app.unmount() } +} + +for (const nextType of ['SHELL', 'FLINK_STREAM', 'SEATUNNEL']) { + check( + `Cancel ${nextType} edit preserves saved task and nested parameters`, + () => { + const { definition, editor, close } = makeEditor() + try { + const original = definition.value.taskDefinitionList[0] + const snapshot = JSON.parse(JSON.stringify(original)) + editor.editTask(20) + changeTaskType(editor.currTask.value, nextType) + editor.currTask.value.taskParams.localParams[0].value = 'draft' + editor.taskCancel() + assert.equal(editor.taskModalVisible.value, false) + assert.equal(definition.value.taskDefinitionList[0], original) + assert.deepEqual(JSON.parse(JSON.stringify(original)), snapshot) + } finally { + close() + } + } + ) +} + +for (const nextType of ['SEATUNNEL', 'FLINK_STREAM', 'SHELL']) { + check( + `Save ${nextType} edit commits mode and preserves task identity`, + () => { + const { definition, editor, close } = makeEditor() + try { + const original = definition.value.taskDefinitionList[0] + const parent = definition.value + editor.editTask(20) + changeTaskType(editor.currTask.value, nextType) + const model = makeForm(editor.currTask.value, 1) + model.taskExecuteType = nextType === 'FLINK_STREAM' ? 'STREAM' : 'BATCH' + model.name = 'edited' + editor.taskConfirm({ data: model }) + const saved = definition.value.taskDefinitionList[0] + assert.equal(editor.workflowDefinition.value, parent) + assert.notEqual(saved, original) + assert.equal(saved.taskType, nextType) + assert.equal(saved.taskExecuteType, model.taskExecuteType) + assert.equal(saved.name, 'edited') + assert.equal(saved.id, 5) + assert.equal(saved.code, 20) + assert.equal(saved.version, 3) + assert.equal(editor.taskModalVisible.value, false) + } finally { + close() + } + } + ) +} + +for (const [taskType, mode, expected] of [ + ['FLINK_STREAM', null, 'STREAM'], + ['SHELL', null, 'BATCH'], + ['SEATUNNEL', null, 'BATCH'], + ['FLINK_STREAM', 'BATCH', 'BATCH'], + ['FLINK_STREAM', 'STREAM', 'STREAM'], + ['SEATUNNEL', 'BATCH', 'BATCH'], + ['SEATUNNEL', 'STREAM', 'STREAM'] +]) { + check(`Confirm ${taskType}/${mode} keeps graph mode ${expected}`, () => { + const { definition, editor, close } = makeEditor(taskType, mode) + try { + editor.editTask(20) + const model = reopen(editor.currTask.value) + assert.equal(model.taskExecuteType, mode) + editor.taskConfirm({ data: model }) + assert.equal(graphModes.get('20'), expected) + assert.equal(definition.value.taskDefinitionList[0].taskExecuteType, mode) + } finally { + close() + } + }) +} + +check('Confirm graph fallback follows the edited type being saved', () => { + const { definition, editor, close } = makeEditor('SHELL', 'BATCH') + try { + editor.editTask(20) + changeTaskType(editor.currTask.value, 'FLINK_STREAM') + const model = makeForm(editor.currTask.value) + model.taskExecuteType = null + editor.taskConfirm({ data: model }) + assert.equal(graphModes.get('20'), 'STREAM') + assert.equal( + definition.value.taskDefinitionList[0].taskType, + 'FLINK_STREAM' + ) + assert.equal(definition.value.taskDefinitionList[0].taskExecuteType, null) + } finally { + close() + } +}) + +for (const [oldType, oldMode, nextType, expected] of [ + ['SHELL', 'BATCH', 'FLINK_STREAM', 'STREAM'], + ['FLINK', 'BATCH', 'FLINK_STREAM', 'STREAM'], + ['FLINK_STREAM', 'STREAM', 'SHELL', 'BATCH'], + ['SEATUNNEL', 'STREAM', 'SHELL', 'BATCH'], + ['FLINK_STREAM', 'STREAM', 'SEATUNNEL', 'BATCH'], + ['SHELL', 'BATCH', 'SEATUNNEL', 'BATCH'] +]) { + check(`${oldType}/${oldMode} -> ${nextType}/${expected}`, () => { + // from=1 exposes the component's dormant task-type selector. + const data = vue.reactive({ + id: '5', + code: 20, + name: 'saved', + taskType: oldType, + taskExecuteType: oldMode, + taskParams: {} + }) + const params = data.taskParams + changeTaskType(data, nextType) + const payload = save(makeForm(data, 1)) + assert.equal(payload.taskType, nextType) + assert.equal(payload.taskExecuteType, expected) + assert.equal(data.code, 20) + assert.equal(data.name, 'saved') + assert.equal(data.taskParams, params) + }) +} +for (const mode of ['BATCH', 'STREAM']) { + check(`same-type SeaTunnel keeps ${mode}`, () => { + const data = { taskType: 'SEATUNNEL', taskExecuteType: mode } + changeTaskType(data, 'SEATUNNEL') + assert.deepEqual(data, { taskType: 'SEATUNNEL', taskExecuteType: mode }) + assert.equal(save(makeForm(data, 1)).taskExecuteType, mode) + }) +} +check('loaded SeaTunnel STREAM initializes without losing its mode', () => { + assert.equal( + makeForm({ id: '5', taskType: 'SEATUNNEL', taskExecuteType: 'STREAM' }) + .taskExecuteType, + 'STREAM' + ) +}) +for (const useCustom of [true, false]) { + check(`SeaTunnel save/reopen with useCustom=${useCustom}`, () => { + let model = makeForm({ taskType: 'SEATUNNEL' }) + model.useCustom = useCustom + model.rawScript = 'env { job.mode = "STREAMING" }' + model.resourceList = useCustom ? [] : ['/seatunnel/config.conf'] + for (const mode of ['BATCH', 'STREAM', 'BATCH']) { + model.taskExecuteType = mode + const payload = save(model) + assert.equal(payload.taskExecuteType, mode) + assert.equal('taskExecuteType' in payload.taskParams, false) + model = reopen(JSON.parse(JSON.stringify(payload))) + assert.equal(model.taskExecuteType, mode) + assert.equal(model.useCustom, useCustom) + assert.equal( + model.rawScript, + useCustom ? 'env { job.mode = "STREAMING" }' : '' + ) + assert.deepEqual( + Array.from(model.resourceList), + useCustom ? [] : ['/seatunnel/config.conf'] + ) + } + }) +} +for (const [taskType, mode] of [ + ['FLINK_STREAM', 'STREAM'], + ['SEATUNNEL', 'BATCH'] +]) { + check(`legacy ${taskType} definition defaults to ${mode}`, () => { + const model = reopen({ id: '5', taskType, taskParams: {} }) + assert.equal(model.taskExecuteType, mode) + assert.equal(save(model).taskExecuteType, mode) + }) +} +for (const item of checks) { + process.stdout.write( + `${item.result} ${item.name}${item.error ? ': ' + item.error : ''}\n` + ) +} +if (process.env.TASK_EXECUTION_REPORT_PATH) { + fs.writeFileSync( + process.env.TASK_EXECUTION_REPORT_PATH, + JSON.stringify({ revision: revision || 'working-tree', checks }, null, 2) + + '\n' + ) +} +process.exitCode = checks.some((item) => item.result === 'FAIL') ? 1 : 0