Skip to content

Commit a80facc

Browse files
committed
Name the JAR manifest attributes the Java coordinator reads
1 parent 178e468 commit a80facc

3 files changed

Lines changed: 17 additions & 10 deletions

File tree

‎task-sdk/src/airflow/sdk/coordinators/java/_dag_importer.py‎

Lines changed: 4 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -23,22 +23,19 @@
2323
from typing import TYPE_CHECKING, ClassVar, Final
2424

2525
from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter
26-
from airflow.sdk.coordinators.java._jar_manifest import read_main_attributes
26+
from airflow.sdk.coordinators.java._jar_manifest import DAG_CODE, MAIN_CLASS, read_main_attributes
2727
from airflow.sdk.importers.base import DagSourceCode
2828

2929
if TYPE_CHECKING:
3030
from airflow.sdk.coordinators.java.coordinator import JavaCoordinator
3131
from airflow.sdk.importers.base import DagDefinition
3232

33-
_DAG_CODE_ATTRIBUTE: Final = "airflow-java-sdk-dag-code"
3433
_MAX_SOURCE_BYTES: Final = 1024 * 1024
3534
_NO_SOURCE: Final = (
3635
"// This JAR embeds no Dag source. Build it with the Airflow Java SDK Gradle plugin, or set\n"
3736
"// airflowBundle.dagSource, to show the source here.\n"
3837
)
39-
_SOURCE_TOO_LARGE: Final = (
40-
f"// The Dag source this JAR embeds is over {_MAX_SOURCE_BYTES} bytes, so it is not shown.\n"
41-
)
38+
_SOURCE_TOO_LARGE: Final = "// The Dag source this JAR embeds is over 1 MiB, so it is not shown.\n"
4239

4340

4441
class JavaDagImporter(CoordinatorDagImporter):
@@ -64,7 +61,7 @@ def might_contain_dag(self, definition: DagDefinition, safe_mode: bool) -> bool:
6461
attributes = read_main_attributes(zf) or {}
6562
except (OSError, zipfile.BadZipFile):
6663
return True
67-
if not (main_class := attributes.get("main-class")):
64+
if not (main_class := attributes.get(MAIN_CLASS)):
6865
return False
6966
return not self.coordinator.main_class or main_class == self.coordinator.main_class
7067

@@ -80,7 +77,7 @@ def get_source_code(self, definition: DagDefinition) -> DagSourceCode:
8077

8178

8279
def _find_source_entry(zf: zipfile.ZipFile) -> zipfile.ZipInfo | None:
83-
if not (entry := (read_main_attributes(zf) or {}).get(_DAG_CODE_ATTRIBUTE)):
80+
if not (entry := (read_main_attributes(zf) or {}).get(DAG_CODE)):
8481
return None
8582
try:
8683
return zf.getinfo(entry)

‎task-sdk/src/airflow/sdk/coordinators/java/_jar_manifest.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,11 @@
2727

2828
MANIFEST_NAME: Final = "META-INF/MANIFEST.MF"
2929

30+
# Attribute keys as parse_main_attributes returns them, lower-cased.
31+
MAIN_CLASS: Final = "main-class"
32+
SUPERVISOR_SCHEMA_VERSION: Final = "airflow-supervisor-schema-version"
33+
DAG_CODE: Final = "airflow-java-sdk-dag-code"
34+
3035
_LINE_END = re.compile(rb"\r\n|\r|\n")
3136

3237

‎task-sdk/src/airflow/sdk/coordinators/java/coordinator.py‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,11 @@
3535
)
3636
from airflow.sdk.coordinators._subprocess import SubprocessCoordinator
3737
from airflow.sdk.coordinators.java._dag_importer import JavaDagImporter
38-
from airflow.sdk.coordinators.java._jar_manifest import read_main_attributes
38+
from airflow.sdk.coordinators.java._jar_manifest import (
39+
MAIN_CLASS,
40+
SUPERVISOR_SCHEMA_VERSION,
41+
read_main_attributes,
42+
)
3943

4044
if TYPE_CHECKING:
4145
from collections.abc import Iterable, Iterator, Sequence
@@ -108,7 +112,7 @@ def from_jar(cls, path: pathlib.Path) -> Self | None:
108112
if attributes is None:
109113
log.debug("JAR does not contain META-INF/MANIFEST.MF; ignored", path=path)
110114
return None
111-
return cls(attributes.get("main-class"), attributes.get("airflow-supervisor-schema-version"))
115+
return cls(attributes.get(MAIN_CLASS), attributes.get(SUPERVISOR_SCHEMA_VERSION))
112116

113117

114118
@attrs.define
@@ -230,7 +234,8 @@ def _build_execute_task_command(self, *, what: TaskInstance) -> tuple[list[str],
230234
return self._build_command(roots, jar.main_class), jar.schema_version
231235

232236
def _build_parse_dag_command(self, *, path: pathlib.Path) -> tuple[list[str], str | None]:
233-
# The classpath and main class match execution's, so a parse runs the code a task runs.
237+
# Same command shape as execution. With one executable JAR per bundle, or main_class set,
238+
# a parse runs the class a task runs.
234239
meta = _JarMetadata.from_jar(path)
235240
if meta is None:
236241
raise ValueError(f"Cannot read the manifest of {path}")

0 commit comments

Comments
 (0)