seanmuth commented on code in PR #73503:
URL: https://github.com/apache/airflow/pull/73503#discussion_r4068107316


##########
task-sdk/src/airflow/sdk/execution_time/supervisor.py:
##########
@@ -729,12 +729,13 @@ def start(
         """
         Fork and start a new subprocess with the specified target function.
 
-        :param use_exec: If True, immediately ``os.execv`` a fresh Python 
interpreter
-            after ``os.fork``: forced on platforms that need it (macOS, whose 
Objective-C
-            frameworks are not fork-safe) and opted into for the task process 
elsewhere via
-            ``[core] execute_tasks_new_python_interpreter`` (a lock a 
supervisor thread
-            held at fork time cannot survive into a fresh address space).
-            ``target`` is rehydrated in the exec'd child from its 
``module:qualname``,
+        :param use_exec: If True, start a fresh Python interpreter via 
``os.posix_spawn``
+            instead of a bare ``os.fork``: forced on platforms that need it 
(macOS, whose
+            Objective-C frameworks are not fork-safe) and opted into for the 
task process
+            elsewhere via ``[core] execute_tasks_new_python_interpreter``. 
Unlike

Review Comment:
   Agreed, that's a pre-existing docstring-placement issue independent of this 
PR -- leaving it as a separate cleanup.



##########
task-sdk/tests/task_sdk/execution_time/test_supervisor.py:
##########
@@ -4703,6 +4703,127 @@ def 
test_start_rejects_non_importable_target_under_exec(self):
             supervisor.WatchedSubprocess.start(target=lambda: None, 
use_exec=True)
 
 
+class TestStartUsesPosixSpawn:
+    """use_exec=True goes through os.posix_spawn, never os.fork -- that's the 
whole point."""
+
+    def _start(self, mocker, **kwargs):
+        spawn = 
mocker.patch("airflow.sdk.execution_time.supervisor.os.posix_spawn", 
return_value=4321)
+        fork = mocker.patch(
+            "airflow.sdk.execution_time.supervisor.os.fork",
+            side_effect=AssertionError("os.fork() must not be called when 
use_exec=True"),
+        )
+        mocker.patch("airflow.sdk.execution_time.supervisor.psutil.Process")
+        supervisor.WatchedSubprocess.start(
+            id=uuid7(), target=supervisor._subprocess_main, use_exec=True, 
**kwargs
+        )
+        return spawn, fork
+
+    def test_does_not_call_fork(self, mocker):
+        """The defining property of the fix: no os.fork() call exists on this 
path at all."""
+        spawn, fork = self._start(mocker)
+        fork.assert_not_called()
+        spawn.assert_called_once()
+
+    def test_spawns_the_bootstrap_with_the_target_env_var(self, mocker):
+        spawn, _ = self._start(mocker)
+        args, kwargs = spawn.call_args
+        path, argv, env = args
+        assert path == sys.executable
+        assert argv == [sys.executable, "-c", supervisor._CHILD_EXEC_BOOTSTRAP]
+        assert env["_AIRFLOW_CHILD_TARGET"] == 
"airflow.sdk.execution_time.supervisor:_subprocess_main"
+
+    def test_file_actions_dup2_the_four_fds(self, mocker):
+        spawn, _ = self._start(mocker)
+        file_actions = spawn.call_args.kwargs["file_actions"]
+        targets = {new_fd for _, _, new_fd in file_actions}
+        assert targets == {0, 1, 2, 3}
+        assert all(action == os.POSIX_SPAWN_DUP2 for action, _, _ in 
file_actions)
+
+    def test_setpgroup_passed_when_new_process_group(self, mocker):
+        spawn, _ = self._start(mocker, new_process_group=True)
+        assert spawn.call_args.kwargs["setpgroup"] == 0
+
+    def test_setpgroup_omitted_when_not_new_process_group(self, mocker):
+        spawn, _ = self._start(mocker, new_process_group=False)
+        assert "setpgroup" not in spawn.call_args.kwargs
+
+    @pytest.mark.skipif(sys.platform == "win32", 
reason="os.fork/os.register_at_fork are POSIX-only")
+    def 
test_hanging_after_fork_handler_wedges_bare_fork_but_not_posix_spawn(self):

Review Comment:
   Fair critique -- that's the only one that's actually deterministic and 
implementation-agnostic. I kept the other five since each one pins down a 
distinct assertion the deterministic test doesn't exercise on its own (env var 
target, the four dup2 mappings, setpgroup present/absent), but I don't feel 
strongly about it -- happy to trim down to just the hang test if you'd rather 
keep the suite lean.



##########
task-sdk/src/airflow/sdk/execution_time/supervisor.py:
##########
@@ -729,12 +729,13 @@ def start(
         """
         Fork and start a new subprocess with the specified target function.
 
-        :param use_exec: If True, immediately ``os.execv`` a fresh Python 
interpreter
-            after ``os.fork``: forced on platforms that need it (macOS, whose 
Objective-C
-            frameworks are not fork-safe) and opted into for the task process 
elsewhere via
-            ``[core] execute_tasks_new_python_interpreter`` (a lock a 
supervisor thread
-            held at fork time cannot survive into a fresh address space).
-            ``target`` is rehydrated in the exec'd child from its 
``module:qualname``,
+        :param use_exec: If True, start a fresh Python interpreter via 
``os.posix_spawn``
+            instead of a bare ``os.fork``: forced on platforms that need it 
(macOS, whose
+            Objective-C frameworks are not fork-safe) and opted into for the 
task process
+            elsewhere via ``[core] execute_tasks_new_python_interpreter``. 
Unlike
+            ``fork()`` followed by ``execv()``, ``posix_spawn`` never runs

Review Comment:
   Good catch, fixed in 91d7bcc13b -- the module docstring, `_child_exec_main`, 
and `_task_process_uses_exec` now describe posix_spawn's file_actions instead 
of the stale execv/set_inheritable text.



##########
task-sdk/src/airflow/sdk/execution_time/supervisor.py:
##########
@@ -755,80 +756,79 @@ def start(
         # Place for child to send requests/read responses, and the server side 
to read/respond
         child_requests, read_requests = socketpair()
 
-        # Open the socketpair before forking off the child, so that it is open 
when we fork.
+        # Open the socketpair before starting the child, so that it is open 
when we do.
         child_logs, read_logs = socketpair()
 
-        pid = os.fork()
-        if pid == 0:
+        if use_exec:
+            # file_actions run as part of the spawn itself -- no forked child 
to run
+            # imperative dup2 code in.
+            file_actions = [

Review Comment:
   Restored in 91d7bcc13b as a comment above the file_actions list.



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