Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
f72c7aa
[Fix-18672][UI] Allow SeaTunnel workflow tasks to select stream execu…
sunlishuo25 Oct 3, 2026
c46368c
[Fix-18672][CI] Register SeaTunnel workflow regression test
sunlishuo25 Oct 3, 2026
e2675d4
[Fix-18672][E2E] Scroll SeaTunnel execution type into view
sunlishuo25 Oct 3, 2026
487ee30
Merge branch 'dev' into Fix-18672
SbloodyS Oct 5, 2026
e18cf87
[Fix-18672][UI] Reset execution mode when task type changes
sunlishuo25 Oct 5, 2026
1c879a4
fix(ui): isolate task editor changes until confirmation
sunlishuo25 Oct 5, 2026
eed7221
Merge branch 'dev' into Fix-18672
SbloodyS Oct 5, 2026
85f1397
Merge local task editor repairs onto maintainer update
sunlishuo25 Oct 5, 2026
81c5273
fix(ui): retain task type graph default on confirmation
sunlishuo25 Oct 5, 2026
d1cb01f
Merge branch 'dev' into Fix-18672
SbloodyS Oct 6, 2026
74650d7
Merge validated task editor repairs with latest maintainer changes
sunlishuo25 Oct 6, 2026
4ee3eef
test: cover SeaTunnel copy graph mode and instance classification
sunlishuo25 Oct 6, 2026
5681f5a
[Fix-18672][E2E] Center task bodies before canvas pointer actions
sunlishuo25 Oct 7, 2026
434b40a
[Fix-18672][E2E] Verify task identity before native pointer actions
sunlishuo25 Oct 7, 2026
565555f
[Fix-18672][E2E] Retry instance reads after table refresh
sunlishuo25 Oct 8, 2026
201bb6c
[Fix-18672][E2E] Wait for save modal transitions before clicking
sunlishuo25 Oct 8, 2026
f1acc46
[Fix-18672][E2E] Wait for task modal readiness before editing
sunlishuo25 Oct 8, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .github/workflows/e2e.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 4 additions & 0 deletions .github/workflows/frontend.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions docs/docs/en/guide/task/seatunnel.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions docs/docs/zh/guide/task/seatunnel.md
Original file line number Diff line number Diff line change
Expand Up @@ -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` 两种模式
Expand Down
Original file line number Diff line number Diff line change
@@ -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.<ShellTaskForm>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.<ShellTaskForm>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<String> 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<String> 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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -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 {
Expand All @@ -49,6 +51,21 @@ public List<Row> 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<Row> 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 {

Expand All @@ -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());
}
Expand Down
Loading
Loading