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


##########
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:
   Done in dca01a5a0c. Native parsing is now opt-in with 
`DagBag(parse_lang_sdk_files=True)`. By default the file is an import error and 
no runtime starts, so `dags reserialize`, the all-bundles fallback and 
`DAG.test()` no longer start runtimes.



##########
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:
   Done in dca01a5a0c. A file's runtime errors are joined into one message.



##########
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:
   This went away in dca01a5a0c, since the importer no longer runs the runtime.



##########
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:
   Done in dca01a5a0c. The patch is `autospec=True` now, and test_cli_util.py 
uses `MagicMock(spec=BaseDagBundle)`.



##########
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:
   Agreed. a68f51b07c on #74035 applies both checks in 
`_handle_parsing_result`, which covers both the Dag processor and the Dag bag. 
The comment here is removed.



##########
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:
   Done in dca01a5a0c. `DagImportResult.dags` stays `list[DAG]`. `DagBag` 
checks whether the claiming importer is a coordinator's, runs the runtime on 
the core side, and bags `SerializedDAG`s. The `SerializedLangSDKDAG` alias is 
gone.



##########
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:
   This went away in dca01a5a0c: no core imports are left in the Task SDK 
importer.



##########
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:
   Done in dca01a5a0c. A new test feeds the whole `java_native.json` through 
one `run()` and checks both Dags are bagged.



##########
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:
   Done in 897a3ef4c6. `tasks list` passes `allow_lang_sdk_dag=True` and lists 
a native Dag's tasks.



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