import json import os import subprocess import sys import threading import urllib.request from typing import Any from unittest.mock import MagicMock import boto3 import pytest from src.worker import ( AbortedStaleTask, MalformedSubprocessOutput, NonEmptyStr, SubprocessCrashError, UnknownSubprocessError, Worker, WorkerRequest, ) from src.worker.clients.step_functions import StaleTaskToken from src.worker.subprocess import WORKER_REQUEST from src.worker.task_protection import TaskProtectionResult class DummyRequest(WorkerRequest): bucket: NonEmptyStr key: NonEmptyStr class DummyTransientError(Exception): pass class DummyFatalError(Exception): pass _QUEUE_NAME = "test-worker-queue" _TIMEOUT_MINUTES = 30 _SUBPROCESS_MODULE = "src.test.dummy" _DUMMY_RESULT: dict[str, Any] = {"is_valid": True} @pytest.fixture(autouse=True) def _reset_moto() -> None: endpoint = os.environ["AWS_ENDPOINT_URL"] request = urllib.request.Request(f"{endpoint}/moto-api/reset", method="POST") urllib.request.urlopen(request, timeout=5) @pytest.fixture def queue_url() -> str: response = boto3.client("sqs").create_queue(QueueName=_QUEUE_NAME) return str(response["QueueUrl"]) def _make_worker( queue_url: str, *, heartbeat_interval_seconds: float = 60.0, message_processing_timeout_minutes: int = _TIMEOUT_MINUTES, exception_levels: dict[type[Exception], int] | None = None, ) -> Worker[DummyRequest]: import logging return Worker( request_schema=DummyRequest, select_subprocess_module=lambda _: _SUBPROCESS_MODULE, sqs_queue_url=queue_url, exception_levels=exception_levels or {DummyTransientError: logging.WARNING, DummyFatalError: logging.ERROR}, heartbeat_interval_seconds=heartbeat_interval_seconds, message_processing_timeout_minutes=message_processing_timeout_minutes, ) def _send(queue_url: str, body: dict[str, object]) -> None: boto3.client("sqs").send_message(QueueUrl=queue_url, MessageBody=json.dumps(body)) def _valid_body(task_token: str = "tok") -> dict[str, object]: return {"task_token": task_token, "bucket": "b", "key": "k"} def _queue_depth(queue_url: str) -> int: response = boto3.client("sqs").get_queue_attributes( QueueUrl=queue_url, AttributeNames=["ApproximateNumberOfMessages"] ) return int(response["Attributes"]["ApproximateNumberOfMessages"]) @pytest.fixture def shutdown_on_subprocess(monkeypatch: pytest.MonkeyPatch) -> threading.Event: shutdown = threading.Event() def stop_after_subprocess(*_: object, **__: object) -> dict[str, Any]: shutdown.set() return _DUMMY_RESULT monkeypatch.setattr(Worker, "_run_subprocess", stop_after_subprocess) return shutdown @pytest.fixture def mock_step_functions(monkeypatch: pytest.MonkeyPatch) -> MagicMock: mock = MagicMock() monkeypatch.setattr("src.worker.worker.step_functions.send_task_success", mock.send_task_success) monkeypatch.setattr("src.worker.worker.step_functions.send_task_failure", mock.send_task_failure) monkeypatch.setattr("src.worker.worker.step_functions.send_task_heartbeat", mock.send_task_heartbeat) return mock @pytest.fixture def mock_protection(monkeypatch: pytest.MonkeyPatch) -> MagicMock: mock = MagicMock() mock.enable.return_value = TaskProtectionResult.ENABLED mock.disable.return_value = True monkeypatch.setattr("src.worker.worker.task_protection.enable", mock.enable) monkeypatch.setattr("src.worker.worker.task_protection.disable", mock.disable) return mock class TestHappyPath: def test_processes_message_and_deletes_it( self, queue_url: str, shutdown_on_subprocess: threading.Event, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown_on_subprocess) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_called_once() mock_step_functions.send_task_failure.assert_not_called() def test_sends_task_token_and_result_to_step_functions( self, queue_url: str, shutdown_on_subprocess: threading.Event, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: _send(queue_url, _valid_body(task_token="my-token")) _make_worker(queue_url).process_messages(shutdown_on_subprocess) call_kwargs = mock_step_functions.send_task_success.call_args.kwargs assert call_kwargs["task_token"] == "my-token" assert call_kwargs["output"] == {"is_valid": True} def test_task_protection_enabled_and_disabled( self, queue_url: str, shutdown_on_subprocess: threading.Event, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown_on_subprocess) mock_protection.enable.assert_called_once_with(expires_in_minutes=_TIMEOUT_MINUTES) mock_protection.disable.assert_called_once() class TestProtectionFailure: @pytest.mark.parametrize( "result", [ pytest.param(TaskProtectionResult.DEPLOYMENT_BLOCKED, id="deployment_blocked"), pytest.param(TaskProtectionResult.TASK_STOPPING_OR_STOPPED, id="task_stopping_or_stopped"), pytest.param(TaskProtectionResult.FAILED, id="failed"), ], ) def test_protection_unavailable_releases_message_and_exits( self, queue_url: str, mock_protection: MagicMock, mock_step_functions: MagicMock, result: TaskProtectionResult, ) -> None: _send(queue_url, _valid_body()) mock_protection.enable.return_value = result _make_worker(queue_url).process_messages(threading.Event()) assert _queue_depth(queue_url) == 1 mock_step_functions.send_task_success.assert_not_called() mock_step_functions.send_task_failure.assert_not_called() mock_protection.disable.assert_not_called() def test_disable_failure_exits_worker( self, queue_url: str, mock_protection: MagicMock, mock_step_functions: MagicMock, monkeypatch: pytest.MonkeyPatch, ) -> None: _send(queue_url, _valid_body()) _send(queue_url, _valid_body()) monkeypatch.setattr(Worker, "_run_subprocess", lambda *_, **__: _DUMMY_RESULT) mock_protection.disable.return_value = False _make_worker(queue_url).process_messages(threading.Event()) mock_step_functions.send_task_success.assert_called_once() assert _queue_depth(queue_url) == 1 class TestMessageErrors: def test_schema_invalid_with_task_token_fast_fails_sfn( self, queue_url: str, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() boto3.client("sqs").send_message(QueueUrl=queue_url, MessageBody=json.dumps({"task_token": "t1"})) mock_protection.disable.side_effect = lambda: shutdown.set() _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_not_called() mock_step_functions.send_task_failure.assert_called_once() call = mock_step_functions.send_task_failure.call_args.kwargs assert call["task_token"] == "t1" assert call["error"] == "MalformedMessage" @pytest.mark.parametrize( "body", [ pytest.param("{not valid json", id="unparseable_json"), pytest.param(json.dumps({"bucket": "b", "key": "k"}), id="missing_task_token"), pytest.param(json.dumps({"task_token": ""}), id="empty_task_token"), pytest.param(json.dumps(["not", "a", "dict"]), id="json_list_not_dict"), ], ) def test_schema_invalid_without_task_token_deleted( self, queue_url: str, mock_protection: MagicMock, mock_step_functions: MagicMock, body: str, ) -> None: shutdown = threading.Event() boto3.client("sqs").send_message(QueueUrl=queue_url, MessageBody=body) mock_protection.disable.side_effect = lambda: shutdown.set() _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_not_called() mock_step_functions.send_task_failure.assert_not_called() def test_schema_invalid_deletes_even_if_sfn_fast_fail_fails( self, queue_url: str, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() boto3.client("sqs").send_message(QueueUrl=queue_url, MessageBody=json.dumps({"task_token": "t1"})) mock_step_functions.send_task_failure.side_effect = RuntimeError("sfn throttled") mock_protection.disable.side_effect = lambda: shutdown.set() _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_failure.assert_called_once() def test_unhandled_exception_sends_task_failure( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() def crash(*_: object, **__: object) -> dict[str, Any]: shutdown.set() error_message = "mediainfo segfault" raise RuntimeError(error_message) monkeypatch.setattr(Worker, "_run_subprocess", crash) _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_not_called() mock_step_functions.send_task_failure.assert_called_once() call = mock_step_functions.send_task_failure.call_args.kwargs assert call["error"] == "RuntimeError" assert "mediainfo segfault" in call["cause"] def test_malformed_subprocess_output_sends_task_failure( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() mock_step_functions.send_task_failure.side_effect = lambda *_, **__: shutdown.set() def malformed(*_: object, **__: object) -> dict[str, Any]: raise MalformedSubprocessOutput("subprocess emitted garbage") monkeypatch.setattr(Worker, "_run_subprocess", malformed) _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown) mock_step_functions.send_task_failure.assert_called_once() assert mock_step_functions.send_task_failure.call_args.kwargs["error"] == "MalformedSubprocessOutput" def test_error_levels_controls_log_level( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, caplog: pytest.LogCaptureFixture, ) -> None: import logging shutdown = threading.Event() mock_step_functions.send_task_failure.side_effect = lambda *_, **__: shutdown.set() def transient(*_: object, **__: object) -> dict[str, Any]: raise DummyTransientError("retryable") monkeypatch.setattr(Worker, "_run_subprocess", transient) _send(queue_url, _valid_body()) with caplog.at_level(logging.WARNING, logger="src.worker.worker"): _make_worker( queue_url, exception_levels={DummyTransientError: logging.WARNING, DummyFatalError: logging.ERROR}, ).process_messages(shutdown) transient_records = [r for r in caplog.records if "DummyTransientError" in r.getMessage()] assert transient_records, "expected DummyTransientError to be logged" assert all(r.levelno == logging.WARNING for r in transient_records) def test_initial_heartbeat_stale_token_skips_validation( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() validate_called = threading.Event() def track(*_: object, **__: object) -> dict[str, Any]: validate_called.set() return _DUMMY_RESULT monkeypatch.setattr(Worker, "_run_subprocess", track) _send(queue_url, _valid_body()) mock_step_functions.send_task_heartbeat.side_effect = StaleTaskToken("retried by SFN") mock_protection.disable.side_effect = lambda: shutdown.set() _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 0 assert not validate_called.is_set() mock_step_functions.send_task_success.assert_not_called() mock_step_functions.send_task_failure.assert_not_called() def test_known_exception_sends_task_failure( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() mock_step_functions.send_task_failure.side_effect = lambda *_, **__: shutdown.set() def fatal(*_: object, **__: object) -> dict[str, Any]: raise DummyFatalError("boom") monkeypatch.setattr(Worker, "_run_subprocess", fatal) _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown) mock_step_functions.send_task_failure.assert_called_once() call = mock_step_functions.send_task_failure.call_args.kwargs assert call["error"] == "DummyFatalError" assert "boom" in call["cause"] def test_send_task_success_exception_logged( self, queue_url: str, shutdown_on_subprocess: threading.Event, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: _send(queue_url, _valid_body()) mock_step_functions.send_task_success.side_effect = RuntimeError("sfn throttled") _make_worker(queue_url).process_messages(shutdown_on_subprocess) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_called_once() def test_send_task_success_stale_token_logged( self, queue_url: str, shutdown_on_subprocess: threading.Event, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: _send(queue_url, _valid_body()) mock_step_functions.send_task_success.side_effect = StaleTaskToken("task already completed") _make_worker(queue_url).process_messages(shutdown_on_subprocess) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_called_once() def test_send_task_failure_stale_token_logged( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() def crash(*_: object, **__: object) -> dict[str, Any]: error_message = "mediainfo segfault" raise RuntimeError(error_message) monkeypatch.setattr(Worker, "_run_subprocess", crash) _send(queue_url, _valid_body()) mock_step_functions.send_task_failure.side_effect = StaleTaskToken("task already completed") mock_protection.disable.side_effect = lambda: shutdown.set() _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_failure.assert_called_once() def test_initial_heartbeat_transient_error_continues( self, queue_url: str, shutdown_on_subprocess: threading.Event, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: _send(queue_url, _valid_body()) mock_step_functions.send_task_heartbeat.side_effect = RuntimeError("sfn throttled") _make_worker(queue_url).process_messages(shutdown_on_subprocess) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_called_once() class TestShutdown: def test_preset_shutdown_skips_polling( self, queue_url: str, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: _send(queue_url, _valid_body()) shutdown = threading.Event() shutdown.set() _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 1 mock_protection.enable.assert_not_called() mock_step_functions.send_task_success.assert_not_called() class TestAbortedStaleTaskFlow: def test_aborted_stale_returns_silently( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() def abort(*_: object, **__: object) -> dict[str, Any]: shutdown.set() raise AbortedStaleTask("token stale mid-run") monkeypatch.setattr(Worker, "_run_subprocess", abort) _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown) assert _queue_depth(queue_url) == 0 mock_step_functions.send_task_success.assert_not_called() mock_step_functions.send_task_failure.assert_not_called() class TestUnknownSubprocessErrorFlow: def test_unknown_subprocess_error_sends_task_failure( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() def unknown(*_: object, **__: object) -> dict[str, Any]: shutdown.set() raise UnknownSubprocessError("unknown error='FooError': boom") monkeypatch.setattr(Worker, "_run_subprocess", unknown) _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown) mock_step_functions.send_task_failure.assert_called_once() assert mock_step_functions.send_task_failure.call_args.kwargs["error"] == "UnknownSubprocessError" class TestSubprocessCrashFlow: def test_subprocess_crash_sends_task_failure( self, queue_url: str, monkeypatch: pytest.MonkeyPatch, mock_protection: MagicMock, mock_step_functions: MagicMock, ) -> None: shutdown = threading.Event() def crash(*_: object, **__: object) -> dict[str, Any]: shutdown.set() raise SubprocessCrashError("subprocess terminated by signal 9") monkeypatch.setattr(Worker, "_run_subprocess", crash) _send(queue_url, _valid_body()) _make_worker(queue_url).process_messages(shutdown) mock_step_functions.send_task_failure.assert_called_once() assert mock_step_functions.send_task_failure.call_args.kwargs["error"] == "SubprocessCrashError" def _spawn_sleeper() -> subprocess.Popen[str]: return subprocess.Popen( [sys.executable, "-c", "import time; time.sleep(30)"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, start_new_session=True, ) class TestTerminateProcessGroup: def test_terminates_running_process_with_sigkill(self) -> None: process = _spawn_sleeper() try: Worker._terminate_process_group(process) process.wait(timeout=2) assert process.returncode < 0 finally: if process.poll() is None: process.kill() process.wait() def test_noop_if_already_exited(self) -> None: process = subprocess.Popen( [sys.executable, "-c", "pass"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, start_new_session=True, ) process.wait() Worker._terminate_process_group(process) class TestRunHeartbeat: def test_terminates_subprocess_on_stale_token(self, monkeypatch: pytest.MonkeyPatch) -> None: process = _spawn_sleeper() try: monkeypatch.setattr( "src.worker.worker.step_functions.send_task_heartbeat", MagicMock(side_effect=StaleTaskToken("stale")), ) stop = threading.Event() aborted_stale_task = threading.Event() Worker._run_heartbeat("tok", stop, process, aborted_stale_task, interval_seconds=0.01) process.wait(timeout=2) assert process.returncode < 0 assert aborted_stale_task.is_set() finally: if process.poll() is None: process.kill() process.wait() def test_exits_cleanly_when_stop_set(self, monkeypatch: pytest.MonkeyPatch) -> None: process = _spawn_sleeper() try: mock_heartbeat = MagicMock() monkeypatch.setattr("src.worker.worker.step_functions.send_task_heartbeat", mock_heartbeat) stop = threading.Event() stop.set() aborted_stale_task = threading.Event() Worker._run_heartbeat("tok", stop, process, aborted_stale_task, interval_seconds=0.01) mock_heartbeat.assert_not_called() assert process.poll() is None assert not aborted_stale_task.is_set() finally: process.kill() process.wait() def test_continues_on_non_stale_exception(self, monkeypatch: pytest.MonkeyPatch) -> None: process = _spawn_sleeper() try: stop = threading.Event() aborted_stale_task = threading.Event() calls = 0 def flaky(*_: object, **__: object) -> None: nonlocal calls calls += 1 if calls >= 3: stop.set() error_message = "transient" raise RuntimeError(error_message) monkeypatch.setattr("src.worker.worker.step_functions.send_task_heartbeat", flaky) Worker._run_heartbeat("tok", stop, process, aborted_stale_task, interval_seconds=0.01) assert calls == 3 assert process.poll() is None # still alive — transient errors don't terminate assert not aborted_stale_task.is_set() finally: process.kill() process.wait() class TestParseSubprocessResult: @pytest.mark.parametrize( ("error", "expected_exc"), [ pytest.param("DummyTransientError", DummyTransientError, id="transient"), pytest.param("DummyFatalError", DummyFatalError, id="fatal"), ], ) def test_maps_subprocess_error_to_exception( self, error: str, expected_exc: type[Exception], queue_url: str, ) -> None: worker = _make_worker(queue_url) process = MagicMock(returncode=1) stdout = json.dumps({"error": error, "message": "boom"}) with pytest.raises(expected_exc, match="boom"): worker._parse_subprocess_result(process, stdout, "", aborted_stale_task=False) def test_unknown_error_raises_unknown_subprocess_error(self, queue_url: str) -> None: worker = _make_worker(queue_url) process = MagicMock(returncode=1) stdout = json.dumps({"error": "MysteryError", "message": "who dis"}) with pytest.raises(UnknownSubprocessError, match="MysteryError"): worker._parse_subprocess_result(process, stdout, "", aborted_stale_task=False) def test_worker_exception_name_in_stdout_is_not_dispatched(self, queue_url: str) -> None: # A subprocess emitting a worker-internal exception name must not be # dispatched as that exception — otherwise a buggy or malicious # subprocess could fake silent control flow (e.g. AbortedStaleTask). import logging worker = Worker( request_schema=DummyRequest, select_subprocess_module=lambda _: _SUBPROCESS_MODULE, sqs_queue_url=queue_url, exception_levels={MalformedSubprocessOutput: logging.WARNING, DummyFatalError: logging.ERROR}, heartbeat_interval_seconds=60.0, message_processing_timeout_minutes=_TIMEOUT_MINUTES, ) process = MagicMock(returncode=1) stdout = json.dumps({"error": "MalformedSubprocessOutput", "message": "spoofed"}) with pytest.raises(UnknownSubprocessError, match="MalformedSubprocessOutput"): worker._parse_subprocess_result(process, stdout, "", aborted_stale_task=False) def test_success_returns_result(self, queue_url: str) -> None: worker = _make_worker(queue_url) process = MagicMock(returncode=0) stdout = json.dumps({"is_valid": True, "extra": "fields"}) result = worker._parse_subprocess_result(process, stdout, "", aborted_stale_task=False) assert result == {"is_valid": True, "extra": "fields"} @pytest.mark.parametrize( "stdout", [ pytest.param("not json at all", id="unparseable_json"), pytest.param('["just", "an", "array"]', id="json_array_not_dict"), ], ) def test_success_with_non_object_json_raises_malformed(self, queue_url: str, stdout: str) -> None: worker = _make_worker(queue_url) process = MagicMock(returncode=0) with pytest.raises(MalformedSubprocessOutput): worker._parse_subprocess_result(process, stdout, "", aborted_stale_task=False) def test_signal_kill_with_aborted_flag_raises_aborted_stale_task(self, queue_url: str) -> None: worker = _make_worker(queue_url) process = MagicMock(returncode=-9) with pytest.raises(AbortedStaleTask, match="task token went stale"): worker._parse_subprocess_result(process, "", "", aborted_stale_task=True) def test_signal_kill_without_aborted_flag_raises_subprocess_crash_error(self, queue_url: str) -> None: worker = _make_worker(queue_url) process = MagicMock(returncode=-9) with pytest.raises(SubprocessCrashError, match="signal 9"): worker._parse_subprocess_result(process, "", "", aborted_stale_task=False) def test_clean_exit_with_aborted_flag_parses_result(self, queue_url: str) -> None: # Race: subprocess finished cleanly right as the heartbeat thread set the flag. # Trust the result; _send_task_success will handle any stale-token rejection. worker = _make_worker(queue_url) process = MagicMock(returncode=0) stdout = json.dumps({"is_valid": False}) result = worker._parse_subprocess_result(process, stdout, "", aborted_stale_task=True) assert result == {"is_valid": False} @pytest.mark.parametrize( "stdout", [ pytest.param("garbage\nmore garbage\n", id="unparseable_json"), pytest.param('["a", "list"]', id="json_array_not_dict"), pytest.param('{"missing_keys": true}', id="missing_required_keys"), pytest.param("", id="empty_stdout"), ], ) def test_malformed_error_envelope_raises_malformed_subprocess_output( self, queue_url: str, stdout: str, ) -> None: worker = _make_worker(queue_url) process = MagicMock(returncode=1) with pytest.raises(MalformedSubprocessOutput, match="unparseable error envelope"): worker._parse_subprocess_result(process, stdout, "", aborted_stale_task=False) class TestSpawnSubprocess: def test_env_var_and_popen_kwargs(self, monkeypatch: pytest.MonkeyPatch, queue_url: str) -> None: request = DummyRequest(task_token="tok", bucket="my-bucket", key="path/to/file.wav") body = '{"task_token":"tok","bucket":"my-bucket","key":"path/to/file.wav","extra_field":"forwarded"}' argv_captured: list[str] = [] kwargs_captured: dict[str, object] = {} def fake_popen(argv: list[str], **kwargs: object) -> MagicMock: argv_captured.extend(argv) kwargs_captured.update(kwargs) return MagicMock() monkeypatch.setattr("src.worker.worker.subprocess.Popen", fake_popen) worker = _make_worker(queue_url) worker._spawn_subprocess(request, body) assert argv_captured[1:] == ["-m", _SUBPROCESS_MODULE] assert kwargs_captured["start_new_session"] is True assert kwargs_captured["stdout"] == subprocess.PIPE assert kwargs_captured["stderr"] == subprocess.PIPE assert kwargs_captured["text"] is True env = kwargs_captured["env"] assert isinstance(env, dict) assert env[WORKER_REQUEST] == body