Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
e9efa47
TS SDK: cover native TypeScript Dags end to end
jason810496 Sep 29, 2026
b86b364
Parse the native TypeScript Dag in the compose e2e
jason810496 Sep 28, 2026
86824c2
Check the native TypeScript Dag's source and triggered run
jason810496 Sep 28, 2026
84897bb
Keep the containerized Java build's Maven repository on the host
jason810496 Sep 28, 2026
e5b7824
Add a native Java Dag to the compose e2e
jason810496 Sep 28, 2026
eced559
Run the Lang-SDK e2e jobs when native Dag parsing changes
jason810496 Sep 28, 2026
0091f00
Share the Dag run lookup and the native Java queue name
jason810496 Sep 28, 2026
3c61f55
Build the k8s TypeScript artifacts only for the native Dag test
jason810496 Sep 28, 2026
8422181
Key the Lang-SDK e2e jobs on native Dag modules only
jason810496 Oct 2, 2026
9514776
Build the TypeScript SDK before packing the k8s example bundle
jason810496 Oct 2, 2026
253bfb8
Pack the TypeScript example against the built SDK
jason810496 Oct 2, 2026
f2a9f33
Explain why the k8s native TypeScript Dag test is skipped
jason810496 Oct 2, 2026
454c2e9
Document how the TypeScript example's native Dag gets parsed
jason810496 Oct 2, 2026
2ede321
Name the skipped-branch e2e test after what it checks
jason810496 Oct 2, 2026
923697e
Parse native Dags in the e2e by class, not by Dag bundle name
jason810496 Oct 3, 2026
931ea73
Drop the dag_bundle_name reason from the k8s native TypeScript skip
jason810496 Oct 3, 2026
2e1ce89
Say which coordinator runs the native e2e Java tasks
jason810496 Oct 3, 2026
9be77c2
Read a native Dag example's served Dags with bundleDags
jason810496 Oct 3, 2026
383ccfe
Keep the k8s lang-SDK values comment within the YAML line limit
jason810496 Oct 5, 2026
18a9601
Rework the lang-SDK k8s harness for mixed-language and native Dag tests
jason810496 Oct 5, 2026
d313c95
Exclude airflowignore_all from the license check
jason810496 Oct 6, 2026
3187346
Cover If, Switch, a task group and TriggerDagRun in the native Java D…
jason810496 Oct 9, 2026
b91e38d
Cover the same native Java Dag features in the k8s native Dag test
jason810496 Oct 9, 2026
c3e39f1
Fix the native TypeScript example's trigger to the current Dag.task API
jason810496 Oct 9, 2026
d610758
Remove client methods a rebase onto main duplicated
jason810496 Oct 10, 2026
5350cf9
Fix the TS Code view test for per-Dag source
jason810496 Oct 10, 2026
6f6e30c
Fix the Java trigger test's XCom key
jason810496 Oct 10, 2026
2c8da56
Add the /coordinator and /hitl subpath aliases typecheck now needs
jason810496 Oct 10, 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
36 changes: 25 additions & 11 deletions .github/workflows/k8s-tests.yml
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ jobs:
# preparing k8s environment with uv takes < 15 seconds with `uv` - there is no point in caching it.
# The lang-SDK coordinator system test is intentionally NOT run in this matrix. It has its own
# dedicated single-combo job (tests-kubernetes-lang-sdk below), so it no longer adds the ~6 minutes
# of Go/Java build + coordinator run to each of these system-test jobs.
# of Go/Java/TypeScript build + coordinator run to each of these system-test jobs.
- name: "\
Run complete K8S tests ${{ matrix.executor }}-${{ env.PYTHON_MAJOR_MINOR_VERSION }}-\
${{env.KUBERNETES_VERSION}}-${{ matrix.use-standard-naming }}"
Expand Down Expand Up @@ -140,11 +140,12 @@ jobs:
run: breeze k8s delete-cluster --all
if: always()

# The Multi-Lang KubernetesExecutor coordinator system test builds a Go bundle and a Java jar and runs
# them as workers under KubernetesExecutor. It is executor/naming/version-independent, so instead of
# bolting it onto all six ``KubernetesExecutor``/``use-standard-naming == false`` system-test jobs above
# (which ran it redundantly and added ~6 minutes each), it runs once here on a single default combo and
# executes only the lang-SDK test (``-k test_lang_sdk_combined_dag_succeeds``) - not the full suite.
# The lang-SDK KubernetesExecutor coordinator system tests build a Go bundle, Java jars and
# TypeScript bundles and run them as workers under KubernetesExecutor. They are
# executor/naming/version-independent, so instead of bolting them onto all six
# ``KubernetesExecutor``/``use-standard-naming == false`` system-test jobs above (which ran them
# redundantly and added several minutes each), they run once here on a single default combo and
# execute only the lang-SDK tests (``-k TestLangSdk``) - not the full suite.
tests-kubernetes-lang-sdk:
timeout-minutes: 60
name: "K8S Lang-SDK:${{ inputs.lang-sdk-kubernetes-combo }}"
Expand Down Expand Up @@ -182,7 +183,7 @@ jobs:
use-uv: ${{ inputs.use-uv }}
make-mnt-writeable-and-cleanup: true
id: breeze
# Provision the lang-SDK Go/Java toolchains on the host so breeze builds the artifacts natively
# Provision the lang-SDK Go/Java/Node toolchains on the host so breeze builds the artifacts natively
# (LANG_SDK_NATIVE_TOOLCHAIN=true) instead of in throwaway toolchain containers, skipping the image
# pulls and reusing the caches restored below. The module/build and Gradle caches are keyed
# explicitly (rather than via the setup-* built-in caching) so the key carries an ``-vN-`` salt:
Expand Down Expand Up @@ -213,15 +214,28 @@ jobs:
path: |
~/.gradle/caches
~/.gradle/wrapper
key: lang-sdk-gradle-v1-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('java-sdk/**/*.gradle*', 'java-sdk/**/gradle-wrapper.properties', 'kubernetes-tests/lang_sdk/java_example/**/*.gradle*') }} # yamllint disable-line rule:line-length
key: lang-sdk-gradle-v1-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('java-sdk/**/*.gradle*', 'java-sdk/**/gradle-wrapper.properties', 'kubernetes-tests/lang_sdk/java_example/**/*.gradle*', 'airflow-e2e-tests/java-native-bundle/**/*.gradle*') }} # yamllint disable-line rule:line-length
restore-keys: |
lang-sdk-gradle-v1-${{ runner.os }}-${{ runner.arch }}-
- name: "Setup pnpm for the lang-SDK TypeScript build"
uses: pnpm/action-setup@ea17c68df8912ef543352723c149a84f56e3d413 # v6.1.0
with:
package_json_file: ts-sdk/package.json
run_install: false
- name: "Setup Node for the lang-SDK TypeScript build"
uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0
with:
node-version: 22
cache: 'pnpm'
cache-dependency-path: |
ts-sdk/pnpm-lock.yaml
kubernetes-tests/lang_sdk/ts_example/pnpm-lock.yaml
# No --upgrade here: the redeploy/upgrade path is already exercised across the main matrix above;
# this job only needs a single deploy to validate the Multi-Lang coordinator, saving ~1.5 minutes.
- name: "Run lang-SDK K8S test ${{ inputs.lang-sdk-kubernetes-combo }}"
# this job only needs a single deploy to validate the lang-SDK coordinators, saving ~1.5 minutes.
- name: "Run lang-SDK K8S tests ${{ inputs.lang-sdk-kubernetes-combo }}"
run: >-
breeze k8s run-complete-tests --no-copy-local-sources
-- -k test_lang_sdk_combined_dag_succeeds
-- -k TestLangSdk
env:
EXECUTOR: "KubernetesExecutor"
USE_STANDARD_NAMING: "false"
Expand Down
1 change: 1 addition & 0 deletions .rat-excludes
Original file line number Diff line number Diff line change
Expand Up @@ -215,6 +215,7 @@ www-hash.txt
**/.prettierrc
**/.airflowignore
**/.airflowignore_glob
**/airflowignore_all


# Vendor includes
Expand Down
10 changes: 7 additions & 3 deletions airflow-e2e-tests/docker/java.yml
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,18 @@
# the pre-built bundle JARs (the Java example under /opt/airflow/java-jars, the
# Scala Spark example under /opt/airflow/scala-jars, and the runner-behaviour
# test fixtures under /opt/airflow/java-test-jars), and configures the worker to
# consume the "java", "scala", and "java-test" Celery queues where @task.stub
# tasks are routed.
# consume the "java", "scala", "java-test", and "java-native" Celery queues where
# Java tasks are routed. The native-Dag JAR is in the Dags folder, so the Dag
# processor runs on the same image to parse it.
---
services:
airflow-dag-processor:
image: airflow-java-worker

airflow-worker:
image: airflow-java-worker
volumes:
- ./java-jars:/opt/airflow/java-jars:ro
- ./scala-jars:/opt/airflow/scala-jars:ro
- ./java-test-jars:/opt/airflow/java-test-jars:ro
command: celery worker -q java,scala,java-test,default
command: celery worker -q java,scala,java-test,java-native,default
2 changes: 2 additions & 0 deletions airflow-e2e-tests/java-native-bundle/.gitignore
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
.gradle
build/
50 changes: 50 additions & 0 deletions airflow-e2e-tests/java-native-bundle/build.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
/*
* 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.
*/

plugins {
id("org.apache.airflow.sdk") version "${projectVersion}"
}

repositories {
mavenLocal()
mavenCentral()
}

dependencies {
annotationProcessor("org.apache.airflow:airflow-sdk-processor:${projectVersion}")
implementation("org.apache.airflow:airflow-sdk:${projectVersion}")
implementation("org.apache.airflow:airflow-sdk-jpl:${projectVersion}")
}

java {
toolchain {
languageVersion.set(JavaLanguageVersion.of(11))
}
sourceCompatibility = JavaVersion.VERSION_11
}

sourceSets {
main {
java.srcDir("src/java")
}
}

airflowBundle {
mainClass = "org.apache.airflow.e2e.NativeBundleBuilder"
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,18 +14,7 @@
# KIND, either express or implied. See the License for the
# specific language governing permissions and limitations
# under the License.
#
# Java worker image for the "java" queue: the stock prod image plus a headless
# JRE so the JavaCoordinator can exec `java`. The Go queue runs on the plain
# prod image (a Go bundle is a self-contained static binary), so keeping this a
# separate image demonstrates routing each coordinator's queue to its own
# pod_template_file with a queue-specific base image.
ARG BASE_IMAGE
FROM ${BASE_IMAGE}

USER root
RUN apt-get update \
&& apt-get install --no-install-recommends -y default-jre-headless \
&& apt-get clean \
&& rm -rf /var/lib/apt/lists/*
USER airflow
org.gradle.configuration-cache=true

projectVersion=1.0.0-SNAPSHOT
33 changes: 33 additions & 0 deletions airflow-e2e-tests/java-native-bundle/settings.gradle
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
/*
* 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.
*/

// Route the plugin lookup to the SDK build published to the local Maven
// repository by conftest._setup_java_sdk_integration.
pluginManagement {
repositories {
mavenLocal()
gradlePluginPortal()
}
}

plugins {
id("org.gradle.toolchains.foojay-resolver-convention") version "0.10.0"
}

rootProject.name = "airflow-e2e-java-native-bundle"
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
/*
* 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.airflow.e2e;

import java.util.List;
import org.apache.airflow.e2e.nativedag.AnnotationDag;
import org.apache.airflow.e2e.nativedag.TargetDag;
import org.apache.airflow.sdk.*;

/**
* Bundle for the native Dag E2E tests: Dags declared entirely in Java, parsed by the Dag processor
* from this JAR.
*
* <p>This is the bundle's main class, so it is the Dag source the Airflow UI shows. Every task sets
* {@code queue}, which {@code queue_to_coordinator} routes to the {@code java-jdk} coordinator.
* The {@code java-native} coordinator only parses this JAR.
*/
public class NativeBundleBuilder {
public static final String QUEUE = "java-native";

public static class Extract implements Task {
@Override
public void execute(Context context, Client client) {
client.setXCom(42L);
}
}

public static class Transform implements Task {
@Override
public void execute(Context context, Client client) {
var extracted = ((Number) client.getXCom("extract")).longValue();
client.setXCom(extracted * 2);
}
}

/** Fails the run unless the value reached it, so a successful run proves the XCom flow. */
public static class Load implements Task {
@Override
public void execute(Context context, Client client) {
var transformed = ((Number) client.getXCom("transform")).longValue();
if (transformed != 84L) {
throw new IllegalStateException("load expected 84 from transform, got " + transformed);
}
}
}

/** Its boolean picks one of {@link ReportMany} and {@link ReportFew}; the other is skipped. */
public static class HasRows implements ConditionTask {
@Override
public boolean decide(Context context, Client client) {
return ((Number) client.getXCom("transform")).longValue() > 0;
}
}

public static class ReportMany implements Task {
@Override
public void execute(Context context, Client client) {}
}

public static class ReportFew implements Task {
@Override
public void execute(Context context, Client client) {}
}

/** Names the one of {@link ReportLong} and {@link ReportShort} that runs; the other is skipped. */
public static class PickReport implements SwitchTask {
@Override
public Class<? extends Task> choose(Context context, Client client) {
return ((Number) client.getXCom("transform")).longValue() > 100 ? ReportLong.class : ReportShort.class;
}
}

public static class ReportLong implements Task {
@Override
public void execute(Context context, Client client) {}
}

public static class ReportShort implements Task {
@Override
public void execute(Context context, Client client) {}
}

public static class Audit implements Task {
@Override
public void execute(Context context, Client client) {}
}

public static DagDef buildDag() {
var dag =
new DagDef("java_native_e2e")
.config("description", "Native Java Dag of the Airflow E2E tests")
.config("catchup", false)
.config("tags", List.of("java-sdk", "native", "e2e"));

var extract = dag.task("extract", Extract.class).config("queue", QUEUE);
var transform = dag.task("transform", Transform.class).config("queue", QUEUE);
var load = dag.task("load", Load.class).config("queue", QUEUE);

transform.after(extract).before(load);

var reportMany = dag.task("report_many", ReportMany.class).config("queue", QUEUE);
var reportFew = dag.task("report_few", ReportFew.class).config("queue", QUEUE);
dag.If("has_rows", HasRows.class).after(transform).config("queue", QUEUE).Then(reportMany).Else(reportFew);

var reportLong = dag.task("report_long", ReportLong.class).config("queue", QUEUE);
var reportShort = dag.task("report_short", ReportShort.class).config("queue", QUEUE);
dag.Switch("pick_report", PickReport.class)
.after(transform)
.config("queue", QUEUE)
.Case(reportLong)
.Case(reportShort);

// A task that starts a run of another Dag; it runs no Java code.
var trigger =
dag.task("trigger_downstream", new TriggerDagRun("java_native_target_e2e")).config("queue", QUEUE);
load.before(trigger);

// Ordering-only edge: the checks group runs after extract, with no data flowing.
var checks = dag.taskGroup("checks");
checks.task("audit", Audit.class).config("queue", QUEUE);
extract.before(checks);

return dag;
}

public static Bundle build() {
return new Bundle().register(buildDag()).register(AnnotationDag.class).register(TargetDag.build());
}

public static void main(String[] args) {
Server.create(args).serve(build());
}
}
Loading
Loading