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


##########
task-sdk/src/airflow/sdk/importers/base.py:
##########
@@ -187,7 +188,7 @@ class DagImportResult:
     """Result of importing DAGs from a definition."""
 
     definition: DagDefinition | None = None
-    dags: list[DAG] = field(default_factory=list)
+    dags: list[DAG | SerializedLangSDKDAG] = field(default_factory=list)

Review Comment:
   This widens the SDK's public `DagImportResult.dags` to a core type. On 
#72369 the thread settled on keeping `sdk.DAG` here, partly to keep core 
serialization types out of the importer API, and #73457 still has the result 
shape open. `SerializedLangSDKDAG` is also only an alias, so every 
`isinstance(dag, SerializedLangSDKDAG)` in the Dag bag really means "any 
`SerializedDAG`", and a third-party importer returning one would skip 
`validate()`, the cluster policies and `sync_bag_to_db` with no signal. Could 
the coordinator importer return the runtime's JSON in a separate field (say 
`serialized_dags: list[dict[str, Any]]`) that `DagBag` deserializes on the core 
side, or the bag key the branch on the claiming importer, as `sync_bag_to_db` 
already does?



##########
airflow-core/src/airflow/dag_processing/dagbag.py:
##########
@@ -359,9 +363,11 @@ def _process_definition(
                 dag.fileloc = fileloc
                 dag.relative_fileloc = self._get_relative_fileloc(fileloc)
                 dag.bundle_name = self.bundle_name
-                dag.validate()
-                _validate_executor_fields(dag, self.bundle_name)
-                _assign_default_team_pools(dag, self.bundle_name)
+                # The Dag processor does not run these on a Lang-SDK Dag, so a 
Dag bag does not either.

Review Comment:
   The parity holds, but `_validate_executor_fields` and 
`_assign_default_team_pools` aren't SDK-only checks, they're core multi-team 
behaviour. With `[core] multi_team` on, a native task in a team bundle keeps 
`default_pool` instead of the team's pool, and an unknown `executor` isn't 
caught at parse time, on the Dag processor path from #74035 as well. Is that a 
known limitation for now? If so it belongs in the list; otherwise 
`_handle_parsing_result` could apply both to the deserialized Dag.



##########
task-sdk/src/airflow/sdk/coordinators/_dag_importer.py:
##########
@@ -78,13 +85,27 @@ def list_dag_definitions(
     def import_definition(
         self, definition: FilesystemDagDefinition, bundle: BaseDagBundle
     ) -> DagImportResult:
-        """Report that only the Dag processor parses *definition*."""
-        return DagImportResult(
-            definition=definition,
-            errors=[
+        from airflow.dag_processing.lang_sdk_processor import 
LangSDKDagFileProcessorProcess  # noqa: SDK002
+        from airflow.serialization.serialized_objects import DagSerialization  
# noqa: SDK002
+
+        source_reference = repr(definition)
+        bundle_path = bundle.path or definition.path.parent
+        relative_loc = definition.get_relative_loc(bundle_path)
+        result = DagImportResult(definition=definition)
+        parsing_result = LangSDKDagFileProcessorProcess.run(

Review Comment:
   Every whole-bundle Dag bag now starts each native file's runtime, and three 
callers throw the result away: `dags reserialize`, the all-bundles fallback in 
`get_bagged_dag`, and the re-sync in `DAG.test()`, since `sync_bag_to_db` drops 
native Dags. So `airflow dags test some_python_dag` in a bundle with N native 
files can start N JVM or Node processes twice, where before it recorded a cheap 
import error, and `reserialize`, which the description gives as the reason for 
this PR, stores nothing new. Could bags that only feed `sync_bag_to_db` skip 
coordinator-claimed files?



##########
airflow-core/src/airflow/utils/cli.py:
##########
@@ -291,6 +292,13 @@ def get_bagged_dag(bundle_names: list | None, dag_id: str, 
dagfile_path: str | N
     """
     from airflow.dag_processing.dagbag import BundleDagBag, sync_bag_to_db
     from airflow.sdk.definitions._internal.dag_parsing_context import 
_airflow_parsing_context_manager
+    from airflow.serialization.definitions.dag import SerializedLangSDKDAG
+
+    def check_python_dag(dag: BaggedDAG) -> DAG:

Review Comment:
   `tasks list` only prints the ids from `dag.tasks`, which works on a 
`SerializedDAG`, so it doesn't need this refusal. Could `task_list` skip 
`check_python_dag` and leave it to `dags test`, `tasks test` and `tasks render`?



##########
airflow-core/tests/unit/dag_processing/test_dagbag.py:
##########
@@ -1534,19 +1544,170 @@ def 
test_dagbag_no_bundle_path_no_syspath_modification(self, tmp_path):
         assert sys.path == syspath_before
 
 
-def test_sync_bag_to_db_leaves_native_files_to_the_dag_processor(tmp_path, 
session, testing_dag_bundle):
-    db.clear_db_import_errors()
-    write_native_file(tmp_path / "dags.native")
-    session.add(ParseImportError(bundle_name="testing", 
filename="dags.native", stacktrace="stored"))
-    session.commit()
+def _serialize_native_dag(dag_id: str, fileloc: str, relative_fileloc: str) -> 
dict:
+    with DAG(dag_id, schedule=None) as dag:
+        BaseOperator(task_id="extract")
+    data = DagSerialization.to_dict(dag)
+    data["dag"].update(fileloc=fileloc, relative_fileloc=relative_fileloc)
+    return data
+
+
[email protected]
+def _native_runtime(**dag_ids_by_file: list[str]):
+    """Make the runtime of a ``.native`` file return the Dag ids given for 
that file's stem."""
+
+    def run(*, path, dag_file_rel_path, **kwargs):
+        dag_ids = dag_ids_by_file[Path(path).stem]
+        return DagFileParsingResult(
+            fileloc=os.fspath(path),
+            serialized_dags=[
+                LazyDeserializedDAG(data=_serialize_native_dag(dag_id, 
os.fspath(path), dag_file_rel_path))
+                for dag_id in dag_ids
+            ],
+        )
+
+    with fake_coordinator(), patch.object(LangSDKDagFileProcessorProcess, 
"run", side_effect=run):

Review Comment:
   This patches `LangSDKDagFileProcessorProcess.run` without `autospec=True`, 
unlike the same patch in test_cli_util.py and test_dag_importer.py, so a 
renamed or wrong kwarg from `import_definition` would still pass. The bundles 
in test_cli_util.py's new tests are bare `MagicMock()`s with the same problem.



##########
task-sdk/src/airflow/sdk/coordinators/_dag_importer.py:
##########
@@ -78,13 +85,27 @@ def list_dag_definitions(
     def import_definition(
         self, definition: FilesystemDagDefinition, bundle: BaseDagBundle
     ) -> DagImportResult:
-        """Report that only the Dag processor parses *definition*."""
-        return DagImportResult(
-            definition=definition,
-            errors=[
+        from airflow.dag_processing.lang_sdk_processor import 
LangSDKDagFileProcessorProcess  # noqa: SDK002
+        from airflow.serialization.serialized_objects import DagSerialization  
# noqa: SDK002
+
+        source_reference = repr(definition)
+        bundle_path = bundle.path or definition.path.parent
+        relative_loc = definition.get_relative_loc(bundle_path)
+        result = DagImportResult(definition=definition)
+        parsing_result = LangSDKDagFileProcessorProcess.run(
+            path=definition.path,
+            bundle_path=bundle_path,
+            bundle_name=bundle.name,
+            dag_file_rel_path=relative_loc,
+            logger=log,
+        )
+        for key, message in (parsing_result.import_errors or {}).items():
+            result.errors.append(
                 DagImportError(

Review Comment:
   All of these errors carry the file's `source_reference`, and 
`_record_import_error` stores by fileloc, so only the last one survives in the 
bag: with the test's `{"main.min.mjs": "one failed", "main.ts": "two failed"}` 
it keeps only `main.ts: two failed`. Joining them into one `DagImportError`, as 
`_handle_parsing_result` does for validation errors, would keep them all.



##########
task-sdk/src/airflow/sdk/coordinators/_dag_importer.py:
##########
@@ -78,13 +85,27 @@ def list_dag_definitions(
     def import_definition(
         self, definition: FilesystemDagDefinition, bundle: BaseDagBundle
     ) -> DagImportResult:
-        """Report that only the Dag processor parses *definition*."""
-        return DagImportResult(
-            definition=definition,
-            errors=[
+        from airflow.dag_processing.lang_sdk_processor import 
LangSDKDagFileProcessorProcess  # noqa: SDK002
+        from airflow.serialization.serialized_objects import DagSerialization  
# noqa: SDK002
+
+        source_reference = repr(definition)
+        bundle_path = bundle.path or definition.path.parent

Review Comment:
   `bundle.path` is typed `Path`, and a `Path` is always truthy, so `or 
definition.path.parent` never runs. The other importers pass `bundle.path` 
straight through.



##########
task-sdk/src/airflow/sdk/coordinators/_dag_importer.py:
##########
@@ -78,13 +85,27 @@ def list_dag_definitions(
     def import_definition(
         self, definition: FilesystemDagDefinition, bundle: BaseDagBundle
     ) -> DagImportResult:
-        """Report that only the Dag processor parses *definition*."""
-        return DagImportResult(
-            definition=definition,
-            errors=[
+        from airflow.dag_processing.lang_sdk_processor import 
LangSDKDagFileProcessorProcess  # noqa: SDK002

Review Comment:
   These two do have to stay inline (task_runner imports this module, and 
`lang_sdk_processor` imports `processor`, which imports `task_runner`), but 
`noqa: SDK002` doesn't say so. A one-line comment naming the cycle would stop 
someone hoisting them later.



##########
task-sdk/tests/task_sdk/coordinators/test_dag_importer.py:
##########
@@ -57,14 +93,47 @@ def test_lists_only_its_artifacts(importer, tmp_path):
     assert [d.path.name for d in definitions] == ["main.min.mjs"]
 
 
-def 
test_import_definition_reports_that_only_the_dag_processor_parses_it(importer, 
tmp_path):
-    bundle_file = tmp_path / "main.min.mjs"
-    bundle_file.write_text("")
-    definition = FilesystemDagDefinition(bundle_file)
-
-    result = importer.import_definition(definition, 
SimpleNamespace(name="testing", path=tmp_path))
-
-    assert result.dags == []
-    assert [error.message for error in result.errors] == [
-        "A native Lang-SDK Dag is parsed only by the Dag processor"
-    ]
+class TestImportDefinition:
+    @staticmethod
+    def _import(importer, tmp_path):
+        definition = FilesystemDagDefinition(tmp_path / "main.min.mjs")
+        return importer.import_definition(definition, 
SimpleNamespace(name="testing", path=tmp_path))
+
+    @patch.object(LangSDKDagFileProcessorProcess, "run", autospec=True)
+    def test_returns_the_dags_the_runtime_serialized(self, mock_run, importer, 
tmp_path):
+        mock_run.return_value = DagFileParsingResult(
+            fileloc=str(tmp_path / "main.min.mjs"),
+            
serialized_dags=[LazyDeserializedDAG(data=_get_payload("conformance_minimal"))],
+            import_errors={"main.min.mjs": "one failed", "main.ts": "two 
failed"},
+        )
+
+        result = self._import(importer, tmp_path)
+
+        [dag] = result.dags
+        assert isinstance(dag, SerializedDAG)
+        assert (dag.dag_id, dag.task_ids) == ("conformance_minimal", ["solo"])
+        assert [(e.source_reference, e.message) for e in result.errors] == [
+            (str(tmp_path / "main.min.mjs"), "one failed"),
+            (str(tmp_path / "main.min.mjs"), "main.ts: two failed"),
+        ]
+        mock_run.assert_called_once_with(
+            path=tmp_path / "main.min.mjs",
+            bundle_path=tmp_path,
+            bundle_name="testing",
+            dag_file_rel_path="main.min.mjs",
+            logger=ANY,
+        )
+
+    @pytest.mark.parametrize("data", _load_payloads())

Review Comment:
   java_native.json records two Dags from one file, but every mocked result 
here holds one Dag, so `import_definition` returning only the first would still 
pass. Feeding both Java Dags as one `DagFileParsingResult` would cover that. 
And since nothing in the tree produces java_native.json, a note on how it was 
recorded would let someone re-record it when `Serde.kt` changes.



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