ashb commented on code in PR #73503:
URL: https://github.com/apache/airflow/pull/73503#discussion_r4065939887
##########
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:
Nit/observation: this should have been in the worker section. Oh well,
(separate issue)
##########
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:
This is probably the only test worth keeping
--
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]