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


##########
task-sdk/src/airflow/sdk/coordinators/java/_dag_importer.py:
##########
@@ -0,0 +1,122 @@
+#
+# 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.
+"""The Dag importer of 
:class:`~airflow.sdk.coordinators.java.JavaCoordinator`."""
+
+from __future__ import annotations
+
+import json
+import zipfile
+from typing import TYPE_CHECKING, Any, ClassVar, Final, cast
+
+from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter
+from airflow.sdk.coordinators.java._jar_manifest import MAIN_CLASS, SOURCES, 
read_main_attributes
+from airflow.sdk.importers.base import DagSourceCode
+
+if TYPE_CHECKING:
+    from airflow.sdk.coordinators.java.coordinator import JavaCoordinator
+    from airflow.sdk.importers.base import DagDefinition
+
+_SOURCES_DIR: Final = "META-INF/airflow/sources/"
+_MAX_SOURCE_BYTES: Final = 1024 * 1024
+_NO_SOURCE: Final = (
+    "// This JAR embeds no Dag source. Build it with the Airflow Java SDK 
Gradle plugin to show the\n"
+    "// source here.\n"
+)
+_SOURCE_TOO_LARGE: Final = "// This Dag source file is over 1 MiB, so it is 
not shown.\n"
+
+
+class JavaDagImporter(CoordinatorDagImporter):
+    """
+    Claim the native Dags of Java bundle JARs.
+
+    A :class:`~airflow.sdk.coordinators.java.JavaCoordinator` parses them. 
Only a JAR whose manifest sets
+    ``Main-Class`` (matching the coordinator's ``main_class`` if set) is 
parsed.
+    """
+
+    coordinator_classpath: ClassVar[str] = 
"airflow.sdk.coordinators.java.JavaCoordinator"
+    artifact_suffix: ClassVar[str] = ".jar"
+    supported_extensions = [".jar"]
+
+    def might_contain_dag(self, definition: DagDefinition, safe_mode: bool) -> 
bool:
+        """
+        Return whether the JAR sets ``Main-Class``, matching the parsing 
coordinator's ``main_class`` if set.
+
+        ``safe_mode`` does not apply, because a JAR without ``Main-Class`` 
cannot run at all. A JAR that
+        cannot be read is kept, so that parsing it reports the error. So is 
every JAR when no coordinator
+        can parse the bundle, so that each parse reports why.
+        """
+        try:
+            with definition.as_file() as path, zipfile.ZipFile(path) as zf:
+                attributes = read_main_attributes(zf) or {}
+        except (OSError, zipfile.BadZipFile):

Review Comment:
   `zf.read` raises more than these two. I tried a scratch zip: a corrupt 
deflate stream gives `zlib.error`, and an unsupported compression method gives 
`NotImplementedError`. Neither is caught here, so it escapes 
`list_dag_definitions` and hits the manager's `except Exception` around 
`_find_files_in_bundle`. That skips the whole bundle on every refresh, so new 
or deleted Python files in it are missed and no import error shows up. 
`ZipImporter` catches `Exception` per archive for the same reason.
   
   Could this be `except Exception: return True`, so the JAR's own parse 
surfaces the error? `_JarMetadata.from_jar` uses the same tuple. There, one bad 
sibling JAR breaks the parse of every other JAR, and the message doesn't say 
which JAR it was.



##########
airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst:
##########
@@ -606,6 +615,67 @@ 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`.
+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"}
+
+Once a ``JavaCoordinator`` is configured, the Dag processor parses the 
executable JARs of every Dag bundle,
+so it needs this ``[sdk]`` configuration and a JRE. With one 
``JavaCoordinator``, it parses them all.
+With several, map each Dag bundle that holds native Java Dags to one of them 
in ``[sdk] dag_bundle_to_coordinator``.
+A JAR in a bundle that has no entry, or an entry that names no 
``JavaCoordinator``, fails to parse with an import error:
+
+.. code-block:: ini
+
+    [sdk]
+    dag_bundle_to_coordinator = {"dags-folder": "java-native"}
+
+A Dag bundle that holds only the JARs that Python Dags' tasks run should list 
``*`` in its ``.airflowignore``.
+Otherwise, with several Java coordinators, its JARs fail to parse.

Review Comment:
   With one Java coordinator the handler bundle's JARs are parsed too: each 
runs in a JVM on every parse loop, and one built with a Java SDK older than 
schema 2026-10-30 (every released one today) becomes an import error, per the 
bullet below. Could this say the ignore file is needed either way? As written 
it reads as optional for the single-coordinator setup the example above 
configures, and the quickstart already says to add it.



##########
task-sdk/src/airflow/sdk/coordinators/java/coordinator.py:
##########
@@ -207,16 +265,44 @@ class JavaCoordinator(SubprocessCoordinator):
     jvm_args: list[str] = attrs.field(factory=list)
     main_class: str = ""
 
+    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
+        # Without main_class, the first executable JAR in path order wins; 
tracked at

Review Comment:
   "The first executable JAR in path order wins" isn't what `_JarInfo.find` 
does when the schema version lives in a separate JAR. It overwrites 
`progress.main_class` for every executable JAR until one with 
`Airflow-Supervisor-Schema-Version` turns up, so a thin bundle with 
`acme-etl.jar` and `acme-tools.jar` (both executable, no schema) ahead of a 
schema-only SDK JAR runs `acme-tools`. Guarding the assignment with 
`progress.main_class is None` would make this comment, the class docstring and 
the `main_class` row in java.rst true.



##########
task-sdk/src/airflow/sdk/coordinators/java/coordinator.py:
##########
@@ -207,16 +265,44 @@ class JavaCoordinator(SubprocessCoordinator):
     jvm_args: list[str] = attrs.field(factory=list)
     main_class: str = ""
 
+    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
+        # Without main_class, the first executable JAR in path 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_jar_command(self, path: pathlib.Path) -> tuple[list[str], str]:
+        """
+        Build the command that runs the executable JAR at *path*, and return 
its schema version.
+
+        The JAR's own ``Main-Class`` runs, so *main_class* must match it when 
set. The bundle's other
+        JARs go on the classpath, so another JAR that sets the same 
``Main-Class`` is rejected.
+
+        :raises ValueError: when the JAR cannot run.
+        """
+        main_class, schema_version = _read_executable_jar(path)
+        if self.main_class and main_class != self.main_class:
+            raise ValueError(
+                f"{path} runs {main_class!r}, but this coordinator's 
main_class is {self.main_class!r}"
+            )
+        roots = self._get_scan_roots()
+        jar = _JarInfo.for_jar(roots, path, main_class, schema_version)
+        return self._build_command(roots, jar.main_class), jar.schema_version

Review Comment:
   The JAR picked here isn't first on the classpath: `_build_command` puts 
every bundle JAR in sorted order. With two fat JARs that each shade the SDK, 
which is the bundle the `mock_start` fixture in test_coordinator.py builds 
(`a.jar` at 2026-06-16, `b.jar` at 2026-10-30), running `b.jar` loads the SDK 
classes from `a.jar`, while the schema version checked against 
`_DAG_PARSING_SCHEMA_VERSION` and sent as `subprocess_schema_version` is 
`b.jar`'s. I checked it with two small JARs: `java -classpath a.jar:b.jar B` 
printed a.jar's `sdk.Sdk` version, and `b.jar:a.jar` printed b.jar's. So 
bumping the SDK in one project of a shared bundle passes the parse gate and 
then runs the old SDK. It's the reason the duplicate `Main-Class` check exists, 
applied to every other class. Could `_build_jar_command` put `path` first and 
the rest of the bundle after it?



##########
task-sdk/tests/task_sdk/coordinators/test_subprocess.py:
##########
@@ -1036,6 +1047,388 @@ def test_named_bundle_mode_locks_the_tree_it_scans(
         mock_lock.assert_called_once_with(bundle_name="artifacts", 
bundle_version="sha-abc")
 
 
[email protected](kw_only=True)
+class _NativeStubCoordinator(_StubSubprocessCoordinator):
+    """Runs native Dag files, recording the roots and the Dag files its 
command builder is given."""
+
+    recorded_dag_files: list[pathlib.Path] = attrs.field(init=False, 
factory=list)
+    scans: int = attrs.field(init=False, default=0)
+
+    def _build_execute_task_command(self, *, what):
+        self.scans += 1
+        return super()._build_execute_task_command(what=what)
+
+    def _build_dag_file_command(self, *, what, path):
+        self.recorded_roots.append(list(self._get_scan_roots()))
+        self.recorded_dag_files.append(path)
+        return [*self.command, os.fspath(path)], self.schema_version
+
+
[email protected](kw_only=True)
+class _BrokenNativeCoordinator(_NativeStubCoordinator):
+    def _build_dag_file_command(self, *, what, path):
+        raise ValueError("no Main-Class")
+
+
[email protected](kw_only=True)
+class _OtherStubCoordinator(_StubSubprocessCoordinator):
+    """A coordinator of another class than the one that runs native Dag 
files."""
+
+
+class _NativeDagImporter(CoordinatorDagImporter):
+    coordinator_classpath = f"{__name__}._NativeStubCoordinator"
+    artifact_suffix = ".native"
+    supported_extensions = [".native"]
+
+    def get_source_code(self, definition) -> DagSourceCode:
+        return DagSourceCode("", "native")
+
+
+class _UnloadableDagImporter(_NativeDagImporter):
+    coordinator_classpath = "nonexistent.module.Coordinator"
+
+
[email protected]
+def _native_dag_files(*keys: str, mapping: dict[str, str] | None = None):
+    """Make ``.native`` files native Dag files, with a _NativeStubCoordinator 
configured for each key."""
+    config = {
+        ("sdk", "coordinators"): json.dumps(
+            {
+                key: {"classpath": f"{__name__}._NativeStubCoordinator", 
"kwargs": {"command": ["/runtime"]}}
+                for key in keys
+            }
+        )
+    }
+    if mapping is not None:
+        config[("sdk", "dag_bundle_to_coordinator")] = json.dumps(mapping)
+    with conf_vars(config):
+        reset_importer_registry()
+        yield
+
+
+@patch(
+    "airflow.sdk.coordinators._dag_importer.COORDINATOR_DAG_IMPORTERS", 
(f"{__name__}._NativeDagImporter",)
+)
+class TestExecuteTaskNativeDagFile:
+    """A task of a native Dag runs its own Dag file at the version of the run, 
with no scan."""
+
+    BUNDLE_INFO = BundleInfo(name="dags", version="v9")
+
+    @pytest.fixture(autouse=True)
+    def _clean_registry(self):
+        reset_importer_registry()
+        yield
+        reset_importer_registry()
+
+    @pytest.fixture
+    def bundle(self, tmp_path):
+        """The task's Dag bundle, resolved at version v9 to a directory that 
holds one native Dag file."""
+        (tmp_path / "sub").mkdir()
+        (tmp_path / "sub" / "dag.native").write_text("")
+        (tmp_path / "dir.native").mkdir()
+        resolved = MagicMock(path=tmp_path, version="v9")

Review Comment:
   These bundles are bare `MagicMock`s, and the `initialize_ti_bundle`, 
`BundleVersionLock` and `start` patches have no `autospec`, unlike 
`TestExecuteTaskBundleWiring` just above, which uses `_make_bundle` (specced on 
`BaseDagBundle`) with autospecced patches. A typo'd bundle attribute in 
`_resolve_task_bundle` or a renamed `start()` keyword would pass here. The 
`mock_start` fixture in java/test_coordinator.py has the same shape.



##########
task-sdk/tests/task_sdk/coordinators/test_subprocess.py:
##########
@@ -1036,6 +1047,388 @@ def test_named_bundle_mode_locks_the_tree_it_scans(
         mock_lock.assert_called_once_with(bundle_name="artifacts", 
bundle_version="sha-abc")
 
 
[email protected](kw_only=True)
+class _NativeStubCoordinator(_StubSubprocessCoordinator):
+    """Runs native Dag files, recording the roots and the Dag files its 
command builder is given."""
+
+    recorded_dag_files: list[pathlib.Path] = attrs.field(init=False, 
factory=list)
+    scans: int = attrs.field(init=False, default=0)
+
+    def _build_execute_task_command(self, *, what):
+        self.scans += 1
+        return super()._build_execute_task_command(what=what)
+
+    def _build_dag_file_command(self, *, what, path):
+        self.recorded_roots.append(list(self._get_scan_roots()))
+        self.recorded_dag_files.append(path)
+        return [*self.command, os.fspath(path)], self.schema_version
+
+
[email protected](kw_only=True)
+class _BrokenNativeCoordinator(_NativeStubCoordinator):
+    def _build_dag_file_command(self, *, what, path):
+        raise ValueError("no Main-Class")
+
+
[email protected](kw_only=True)
+class _OtherStubCoordinator(_StubSubprocessCoordinator):
+    """A coordinator of another class than the one that runs native Dag 
files."""
+
+
+class _NativeDagImporter(CoordinatorDagImporter):
+    coordinator_classpath = f"{__name__}._NativeStubCoordinator"
+    artifact_suffix = ".native"
+    supported_extensions = [".native"]
+
+    def get_source_code(self, definition) -> DagSourceCode:
+        return DagSourceCode("", "native")
+
+
+class _UnloadableDagImporter(_NativeDagImporter):
+    coordinator_classpath = "nonexistent.module.Coordinator"
+
+
[email protected]
+def _native_dag_files(*keys: str, mapping: dict[str, str] | None = None):
+    """Make ``.native`` files native Dag files, with a _NativeStubCoordinator 
configured for each key."""
+    config = {
+        ("sdk", "coordinators"): json.dumps(
+            {
+                key: {"classpath": f"{__name__}._NativeStubCoordinator", 
"kwargs": {"command": ["/runtime"]}}
+                for key in keys
+            }
+        )
+    }
+    if mapping is not None:
+        config[("sdk", "dag_bundle_to_coordinator")] = json.dumps(mapping)
+    with conf_vars(config):
+        reset_importer_registry()
+        yield
+
+
+@patch(
+    "airflow.sdk.coordinators._dag_importer.COORDINATOR_DAG_IMPORTERS", 
(f"{__name__}._NativeDagImporter",)
+)
+class TestExecuteTaskNativeDagFile:
+    """A task of a native Dag runs its own Dag file at the version of the run, 
with no scan."""
+
+    BUNDLE_INFO = BundleInfo(name="dags", version="v9")
+
+    @pytest.fixture(autouse=True)
+    def _clean_registry(self):
+        reset_importer_registry()
+        yield
+        reset_importer_registry()
+
+    @pytest.fixture
+    def bundle(self, tmp_path):
+        """The task's Dag bundle, resolved at version v9 to a directory that 
holds one native Dag file."""
+        (tmp_path / "sub").mkdir()
+        (tmp_path / "sub" / "dag.native").write_text("")
+        (tmp_path / "dir.native").mkdir()
+        resolved = MagicMock(path=tmp_path, version="v9")
+        resolved.name = "dags"
+        return resolved
+
+    @pytest.fixture
+    def mock_initialize(self, bundle):
+        with 
patch("airflow.sdk.coordinators._subprocess.initialize_ti_bundle", 
return_value=bundle) as m:
+            yield m
+
+    @pytest.fixture
+    def mock_lock(self):
+        with patch("airflow.sdk.coordinators._subprocess.BundleVersionLock") 
as mock_lock:
+            yield mock_lock
+
+    @pytest.fixture
+    def mock_start(self):
+        with patch.object(_PopenActivitySubprocess, "start") as mock_start:
+            mock_start.return_value.wait.return_value = 0
+            yield mock_start
+
+    def _execute(self, coordinator, client, *, rel_path="sub/dag.native", 
bundle_info=BUNDLE_INFO):
+        return coordinator.execute_task(
+            what=_make_ti(),
+            dag_rel_path=rel_path,
+            bundle_info=bundle_info,
+            client=client,
+            subprocess_logs_to_stdout=False,
+        )
+
+    @pytest.mark.usefixtures("mock_initialize")
+    def test_runs_the_dag_file_at_the_version_of_the_run_without_a_scan(
+        self, mock_start, mock_lock, mock_initialize, mock_client, tmp_path
+    ):
+        coordinator = _NativeStubCoordinator(command=["/runtime"], 
schema_version="2026-06-16")
+
+        with _native_dag_files("native"):
+            result = self._execute(coordinator, mock_client)
+
+        dag_file = tmp_path / "sub" / "dag.native"
+        mock_initialize.assert_called_once_with(self.BUNDLE_INFO)
+        assert coordinator.recorded_dag_files == [dag_file]
+        assert coordinator.recorded_roots == [[tmp_path]]
+        assert coordinator.scans == 0
+        assert coordinator._active_scan_roots is None
+        mock_lock.assert_called_once_with(bundle_name="dags", 
bundle_version="v9")
+        mock_lock.return_value.__exit__.assert_called_once()
+        assert mock_start.call_args.kwargs["command"] == ["/runtime", 
os.fspath(dag_file)]
+        assert mock_start.call_args.kwargs["subprocess_schema_version"] == 
"2026-06-16"
+        assert result.exit_code == 0
+
+    @pytest.mark.usefixtures("mock_start", "mock_lock")
+    def test_does_not_read_the_artifact_bundle_of_the_coordinator(self, 
mock_initialize, mock_client):
+        coordinator = _NativeStubCoordinator(command=["/runtime"], 
task_handler_bundle_name="artifacts")
+
+        with _native_dag_files("native"):
+            self._execute(coordinator, mock_client)
+
+        mock_initialize.assert_called_once_with(self.BUNDLE_INFO)
+
+    @pytest.mark.usefixtures("mock_start", "mock_lock")
+    def test_pins_a_run_without_a_version_to_the_current_version(
+        self, mock_initialize, mock_lock, mock_client, tmp_path
+    ):
+        pinned_tree = tmp_path / "versions" / "sha-abc"
+        (pinned_tree / "sub").mkdir(parents=True)
+        (pinned_tree / "sub" / "dag.native").write_text("")
+        unpinned = MagicMock(path=tmp_path, version=None)
+        unpinned.get_current_version.return_value = 
BundleVersion(version="sha-abc", data=None)
+        pinned = MagicMock(path=pinned_tree, version="sha-abc")
+        pinned.name = "dags"
+        mock_initialize.side_effect = [unpinned, pinned]
+        coordinator = _NativeStubCoordinator(command=["/runtime"])
+
+        with _native_dag_files("native"):
+            self._execute(coordinator, mock_client, 
bundle_info=BundleInfo(name="dags"))
+
+        assert coordinator.recorded_dag_files == [pinned_tree / "sub" / 
"dag.native"]
+        mock_lock.assert_called_once_with(bundle_name="dags", 
bundle_version="sha-abc")
+
+    @pytest.mark.usefixtures("mock_initialize", "mock_lock")
+    @pytest.mark.parametrize(
+        "mapping",
+        [
+            pytest.param(None, id="no-entry-for-the-bundle"),
+            pytest.param({"dags": "first"}, 
id="entry-for-another-coordinator"),
+        ],
+    )
+    def test_runs_on_a_coordinator_that_does_not_parse_the_bundle(
+        self, mock_start, mock_client, tmp_path, mapping
+    ):
+        queue_coordinator = _NativeStubCoordinator(command=["/runtime"])
+
+        with _native_dag_files("first", "second", mapping=mapping):
+            importer = find_claiming_importer("sub/dag.native", "dags")
+            parsing_error = (
+                pytest.raises(InvalidCoordinatorError, match="Dag bundle 
'dags' has 2 _NativeStubCoordinator")
+                if mapping is None
+                else contextlib.nullcontext()
+            )
+            with parsing_error:
+                parsing_coordinator = importer.get_parsing_coordinator()
+            self._execute(queue_coordinator, mock_client)
+
+        assert mapping is None or parsing_coordinator is not queue_coordinator

Review Comment:
   `queue_coordinator` is a local instance the manager never builds, so 
`parsing_coordinator is not queue_coordinator` can't fail, and with 
`mapping=None` the name is never bound. What the test name promises is that 
`execute_task` never asks for the parsing coordinator. Patching 
`CoordinatorDagImporter.get_parsing_coordinator` with `autospec=True` and 
asserting `assert_not_called()` would pin that for both params, and the 
`pytest.raises`/`nullcontext` block could go.



##########
task-sdk/tests/task_sdk/execution_time/test_supervisor.py:
##########
@@ -270,6 +273,260 @@ def test_supervise(
                 supervise_task(**kw)
 
 
+def _response_error(status: int, detail: dict[str, Any]) -> 
ServerResponseError:
+    error = ServerResponseError.from_response(
+        httpx.Response(status, json={"detail": detail}, 
request=httpx.Request("PATCH", "http://server/x";))
+    )
+    assert error is not None
+    return error
+
+
+class TestSuperviseTaskLaunchError:
+    """A task whose runtime cannot be launched fails its try in the 
supervisor, with the reason."""
+
+    LAUNCH_EVENT = "Cannot launch the task's runtime"
+
+    @pytest.fixture
+    def ti(self):
+        return TaskInstance(
+            id=uuid7(),
+            task_id="extract",
+            dag_id="etl",
+            run_id="r",
+            try_number=1,
+            dag_version_id=uuid7(),
+            queue="java",
+        )
+
+    @pytest.fixture
+    def client(self, make_ti_context):
+        client = MagicMock(spec=sdk_client.Client)
+        client.task_instances.start.return_value = 
make_ti_context(should_retry=False)
+        return client
+
+    @pytest.fixture
+    def upload_to_remote(self):
+        with patch("airflow.sdk.log.upload_to_remote", autospec=True) as 
mock_upload:
+            yield mock_upload
+
+    def _supervise(
+        self, client, ti, error: BaseException | Callable[..., Any], *, 
log_path: str | None = None
+    ) -> int:
+        coordinator = MagicMock(spec=BaseCoordinator)
+        coordinator.execute_task.side_effect = error
+        with patch.object(supervisor, "get_coordinator_manager", 
autospec=True) as mock_manager:
+            mock_manager.return_value.for_queue.return_value = coordinator
+            return supervise_task(
+                ti=ti,
+                bundle_info=BundleInfo(name="dags", version="v1"),
+                dag_rel_path="etl.jar",
+                token="",
+                client=client,
+                log_path=log_path,
+            )
+
+    def _read_task_log(self, tmp_path, log_path: str) -> list[dict[str, Any]]:
+        return [json.loads(line) for line in (tmp_path / 
log_path).read_text().splitlines() if line]
+
+    def test_fails_the_try_with_the_reason(self, client, ti, upload_to_remote):
+        exit_code = self._supervise(client, ti, TaskLaunchError("no jar for 
'extract'"))
+
+        assert exit_code == 1
+        client.task_instances.start.assert_called_once_with(ti.id, 
os.getpid(), mock.ANY)
+        client.task_instances.finish.assert_called_once_with(
+            id=ti.id,
+            state=TaskInstanceState.FAILED,
+            when=mock.ANY,
+            rendered_map_index=None,
+            retry_reason="no jar for 'extract'",
+        )
+        client.task_instances.retry.assert_not_called()
+
+    def test_retries_the_try_when_the_task_has_retries_left(
+        self, client, ti, make_ti_context, upload_to_remote
+    ):
+        client.task_instances.start.return_value = 
make_ti_context(should_retry=True, max_tries=2)
+
+        exit_code = self._supervise(client, ti, TaskLaunchError("no jar for 
'extract'"))
+
+        assert exit_code == 1
+        client.task_instances.retry.assert_called_once_with(
+            ti.id, end_date=mock.ANY, rendered_map_index=None, 
retry_reason="no jar for 'extract'"
+        )
+        client.task_instances.finish.assert_not_called()
+
+    def 
test_fails_the_try_of_an_error_that_is_not_retryable_even_with_retries_left(
+        self, client, ti, make_ti_context, upload_to_remote
+    ):
+        client.task_instances.start.return_value = 
make_ti_context(should_retry=True, max_tries=2)
+
+        exit_code = self._supervise(client, ti, TaskLaunchError("wrong 
coordinator", retryable=False))
+
+        assert exit_code == 1
+        client.task_instances.finish.assert_called_once_with(
+            id=ti.id,
+            state=TaskInstanceState.FAILED,
+            when=mock.ANY,
+            rendered_map_index=None,
+            retry_reason="wrong coordinator",
+        )
+        client.task_instances.retry.assert_not_called()
+
+    def test_writes_the_reason_to_the_task_log(self, client, ti, tmp_path, 
upload_to_remote):
+        with conf_vars({("logging", "base_log_folder"): str(tmp_path)}):
+            self._supervise(client, ti, TaskLaunchError("no jar for 
'extract'"), log_path="etl/extract.log")
+
+        entries = self._read_task_log(tmp_path, "etl/extract.log")
+        assert [(e["level"], e["event"], e["reason"]) for e in entries if 
"reason" in e] == [
+            ("error", self.LAUNCH_EVENT, "no jar for 'extract'")
+        ]
+
+    def test_logs_the_cause_of_the_error(self, client, ti, tmp_path, 
upload_to_remote):
+        try:
+            try:
+                raise OSError("disk gone")
+            except OSError as cause:
+                raise TaskLaunchError("cannot read the bundle") from cause
+        except TaskLaunchError as e:
+            error = e
+
+        with conf_vars({("logging", "base_log_folder"): str(tmp_path)}):
+            self._supervise(client, ti, error, log_path="etl/extract.log")
+
+        (entry,) = [e for e in self._read_task_log(tmp_path, 
"etl/extract.log") if "reason" in e]
+        assert "disk gone" in json.dumps(entry["exception"])
+
+    def test_truncates_the_reason_for_the_server_but_not_the_log(
+        self, client, ti, tmp_path, upload_to_remote
+    ):
+        reason = "x" * 700
+
+        with conf_vars({("logging", "base_log_folder"): str(tmp_path)}):
+            self._supervise(client, ti, TaskLaunchError(reason), 
log_path="etl/extract.log")
+
+        assert client.task_instances.finish.call_args.kwargs["retry_reason"] 
== "x" * 500
+        (entry,) = [e for e in self._read_task_log(tmp_path, 
"etl/extract.log") if "reason" in e]
+        assert entry["reason"] == reason
+
+    @pytest.mark.enable_redact
+    @pytest.mark.parametrize(
+        ("message", "expected"),
+        [
+            pytest.param(
+                "cannot read https://user:hunter2-hunter2@git/repo";,
+                "cannot read https://user:***@git/repo";,
+                id="short",
+            ),
+            pytest.param("x" * 490 + "hunter2-hunter2", "x" * 490 + "***", 
id="across-the-cut"),
+        ],
+    )
+    def test_redacts_a_masked_secret_in_the_reason_before_cutting_it(
+        self, client, ti, tmp_path, upload_to_remote, message, expected
+    ):
+        from airflow.sdk._shared.secrets_masker import mask_secret

Review Comment:
   Can this import move to the top of the module? Nothing here needs it 
deferred.



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