/
/
1"""Tests for the AsyncProcess helper."""
2
3from __future__ import annotations
4
5import asyncio
6import os
7import sys
8import time
9from collections.abc import AsyncGenerator
10from unittest.mock import MagicMock
11
12import pytest
13
14from music_assistant.helpers import process as process_module
15from music_assistant.helpers.process import AsyncProcess
16
17# Comfortably beyond the OS pipe capacity plus asyncio's default high-water mark,
18# so the bytes are guaranteed to still be queued in our own write buffer.
19_MORE_THAN_THE_PIPE_HOLDS = b"\x00" * (4 * 1024 * 1024)
20
21# Ignores SIGINT and keeps stdout open, so nothing but SIGKILL ends it and its
22# pipe never reaches EOF. The marker tells the test the handler is installed.
23_WEDGED_CHILD = (
24 "import signal, sys, time; signal.signal(signal.SIGINT, signal.SIG_IGN); "
25 "sys.stdout.write('ready\\n'); sys.stdout.flush(); time.sleep(30)"
26)
27
28
29@pytest.fixture(name="piped_process")
30async def piped_process_fixture() -> AsyncGenerator[tuple[AsyncProcess, int]]:
31 """
32 Yield an AsyncProcess writing to a real pipe, plus its unread read end.
33
34 The write side is a genuine asyncio pipe transport, so the flow control
35 stdin_quiesced relies on behaves exactly as it does against a real process,
36 while nothing consumes the read end until a test chooses to.
37 """
38 read_fd, write_fd = os.pipe()
39 loop = asyncio.get_running_loop()
40 transport, protocol = await loop.connect_write_pipe(
41 asyncio.streams.FlowControlMixin, os.fdopen(write_fd, "wb", 0)
42 )
43 writer = asyncio.StreamWriter(transport, protocol, None, loop)
44 proc = AsyncProcess(["cat"], stdin=True)
45 proc.proc = MagicMock(stdin=writer)
46 try:
47 yield proc, read_fd
48 finally:
49 transport.abort()
50 # abort() only schedules the callback that closes the write side, so let
51 # it run before the loop goes away rather than relying on the ordering.
52 await asyncio.sleep(0)
53 os.close(read_fd)
54
55
56def _queue_without_waiting(proc: AsyncProcess, data: bytes) -> None:
57 """
58 Hand data to stdin without awaiting it, leaving it in the write buffer.
59
60 ``write`` would block on its own drain while the pipe is unread, which is the
61 state these tests need to set up rather than something to wait through.
62 """
63 assert proc.proc is not None
64 assert proc.proc.stdin is not None
65 proc.proc.stdin.write(data)
66
67
68async def _consume(read_fd: int, total: int) -> None:
69 """Read `total` bytes off the pipe so the write buffer can empty."""
70 remaining = total
71 while remaining > 0:
72 remaining -= len(await asyncio.to_thread(os.read, read_fd, 65536))
73
74
75@pytest.mark.asyncio
76async def test_stdin_quiesced_waits_for_the_write_buffer_to_empty(
77 piped_process: tuple[AsyncProcess, int],
78) -> None:
79 """The block is entered only once every queued byte has reached the reader."""
80 proc, read_fd = piped_process
81 _queue_without_waiting(proc, _MORE_THAN_THE_PIPE_HOLDS)
82 assert proc.proc is not None
83 assert proc.proc.stdin is not None
84 assert proc.proc.stdin.transport.get_write_buffer_size() > 0
85 reader = asyncio.create_task(_consume(read_fd, len(_MORE_THAN_THE_PIPE_HOLDS)))
86
87 async with proc.stdin_quiesced() as quiesced:
88 assert quiesced is True
89 assert proc.proc.stdin.transport.get_write_buffer_size() == 0
90
91 await reader
92
93
94@pytest.mark.asyncio
95async def test_stdin_quiesced_reports_a_reader_that_never_catches_up(
96 piped_process: tuple[AsyncProcess, int],
97) -> None:
98 """A reader that never consumes the pipe reports failure instead of hanging."""
99 proc, _ = piped_process
100 _queue_without_waiting(proc, _MORE_THAN_THE_PIPE_HOLDS)
101
102 async with proc.stdin_quiesced(timeout=0.2) as quiesced:
103 assert quiesced is False
104
105
106@pytest.mark.asyncio
107async def test_stdin_quiesced_keeps_writes_out_of_the_block(
108 piped_process: tuple[AsyncProcess, int],
109) -> None:
110 """
111 No write can land while the block runs, so nothing queues up behind a drain.
112
113 This is the guarantee the block exists for: a caller telling the process
114 something about the bytes it has been handed needs stdin to stay as it left it
115 until the process has answered.
116 """
117 proc, read_fd = piped_process
118 reader = asyncio.create_task(_consume(read_fd, len(b"late")))
119
120 async with proc.stdin_quiesced() as quiesced:
121 assert quiesced is True
122 writing = asyncio.create_task(proc.write(b"late"))
123 await asyncio.sleep(0)
124
125 assert not writing.done()
126 assert proc.proc is not None
127 assert proc.proc.stdin is not None
128 assert proc.proc.stdin.transport.get_write_buffer_size() == 0
129
130 await writing
131 await reader
132
133
134@pytest.mark.asyncio
135async def test_stdin_quiesced_restores_normal_writing(
136 piped_process: tuple[AsyncProcess, int],
137) -> None:
138 """Writing carries on unaffected afterwards, on the restored buffer limits."""
139 proc, read_fd = piped_process
140 assert proc.proc is not None
141 assert proc.proc.stdin is not None
142 limits = proc.proc.stdin.transport.get_write_buffer_limits()
143 reader = asyncio.create_task(_consume(read_fd, len(b"first") + len(b"second")))
144 await proc.write(b"first")
145
146 async with proc.stdin_quiesced() as quiesced:
147 assert quiesced is True
148
149 assert proc.proc.stdin.transport.get_write_buffer_limits() == limits
150 await proc.write(b"second")
151 await reader
152
153
154@pytest.mark.asyncio
155async def test_stdin_quiesced_is_a_noop_without_a_process() -> None:
156 """A process that was never started quiesces trivially rather than raising."""
157 proc = AsyncProcess(["cat"], stdin=True)
158
159 async with proc.stdin_quiesced() as quiesced:
160 assert quiesced is True
161
162
163@pytest.mark.asyncio
164async def test_iter_stdout_drains_lines_buffered_after_exit() -> None:
165 """
166 Output written just before the process exits is still delivered.
167
168 A short-lived process can write everything and be reaped before the reader
169 runs, so keying the stdout reader off the returncode would drop exactly the
170 output that explains why it exited.
171 """
172 proc = AsyncProcess(
173 ["sh", "-c", "for i in $(seq 1 50); do echo line$i; done"],
174 stdout=True,
175 stderr=asyncio.subprocess.STDOUT,
176 )
177 await proc.start()
178 await proc.wait()
179
180 lines = [line async for line in proc.iter_stdout()]
181
182 assert lines == [f"line{index}" for index in range(1, 51)]
183 await proc.close()
184
185
186@pytest.mark.asyncio
187async def test_read_stdout_stops_once_the_process_is_closed() -> None:
188 """A closed process reports EOF instead of waiting on a stream it no longer owns."""
189 proc = AsyncProcess(["sh", "-c", "sleep 30"], stdout=True, stderr=asyncio.subprocess.STDOUT)
190 await proc.start()
191 await proc.close()
192
193 assert await proc.read_stdout() == b""
194
195
196@pytest.mark.asyncio
197async def test_second_close_returns_without_waiting_out_the_stream_locks() -> None:
198 """
199 Closing an already-closed process is cheap.
200
201 close() keeps the stdin/stdout locks it takes, so a second call used to sit
202 through both 5s acquire timeouts - a delay paid on every supervised restart
203 that closes the process before its own cleanup runs.
204 """
205 proc = AsyncProcess(["sh", "-c", "sleep 30"], stdout=True, stderr=asyncio.subprocess.STDOUT)
206 await proc.start()
207 await proc.close()
208
209 started = time.monotonic()
210 await proc.close()
211
212 assert time.monotonic() - started < 1
213
214
215@pytest.mark.asyncio
216async def test_close_reaps_a_child_that_never_closes_its_pipes(
217 monkeypatch: pytest.MonkeyPatch,
218) -> None:
219 """
220 A child holding its pipes open must not keep close() from reaping it.
221
222 Draining stdout is what lets a healthy process flush before it is reaped, so
223 an unbounded drain waits out a wedged child forever and the terminate/SIGKILL
224 escalation is never reached.
225 """
226 monkeypatch.setattr(process_module, "PIPE_DRAIN_TIMEOUT", 0.2)
227 proc = AsyncProcess(
228 [sys.executable, "-c", _WEDGED_CHILD], stdout=True, stderr=asyncio.subprocess.STDOUT
229 )
230 await proc.start()
231 assert await proc.read_stdout() == b"ready\n"
232
233 async with asyncio.timeout(20):
234 await proc.close()
235
236 assert proc.returncode is not None
237