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


##########
task-sdk/src/airflow/sdk/coordinators/_subprocess.py:
##########
@@ -467,11 +496,61 @@ def _build_execute_task_command(self, *, what: 
TaskInstance) -> tuple[list[str],
         """
         raise NotImplementedError
 
+    def _build_parse_dag_command(self, *, path: pathlib.Path) -> 
tuple[list[str], str | None]:
+        """
+        Build the command that parses the Dag file at *path* and resolve its 
wire-schema version.
+
+        Subclasses can retrieve the directories to scan for artifacts with
+        :meth:`_get_scan_roots`; for a parse they are the Dag bundle's root. 
The contract is
+        that of :meth:`_build_execute_task_command`: *command* MUST NOT 
include the
+        ``--comm`` / ``--logs`` flags.
+        """
+        raise NotImplementedError
+
+    def parse_dag(
+        self,
+        *,
+        path: pathlib.Path,
+        bundle_path: pathlib.Path,
+        comm_address: tuple[str, int],
+        logs_address: tuple[str, int],
+        report_schema_version: Callable[[str | None], None],
+    ) -> NoReturn:
+        """
+        Replace the current process with the runtime that parses the Dag file 
at *path*.
+
+        Call this in a child process whose standard streams are already set 
up; the runtime
+        inherits them and no other file descriptor. The command and its 
supervisor wire-schema
+        version are resolved against *bundle_path*, and the version is passed 
to
+        *report_schema_version* just before the exec. The runtime connects 
back to
+        *comm_address* and *logs_address*.
+
+        :raises Exception: when the command cannot be resolved or started; the 
process is then
+            unchanged.
+        """
+        with self._set_scan_roots([bundle_path]):
+            command, schema_version = self._build_parse_dag_command(path=path)
+        if schema_version is not None:
+            get_schema_version_migrator().resolve_version(schema_version)
+        argv = [
+            *command,
+            f"--comm={comm_address[0]}:{comm_address[1]}",
+            f"--logs={logs_address[0]}:{logs_address[1]}",
+        ]
+        report_schema_version(schema_version)
+        # Python ignores these at startup and exec keeps ignored signals; 
subprocess.Popen resets
+        # them the same way for the task runtime.
+        for name in ("SIGPIPE", "SIGXFZ", "SIGXFSZ"):

Review Comment:
   Yes, it's a typo, let me remove it, thanks.



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