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


##########
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:
   The parse of a JAR only reads its `Main-Class` and then builds the command 
from the whole bundle, so two JARs that share one (say `etl-old.jar` left next 
to `etl-new.jar`) both run `-classpath etl-new.jar:etl-old.jar com.A`. The JVM 
loads one JAR's classes, while `_JarInfo.find` can report the other's schema 
version because its walk is unsorted, and neither `main_class` nor 
`.airflowignore` helps since the classpath and `find` don't read them. Could a 
`Main-Class` that more than one JAR in the bundle sets be an import error? 
There's also no parse test with two executable JARs: swapping `meta.main_class` 
for `self.main_class` on this line keeps the suite green.



##########
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:
   At this layer the manager still finds JARs through the zip path, so opting 
Java in here has two side effects until the importer-driven discovery is in 
(#74020, or #73841 in this stack). The processor never asks 
`might_contain_dag`, so each dependency JAR of a thin bundle becomes a "Cannot 
start the Lang-SDK runtime" import error. And `_get_observed_filelocs` expands 
the JAR into its members, so its Dags are marked stale and its import errors 
cleared on every bundle refresh until the next parse. Could the description 
name #74020 as a merge prerequisite next to #71190?



##########
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:
   No released Java SDK answers the parse request yet: in java-sdk 1.0.0-beta1, 
`handleIncoming` only acts on `StartupDetails` and `ErrorResponse`. So once a 
coordinator without `jars_root` serves a bundle, every executable JAR built 
with it, task-handler JARs for Python Dags included, ends in an import error. 
Could this section name the first Java SDK version that answers the parse, and 
point older bundles at `jars_root`?



##########
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:
   This runs the real `parse_dag`, which resets SIGPIPE and SIGXFSZ to 
`SIG_DFL` before the exec, and only `os.execvpe` is patched, so the pytest 
worker keeps those dispositions afterwards and a later write to a closed pipe 
kills the worker instead of raising `BrokenPipeError`. The matching test in 
test_subprocess.py patches `signal.signal` with autospec; could this one do the 
same?



##########
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:
   Catching every `OSError` here changes the released `jars_root` task path 
too: in task-sdk 1.3.2 an unreadable JAR failed the task with a 
`PermissionError` naming the file, and now it fails with "cannot find a JAR 
with Main-Class matching ..." while the cause only reaches the module log. 
Could this catch `FileNotFoundError` and `IsADirectoryError`, which the new 
tests cover, and let `PermissionError` through? On the parse side, "Cannot read 
the manifest of X" at line 240 also drops the reason (not a zip, no manifest), 
which is what a user needs on the import errors page.



##########
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:
   Plenty of dependencies set `Main-Class` too (the PostgreSQL JDBC driver's is 
`org.postgresql.util.PGJDBCMain`, H2's is `org.h2.tools.Console`), so with 
`main_class` unset each one is parsed: its main runs in the Dag processor, 
prints a banner, exits, and shows up as an import error. Could this say a thin 
bundle should set `main_class` or list its dependency JARs in `.airflowignore`? 
The same sentence is in the `_dag_importer.py` docstring.



-- 
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