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