/
/
1"""Tests for the AsyncProcess helper."""
2
3from __future__ import annotations
4
5import asyncio
6import os
7from collections.abc import AsyncGenerator
8from unittest.mock import MagicMock
9
10import pytest
11
12from music_assistant.helpers.process import AsyncProcess
13
14# Comfortably beyond the OS pipe capacity plus asyncio's default high-water mark,
15# so the bytes are guaranteed to still be queued in our own write buffer.
16_MORE_THAN_THE_PIPE_HOLDS = b"\x00" * (4 * 1024 * 1024)
17
18
19@pytest.fixture(name="piped_process")
20async def piped_process_fixture() -> AsyncGenerator[tuple[AsyncProcess, int]]:
21 """
22 Yield an AsyncProcess writing to a real pipe, plus its unread read end.
23
24 The write side is a genuine asyncio pipe transport, so the flow control
25 stdin_quiesced relies on behaves exactly as it does against a real process,
26 while nothing consumes the read end until a test chooses to.
27 """
28 read_fd, write_fd = os.pipe()
29 loop = asyncio.get_running_loop()
30 transport, protocol = await loop.connect_write_pipe(
31 asyncio.streams.FlowControlMixin, os.fdopen(write_fd, "wb", 0)
32 )
33 writer = asyncio.StreamWriter(transport, protocol, None, loop)
34 proc = AsyncProcess(["cat"], stdin=True)
35 proc.proc = MagicMock(stdin=writer)
36 try:
37 yield proc, read_fd
38 finally:
39 transport.abort()
40 # abort() only schedules the callback that closes the write side, so let
41 # it run before the loop goes away rather than relying on the ordering.
42 await asyncio.sleep(0)
43 os.close(read_fd)
44
45
46def _queue_without_waiting(proc: AsyncProcess, data: bytes) -> None:
47 """
48 Hand data to stdin without awaiting it, leaving it in the write buffer.
49
50 ``write`` would block on its own drain while the pipe is unread, which is the
51 state these tests need to set up rather than something to wait through.
52 """
53 assert proc.proc is not None
54 assert proc.proc.stdin is not None
55 proc.proc.stdin.write(data)
56
57
58async def _consume(read_fd: int, total: int) -> None:
59 """Read `total` bytes off the pipe so the write buffer can empty."""
60 remaining = total
61 while remaining > 0:
62 remaining -= len(await asyncio.to_thread(os.read, read_fd, 65536))
63
64
65@pytest.mark.asyncio
66async def test_stdin_quiesced_waits_for_the_write_buffer_to_empty(
67 piped_process: tuple[AsyncProcess, int],
68) -> None:
69 """The block is entered only once every queued byte has reached the reader."""
70 proc, read_fd = piped_process
71 _queue_without_waiting(proc, _MORE_THAN_THE_PIPE_HOLDS)
72 assert proc.proc is not None
73 assert proc.proc.stdin is not None
74 assert proc.proc.stdin.transport.get_write_buffer_size() > 0
75 reader = asyncio.create_task(_consume(read_fd, len(_MORE_THAN_THE_PIPE_HOLDS)))
76
77 async with proc.stdin_quiesced() as quiesced:
78 assert quiesced is True
79 assert proc.proc.stdin.transport.get_write_buffer_size() == 0
80
81 await reader
82
83
84@pytest.mark.asyncio
85async def test_stdin_quiesced_reports_a_reader_that_never_catches_up(
86 piped_process: tuple[AsyncProcess, int],
87) -> None:
88 """A reader that never consumes the pipe reports failure instead of hanging."""
89 proc, _ = piped_process
90 _queue_without_waiting(proc, _MORE_THAN_THE_PIPE_HOLDS)
91
92 async with proc.stdin_quiesced(timeout=0.2) as quiesced:
93 assert quiesced is False
94
95
96@pytest.mark.asyncio
97async def test_stdin_quiesced_keeps_writes_out_of_the_block(
98 piped_process: tuple[AsyncProcess, int],
99) -> None:
100 """
101 No write can land while the block runs, so nothing queues up behind a drain.
102
103 This is the guarantee the block exists for: a caller telling the process
104 something about the bytes it has been handed needs stdin to stay as it left it
105 until the process has answered.
106 """
107 proc, read_fd = piped_process
108 reader = asyncio.create_task(_consume(read_fd, len(b"late")))
109
110 async with proc.stdin_quiesced() as quiesced:
111 assert quiesced is True
112 writing = asyncio.create_task(proc.write(b"late"))
113 await asyncio.sleep(0)
114
115 assert not writing.done()
116 assert proc.proc is not None
117 assert proc.proc.stdin is not None
118 assert proc.proc.stdin.transport.get_write_buffer_size() == 0
119
120 await writing
121 await reader
122
123
124@pytest.mark.asyncio
125async def test_stdin_quiesced_restores_normal_writing(
126 piped_process: tuple[AsyncProcess, int],
127) -> None:
128 """Writing carries on unaffected afterwards, on the restored buffer limits."""
129 proc, read_fd = piped_process
130 assert proc.proc is not None
131 assert proc.proc.stdin is not None
132 limits = proc.proc.stdin.transport.get_write_buffer_limits()
133 reader = asyncio.create_task(_consume(read_fd, len(b"first") + len(b"second")))
134 await proc.write(b"first")
135
136 async with proc.stdin_quiesced() as quiesced:
137 assert quiesced is True
138
139 assert proc.proc.stdin.transport.get_write_buffer_limits() == limits
140 await proc.write(b"second")
141 await reader
142
143
144@pytest.mark.asyncio
145async def test_stdin_quiesced_is_a_noop_without_a_process() -> None:
146 """A process that was never started quiesces trivially rather than raising."""
147 proc = AsyncProcess(["cat"], stdin=True)
148
149 async with proc.stdin_quiesced() as quiesced:
150 assert quiesced is True
151