jason810496 commented on code in PR #74036:
URL: https://github.com/apache/airflow/pull/74036#discussion_r4163183016


##########
task-sdk/src/airflow/sdk/coordinators/java/coordinator.py:
##########
@@ -211,16 +218,32 @@ class JavaCoordinator(SubprocessCoordinator):
 
     _explicit_root_kwarg = "jars_root"
 
+    @classmethod
+    def get_dag_importer_class(cls) -> type[JavaDagImporter]:
+        return JavaDagImporter
+
+    def _build_command(self, roots: Sequence[pathlib.Path], main_class: str) 
-> list[str]:
+        return [self.java_executable, "-classpath", 
_calculate_classpath(roots), *self.jvm_args, main_class]
+
     def _build_execute_task_command(self, *, what: TaskInstance) -> 
tuple[list[str], str | None]:
         # Without main_class, the first executable JAR in walk order wins; 
tracked at
         # https://github.com/apache/airflow/issues/71134
         roots = self._get_scan_roots()
         jar = _JarInfo.find(roots, self.main_class)
-        command = [
-            self.java_executable,
-            "-classpath",
-            _calculate_classpath(roots),
-            *self.jvm_args,
-            jar.main_class,
-        ]
-        return command, jar.schema_version
+        return self._build_command(roots, jar.main_class), jar.schema_version
+
+    def _build_parse_dag_command(self, *, path: pathlib.Path) -> 
tuple[list[str], str | None]:
+        # Same command shape as execution. With one executable JAR per bundle, 
or main_class set,
+        # a parse runs the class a task runs.
+        meta = _JarMetadata.from_jar(path)
+        if meta is None:
+            raise ValueError(f"Cannot read the manifest of {path}")
+        if not meta.main_class:
+            raise ValueError(f"{path} is not an executable JAR: its manifest 
sets no Main-Class")
+        if self.main_class and meta.main_class != self.main_class:
+            raise ValueError(
+                f"{path} runs {meta.main_class!r}, but this coordinator's 
main_class is {self.main_class!r}"
+            )
+        roots = self._get_scan_roots()
+        jar = _JarInfo.find(roots, meta.main_class)

Review Comment:
   Done in 5f546e305c. One helper serves both parse and execute: it rejects two 
JARs with the same `Main-Class`, and the JAR scan is sorted. A task now runs 
the JAR its Dag was parsed from, which 620f89fcec passes to the command builder.



##########
airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst:
##########
@@ -589,6 +589,46 @@ Durations and date-times are ISO-8601 strings in 
annotations (``retryDelay = "PT
 ``java.time.OffsetDateTime`` values in ``config`` calls.  An unknown key or a 
mismatched value type
 fails the build for an annotation, and the ``config`` call itself for an 
object.
 
+.. _java-sdk/native-dag-parsing:
+
+Parsing native Java Dags
+~~~~~~~~~~~~~~~~~~~~~~~~
+
+To have Airflow parse the Dags a bundle JAR declares, put the JAR in a Dag 
bundle and configure a
+:class:`~airflow.sdk.coordinators.java.JavaCoordinator` that reads that 
bundle: leave ``jars_root``
+unset to read each Dag's own bundle, or set ``dag_bundle_name``. The Dag 
processor runs the JAR's main
+class to list its Dags, so it needs a Java executable, as the workers do:
+
+.. code-block:: ini
+
+    [sdk]
+    coordinators = {
+      "java-native": {
+        "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
+        "kwargs": {"java_executable": "/usr/lib/jvm/java-17-openjdk/bin/java"}
+      }
+    }
+    queue_to_coordinator = {"java-native": "java-native"}
+
+A coordinator with ``jars_root`` set only runs tasks; it parses no Dags.
+
+* Every JAR in the bundle whose manifest sets ``Main-Class`` is parsed. Each 
Dag its main class

Review Comment:
   Done in bc7876e38f. `jars_root` is gone from this stack. Parsing a JAR built 
with a supervisor schema older than 2026-10-30 now fails with a clear error, 
and the docs say to list such JARs in `.airflowignore`. Tasks still find them.



##########
task-sdk/tests/task_sdk/coordinators/java/test_coordinator.py:
##########
@@ -425,3 +450,90 @@ def test_returns_execution_result(self, jars_root, 
mock_client):
 
         assert isinstance(result, BaseCoordinator.ExecutionResult)
         assert result.exit_code == 0
+
+
+class TestBuildParseDagCommand:
+    def _build(self, coordinator: JavaCoordinator, root: pathlib.Path, jar: 
pathlib.Path):
+        with coordinator._set_scan_roots([root]):
+            return coordinator._build_parse_dag_command(path=jar)
+
+    def test_fat_jar(self, tmp_path):
+        jar = _make_jar(tmp_path / "app.jar", main_class="com.example.Dags", 
schema_version="2026-06-16")
+        coordinator = JavaCoordinator(java_executable="/opt/java/bin/java", 
jvm_args=["-Xmx256m"])
+
+        command, schema_version = self._build(coordinator, tmp_path, jar)
+
+        assert command == ["/opt/java/bin/java", "-classpath", jar.as_posix(), 
"-Xmx256m", "com.example.Dags"]
+        assert schema_version == "2026-06-16"
+
+    def test_thin_jar_takes_the_schema_version_from_a_sibling(self, tmp_path):
+        jar = _make_jar(tmp_path / "app.jar", main_class="com.example.Dags")
+        (tmp_path / "libs").mkdir()
+        sdk = _make_jar(tmp_path / "libs" / "airflow-sdk.jar", 
main_class=None, schema_version="2026-06-16")
+
+        command, schema_version = self._build(JavaCoordinator(), tmp_path, jar)
+
+        assert command[2].split(os.pathsep) == [jar.as_posix(), sdk.as_posix()]
+        assert command[-1] == "com.example.Dags"
+        assert schema_version == "2026-06-16"
+
+    def test_matches_the_execute_command(self, tmp_path):
+        jar = _make_jar(tmp_path / "app.jar", main_class="com.example.Dags", 
schema_version="2026-06-16")
+        _make_jar(tmp_path / "dep.jar", main_class=None)
+        coordinator = JavaCoordinator(jvm_args=["-Xmx256m"])
+
+        parse = self._build(coordinator, tmp_path, jar)
+        with coordinator._set_scan_roots([tmp_path]):
+            execute = coordinator._build_execute_task_command(what=_make_ti())
+
+        assert parse == execute
+
+    @pytest.mark.parametrize(
+        ("attributes", "match"),
+        [
+            pytest.param(None, "Cannot read the manifest", id="no-manifest"),
+            pytest.param({"Manifest-Version": "1.0"}, "sets no Main-Class", 
id="no-main-class"),
+            pytest.param({"Main-Class": "com.example.Other"}, "main_class is 
'com.example.Dags'", id="pin"),
+        ],
+    )
+    def test_rejects_a_jar_it_cannot_run(self, tmp_path, attributes, match):
+        jar = make_jar(tmp_path / "app.jar", attributes=attributes, 
entries={"a.class": b""})
+
+        with pytest.raises(ValueError, match=match):
+            self._build(JavaCoordinator(main_class="com.example.Dags"), 
tmp_path, jar)
+
+    def test_rejects_a_non_zip(self, tmp_path):
+        jar = tmp_path / "broken.jar"
+        jar.write_bytes(b"not a zip")
+
+        with pytest.raises(ValueError, match="Cannot read the manifest"):
+            self._build(JavaCoordinator(), tmp_path, jar)
+
+    def test_rejects_a_jar_deleted_after_discovery(self, tmp_path):
+        with pytest.raises(ValueError, match="Cannot read the manifest"):
+            self._build(JavaCoordinator(), tmp_path, tmp_path / "gone.jar")
+
+    @patch("airflow.sdk.coordinators._subprocess.os.execvpe", autospec=True, 
side_effect=OSError("exec"))
+    def test_parse_dag_execs_the_jvm(self, mock_execvpe, tmp_path):

Review Comment:
   Done in ecb55b91bc. The test now also patches `signal.signal` and 
`_set_close_on_exec_above_stderr`.



##########
task-sdk/src/airflow/sdk/coordinators/java/coordinator.py:
##########
@@ -211,16 +218,32 @@ class JavaCoordinator(SubprocessCoordinator):
 
     _explicit_root_kwarg = "jars_root"
 
+    @classmethod
+    def get_dag_importer_class(cls) -> type[JavaDagImporter]:

Review Comment:
   Agreed. The stack's discovery layer is now a mirror of #74020. Rather than 
reorder the stack, the PR description says to merge after #74020.



##########
airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst:
##########
@@ -589,6 +589,46 @@ Durations and date-times are ISO-8601 strings in 
annotations (``retryDelay = "PT
 ``java.time.OffsetDateTime`` values in ``config`` calls.  An unknown key or a 
mismatched value type
 fails the build for an annotation, and the ``config`` call itself for an 
object.
 
+.. _java-sdk/native-dag-parsing:
+
+Parsing native Java Dags
+~~~~~~~~~~~~~~~~~~~~~~~~
+
+To have Airflow parse the Dags a bundle JAR declares, put the JAR in a Dag 
bundle and configure a
+:class:`~airflow.sdk.coordinators.java.JavaCoordinator` that reads that 
bundle: leave ``jars_root``
+unset to read each Dag's own bundle, or set ``dag_bundle_name``. The Dag 
processor runs the JAR's main
+class to list its Dags, so it needs a Java executable, as the workers do:
+
+.. code-block:: ini
+
+    [sdk]
+    coordinators = {
+      "java-native": {
+        "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
+        "kwargs": {"java_executable": "/usr/lib/jvm/java-17-openjdk/bin/java"}
+      }
+    }
+    queue_to_coordinator = {"java-native": "java-native"}
+
+A coordinator with ``jars_root`` set only runs tasks; it parses no Dags.
+
+* Every JAR in the bundle whose manifest sets ``Main-Class`` is parsed. Each 
Dag its main class
+  declares, through ``Bundle.register`` of a ``DagDef`` or an ``@Builder.Dag`` 
class, is stored with
+  that JAR as its file. Task handlers for a Python Dag are not Dags, and a JAR 
without ``Main-Class``,

Review Comment:
   Done in 204ddea7bd. The docs and the docstring now say many dependencies set 
`Main-Class`, so a thin bundle should set `main_class` or list its dependency 
JARs in `.airflowignore`.



##########
task-sdk/src/airflow/sdk/coordinators/java/coordinator.py:
##########
@@ -93,24 +98,21 @@ def _calculate_classpath(jars_root: Sequence[pathlib.Path]) 
-> str:
 
 @attrs.define
 class _JarMetadata:
-    main_class: str
-    schema_version: str
+    main_class: str | None
+    schema_version: str | None
 
     @classmethod
     def from_jar(cls, path: pathlib.Path) -> Self | None:
         try:
             with zipfile.ZipFile(path) as zf:
-                try:
-                    manifest_info = zf.getinfo("META-INF/MANIFEST.MF")
-                except KeyError:
-                    log.debug("JAR does not contain META-INF/MANIFEST.MF; 
ignored", path=path)
-                    return None
-                with zf.open(manifest_info) as f:
-                    manifest = email.message_from_binary_file(f)
-            return cls(manifest["Main-Class"], 
manifest["Airflow-Supervisor-Schema-Version"])
-        except zipfile.BadZipFile:
+                attributes = read_main_attributes(zf)
+        except (OSError, zipfile.BadZipFile):

Review Comment:
   Done in 62c038ff91. The catch is narrowed to `FileNotFoundError`, 
`IsADirectoryError` and `BadZipFile`. The parse path now says whether the JAR 
is invalid, has no manifest, or has no `Main-Class`.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to