Skip to content

Commit ec50558

Browse files
committed
test(worker): check logging on a real worker process
Keep the fix to the if-statement in start_listen, and replace the mocked tests with a test that starts a real `taskiq worker` for each available start method and waits for "Listening started." in its output.
1 parent 9a9f04d commit ec50558

3 files changed

Lines changed: 127 additions & 69 deletions

File tree

‎taskiq/cli/worker/run.py‎

Lines changed: 5 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -70,25 +70,6 @@ def get_receiver_type(args: WorkerArgs) -> type[Receiver]:
7070
return receiver_type
7171

7272

73-
def configure_child_logging(args: WorkerArgs) -> None:
74-
"""
75-
Configure logging in a worker process.
76-
77-
A process started with the ``fork`` start method inherits
78-
the logging configuration of the main process.
79-
Processes started with ``spawn`` or ``forkserver``
80-
(the default on Linux since Python 3.14) begin with
81-
a fresh interpreter, so logging has to be configured again.
82-
83-
:param args: CLI arguments.
84-
"""
85-
if args.configure_logging and get_start_method() != "fork":
86-
logging.basicConfig(
87-
level=args.log_level,
88-
format=args.log_format,
89-
)
90-
91-
9273
def start_listen(args: WorkerArgs) -> None:
9374
"""
9475
This function starts actual listening process.
@@ -105,7 +86,11 @@ def start_listen(args: WorkerArgs) -> None:
10586
"""
10687
shutdown_event = asyncio.Event()
10788
hardkill_counter = 0
108-
configure_child_logging(args)
89+
if args.configure_logging and get_start_method() != "fork":
90+
logging.basicConfig(
91+
level=args.log_level,
92+
format=args.log_format,
93+
)
10994

11095
def interrupt_handler(signum: int, _frame: Any) -> None:
11196
"""

‎tests/cli/worker/test_child_logging.py‎

Lines changed: 0 additions & 49 deletions
This file was deleted.
Lines changed: 122 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,122 @@
1+
import multiprocessing
2+
import queue
3+
import signal
4+
import subprocess
5+
import sys
6+
import threading
7+
import time
8+
from pathlib import Path
9+
10+
import pytest
11+
12+
BROKER_MODULE = """
13+
import asyncio
14+
from collections.abc import AsyncGenerator
15+
16+
from taskiq import AsyncBroker, BrokerMessage
17+
18+
19+
class IdleBroker(AsyncBroker):
20+
async def kick(self, message: BrokerMessage) -> None:
21+
pass
22+
23+
async def listen(self) -> AsyncGenerator[bytes, None]:
24+
while True:
25+
await asyncio.sleep(3600)
26+
yield b""
27+
28+
29+
broker = IdleBroker()
30+
"""
31+
32+
# Runs the taskiq CLI, optionally forcing a multiprocessing start method first.
33+
CLI_RUNNER = """
34+
import multiprocessing
35+
import sys
36+
37+
from taskiq.__main__ import main
38+
39+
if __name__ == "__main__":
40+
start_method = sys.argv.pop(1)
41+
if start_method != "default":
42+
multiprocessing.set_start_method(start_method)
43+
main()
44+
"""
45+
46+
STARTUP_TIMEOUT = 30
47+
48+
49+
def worker_start_methods() -> list[str]:
50+
# The worker always uses "spawn" on macOS: forcing another method would conflict.
51+
if sys.platform == "darwin":
52+
return ["default"]
53+
# The default method is one of these, so it's covered without a separate run.
54+
return multiprocessing.get_all_start_methods()
55+
56+
57+
def wait_for_output(process: subprocess.Popen[str], text: str) -> bool:
58+
lines: queue.Queue[str] = queue.Queue()
59+
60+
def read_output() -> None:
61+
assert process.stdout is not None
62+
for line in process.stdout:
63+
lines.put(line)
64+
65+
threading.Thread(target=read_output, daemon=True).start()
66+
deadline = time.monotonic() + STARTUP_TIMEOUT
67+
while (remaining := deadline - time.monotonic()) > 0:
68+
try:
69+
if text in lines.get(timeout=remaining):
70+
return True
71+
except queue.Empty:
72+
return False
73+
return False
74+
75+
76+
def stop(process: subprocess.Popen[str]) -> None:
77+
process.send_signal(signal.SIGINT)
78+
try:
79+
process.wait(timeout=10)
80+
except subprocess.TimeoutExpired:
81+
process.kill()
82+
process.wait()
83+
84+
85+
@pytest.mark.skipif(
86+
sys.platform == "win32",
87+
reason="Worker processes can't be stopped with SIGINT on Windows.",
88+
)
89+
@pytest.mark.parametrize("start_method", worker_start_methods())
90+
def test_worker_process_logs_listening_started(
91+
tmp_path: Path,
92+
start_method: str,
93+
) -> None:
94+
"""
95+
Worker processes emit their own log records.
96+
97+
A worker started with "spawn" or "forkserver" (the default on Linux
98+
since Python 3.14) doesn't inherit the main process logging config.
99+
"""
100+
(tmp_path / "idle_broker.py").write_text(BROKER_MODULE)
101+
(tmp_path / "run_cli.py").write_text(CLI_RUNNER)
102+
process = subprocess.Popen( # noqa: S603
103+
[
104+
sys.executable,
105+
"run_cli.py",
106+
start_method,
107+
"worker",
108+
"idle_broker:broker",
109+
"--workers",
110+
"1",
111+
"--log-level",
112+
"INFO",
113+
],
114+
cwd=tmp_path,
115+
stdout=subprocess.PIPE,
116+
stderr=subprocess.STDOUT,
117+
text=True,
118+
)
119+
try:
120+
assert wait_for_output(process, "Listening started.")
121+
finally:
122+
stop(process)

0 commit comments

Comments
 (0)