/
/
1"""
2Unit tests for the Sendspin -> AirPlay bridge timing.
3
4Cover ten things, with the Sendspin clock mocked via ``ManualClock`` so the
5tests are deterministic and independent of the host wall-clock:
6
7* the clock-domain conversion turning a Sendspin audible instant (Sendspin's own
8 monotonic clock) into the unix epoch ms used by the START command, and back;
9* the startup lead reported to Sendspin, which decides how far ahead of the
10 audible instant it schedules the first chunk;
11* the start anchor: byte 0 is anchored to the first chunk Sendspin delivers, so a
12 fresh track keeps position 0 and a late joiner lands at the group's live position;
13* anchoring against the binary: the commanded instant honours the join headroom,
14 the receiver's clock-ready projection and the content Sendspin already
15 scheduled, and the content is then mapped onto the instant the binary acked;
16* the timeline alignment that keeps every chunk at the byte offset its timestamp
17 claims, so a discontinuity in the Sendspin timeline does not shift the device
18 off the group's clock for the rest of the stream;
19* the playout shift the binary reports after a PCM starvation, which moves the
20 anchor so the device does not stay behind the group once it re-anchors itself;
21* the write pacing that keeps the device buffered a bounded amount ahead of real
22 time so a late-join catch-up backlog is not dumped into the CLI;
23* the warm handover: a running, connected stream is kept (not torn down) across
24 a new Sendspin stream and rides the persistent-stdin flush-refill (FLUSH +
25 re-anchoring START) instead of a cold reconnect -- with flush-timeout and
26 superseded-task fallback, and the supersession handling that keeps a stale
27 start from spawning a process or touching the stream a newer one owns;
28* the recovery from a transport lost mid-stream: the dead CLI is released and
29 re-anchored on the group's live timeline. Every give-up then takes the speaker
30 out of the Sendspin session, so the player stops reporting playback nobody can
31 hear, and a bounded re-join brings back one that was only briefly away.
32"""
33
34import asyncio
35from collections.abc import Coroutine
36from typing import cast
37from unittest.mock import AsyncMock, MagicMock, patch
38
39import pytest
40from aiosendspin.clock import ManualClock
41from aiosendspin.server.roles import AudioChunk
42
43from music_assistant.providers.airplay.constants import (
44 AIRPLAY_CLOCK_READY_LEAD_MS,
45 AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS,
46 AIRPLAY_SPLICE_LEAD_MARGIN_MS,
47 ClockReadiness,
48 StreamingProtocol,
49)
50from music_assistant.providers.airplay.sendspin_bridge import (
51 BRIDGE_COLD_START_LEAD_MS,
52 BRIDGE_MIN_BUFFER_MS,
53 BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS,
54 BRIDGE_WARM_START_LEAD_MS,
55 MAX_DEVICE_BUFFER_SECONDS,
56 MAX_HELD_AUDIO_US,
57 MAX_HELD_CHUNKS,
58 PAD_BLOCK_FRAMES,
59 SILENCE_BLOCK,
60 SendspinAirPlayBridge,
61 SendspinBridgeManager,
62 device_buffer_ahead_seconds,
63 sendspin_audible_instant_to_unix_ms,
64 unix_ms_to_sendspin_audible_instant,
65)
66from music_assistant.providers.sendspin.bridge_role import (
67 BRIDGE_BYTES_PER_SAMPLE,
68 BRIDGE_CHANNELS,
69 BRIDGE_SAMPLE_RATE,
70 BridgePlayerRole,
71)
72
73BRIDGE_BYTES_PER_SECOND = BRIDGE_SAMPLE_RATE * BRIDGE_CHANNELS * BRIDGE_BYTES_PER_SAMPLE
74BRIDGE_BYTES_PER_FRAME = BRIDGE_CHANNELS * BRIDGE_BYTES_PER_SAMPLE
75
76# A large, arbitrary Sendspin monotonic-clock epoch (microseconds). Real
77# monotonic clocks start from an unspecified point (e.g. host boot), so the
78# conversion must never depend on this value.
79SENDSPIN_EPOCH_US = 5_000_000_000_000 # ~57.8 days of monotonic uptime
80UNIX_NOW_S = 1_784_000_000.0 # fixed unix wall-clock reading for the tests
81UNIX_NOW_MS = int(UNIX_NOW_S * 1000)
82COLD_LEAD_MS = BRIDGE_COLD_START_LEAD_MS
83WARM_LEAD_MS = BRIDGE_WARM_START_LEAD_MS
84# Patched with zero delays so the re-join backoff runs instantly while the
85# attempt-count and give-up logic around it stays real.
86_NO_REJOIN_DELAYS = "music_assistant.providers.airplay.sendspin_bridge.BRIDGE_REJOIN_ATTEMPT_DELAYS"
87
88
89def _audible_instant_us(clock: ManualClock, lead_ms: int) -> int:
90 """Return a sample Sendspin audible instant that far ahead of now (exercises the mapping)."""
91 return clock.now_us() + lead_ms * 1_000
92
93
94def _unix_at(sendspin_us: int) -> float:
95 """Model a constant-offset, same-rate Sendspin<->unix relationship."""
96 return UNIX_NOW_S + (sendspin_us - SENDSPIN_EPOCH_US) / 1_000_000
97
98
99def test_maps_future_delta_to_unix_now_plus_lead() -> None:
100 """An instant a lead ahead maps to unix_now + that lead (in ms)."""
101 clock = ManualClock(now_us_value=SENDSPIN_EPOCH_US)
102 drop_until = _audible_instant_us(clock, COLD_LEAD_MS)
103
104 start_unix_ms = sendspin_audible_instant_to_unix_ms(drop_until, clock.now_us(), UNIX_NOW_S)
105
106 assert start_unix_ms == int(UNIX_NOW_S * 1000) + COLD_LEAD_MS
107
108
109def test_standing_clock_offset_cancels_out() -> None:
110 """
111 The absolute Sendspin epoch must not affect the result.
112
113 Two wildly different monotonic epochs, with the same future delta and the
114 same unix reading, must yield the exact same start instant. This is what
115 makes the naive ``now/now`` subtraction correct: only the delta transfers
116 between the clocks, so any standing offset cancels.
117 """
118 clock_a = ManualClock(now_us_value=SENDSPIN_EPOCH_US)
119 clock_b = ManualClock(now_us_value=SENDSPIN_EPOCH_US + 987_654_321_000)
120
121 result_a = sendspin_audible_instant_to_unix_ms(
122 _audible_instant_us(clock_a, COLD_LEAD_MS), clock_a.now_us(), UNIX_NOW_S
123 )
124 result_b = sendspin_audible_instant_to_unix_ms(
125 _audible_instant_us(clock_b, COLD_LEAD_MS), clock_b.now_us(), UNIX_NOW_S
126 )
127
128 assert result_a == result_b
129
130
131def test_derived_start_equals_sendspin_audible_instant_in_unix() -> None:
132 """
133 The derived start lands on the unix time that coincides with the Sendspin instant.
134
135 Models the two clocks as running at the same rate with a constant offset
136 (unix = anchor + (sendspin_us - epoch)/1e6). The bridge only ever reads the
137 two clocks together, so the result must land exactly on the unix time that
138 coincides with the Sendspin audible instant, for any offset and any lead.
139 """
140 for lead_ms in (WARM_LEAD_MS, COLD_LEAD_MS):
141 clock = ManualClock(now_us_value=SENDSPIN_EPOCH_US)
142 drop_until = _audible_instant_us(clock, lead_ms)
143 # Some real time passes between setting the anchor and starting the CLI.
144 clock.advance_us(40_000) # 40 ms of setup churn (cleanup, task hop)
145 sendspin_now = clock.now_us()
146 unix_now = _unix_at(sendspin_now)
147
148 start_unix_ms = sendspin_audible_instant_to_unix_ms(drop_until, sendspin_now, unix_now)
149
150 assert start_unix_ms == int(_unix_at(drop_until) * 1000)
151
152
153def test_scheduling_gap_between_reads_shrinks_lead_not_target() -> None:
154 """
155 A gap before CLI start shrinks the remaining lead but keeps the audible instant fixed.
156
157 Computing the anchor immediately vs after a 400 ms gap must resolve to the
158 same unix instant, because the future delta is recomputed against the same
159 (advanced) Sendspin clock and unix reading.
160 """
161 clock = ManualClock(now_us_value=SENDSPIN_EPOCH_US)
162 drop_until = _audible_instant_us(clock, COLD_LEAD_MS)
163
164 immediate = sendspin_audible_instant_to_unix_ms(drop_until, clock.now_us(), UNIX_NOW_S)
165
166 gap_s = 0.4
167 clock.advance_us(int(gap_s * 1_000_000))
168 delayed = sendspin_audible_instant_to_unix_ms(drop_until, clock.now_us(), UNIX_NOW_S + gap_s)
169
170 assert immediate == delayed
171 # And the remaining lead really did shrink by the gap.
172 remaining_lead_ms = delayed - int((UNIX_NOW_S + gap_s) * 1000)
173 assert remaining_lead_ms == COLD_LEAD_MS - int(gap_s * 1000)
174
175
176def test_anchor_already_in_the_past_maps_to_a_past_unix_instant() -> None:
177 """
178 An audible instant behind 'now' yields a unix ms before the unix reading.
179
180 This is the setup-outran-the-lead edge case: the value stays a faithful
181 projection (negative lead) rather than being clamped here, so the anchor
182 math can see it and raise the start to the join floor itself.
183 """
184 clock = ManualClock(now_us_value=SENDSPIN_EPOCH_US)
185 audible_in_the_past = clock.now_us() - 300_000 # 300 ms ago
186
187 start_unix_ms = sendspin_audible_instant_to_unix_ms(
188 audible_in_the_past, clock.now_us(), UNIX_NOW_S
189 )
190
191 assert start_unix_ms == int(UNIX_NOW_S * 1000) - 300
192 assert start_unix_ms < int(UNIX_NOW_S * 1000)
193
194
195# --- Start anchor: fresh keeps position 0, late join lands at live position ---
196
197
198def _make_bridge(
199 clock_now_us: int,
200 protocol: StreamingProtocol = StreamingProtocol.AIRPLAY2,
201 sync_adjust: int = 0,
202) -> SendspinAirPlayBridge:
203 """Build a bridge with mocked provider/player/server and a ManualClock."""
204 provider = MagicMock()
205 provider.mass = MagicMock()
206 # Real values: the decision is handed to the CLI verbatim and the group's is
207 # compared against it, both of which a MagicMock would answer truthily
208 # whatever was resolved. None models a group with no live decision.
209 provider.bridge_manager.resolve_shared_ptp = MagicMock(return_value=False)
210 provider.bridge_manager.group_shared_ptp = MagicMock(return_value=None)
211 airplay_player = MagicMock()
212 airplay_player.player_id = "apc43875e9e53a"
213 airplay_player.display_name = "Test Player"
214 airplay_player.protocol = protocol
215 # A real int: the anchor math guards sync_adjust with isinstance(..., int), so
216 # a MagicMock would silently read as 0 and pass the test for the wrong reason.
217 airplay_player.config.get_value = MagicMock(return_value=sync_adjust)
218 sendspin_server = MagicMock()
219 sendspin_server.clock = ManualClock(now_us_value=clock_now_us)
220 bridge = SendspinAirPlayBridge(provider, airplay_player, sendspin_server)
221 bridge._is_streaming = True
222 return bridge
223
224
225def _pcm_chunk(timestamp_us: int, duration_us: int = 100_000) -> AudioChunk:
226 """Build a silent PCM AudioChunk at a Sendspin timestamp."""
227 frames = int(duration_us * BRIDGE_SAMPLE_RATE / 1_000_000)
228 data = b"\x00" * (frames * BRIDGE_CHANNELS * BRIDGE_BYTES_PER_SAMPLE)
229 return AudioChunk(
230 data=data, timestamp_us=timestamp_us, duration_us=duration_us, byte_count=len(data)
231 )
232
233
234def test_fresh_start_anchors_to_first_chunk_and_keeps_intro() -> None:
235 """
236 A fresh track's opening is kept: byte 0 anchors to the first chunk, not now+lead.
237
238 Models the clip scenario where the first delivered chunk (file position 0)
239 is scheduled earlier than ``clock.now() + the bridge lead``. Anchoring to
240 ``now + lead`` would drop everything before it -- the intro. The chunk
241 timestamp must win, and its audio must reach the CLI, not be dropped.
242 """
243 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
244 now_plus_lead = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
245 first_chunk_ts = SENDSPIN_EPOCH_US + 250_000 # position 0, only 250 ms ahead of now
246 assert first_chunk_ts < now_plus_lead
247
248 with patch.object(bridge, "_start_protocol_from_chunk", MagicMock()):
249 bridge._on_audio_chunk(_pcm_chunk(first_chunk_ts))
250
251 assert bridge._drop_until_us == first_chunk_ts
252 # Held while the anchor is negotiated, then queued -- never discarded.
253 _settle_anchor(bridge)
254 assert not bridge._write_queue.empty()
255
256
257def test_late_join_anchors_to_catchup_target_live_position() -> None:
258 """
259 A late joiner lands at the group's current position, not at track zero.
260
261 After minutes of playback the first delivered chunk is the catch-up target
262 (playhead + the bridge lead), far from the track start. The anchor must follow that
263 chunk so the joiner maps onto the live timeline instead of restarting at 0.
264 """
265 playhead_us = SENDSPIN_EPOCH_US + 600_000_000 # 600 s into the session
266 bridge = _make_bridge(clock_now_us=playhead_us)
267 catchup_target_ts = playhead_us + COLD_LEAD_MS * 1_000
268
269 with patch.object(bridge, "_start_protocol_from_chunk", MagicMock()):
270 bridge._on_audio_chunk(_pcm_chunk(catchup_target_ts))
271
272 assert bridge._drop_until_us == catchup_target_ts
273 # The anchor tracks the advanced playhead, not a fresh now+lead-from-zero.
274 assert bridge._drop_until_us > SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
275
276
277# --- Timeline alignment: a discontinuity must not shift the device off the group clock ---
278
279
280def _drain_queued_bytes(bridge: SendspinAirPlayBridge) -> int:
281 """Return the number of audio bytes handed to the CLI writer, emptying the queue."""
282 total = 0
283 while not bridge._write_queue.empty():
284 data = bridge._write_queue.get_nowait()
285 if data is not None:
286 total += len(data)
287 return total
288
289
290def _settle_anchor(bridge: SendspinAirPlayBridge) -> None:
291 """Model the binary acking exactly the anchor asked for: replay what was held."""
292 bridge._anchor_settled = True
293 held = list(bridge._held_chunks)
294 bridge._held_chunks.clear()
295 bridge._held_us = 0
296 for chunk in held:
297 bridge._align_chunk(chunk)
298
299
300def _start_stream_at(bridge: SendspinAirPlayBridge, first_chunk_ts: int) -> None:
301 """Feed the anchoring first chunk so the bridge is aligned and streaming."""
302 with patch.object(bridge, "_start_protocol_from_chunk", MagicMock()):
303 bridge._on_audio_chunk(_pcm_chunk(first_chunk_ts))
304 # The mocked task reports done() truthy by default, which the chunk handler
305 # reads as a failed protocol start; model a start still in flight instead.
306 cast("MagicMock", bridge._airplay_stream_start_task).done.return_value = False
307 _settle_anchor(bridge)
308
309
310def _expected_frames(bridge: SendspinAirPlayBridge, timeline_end_us: int) -> int:
311 """Frames the CLI stream must hold for its cursor to sit at a timeline instant."""
312 return round((timeline_end_us - bridge._drop_until_us) * BRIDGE_SAMPLE_RATE / 1_000_000)
313
314
315def test_timeline_gap_is_padded_with_silence() -> None:
316 """
317 A hole in the Sendspin timeline is filled so the device stays on the group clock.
318
319 Sendspin rebases the shared timeline forward when audio production stalls,
320 delivering no audio for the skipped span. The CLI plays its byte stream at a
321 fixed rate from an anchor that is never revised, so writing the next chunk
322 straight after the previous one would leave this device permanently ahead of
323 the group by the size of the hole.
324 """
325 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
326 first_ts = SENDSPIN_EPOCH_US + 250_000
327 _start_stream_at(bridge, first_ts)
328 _drain_queued_bytes(bridge)
329
330 gap_us = 415_711
331 next_ts = first_ts + 100_000 + gap_us
332 bridge._on_audio_chunk(_pcm_chunk(next_ts))
333
334 expected = _expected_frames(bridge, next_ts + 100_000)
335 assert bridge._queued_frames == expected
336 assert (
337 _drain_queued_bytes(bridge)
338 == (expected - _expected_frames(bridge, first_ts + 100_000))
339 * BRIDGE_CHANNELS
340 * BRIDGE_BYTES_PER_SAMPLE
341 )
342
343
344def test_overlapping_chunk_head_is_trimmed() -> None:
345 """A chunk reaching back behind the write cursor keeps only its unwritten tail."""
346 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
347 first_ts = SENDSPIN_EPOCH_US + 250_000
348 _start_stream_at(bridge, first_ts)
349 _drain_queued_bytes(bridge)
350
351 overlap_us = 40_000
352 next_ts = first_ts + 100_000 - overlap_us
353 bridge._on_audio_chunk(_pcm_chunk(next_ts))
354
355 assert bridge._queued_frames == _expected_frames(bridge, next_ts + 100_000)
356
357
358def test_chunk_entirely_behind_the_cursor_is_dropped() -> None:
359 """Audio already written is not queued a second time."""
360 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
361 first_ts = SENDSPIN_EPOCH_US + 250_000
362 _start_stream_at(bridge, first_ts)
363 _drain_queued_bytes(bridge)
364 cursor_frames = bridge._queued_frames
365
366 bridge._on_audio_chunk(_pcm_chunk(first_ts + 10_000, duration_us=50_000))
367
368 assert bridge._queued_frames == cursor_frames
369 assert _drain_queued_bytes(bridge) == 0
370
371
372@pytest.mark.parametrize("server_side", [False, True])
373def test_stream_start_resets_the_write_cursor(server_side: bool) -> None:
374 """
375 Both stream-start entry points rewind the cursor so the next chunk re-anchors byte 0.
376
377 A cursor carried over from the previous stream would place the first chunk of
378 the new one far behind the write position and get it trimmed away as already
379 written, and a settled-anchor flag carried over would let the new stream's
380 chunks be placed against the previous stream's anchor.
381 """
382 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
383 _start_stream_at(bridge, SENDSPIN_EPOCH_US + 250_000)
384 bridge._held_chunks.append(_pcm_chunk(SENDSPIN_EPOCH_US + 250_000))
385 assert bridge._queued_frames > 0
386
387 if server_side:
388 bridge._on_stream_start(MagicMock())
389 else:
390 bridge._on_bridge_stream_start()
391
392 assert bridge._queued_frames == 0
393 assert bridge._drop_until_us == 0
394 assert bridge._anchor_settled is False
395 assert not bridge._held_chunks
396
397
398def test_contiguous_chunks_are_written_untouched() -> None:
399 """Normal playback queues exactly its own audio -- no padding, no trimming."""
400 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
401 first_ts = SENDSPIN_EPOCH_US + 250_000
402 _start_stream_at(bridge, first_ts)
403 _drain_queued_bytes(bridge)
404
405 for index in range(1, 20):
406 bridge._on_audio_chunk(_pcm_chunk(first_ts + index * 100_000))
407
408 assert bridge._queued_frames == _expected_frames(bridge, first_ts + 20 * 100_000)
409 assert _drain_queued_bytes(bridge) == 19 * 100_000 * BRIDGE_BYTES_PER_SECOND // 1_000_000
410
411
412# --- Write pacing: bound the device buffer so a catch-up backlog is not dumped ---
413
414
415def test_device_buffer_ahead_seconds_tracks_write_cursor() -> None:
416 """The buffered-ahead measure follows byte 0 = start anchor, +1 s per second written."""
417 start_unix_ms = 1_784_000_000_000
418 now = start_unix_ms / 1000
419
420 assert device_buffer_ahead_seconds(start_unix_ms, 0, BRIDGE_BYTES_PER_SECOND, now) == 0.0
421 one_second = BRIDGE_BYTES_PER_SECOND
422 ahead = device_buffer_ahead_seconds(start_unix_ms, one_second, BRIDGE_BYTES_PER_SECOND, now)
423 assert abs(ahead - 1.0) < 1e-9
424
425
426def test_late_join_backlog_trips_pacing_bound_but_steady_feed_does_not() -> None:
427 """A ~27 s catch-up backlog exceeds the bound; a few seconds of steady audio stays under it."""
428 start_unix_ms = 1_784_000_000_000
429 now = start_unix_ms / 1000
430
431 backlog_ahead = device_buffer_ahead_seconds(
432 start_unix_ms, 27 * BRIDGE_BYTES_PER_SECOND, BRIDGE_BYTES_PER_SECOND, now
433 )
434 assert backlog_ahead > MAX_DEVICE_BUFFER_SECONDS
435
436 steady_ahead = device_buffer_ahead_seconds(
437 start_unix_ms, 3 * BRIDGE_BYTES_PER_SECOND, BRIDGE_BYTES_PER_SECOND, now
438 )
439 assert steady_ahead < MAX_DEVICE_BUFFER_SECONDS
440
441
442async def test_failed_cli_write_does_not_advance_pacing_cursor() -> None:
443 """A dropped write cannot move the pacing cursor past audio the CLI never received."""
444 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
445 stream = MagicMock()
446 stream.write_audio = AsyncMock(side_effect=[OSError("write failed"), None])
447 stream.write_audio_eof = AsyncMock()
448 bridge._airplay_stream = stream
449 bridge._airplay_stream_ready.set()
450 bridge._start_unix_ms = int(UNIX_NOW_S * 1000)
451 bridge._write_queue.put_nowait(b"first")
452 bridge._write_queue.put_nowait(b"second")
453 bridge._write_queue.put_nowait(None)
454
455 with patch(
456 "music_assistant.providers.airplay.sendspin_bridge.device_buffer_ahead_seconds",
457 return_value=0.0,
458 ) as buffer_ahead:
459 await bridge._cli_writer()
460
461 assert [call.args[1] for call in buffer_ahead.call_args_list] == [0, 0]
462 assert stream.write_audio.await_count == 2
463
464
465# --- Commanded cold start and warm handover ------------------------------------
466
467
468def _make_anchor_stream(
469 *,
470 ready_at_unix_ms: int | None = None,
471 ack: int | None = None,
472 warm_lead_ms: int = 0,
473 flushed_head_unix_ms: int = 0,
474 audio_pending_ms: int = 0,
475) -> MagicMock:
476 """
477 Build an AirPlayStream mock the anchor math can run against.
478
479 The bridge reads these off the stream and does arithmetic on them, so they
480 must be real numbers: the anchor compares ``warm_lead_ms`` /
481 ``flushed_head_unix_ms`` / ``audio_pending_ms`` with ``> 0`` and the shift
482 fold subtracts ``cumulative_shift_seconds``, none of which a bare MagicMock
483 can answer (every one of them is truthy).
484
485 :param ack: Instant the binary acks the START at. None acks the commanded
486 instant, as a feasible one is.
487 """
488
489 async def _ack_start(start_unix_ms: int = 0, **_kwargs: object) -> int:
490 return start_unix_ms if ack is None else ack
491
492 stream = MagicMock()
493 stream.cumulative_shift_seconds = 0.0
494 stream.connect = AsyncMock()
495 stream.wait_for_connection = AsyncMock()
496 stream.stop = AsyncMock()
497 stream.flush = AsyncMock(return_value=True)
498 stream.wait_clock_ready = AsyncMock(return_value=(ClockReadiness.PROJECTED, ready_at_unix_ms))
499 stream.start = AsyncMock(side_effect=_ack_start)
500 stream.warm_lead_ms = warm_lead_ms
501 stream.flushed_head_unix_ms = flushed_head_unix_ms
502 stream.audio_pending_ms = audio_pending_ms
503 return stream
504
505
506async def test_cold_start_connects_then_anchors_first_start() -> None:
507 """A fresh bridge stream anchors its first START only after the CLI connects."""
508 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
509 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
510 bridge._airplay_stream_start_task = asyncio.current_task()
511 commanded = UNIX_NOW_MS + COLD_LEAD_MS
512 stream = _make_anchor_stream(ack=commanded)
513
514 with (
515 patch(
516 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
517 return_value=stream,
518 ),
519 patch(
520 "music_assistant.providers.airplay.sendspin_bridge.time.time",
521 return_value=UNIX_NOW_S,
522 ),
523 ):
524 await bridge._start_protocol_from_chunk()
525
526 stream.connect.assert_awaited_once_with(False)
527 stream.wait_for_connection.assert_awaited_once_with()
528 stream.start.assert_awaited_once_with(commanded, join=True)
529 assert bridge._airplay_stream is stream
530 assert bridge.airplay_player.stream is stream
531 assert bridge._started is True
532 assert bridge._airplay_stream_ready.is_set()
533
534
535async def test_a_fresh_process_releases_a_foreign_mute_latch_before_it_connects() -> None:
536 """
537 A cold start releases a foreign mute latch before it connects.
538
539 Connecting is what carries that state to the device, so releasing the latch
540 after it would not be heard until the next command.
541 """
542 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
543 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
544 bridge._airplay_stream_start_task = asyncio.current_task()
545 stream = _make_anchor_stream()
546 order: list[str] = []
547 cast("MagicMock", bridge.airplay_player).release_foreign_mute_latch = MagicMock(
548 side_effect=lambda: order.append("release_mute")
549 )
550 stream.connect = AsyncMock(side_effect=lambda *_a, **_kw: order.append("connect"))
551
552 with (
553 patch(
554 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
555 return_value=stream,
556 ),
557 patch(
558 "music_assistant.providers.airplay.sendspin_bridge.time.time",
559 return_value=UNIX_NOW_S,
560 ),
561 ):
562 await bridge._start_protocol_from_chunk()
563
564 assert order == ["release_mute", "connect"]
565
566
567async def test_a_kept_process_keeps_a_mute_latch_owned_by_another_control() -> None:
568 """
569 A warm handover does not release a foreign mute latch.
570
571 Only a connect re-sends VOLUME=, so releasing the latch over a kept process
572 would clear it while the device is still muted on its own end, leaving the
573 stream silent with nothing to say so.
574 """
575 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
576 kept_stream = _make_anchor_stream()
577 bridge._airplay_stream = kept_stream
578 bridge.airplay_player.stream = kept_stream
579 bridge._started = True
580
581 # the whole warm restart, from the Sendspin stream-start callback to the handover
582 bridge._on_bridge_stream_start()
583 assert bridge._airplay_stream is kept_stream
584 bridge._drop_until_us = SENDSPIN_EPOCH_US
585 bridge._airplay_stream_start_task = asyncio.current_task()
586
587 with patch(
588 "music_assistant.providers.airplay.sendspin_bridge.time.time",
589 return_value=UNIX_NOW_S,
590 ):
591 await bridge._start_protocol_from_chunk()
592
593 kept_stream.flush.assert_awaited_once_with()
594 cast("MagicMock", bridge.airplay_player).release_foreign_mute_latch.assert_not_called()
595
596
597async def test_a_failed_warm_handover_releases_the_latch_before_its_cold_retry() -> None:
598 """
599 The cold retry after a failed warm handover still releases a foreign mute latch.
600
601 That retry spawns a fresh process, which is sent whatever volume and mute it
602 finds on connect, so a mute latch the parent no longer owns would start it silent.
603 """
604 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
605 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
606 bridge._airplay_stream_start_task = asyncio.current_task()
607 kept_stream = _make_anchor_stream()
608 kept_stream.flush = AsyncMock(return_value=False)
609 bridge._airplay_stream = kept_stream
610 cold_stream = _make_anchor_stream()
611
612 with (
613 patch(
614 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
615 return_value=cold_stream,
616 ),
617 patch(
618 "music_assistant.providers.airplay.sendspin_bridge.time.time",
619 return_value=UNIX_NOW_S,
620 ),
621 ):
622 await bridge._start_protocol_from_chunk()
623
624 cold_stream.connect.assert_awaited_once_with(False)
625 cast("MagicMock", bridge.airplay_player).release_foreign_mute_latch.assert_called_once_with()
626
627
628async def test_a_superseded_cold_start_never_reaches_the_receiver() -> None:
629 """
630 A cold start that already lost the race bails out before it spawns anything.
631
632 Connecting first would pay a full process spawn and session setup only to
633 kill it again, put a second session on a receiver the newer start is about
634 to claim, and overwrite the shared-clock decision of the process that start
635 is really running.
636 """
637 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
638 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
639 # the decision the newer start recorded for the process it is spawning
640 bridge._use_shared_ptp = True
641 # a different task owns the bridge: this cold start is stale
642 bridge._airplay_stream_start_task = MagicMock()
643 stream = _make_anchor_stream()
644
645 with (
646 patch(
647 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
648 return_value=stream,
649 ),
650 patch(
651 "music_assistant.providers.airplay.sendspin_bridge.time.time",
652 return_value=UNIX_NOW_S,
653 ),
654 ):
655 await bridge._start_protocol_from_chunk()
656
657 stream.connect.assert_not_awaited()
658 stream.stop.assert_not_awaited()
659 assert bridge._use_shared_ptp is True
660
661
662async def test_a_superseded_start_leaves_the_kept_stream_untouched() -> None:
663 """
664 A start that lost the race never flushes the stream the newer one kept.
665
666 Arming the bridge keeps a warm-eligible stream alive, so the stale and the
667 newer start find the same instance; flushing it here would cut into the
668 audio the newer start is anchoring on it.
669 """
670 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
671 kept_stream = _make_anchor_stream()
672 bridge._airplay_stream = kept_stream
673 # a different task owns the bridge: this start is stale
674 bridge._airplay_stream_start_task = MagicMock()
675
676 await bridge._start_protocol_from_chunk()
677
678 kept_stream.flush.assert_not_awaited()
679 kept_stream.stop.assert_not_awaited()
680 assert bridge._airplay_stream is kept_stream
681
682
683async def test_a_start_superseded_during_the_warm_fallback_spawns_nothing() -> None:
684 """
685 Losing the race while releasing the kept stream still stops short of the receiver.
686
687 A failed warm handover tears the kept stream down before it falls back to a
688 cold start, and that teardown is long enough for a newer start to claim the
689 bridge in the meantime.
690 """
691 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
692 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
693 bridge._airplay_stream_start_task = asyncio.current_task()
694 kept_stream = _make_anchor_stream()
695 kept_stream.flush = AsyncMock(return_value=False)
696 bridge._airplay_stream = kept_stream
697 cold_stream = _make_anchor_stream()
698
699 async def stop(**_kwargs: object) -> None:
700 # a newer stream start claimed the bridge while the kept stream went down
701 bridge._airplay_stream_start_task = MagicMock()
702
703 kept_stream.stop = AsyncMock(side_effect=stop)
704
705 with patch(
706 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
707 return_value=cold_stream,
708 ):
709 await bridge._start_protocol_from_chunk()
710
711 cold_stream.connect.assert_not_awaited()
712 cold_stream.stop.assert_not_awaited()
713
714
715async def test_cold_start_superseded_while_connecting_stops_its_transport() -> None:
716 """A cold stream superseded while its process comes up is torn down again."""
717 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
718 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
719 bridge._airplay_stream_start_task = asyncio.current_task()
720 stream = _make_anchor_stream()
721
722 async def wait_for_connection() -> None:
723 # a newer stream start claimed the bridge while the process came up
724 bridge._airplay_stream_start_task = MagicMock()
725
726 stream.wait_for_connection = AsyncMock(side_effect=wait_for_connection)
727
728 assert await bridge._start_cold_stream(stream) is False
729
730 stream.start.assert_not_awaited()
731 stream.stop.assert_awaited_once_with(force=True)
732
733
734async def test_cold_start_superseded_during_the_anchor_stops_its_transport() -> None:
735 """
736 A cold stream superseded while the binary holds its ack is torn down.
737
738 The anchor publishes the stream before commanding START, so a supersession
739 inside it would otherwise leave a live cliairplay attached to the receiver
740 with nobody owning it.
741 """
742 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
743 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
744 bridge._airplay_stream_start_task = asyncio.current_task()
745 stream = _make_anchor_stream()
746
747 async def start(start_unix_ms: int, **_kwargs: object) -> int:
748 # A newer stream start claimed the bridge while the binary held its ack.
749 bridge._airplay_stream_start_task = MagicMock()
750 return start_unix_ms
751
752 stream.start = AsyncMock(side_effect=start)
753
754 with patch(
755 "music_assistant.providers.airplay.sendspin_bridge.time.time",
756 return_value=UNIX_NOW_S,
757 ):
758 assert await bridge._start_cold_stream(stream) is False
759
760 stream.stop.assert_awaited_once_with(force=True)
761 assert bridge._airplay_stream is None
762 assert bridge.airplay_player.stream is None
763
764
765async def test_superseded_cold_stream_teardown_spares_the_newer_owner() -> None:
766 """A newer start's published stream survives the stale cold stream's teardown."""
767 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
768 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
769 bridge._airplay_stream_start_task = asyncio.current_task()
770 stream = _make_anchor_stream()
771 newer_stream = _make_anchor_stream()
772
773 async def start(start_unix_ms: int, **_kwargs: object) -> int:
774 bridge._airplay_stream_start_task = MagicMock()
775 bridge._airplay_stream = newer_stream
776 bridge.airplay_player.stream = newer_stream
777 return start_unix_ms
778
779 stream.start = AsyncMock(side_effect=start)
780
781 with patch(
782 "music_assistant.providers.airplay.sendspin_bridge.time.time",
783 return_value=UNIX_NOW_S,
784 ):
785 assert await bridge._start_cold_stream(stream) is False
786
787 stream.stop.assert_awaited_once_with(force=True)
788 newer_stream.stop.assert_not_awaited()
789 assert bridge._airplay_stream is newer_stream
790 assert bridge.airplay_player.stream is newer_stream
791
792
793# --- Anchoring: command an instant the device can hit, then honour the ack ---
794
795
796def _prepare_anchor(
797 bridge: SendspinAirPlayBridge, stream: MagicMock, first_chunk_lead_ms: int
798) -> None:
799 """Wire a bridge so ``_anchor_stream`` can be awaited directly on ``stream``."""
800 bridge._drop_until_us = bridge.sendspin_server.clock.now_us() + first_chunk_lead_ms * 1_000
801 bridge._airplay_stream = stream
802 bridge._airplay_stream_start_task = asyncio.current_task()
803
804
805async def _anchor(bridge: SendspinAirPlayBridge, stream: MagicMock, *, warm: bool = False) -> bool:
806 """Run ``_anchor_stream`` with the unix clock pinned to UNIX_NOW_S."""
807 with patch(
808 "music_assistant.providers.airplay.sendspin_bridge.time.time",
809 return_value=UNIX_NOW_S,
810 ):
811 return await bridge._anchor_stream(stream, warm=warm)
812
813
814def _commanded_instant(stream: MagicMock) -> int:
815 """Return the instant the START command carried."""
816 return int(stream.start.await_args.args[0])
817
818
819async def test_anchor_floors_at_the_join_headroom() -> None:
820 """
821 A Sendspin lead shorter than the join floor is raised to it.
822
823 AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS is the least the anchor may sit ahead of
824 now, so a shorter lead is floored rather than honoured.
825 """
826 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
827 stream = _make_anchor_stream()
828 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
829
830 assert await _anchor(bridge, stream) is True
831
832 assert _commanded_instant(stream) == UNIX_NOW_MS + AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS
833
834
835async def test_anchor_reports_leftover_audio_pending_on_stdin() -> None:
836 """
837 Audio still pending when the anchor is commanded is named as the offset it causes.
838
839 The writer is gated until the anchor settles, so anything pending was left
840 behind by an earlier stream and the START anchors it as this one's first
841 sample. The cursor only counts what the bridge queued itself, so no later
842 realignment can see the resulting offset, which leaves this warning as the
843 one place it surfaces.
844 """
845 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
846 stream = _make_anchor_stream(audio_pending_ms=92)
847 _prepare_anchor(bridge, stream, first_chunk_lead_ms=WARM_LEAD_MS)
848
849 with patch.object(bridge.logger, "warning") as warning:
850 assert await _anchor(bridge, stream, warm=True) is True
851
852 pending = [call for call in warning.call_args_list if "pending" in call.args[0]]
853 assert len(pending) == 1
854 assert pending[0].args[2] == 92
855
856
857async def test_anchor_stays_quiet_when_stdin_was_left_empty() -> None:
858 """An anchor commanded against empty stdin reports nothing."""
859 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
860 stream = _make_anchor_stream()
861 _prepare_anchor(bridge, stream, first_chunk_lead_ms=WARM_LEAD_MS)
862
863 with patch.object(bridge.logger, "warning") as warning:
864 assert await _anchor(bridge, stream, warm=True) is True
865
866 assert not [call for call in warning.call_args_list if "pending" in call.args[0]]
867
868
869async def test_anchor_follows_the_clock_ready_projection() -> None:
870 """A receiver that projects a later readiness pushes the anchor out past it."""
871 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
872 ready_at = UNIX_NOW_MS + 3200
873 stream = _make_anchor_stream(ready_at_unix_ms=ready_at)
874 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
875
876 assert await _anchor(bridge, stream) is True
877
878 assert _commanded_instant(stream) == ready_at + AIRPLAY_CLOCK_READY_LEAD_MS
879
880
881async def test_anchor_still_starts_a_receiver_whose_clock_stalled() -> None:
882 """
883 A stalled receiver is anchored anyway, unlike a late joiner, and warned about.
884
885 A joiner is dropped because the session plays on without it, while here
886 dropping would stop the speaker, and the binary's stall report is a
887 diagnosis a receiver can still come good from.
888 """
889 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
890 stream = _make_anchor_stream()
891 stream.wait_clock_ready = AsyncMock(return_value=(ClockReadiness.STALLED, 0))
892 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
893
894 with patch.object(bridge.logger, "warning") as warning:
895 assert await _anchor(bridge, stream) is True
896
897 assert _commanded_instant(stream) == UNIX_NOW_MS + AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS
898 assert bridge._started is True
899 assert len([call for call in warning.call_args_list if "PTP clock" in call.args[0]]) == 1
900
901
902async def test_anchor_never_precedes_content_already_scheduled() -> None:
903 """
904 A buffered source keeps its intro: the anchor lands on the first chunk we hold.
905
906 Sendspin can schedule the first sample much further out than the device
907 needs. Anchoring on the floor instead would place byte 0 in the middle of
908 the audio already delivered and throw away everything before it.
909 """
910 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
911 first_chunk_us = SENDSPIN_EPOCH_US + 6_000_000
912 stream = _make_anchor_stream(ack=UNIX_NOW_MS + 6000)
913 _prepare_anchor(bridge, stream, first_chunk_lead_ms=6000)
914 bridge._held_chunks.append(_pcm_chunk(first_chunk_us))
915
916 assert await _anchor(bridge, stream) is True
917
918 assert _commanded_instant(stream) == UNIX_NOW_MS + 6000
919 # Nothing skipped, and the held opening reached the writer intact.
920 assert bridge._drop_until_us == first_chunk_us
921 assert _drain_queued_bytes(bridge) == 100_000 * BRIDGE_BYTES_PER_SECOND // 1_000_000
922
923
924async def test_writer_stays_blocked_until_the_start_is_acked() -> None:
925 """
926 The writer is released only once the content is mapped onto the acked instant.
927
928 Feeding the CLI before the ack would place bytes against an anchor the
929 binary has not confirmed and may still correct forward.
930 """
931 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
932 stream = _make_anchor_stream()
933 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
934 ack_gate = asyncio.Event()
935 start_called = asyncio.Event()
936
937 async def start(start_unix_ms: int, **_kwargs: object) -> int:
938 start_called.set()
939 await ack_gate.wait()
940 return start_unix_ms
941
942 stream.start = AsyncMock(side_effect=start)
943 anchor_task = asyncio.create_task(_anchor(bridge, stream))
944 bridge._airplay_stream_start_task = cast("asyncio.Task[None]", anchor_task)
945 await start_called.wait()
946
947 assert not bridge._airplay_stream_ready.is_set()
948 ack_gate.set()
949 assert await anchor_task is True
950 assert bridge._airplay_stream_ready.is_set()
951
952
953async def test_chunks_arriving_before_the_ack_are_held_and_replayed_in_order() -> None:
954 """Audio delivered while the anchor is outstanding is queued once, in order."""
955 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
956 stream = _make_anchor_stream()
957 _prepare_anchor(bridge, stream, first_chunk_lead_ms=2500)
958 first_chunk_us = SENDSPIN_EPOCH_US + 2_500_000
959 ack_gate = asyncio.Event()
960 start_called = asyncio.Event()
961
962 async def start(start_unix_ms: int, **_kwargs: object) -> int:
963 start_called.set()
964 await ack_gate.wait()
965 return start_unix_ms
966
967 stream.start = AsyncMock(side_effect=start)
968 anchor_task = asyncio.create_task(_anchor(bridge, stream))
969 bridge._airplay_stream_start_task = cast("asyncio.Task[None]", anchor_task)
970 await start_called.wait()
971
972 for index in range(4):
973 bridge._on_audio_chunk(_pcm_chunk(first_chunk_us + index * 100_000))
974 assert len(bridge._held_chunks) == 4
975 assert bridge._write_queue.empty()
976
977 ack_gate.set()
978 assert await anchor_task is True
979
980 # Four contiguous 100 ms chunks, replayed without padding or trimming.
981 assert bridge._queued_frames == _expected_frames(bridge, first_chunk_us + 400_000)
982 assert _drain_queued_bytes(bridge) == 4 * 100_000 * BRIDGE_BYTES_PER_SECOND // 1_000_000
983
984
985def test_held_backlog_is_capped_and_drops_the_oldest() -> None:
986 """An anchor that never settles cannot grow the hold without bound."""
987 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
988 chunk_us = 100_000
989 over_cap = MAX_HELD_AUDIO_US // chunk_us + 50
990
991 for index in range(over_cap):
992 bridge._hold_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + index * chunk_us))
993
994 assert sum(chunk.duration_us for chunk in bridge._held_chunks) <= MAX_HELD_AUDIO_US
995 # The running total the cap is measured against tracks the deque exactly.
996 assert bridge._held_us == sum(chunk.duration_us for chunk in bridge._held_chunks)
997 # The oldest went, the newest stayed.
998 assert bridge._held_chunks[0].timestamp_us > SENDSPIN_EPOCH_US
999 assert bridge._held_chunks[-1].timestamp_us == SENDSPIN_EPOCH_US + (over_cap - 1) * chunk_us
1000
1001
1002def test_held_backlog_is_capped_by_chunk_count() -> None:
1003 """
1004 A run of zero-duration chunks is bounded by the count cap.
1005
1006 Such chunks carry no duration at all, so the µs cap can never trip on them
1007 however many arrive.
1008 """
1009 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1010
1011 for index in range(MAX_HELD_CHUNKS + 50):
1012 bridge._hold_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + index, duration_us=0))
1013
1014 assert len(bridge._held_chunks) == MAX_HELD_CHUNKS
1015 assert bridge._held_us == 0
1016 assert bridge._held_chunks[-1].timestamp_us == SENDSPIN_EPOCH_US + MAX_HELD_CHUNKS + 49
1017
1018
1019async def test_corrected_ack_rebases_the_content_onto_the_acked_instant() -> None:
1020 """An instant the binary moved forward re-bases the anchor, cursor and pacing base."""
1021 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1022 correction_ms = 3000
1023 acked = UNIX_NOW_MS + AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS + correction_ms
1024 stream = _make_anchor_stream(ack=acked)
1025 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
1026 bridge._queued_frames = 12_345
1027
1028 assert await _anchor(bridge, stream) is True
1029
1030 assert _commanded_instant(stream) == UNIX_NOW_MS + AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS
1031 assert bridge._drop_until_us == SENDSPIN_EPOCH_US + (acked - UNIX_NOW_MS) * 1_000
1032 assert bridge._queued_frames == 0
1033 assert bridge._start_unix_ms == acked
1034
1035
1036async def test_content_before_the_acked_instant_is_dropped_and_trimmed() -> None:
1037 """Held audio the acked anchor moved past is dropped, the straddling chunk trimmed."""
1038 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1039 acked = UNIX_NOW_MS + 2600
1040 stream = _make_anchor_stream(ack=acked)
1041 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
1042 anchor_us = SENDSPIN_EPOCH_US + 2_600_000
1043 for offset in (-200_000, -50_000, 50_000):
1044 bridge._held_chunks.append(_pcm_chunk(anchor_us + offset))
1045
1046 assert await _anchor(bridge, stream) is True
1047
1048 # Fully-behind chunk gone, the straddling one keeps its 50 ms tail, and the
1049 # last chunk continues contiguously: 150 ms of audio in total.
1050 assert bridge._queued_frames == _expected_frames(bridge, anchor_us + 150_000)
1051 assert _drain_queued_bytes(bridge) == bridge._queued_frames * BRIDGE_BYTES_PER_FRAME
1052
1053
1054async def test_sync_adjust_shifts_the_command_but_not_the_content_mapping() -> None:
1055 """The device's own offset rides on the command; the group timeline stays untouched."""
1056 adjust_ms = 300
1057 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US, sync_adjust=adjust_ms)
1058 anchor_ms = UNIX_NOW_MS + AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS
1059 stream = _make_anchor_stream(ack=anchor_ms + adjust_ms)
1060 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
1061
1062 assert await _anchor(bridge, stream) is True
1063
1064 assert _commanded_instant(stream) == anchor_ms + adjust_ms
1065 # Pacing tracks the real wall-clock instant of byte 0 (adjust included)...
1066 assert bridge._start_unix_ms == anchor_ms + adjust_ms
1067 # ...while the content is placed on the group timeline, without it.
1068 assert bridge._drop_until_us == SENDSPIN_EPOCH_US + AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS * 1_000
1069
1070
1071async def test_ack_earlier_than_commanded_is_trusted_verbatim() -> None:
1072 """
1073 An acked instant before the commanded one is used as-is, never clamped up.
1074
1075 Clamping to the commanded instant would map the content onto a moment the
1076 binary is not rendering at and put the device ahead of the rest of the group.
1077 """
1078 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1079 acked = UNIX_NOW_MS + AIRPLAY_LATE_JOIN_MIN_HEADROOM_MS - 400
1080 stream = _make_anchor_stream(ack=acked)
1081 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
1082
1083 assert await _anchor(bridge, stream) is True
1084
1085 assert _commanded_instant(stream) > acked
1086 assert bridge._start_unix_ms == acked
1087 assert bridge._drop_until_us == SENDSPIN_EPOCH_US + (acked - UNIX_NOW_MS) * 1_000
1088
1089
1090def test_first_chunk_after_an_anchor_is_not_reported_as_drift() -> None:
1091 """The gap between a fresh anchor and its first chunk is placement, not drift."""
1092 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1093 bridge._drop_until_us = SENDSPIN_EPOCH_US
1094 bridge._queued_frames = 0
1095
1096 bridge._align_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 900_000))
1097
1098 cast("MagicMock", bridge.logger).warning.assert_not_called()
1099
1100 # A cursor that already advanced can drift, and that is still reported.
1101 bridge._align_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 2_000_000))
1102
1103 cast("MagicMock", bridge.logger).warning.assert_called_once()
1104
1105
1106def test_a_long_timeline_gap_is_padded_from_one_shared_block() -> None:
1107 """
1108 A long hole is queued as repeats of one silence block, not as a single buffer.
1109
1110 Sendspin rebases the shared timeline forward when audio production stalls,
1111 which can open a hole of tens of seconds. Building that as one bytes object
1112 puts megabytes on the event loop in a single synchronous allocation.
1113 """
1114 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1115 first_ts = SENDSPIN_EPOCH_US + 250_000
1116 _start_stream_at(bridge, first_ts)
1117 _drain_queued_bytes(bridge)
1118
1119 gap_us = 60_000_000
1120 next_ts = first_ts + 100_000 + gap_us
1121 bridge._on_audio_chunk(_pcm_chunk(next_ts))
1122
1123 queued: list[bytes] = []
1124 while not bridge._write_queue.empty():
1125 block = bridge._write_queue.get_nowait()
1126 assert block is not None
1127 queued.append(block)
1128
1129 pad, data = queued[:-1], queued[-1]
1130 gap_frames = round(gap_us * BRIDGE_SAMPLE_RATE / 1_000_000)
1131 assert len(pad) == gap_frames // PAD_BLOCK_FRAMES
1132 # Identity, not equality: the whole hole costs one allocation.
1133 assert all(block is SILENCE_BLOCK for block in pad)
1134 # The hole is still filled exactly, so the device stays on the group's clock.
1135 assert sum(len(block) for block in pad) == gap_frames * BRIDGE_BYTES_PER_FRAME
1136 assert len(data) == 100_000 * BRIDGE_BYTES_PER_SECOND // 1_000_000
1137
1138
1139# --- Playout shift: the binary's own mid-stream re-anchors move the mapping ---
1140
1141
1142def _shifted_bridge(shift_seconds: float) -> tuple[SendspinAirPlayBridge, MagicMock, int]:
1143 """
1144 Start a streaming bridge whose CLI reports a mid-stream playout shift.
1145
1146 :param shift_seconds: Cumulative shift the binary reports since its START.
1147 :return: The bridge, its stream mock and the first chunk's timestamp.
1148 """
1149 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1150 first_ts = SENDSPIN_EPOCH_US + 250_000
1151 _start_stream_at(bridge, first_ts)
1152 bridge._start_unix_ms = UNIX_NOW_MS
1153 _drain_queued_bytes(bridge)
1154 stream = MagicMock()
1155 # A real float: the fold subtracts it from the applied baseline, which a bare
1156 # MagicMock cannot answer.
1157 stream.cumulative_shift_seconds = shift_seconds
1158 bridge._airplay_stream = stream
1159 return bridge, stream, first_ts
1160
1161
1162def test_reported_reanchor_moves_the_anchor_and_the_pacing_start() -> None:
1163 """
1164 A starvation re-anchor makes every byte audible later, so the anchor follows it.
1165
1166 cliairplay shifts its playout forward when it runs out of PCM on stdin.
1167 Leaving the bridge's mapping where it was would place all following content
1168 ahead of where the device actually plays it, for the rest of the stream.
1169 """
1170 shift_s = 1.5
1171 bridge, _, first_ts = _shifted_bridge(shift_s)
1172 anchor_before = bridge._drop_until_us
1173
1174 bridge._on_audio_chunk(_pcm_chunk(first_ts + 100_000))
1175
1176 assert bridge._drop_until_us == anchor_before + round(shift_s * 1_000_000)
1177 # The write pacing measures from the same instant, so it moves with it.
1178 assert bridge._start_unix_ms == UNIX_NOW_MS + round(shift_s * 1000)
1179
1180
1181def test_reported_reanchor_skips_content_until_the_device_catches_up() -> None:
1182 """The device plays the shift late, so exactly that much content is dropped."""
1183 shift_s = 0.5
1184 bridge, _, first_ts = _shifted_bridge(shift_s)
1185
1186 fed_us = 1_000_000
1187 for index in range(fed_us // 100_000):
1188 bridge._on_audio_chunk(_pcm_chunk(first_ts + 100_000 * (index + 1)))
1189
1190 assert bridge._queued_frames == _expected_frames(bridge, first_ts + 100_000 + fed_us)
1191 assert _drain_queued_bytes(bridge) == round(
1192 (fed_us / 1_000_000 - shift_s) * BRIDGE_BYTES_PER_SECOND
1193 )
1194
1195
1196def test_absorbing_a_reanchor_is_not_reported_as_timeline_drift() -> None:
1197 """The trim that works a folded shift off is a correction, not Sendspin drift."""
1198 bridge, _, first_ts = _shifted_bridge(0.5)
1199 logger = cast("MagicMock", bridge.logger)
1200
1201 # Straddles the shifted anchor, so it is trimmed rather than dropped outright.
1202 bridge._on_audio_chunk(_pcm_chunk(first_ts + 550_000))
1203
1204 assert bridge._queued_frames == _expected_frames(bridge, first_ts + 650_000)
1205 # Only the re-anchor itself is reported; the trim it caused stays quiet.
1206 logger.warning.assert_called_once()
1207 logger.warning.reset_mock()
1208
1209 # Back on the timeline: the shift is worked off and nothing is realigned.
1210 bridge._on_audio_chunk(_pcm_chunk(first_ts + 650_000))
1211 assert not bridge._absorbing_shift
1212 logger.warning.assert_not_called()
1213
1214 # A real discontinuity is reported again.
1215 bridge._on_audio_chunk(_pcm_chunk(first_ts + 1_500_000))
1216 logger.warning.assert_called_once()
1217
1218
1219def test_a_cursor_off_the_frame_grid_still_absorbs_quietly() -> None:
1220 """
1221 The trim stays quiet however the cursor happens to sit when the shift lands.
1222
1223 ``_align_chunk`` re-targets each chunk against the anchor independently, so
1224 the cursor routinely rests a frame either side of the timeline. The trim a
1225 fold asks for is a correction at any of those offsets, never Sendspin drift.
1226 """
1227 bridge, stream, first_ts = _shifted_bridge(0.0)
1228 logger = cast("MagicMock", bridge.logger)
1229
1230 # Play on a while first, so the cursor sits well past the anchor the way it
1231 # does mid-stream when a starvation hits.
1232 for index in range(1, 11):
1233 bridge._on_audio_chunk(_pcm_chunk(first_ts + 100_000 * index))
1234 logger.warning.assert_not_called()
1235
1236 # A frame past the timeline: the trim the fold asks for then runs one frame
1237 # deeper than the shift itself.
1238 bridge._queued_frames += 1
1239 stream.cumulative_shift_seconds = 0.5
1240 bridge._on_audio_chunk(_pcm_chunk(first_ts + 1_100_000))
1241
1242 # Only the re-anchor itself is reported.
1243 logger.warning.assert_called_once()
1244
1245
1246def test_an_absorption_spanning_chunks_reports_only_the_reanchor() -> None:
1247 """
1248 A shift takes several chunks to trim off, and stays one report throughout.
1249
1250 The binary keeps reporting the same running total while the trim works, so
1251 every chunk until the cursor is back on the timeline realigns by design.
1252 """
1253 bridge, stream, first_ts = _shifted_bridge(0.0)
1254 logger = cast("MagicMock", bridge.logger)
1255 for index in range(1, 11):
1256 bridge._on_audio_chunk(_pcm_chunk(first_ts + 100_000 * index))
1257 logger.warning.assert_not_called()
1258
1259 stream.cumulative_shift_seconds = 0.5
1260 # Five chunks of content are trimmed away before the cursor catches up.
1261 for index in range(11, 17):
1262 bridge._on_audio_chunk(_pcm_chunk(first_ts + 100_000 * index))
1263
1264 logger.warning.assert_called_once()
1265 assert not bridge._absorbing_shift
1266
1267
1268def test_an_unchanged_reanchor_total_is_folded_only_once() -> None:
1269 """The binary reports a running total, so only what is new moves the anchor."""
1270 bridge, _, first_ts = _shifted_bridge(0.5)
1271 bridge._on_audio_chunk(_pcm_chunk(first_ts + 600_000))
1272 anchor_after_fold = bridge._drop_until_us
1273 assert anchor_after_fold == first_ts + 500_000
1274
1275 bridge._on_audio_chunk(_pcm_chunk(first_ts + 700_000))
1276
1277 assert bridge._drop_until_us == anchor_after_fold
1278
1279
1280def test_a_reset_reanchor_total_rebaselines_without_moving_the_anchor() -> None:
1281 """
1282 A total that went backwards means a START already replaced the mapping.
1283
1284 The binary zeroes its running total on every START, so the drop is not the
1285 device un-shifting: the bridge takes the new baseline and leaves the anchor
1286 to the START that set it.
1287 """
1288 bridge, stream, first_ts = _shifted_bridge(0.5)
1289 bridge._on_audio_chunk(_pcm_chunk(first_ts + 600_000))
1290 anchor_after_fold = bridge._drop_until_us
1291 assert anchor_after_fold == first_ts + 500_000
1292
1293 stream.cumulative_shift_seconds = 0.0
1294 bridge._on_audio_chunk(_pcm_chunk(first_ts + 700_000))
1295
1296 assert bridge._drop_until_us == anchor_after_fold
1297 assert bridge._applied_shift_seconds == 0.0
1298
1299
1300def test_a_reset_reanchor_total_ends_the_absorption() -> None:
1301 """
1302 A total that went backwards leaves no correction outstanding to stay quiet for.
1303
1304 The START that zeroed the total replaced the mapping the trim was working
1305 against, so a realignment after it is Sendspin timeline drift again and
1306 worth reporting.
1307 """
1308 bridge, stream, first_ts = _shifted_bridge(0.5)
1309 logger = cast("MagicMock", bridge.logger)
1310 # Mid-absorption: the trim has not caught the cursor up to the timeline yet.
1311 bridge._on_audio_chunk(_pcm_chunk(first_ts + 550_000))
1312 logger.warning.reset_mock()
1313
1314 stream.cumulative_shift_seconds = 0.0
1315 bridge._on_audio_chunk(_pcm_chunk(first_ts + 600_000))
1316
1317 assert not bridge._absorbing_shift
1318 logger.warning.assert_called_once()
1319
1320
1321async def test_anchoring_clears_the_folded_shift_baseline() -> None:
1322 """A START re-anchors the binary from scratch, so the fold starts over with it."""
1323 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1324 stream = _make_anchor_stream(ack=UNIX_NOW_MS + COLD_LEAD_MS)
1325 _prepare_anchor(bridge, stream, COLD_LEAD_MS)
1326 bridge._applied_shift_seconds = 1.5
1327 bridge._absorbing_shift = True
1328
1329 assert await _anchor(bridge, stream)
1330
1331 assert bridge._applied_shift_seconds == 0.0
1332 assert not bridge._absorbing_shift
1333
1334
1335@pytest.mark.parametrize(
1336 ("warm_lead_ms", "flushed_head_offset_ms", "adjust_ms", "expected_anchor_offset_ms"),
1337 [
1338 (4000, 0, 0, 4000 + AIRPLAY_SPLICE_LEAD_MARGIN_MS),
1339 (4000, 0, -600, 4600 + AIRPLAY_SPLICE_LEAD_MARGIN_MS),
1340 (0, 5000, 0, 5000 + AIRPLAY_SPLICE_LEAD_MARGIN_MS),
1341 ],
1342)
1343async def test_warm_anchor_clears_the_receivers_queued_audio(
1344 warm_lead_ms: int,
1345 flushed_head_offset_ms: int,
1346 adjust_ms: int,
1347 expected_anchor_offset_ms: int,
1348) -> None:
1349 """
1350 A warm re-anchor lands beyond the audio the receiver still has queued.
1351
1352 A negative sync_adjust moves the commanded instant earlier and eats into that
1353 lead, so it is added back to the requirement.
1354 """
1355 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US, sync_adjust=adjust_ms)
1356 stream = _make_anchor_stream(
1357 warm_lead_ms=warm_lead_ms,
1358 flushed_head_unix_ms=UNIX_NOW_MS + flushed_head_offset_ms if flushed_head_offset_ms else 0,
1359 )
1360 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
1361
1362 assert await _anchor(bridge, stream, warm=True) is True
1363
1364 assert _commanded_instant(stream) == UNIX_NOW_MS + expected_anchor_offset_ms + adjust_ms
1365
1366
1367async def test_superseded_during_the_ack_mutates_nothing() -> None:
1368 """A newer stream start taking over while the ack is outstanding wins untouched."""
1369 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1370 stream = _make_anchor_stream()
1371 _prepare_anchor(bridge, stream, first_chunk_lead_ms=250)
1372 original_drop_until = bridge._drop_until_us
1373 bridge._held_chunks.append(_pcm_chunk(original_drop_until))
1374
1375 async def start(start_unix_ms: int, **_kwargs: object) -> int:
1376 # A newer stream start claimed the bridge while the binary held its ack.
1377 bridge._airplay_stream_start_task = MagicMock()
1378 return start_unix_ms
1379
1380 stream.start = AsyncMock(side_effect=start)
1381
1382 assert await _anchor(bridge, stream) is False
1383
1384 assert bridge._drop_until_us == original_drop_until
1385 assert bridge._start_unix_ms == 0
1386 assert bridge._started is False
1387 assert bridge._anchor_settled is False
1388 assert len(bridge._held_chunks) == 1
1389 assert not bridge._airplay_stream_ready.is_set()
1390
1391
1392def test_unix_to_sendspin_instant_round_trips() -> None:
1393 """The two clock-domain helpers are exact inverses of each other."""
1394 clock = ManualClock(now_us_value=SENDSPIN_EPOCH_US)
1395 for lead_ms in (-300, 0, WARM_LEAD_MS, COLD_LEAD_MS):
1396 audible_us = _audible_instant_us(clock, lead_ms)
1397 unix_ms = sendspin_audible_instant_to_unix_ms(audible_us, clock.now_us(), UNIX_NOW_S)
1398
1399 assert unix_ms == UNIX_NOW_MS + lead_ms
1400 assert (
1401 unix_ms_to_sendspin_audible_instant(unix_ms, clock.now_us(), UNIX_NOW_S) == audible_us
1402 )
1403
1404
1405# --- Warm handover: a kept stream survives a new stream start and rides flush-refill ---
1406
1407
1408def _make_kept_stream(
1409 *, running: bool = True, connected: bool = True, ended_cleanly: bool = False
1410) -> MagicMock:
1411 """
1412 Build a mock AirPlayStream reporting the given running/connected state.
1413
1414 :param running: Whether the cli process behind the stream is still alive.
1415 :param connected: Whether the device connection has been established.
1416 :param ended_cleanly: Whether the binary reported the end of the stream
1417 itself. A real bool: a bare MagicMock reads as a clean end, which the
1418 loss check treats as no loss at all.
1419 """
1420 stream = MagicMock()
1421 stream.running = running
1422 stream.connected = connected
1423 stream.ended_cleanly = ended_cleanly
1424 # A real float: the shift fold subtracts it from the applied baseline, which
1425 # a bare MagicMock cannot answer.
1426 stream.cumulative_shift_seconds = 0.0
1427 return stream
1428
1429
1430def test_on_bridge_stream_start_keeps_warm_eligible_stream() -> None:
1431 """A running, connected AirPlay 2 stream survives a new Sendspin stream start."""
1432 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1433 kept_stream = _make_kept_stream()
1434 bridge._airplay_stream = kept_stream
1435 bridge.airplay_player.stream = kept_stream
1436 bridge._started = True
1437
1438 bridge._on_bridge_stream_start()
1439
1440 assert bridge._airplay_stream is kept_stream
1441 assert bridge.airplay_player.stream is kept_stream
1442 assert bridge._stream_is_warm_eligible()
1443
1444
1445def test_on_bridge_stream_start_keeps_raop_stream() -> None:
1446 """A started legacy RAOP stream is eligible for warm Sendspin flush-refill."""
1447 bridge = _make_bridge(
1448 clock_now_us=SENDSPIN_EPOCH_US,
1449 protocol=StreamingProtocol.RAOP,
1450 )
1451 old_stream = _make_kept_stream()
1452 bridge._airplay_stream = old_stream
1453 bridge.airplay_player.stream = old_stream
1454 bridge._started = True
1455
1456 bridge._on_bridge_stream_start()
1457
1458 assert bridge._airplay_stream is old_stream
1459 assert bridge.airplay_player.stream is old_stream
1460
1461
1462def test_sendspin_callbacks_keep_raop_stream_until_warm_handover() -> None:
1463 """Both Sendspin start callbacks preserve a reusable legacy RAOP session."""
1464 bridge = _make_bridge(
1465 clock_now_us=SENDSPIN_EPOCH_US,
1466 protocol=StreamingProtocol.RAOP,
1467 )
1468 kept_stream = _make_kept_stream()
1469 bridge._airplay_stream = kept_stream
1470 bridge.airplay_player.stream = kept_stream
1471 bridge._started = True
1472
1473 bridge._on_stream_start(MagicMock())
1474 bridge._on_bridge_stream_start()
1475
1476 assert bridge._airplay_stream is kept_stream
1477 assert bridge.airplay_player.stream is kept_stream
1478
1479
1480def test_on_bridge_stream_start_replaces_uncommitted_stream() -> None:
1481 """A connected stream cannot be retained before its first START."""
1482 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1483 old_stream = _make_kept_stream()
1484 bridge._airplay_stream = old_stream
1485 bridge.airplay_player.stream = old_stream
1486
1487 bridge._on_bridge_stream_start()
1488
1489 assert bridge._airplay_stream is None
1490 assert bridge.airplay_player.stream is None # type: ignore[unreachable]
1491
1492
1493def test_on_stream_start_keeps_warm_eligible_stream() -> None:
1494 """The Sendspin-server-side stream-start callback also keeps a warm-eligible stream."""
1495 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1496 kept_stream = _make_kept_stream()
1497 bridge._airplay_stream = kept_stream
1498 bridge.airplay_player.stream = kept_stream
1499 bridge._started = True
1500
1501 bridge._on_stream_start(MagicMock())
1502
1503 assert bridge._airplay_stream is kept_stream
1504 assert bridge.airplay_player.stream is kept_stream
1505
1506
1507async def test_warm_stream_flushes_and_reanchors_on_kept_instance() -> None:
1508 """A warm handover flushes and re-anchors START on the same stream instance."""
1509 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1510 bridge._drop_until_us = SENDSPIN_EPOCH_US + WARM_LEAD_MS * 1_000
1511 commanded = UNIX_NOW_MS + WARM_LEAD_MS
1512 kept_stream = _make_anchor_stream(ack=commanded)
1513 bridge._airplay_stream = kept_stream
1514 bridge._airplay_stream_start_task = asyncio.current_task()
1515
1516 with patch(
1517 "music_assistant.providers.airplay.sendspin_bridge.time.time",
1518 return_value=UNIX_NOW_S,
1519 ):
1520 committed = await bridge._start_warm_stream(kept_stream)
1521
1522 assert committed is True
1523 assert bridge._airplay_stream is kept_stream # no new instance was built
1524 kept_stream.flush.assert_awaited_once_with()
1525 kept_stream.start.assert_awaited_once_with(commanded, join=True)
1526 assert bridge._started is True
1527 assert bridge._airplay_stream_ready.is_set()
1528
1529
1530async def test_warm_stream_flush_timeout_falls_back_to_cold() -> None:
1531 """A flush that is never acknowledged never re-anchors and falls back to cold."""
1532 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1533 kept_stream = _make_anchor_stream()
1534 kept_stream.flush = AsyncMock(return_value=False)
1535 bridge._airplay_stream = kept_stream
1536 bridge._airplay_stream_start_task = asyncio.current_task()
1537
1538 committed = await bridge._start_warm_stream(kept_stream)
1539
1540 assert committed is False
1541 kept_stream.start.assert_not_awaited()
1542 assert bridge._started is False
1543
1544
1545async def test_warm_stream_superseded_before_start_does_not_anchor() -> None:
1546 """If a newer stream start already owns the bridge, the stale flush never re-anchors."""
1547 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1548 kept_stream = _make_anchor_stream()
1549 bridge._airplay_stream = kept_stream
1550 # Simulate a newer stream start having already replaced the tracked task.
1551 bridge._airplay_stream_start_task = MagicMock()
1552
1553 committed = await bridge._start_warm_stream(kept_stream)
1554
1555 assert committed is False
1556 kept_stream.start.assert_not_awaited()
1557
1558
1559async def test_warm_stream_cancellation_propagates() -> None:
1560 """Cancellation while flushing propagates without re-anchoring."""
1561 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1562 kept_stream = _make_anchor_stream()
1563 flush_waiting = asyncio.Event()
1564
1565 async def flush(*_args: object, **_kwargs: object) -> bool:
1566 bridge._airplay_stream_start_task = asyncio.current_task()
1567 flush_waiting.set()
1568 await asyncio.Event().wait()
1569 return True
1570
1571 kept_stream.flush = AsyncMock(side_effect=flush)
1572 bridge._airplay_stream = kept_stream
1573
1574 warm_task = asyncio.create_task(bridge._start_warm_stream(kept_stream))
1575 await flush_waiting.wait()
1576 warm_task.cancel()
1577 with pytest.raises(asyncio.CancelledError):
1578 await warm_task
1579
1580 kept_stream.start.assert_not_awaited()
1581
1582
1583async def test_warm_handover_superseded_during_the_anchor_keeps_the_stream() -> None:
1584 """
1585 A superseded warm handover leaves the kept stream to the newer start.
1586
1587 The anchor reports the same False for "superseded" as for a genuine failure,
1588 so without an ownership re-check the stale task would stop the very transport
1589 the newer start decided to keep and put a second cliairplay on the receiver.
1590 """
1591 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1592 bridge._drop_until_us = SENDSPIN_EPOCH_US + WARM_LEAD_MS * 1_000
1593 bridge._airplay_stream_start_task = asyncio.current_task()
1594 kept_stream = _make_anchor_stream()
1595 bridge._airplay_stream = kept_stream
1596 bridge.airplay_player.stream = kept_stream
1597 cold_stream = _make_anchor_stream()
1598
1599 async def start(start_unix_ms: int, **_kwargs: object) -> int:
1600 # A newer stream start claimed the bridge while the binary held its ack.
1601 bridge._airplay_stream_start_task = MagicMock()
1602 return start_unix_ms
1603
1604 kept_stream.start = AsyncMock(side_effect=start)
1605
1606 with (
1607 patch(
1608 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
1609 return_value=cold_stream,
1610 ),
1611 patch(
1612 "music_assistant.providers.airplay.sendspin_bridge.time.time",
1613 return_value=UNIX_NOW_S,
1614 ),
1615 ):
1616 await bridge._start_protocol_from_chunk()
1617
1618 kept_stream.stop.assert_not_awaited()
1619 cold_stream.connect.assert_not_awaited()
1620 assert bridge._airplay_stream is kept_stream
1621 assert bridge.airplay_player.stream is kept_stream
1622
1623
1624async def test_superseded_start_failure_leaves_the_newer_stream_alone() -> None:
1625 """
1626 A stale start that fails must not take the newer stream's state with it.
1627
1628 The receiver is busy precisely because the newer start just claimed it, so a
1629 superseded cold connect failing is the ordinary outcome. Running the recovery
1630 would drop the newer stream's held backlog, release its writer before its own
1631 anchor is settled and schedule its teardown.
1632 """
1633 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1634 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
1635 bridge._airplay_stream_start_task = asyncio.current_task()
1636 bridge._hold_chunk(_pcm_chunk(SENDSPIN_EPOCH_US))
1637 newer_stream = _make_anchor_stream()
1638 stale_stream = _make_anchor_stream()
1639
1640 async def connect(_use_shared_ptp: bool | None) -> None:
1641 # The newer start won the receiver, so this one cannot have it.
1642 bridge._airplay_stream_start_task = MagicMock()
1643 bridge._airplay_stream = newer_stream
1644 bridge.airplay_player.stream = newer_stream
1645 raise OSError("device busy")
1646
1647 stale_stream.connect = AsyncMock(side_effect=connect)
1648
1649 with patch(
1650 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
1651 return_value=stale_stream,
1652 ):
1653 await bridge._start_protocol_from_chunk()
1654
1655 assert bridge._is_streaming is True
1656 assert len(bridge._held_chunks) == 1
1657 assert not bridge._airplay_stream_ready.is_set()
1658 assert bridge._airplay_stream is newer_stream
1659 newer_stream.stop.assert_not_awaited()
1660 # No teardown was scheduled for the newer stream's resources.
1661 cast("MagicMock", bridge.mass).create_task.assert_not_called()
1662
1663
1664# --- Startup lead: how far ahead Sendspin schedules the first chunk ---
1665
1666
1667def _make_timed_bridge() -> tuple[SendspinAirPlayBridge, MagicMock]:
1668 """Return a bridge with a mocked bridge role attached, plus that role."""
1669 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1670 role = MagicMock()
1671 bridge._bridge_role = role
1672 return bridge, role
1673
1674
1675def test_bridge_timing_reports_the_cold_lead_without_a_warm_stream() -> None:
1676 """Without a reusable stream the lead has to cover a full process spawn and connect."""
1677 bridge, role = _make_timed_bridge()
1678
1679 bridge._refresh_bridge_timing()
1680
1681 role.set_timing.assert_called_once_with(
1682 required_lead_time_ms=BRIDGE_COLD_START_LEAD_MS, min_buffer_ms=BRIDGE_MIN_BUFFER_MS
1683 )
1684
1685
1686def test_bridge_timing_reports_the_warm_lead_for_a_reusable_stream() -> None:
1687 """A kept, connected, already-anchored stream pays no connect, so it needs less lead."""
1688 bridge, role = _make_timed_bridge()
1689 bridge._airplay_stream = _make_kept_stream()
1690 bridge._started = True
1691
1692 bridge._refresh_bridge_timing()
1693
1694 role.set_timing.assert_called_once_with(
1695 required_lead_time_ms=BRIDGE_WARM_START_LEAD_MS, min_buffer_ms=BRIDGE_MIN_BUFFER_MS
1696 )
1697
1698
1699def test_bridge_timing_is_a_noop_without_a_bridge_role() -> None:
1700 """Timing can be refreshed before registration completes, with nothing to push it to."""
1701 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1702 assert bridge._bridge_role is None
1703
1704 bridge._refresh_bridge_timing()
1705
1706
1707def test_stream_start_reads_the_lead_before_rewinding_the_stream_state() -> None:
1708 """
1709 The warm/cold decision is taken while the previous stream's state is intact.
1710
1711 ``_on_stream_start`` rewinds the per-stream state; reading the timing after
1712 that rewind would report the cold lead for every start, including the warm
1713 handovers that need none of that budget.
1714 """
1715 bridge, role = _make_timed_bridge()
1716 kept_stream = _make_kept_stream()
1717 bridge._airplay_stream = kept_stream
1718 bridge.airplay_player.stream = kept_stream
1719 bridge._started = True
1720 observed: list[tuple[bool, object]] = []
1721 refresh = bridge._refresh_bridge_timing
1722
1723 def record() -> None:
1724 observed.append((bridge._started, bridge._airplay_stream))
1725 refresh()
1726
1727 with patch.object(bridge, "_refresh_bridge_timing", record):
1728 bridge._on_stream_start(MagicMock())
1729
1730 assert observed == [(True, kept_stream)]
1731 role.set_timing.assert_called_once_with(
1732 required_lead_time_ms=BRIDGE_WARM_START_LEAD_MS, min_buffer_ms=BRIDGE_MIN_BUFFER_MS
1733 )
1734
1735
1736# --- Mid-stream transport loss: re-anchoring, and giving up when it keeps dropping ---
1737
1738
1739def _make_completed_start_task(*, failed: bool = False) -> MagicMock:
1740 """
1741 Build a start-task mock the chunk handler reads as a finished protocol start.
1742
1743 Every predicate must answer a real bool: a bare MagicMock reports itself as
1744 cancelled, which the handler reads as a failed start.
1745 """
1746 task = MagicMock()
1747 task.done.return_value = True
1748 task.cancelled.return_value = failed
1749 task.exception.return_value = None
1750 return task
1751
1752
1753def _make_anchored_bridge(
1754 *, running: bool, ended_cleanly: bool = False
1755) -> tuple[SendspinAirPlayBridge, MagicMock]:
1756 """
1757 Return a bridge anchored on a transport in the given running state, plus that transport.
1758
1759 :param running: Whether the cli process behind the transport is still alive.
1760 :param ended_cleanly: Whether the binary reported the end of the stream
1761 itself, which stops the transport without losing it.
1762 """
1763 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1764 stream = _make_kept_stream(running=running, ended_cleanly=ended_cleanly)
1765 bridge._airplay_stream = stream
1766 bridge.airplay_player.stream = stream
1767 bridge._airplay_stream_start_task = _make_completed_start_task()
1768 bridge._started = True
1769 bridge._anchor_settled = True
1770 bridge._drop_until_us = SENDSPIN_EPOCH_US
1771 return bridge, stream
1772
1773
1774def test_lost_transport_rearms_a_cold_start_on_the_current_chunk() -> None:
1775 """
1776 A transport that died mid-stream is released and re-anchored on the live timeline.
1777
1778 The CLI accepts and discards writes once its process is gone, so the loss is
1779 only visible on the stream itself. The chunk that exposes it is also the one
1780 the fresh transport anchors to, which is where the group is playing now.
1781 """
1782 bridge, _ = _make_anchored_bridge(running=False)
1783 chunk_ts = SENDSPIN_EPOCH_US + 30_000_000
1784
1785 bridge._on_audio_chunk(_pcm_chunk(chunk_ts))
1786
1787 assert bridge._airplay_stream is None
1788 assert bridge.airplay_player.stream is None
1789 assert bridge._started is False
1790 assert bridge._anchor_settled is False
1791 # a fresh start is armed and anchored where the group is playing right now
1792 assert bridge._drop_until_us == chunk_ts
1793 # the chunk is held until the new anchor is acked, not placed against the dead one
1794 assert len(bridge._held_chunks) == 1
1795
1796
1797def test_lost_transport_and_its_writer_are_torn_down() -> None:
1798 """The dead transport and the writer feeding it are handed to the cleanup path."""
1799 bridge, dead_stream = _make_anchored_bridge(running=False)
1800 writer_task = MagicMock()
1801 bridge._writer_task = writer_task
1802 start_task = bridge._airplay_stream_start_task
1803
1804 with patch.object(bridge, "_cleanup_old_stream", MagicMock()) as cleanup:
1805 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
1806
1807 assert cleanup.call_args.args[:3] == (dead_stream, writer_task, start_task)
1808
1809
1810def test_live_transport_keeps_streaming_untouched() -> None:
1811 """A running transport is left alone: chunks keep flowing to the same stream."""
1812 bridge, stream = _make_anchored_bridge(running=True)
1813 start_task = bridge._airplay_stream_start_task
1814
1815 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
1816
1817 assert bridge._airplay_stream is stream
1818 assert bridge._airplay_stream_start_task is start_task
1819 assert bridge._started is True
1820 assert not bridge._write_queue.empty()
1821
1822
1823def test_transport_is_not_judged_while_a_start_is_in_flight() -> None:
1824 """
1825 A start owns its transport, so a stream it is tearing down is not a loss.
1826
1827 A warm handover that fails stops the kept stream before dropping it, leaving
1828 a window where the bridge still points at a stopped stream. Restarting from
1829 that window would fight the start already falling back to a cold reconnect.
1830 """
1831 bridge, stopped_stream = _make_anchored_bridge(running=False)
1832 cast("MagicMock", bridge._airplay_stream_start_task).done.return_value = False
1833
1834 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
1835
1836 assert bridge._airplay_stream is stopped_stream
1837 assert bridge._started is True
1838
1839
1840def test_unanchored_transport_is_not_treated_as_a_loss() -> None:
1841 """
1842 A stream that never anchored is the start's to report, not a mid-stream loss.
1843
1844 Recovery re-joins the group where the current chunk sits, which only means
1845 anything once an anchor existed. A start that finished without one has
1846 already taken the bridge out of streaming through its own failure path.
1847 """
1848 bridge, stopped_stream = _make_anchored_bridge(running=False)
1849 start_task = bridge._airplay_stream_start_task
1850 bridge._started = False
1851
1852 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
1853
1854 assert bridge._airplay_stream is stopped_stream
1855 assert bridge._airplay_stream_start_task is start_task
1856
1857
1858def test_a_stream_the_native_path_took_over_is_not_recovered() -> None:
1859 """
1860 A transport the bridge no longer owns is not the bridge's to restart.
1861
1862 The native path stops (or replaces) the player's stream without telling the
1863 bridge, which reads its own stopped stream as a crash. Recovering would put
1864 a second cli process on the same receiver and let the cold start publish its
1865 stream over the native session's.
1866 """
1867 bridge, stopped_stream = _make_anchored_bridge(running=False)
1868 # the native path took the player over and left the bridge holding a stream
1869 # that is no longer the player's
1870 cast("MagicMock", bridge.airplay_player).stream = _make_kept_stream()
1871
1872 with patch.object(bridge, "_restart_transport", MagicMock()) as restart:
1873 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
1874
1875 restart.assert_not_called()
1876 assert bridge._airplay_stream is stopped_stream
1877
1878
1879def test_a_stream_the_binary_ended_itself_is_not_a_loss() -> None:
1880 """
1881 A cli process that reported the end of the stream did not lose its transport.
1882
1883 The stderr loop also ends on a clean [STATUS] eof or the binary's idle cap,
1884 which stops the stream exactly like a crash does. Restarting one of those
1885 spawns a process for audio that is already over, and two such restarts
1886 inside the guard window take the speaker out of the group for good. Its
1887 counterpart is test_lost_transport_rearms_a_cold_start_on_the_current_chunk,
1888 where the same stopped stream ended without saying so.
1889 """
1890 bridge, ended_stream = _make_anchored_bridge(running=False, ended_cleanly=True)
1891
1892 with patch.object(bridge, "_restart_transport", MagicMock()) as restart:
1893 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
1894
1895 restart.assert_not_called()
1896 assert bridge._airplay_stream is ended_stream
1897 assert bridge._started is True
1898
1899
1900def test_restarting_the_transport_drops_a_deferred_teardown() -> None:
1901 """
1902 A teardown deferred by an earlier stream end must not fire into the new transport.
1903
1904 The restart arms a transport that pending timer knows nothing about, so it is
1905 cancelled along with the stream it was scheduled for.
1906 """
1907 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1908
1909 bridge._restart_transport()
1910
1911 cast("MagicMock", bridge.mass).cancel_timer.assert_called_once_with(bridge._teardown_timer_id)
1912
1913
1914def test_a_grace_timer_that_already_fired_spares_the_restarted_stream() -> None:
1915 """
1916 A teardown whose timer fired before the restart cancelled it leaves the new stream alone.
1917
1918 cancel_timer cannot recall a handle that already fired, so a stream arriving
1919 at the very end of the grace window still gets the call. Reading the live
1920 fields there would cancel that stream's writer and drain its queue, leaving
1921 the speaker silent for the whole track -- and the warm restart keeps the same
1922 stream object, so telling the two apart by the stream alone cannot work.
1923 """
1924 bridge, stream = _make_anchored_bridge(running=True)
1925 bridge._writer_task = MagicMock()
1926
1927 bridge._on_bridge_stream_end()
1928 # the next stream arrives and rides the kept process, cancelling a timer that
1929 # has already fired
1930 bridge._restart_transport()
1931 new_writer_task = bridge._writer_task
1932 bridge._write_queue.put_nowait(b"\x00" * BRIDGE_BYTES_PER_FRAME)
1933
1934 bridge._deferred_cleanup()
1935
1936 assert bridge._airplay_stream is stream
1937 assert bridge._writer_task is new_writer_task
1938 assert not bridge._write_queue.empty()
1939
1940
1941async def test_the_cleanup_a_start_waits_on_cannot_cancel_it() -> None:
1942 """
1943 A start waiting for the pending teardown is not among the handles it cancels.
1944
1945 _start_protocol_from_chunk and _cli_writer both await _cleanup_task before
1946 touching the transport. A teardown reading the live fields when it finally
1947 ran would find the waiting start there and cancel it, killing the stream it
1948 was clearing the way for.
1949 """
1950 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1951 bridge._is_streaming = False
1952 stream = _make_kept_stream()
1953 stream.stop = AsyncMock()
1954 bridge._airplay_stream = stream
1955 bridge._airplay_stream_start_task = _make_completed_start_task()
1956
1957 bridge._schedule_cleanup()
1958 teardown = cast("MagicMock", bridge.mass).create_task.call_args.args[0]
1959 # the start that arrives next publishes itself and then awaits the teardown
1960 start = cast("asyncio.Task[None]", asyncio.current_task())
1961 bridge._airplay_stream_start_task = start
1962
1963 await teardown
1964
1965 # the teardown ran against what the bridge held when it was scheduled
1966 stream.stop.assert_awaited_once_with(force=True)
1967 assert start.cancelling() == 0
1968
1969
1970def test_a_new_sendspin_stream_restores_the_recovery_budget() -> None:
1971 """
1972 Every Sendspin stream starts with a full recovery budget.
1973
1974 A loss on the previous stream says nothing about the device's health on this
1975 one; carrying the stamp over would abandon a speaker on its very first loss.
1976 """
1977 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1978 bridge._last_transport_recovery = 100.0
1979
1980 bridge._on_stream_start(MagicMock())
1981
1982 assert bridge._last_transport_recovery is None
1983
1984
1985def test_a_stream_start_on_an_unavailable_player_still_restores_the_budget() -> None:
1986 """
1987 The recovery budget is settled before any early return can skip it.
1988
1989 The stream-start callback bails out when the player is unavailable, but the
1990 role-side entry point has no such gate; leaving the verdict of the previous
1991 stream in place would abandon the speaker on the next stream's first loss.
1992 """
1993 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
1994 cast("MagicMock", bridge.airplay_player).available = False
1995 bridge._last_transport_recovery = 100.0
1996
1997 bridge._on_stream_start(MagicMock())
1998
1999 assert bridge._last_transport_recovery is None
2000
2001
2002def test_second_transport_loss_within_the_guard_window_gives_up() -> None:
2003 """A device dropping its transport again right away is abandoned, not re-anchored."""
2004 bridge, _ = _make_anchored_bridge(running=False)
2005
2006 with (
2007 patch(
2008 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2009 side_effect=[100.0, 100.0 + BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS - 1],
2010 ),
2011 patch.object(bridge, "_restart_transport", MagicMock()) as restart,
2012 patch.object(bridge, "_abandon_streaming", MagicMock()) as abandon,
2013 ):
2014 assert bridge._recover_transport() is True
2015 assert bridge._recover_transport() is False
2016
2017 restart.assert_called_once_with()
2018 abandon.assert_called_once_with()
2019
2020
2021def test_transport_loss_after_the_guard_window_recovers_again() -> None:
2022 """A single blip hours apart is a new incident, not a flapping device."""
2023 bridge, _ = _make_anchored_bridge(running=False)
2024
2025 with (
2026 patch(
2027 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2028 side_effect=[100.0, 100.0 + BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS + 1],
2029 ),
2030 patch.object(bridge, "_restart_transport", MagicMock()) as restart,
2031 patch.object(bridge, "_abandon_streaming", MagicMock()) as abandon,
2032 ):
2033 assert bridge._recover_transport() is True
2034 assert bridge._recover_transport() is True
2035
2036 assert restart.call_count == 2
2037 abandon.assert_not_called()
2038
2039
2040def test_giving_up_does_not_queue_the_chunk_that_exposed_the_loss() -> None:
2041 """
2042 The chunk that trips the give-up is dropped, not written into the dead stream.
2043
2044 Giving up leaves the anchor and the stream reference untouched, so a chunk
2045 that carried on through the handler would still be placed and queued.
2046 """
2047 bridge, _ = _make_anchored_bridge(running=False)
2048 bridge._last_transport_recovery = 100.0
2049
2050 with (
2051 patch(
2052 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2053 return_value=100.0 + BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS - 1,
2054 ),
2055 patch.object(bridge, "_restart_transport", MagicMock()) as restart,
2056 ):
2057 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
2058
2059 restart.assert_not_called()
2060 assert bridge._write_queue.empty()
2061
2062
2063async def test_a_failed_protocol_start_leaves_the_session() -> None:
2064 """
2065 A cold start that raised takes the speaker out of the group it cannot play in.
2066
2067 Whether the start was the stream's first or a replacement for a transport
2068 that died, the outcome is the same silence; leaving is what stops the player
2069 reporting playback nobody can hear.
2070 """
2071 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2072 bridge._airplay_stream_start_task = asyncio.current_task()
2073 stream = _make_anchor_stream()
2074 stream.connect = AsyncMock(side_effect=OSError("no route to device"))
2075
2076 with (
2077 patch(
2078 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
2079 return_value=stream,
2080 ),
2081 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
2082 ):
2083 await bridge._start_protocol_from_chunk()
2084
2085 assert bridge._is_streaming is False
2086 leave.assert_called_once_with()
2087
2088
2089async def test_losing_a_speaker_for_good_runs_the_whole_chain() -> None:
2090 """
2091 End to end: a transport dies, the reconnect is refused, the speaker leaves the group.
2092
2093 Every step here is the real one -- detection, the recovery decision, the
2094 re-arm and the cold start -- so a give-up swallowed anywhere along that
2095 chain shows up as a speaker that stays silently "playing" instead of
2096 dropping out.
2097 """
2098 bridge, _ = _make_anchored_bridge(running=False)
2099 stream = _make_anchor_stream()
2100 stream.connect = AsyncMock(side_effect=OSError("device gone"))
2101
2102 with (
2103 patch(
2104 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
2105 return_value=stream,
2106 ),
2107 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
2108 ):
2109 # the chunk that exposes the loss re-arms and anchors a replacement
2110 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
2111 assert bridge._airplay_stream_start_task is not None
2112 leave.assert_not_called()
2113 # run the replacement start the chunk handler scheduled
2114 start = cast("MagicMock", bridge.mass).create_task.call_args.args[0]
2115 bridge._airplay_stream_start_task = asyncio.current_task()
2116 await start
2117
2118 assert bridge._is_streaming is False
2119 leave.assert_called_once_with()
2120
2121
2122def test_abandoning_streaming_stops_the_feed_and_leaves_the_session() -> None:
2123 """Giving up stops accepting chunks, unblocks the writer and leaves the session."""
2124 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2125 bridge._held_chunks.append(_pcm_chunk(SENDSPIN_EPOCH_US))
2126 bridge._held_us = 100_000
2127
2128 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
2129 bridge._abandon_streaming()
2130
2131 assert bridge._is_streaming is False
2132 assert not bridge._held_chunks
2133 assert bridge._held_us == 0
2134 assert bridge._airplay_stream_ready.is_set()
2135 # scheduled, not merely constructed: an unscheduled coroutine never leaves
2136 leave.assert_called_once_with()
2137 scheduled = [call.args[0] for call in cast("MagicMock", bridge.mass).create_task.call_args_list]
2138 assert leave.return_value in scheduled
2139
2140
2141def test_a_flapping_device_is_taken_out_of_the_sendspin_session() -> None:
2142 """
2143 Only a device that cannot hold a transport is dropped from the group.
2144
2145 Its silence is real and permanent, so the visible player must stop reporting
2146 playback; the rest of the group keeps going without it.
2147 """
2148 bridge, _ = _make_anchored_bridge(running=False)
2149
2150 with (
2151 patch(
2152 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2153 side_effect=[100.0, 100.0 + BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS - 1],
2154 ),
2155 patch.object(bridge, "_restart_transport", MagicMock()),
2156 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
2157 ):
2158 assert bridge._recover_transport() is True
2159 # the first loss is recoverable, so the speaker keeps its place
2160 leave.assert_not_called()
2161 assert bridge._recover_transport() is False
2162
2163 # scheduled, not merely constructed: an unscheduled coroutine never leaves
2164 leave.assert_called_once_with()
2165 scheduled = [call.args[0] for call in cast("MagicMock", bridge.mass).create_task.call_args_list]
2166 assert leave.return_value in scheduled
2167
2168
2169def test_failed_start_task_gives_up_on_the_stream() -> None:
2170 """A protocol start that failed stops the feed and drops out of the group."""
2171 bridge, _ = _make_anchored_bridge(running=True)
2172 bridge._airplay_stream_start_task = _make_completed_start_task(failed=True)
2173
2174 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
2175 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + 1_000_000))
2176
2177 assert bridge._is_streaming is False
2178 leave.assert_called_once_with()
2179
2180
2181async def test_writer_readiness_timeout_gives_up_on_the_stream() -> None:
2182 """
2183 A protocol that never becomes ready stops the feed and drops out of the group.
2184
2185 A transport that hangs instead of failing renders the same silence as one
2186 that refused the connection, so it is given up on the same way.
2187 """
2188 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2189 bridge._writer_task = asyncio.current_task()
2190 bridge._airplay_stream_ready = MagicMock(wait=AsyncMock(side_effect=TimeoutError))
2191
2192 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
2193 await bridge._cli_writer()
2194
2195 assert bridge._is_streaming is False
2196 leave.assert_called_once_with()
2197
2198
2199async def test_a_stale_writer_cannot_give_up_on_a_newer_stream() -> None:
2200 """
2201 Only the writer still feeding the bridge may abandon it.
2202
2203 A writer left behind by a slow teardown speaks for a stream that is already
2204 gone; letting it give up would stop, and un-group, its successor.
2205 """
2206 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2207 # a newer stream owns the bridge; this writer is the previous stream's
2208 bridge._writer_task = MagicMock()
2209 bridge._airplay_stream_ready = MagicMock(wait=AsyncMock(side_effect=TimeoutError))
2210
2211 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
2212 await bridge._cli_writer()
2213
2214 assert bridge._is_streaming is True
2215 leave.assert_not_called()
2216
2217
2218def _make_grouped_client(*, group_members: int = 2, has_active_stream: bool = False) -> MagicMock:
2219 """
2220 Build a bridge client mock that reads as having left a shared group.
2221
2222 Quiescing moves the client on to a solo group, exactly as the real one does,
2223 so a caller reading ``client.group`` after leaving no longer sees the group
2224 that was left.
2225
2226 :param group_members: Members in the group the client lands in after
2227 leaving; more than one means it was grouped again meanwhile.
2228 :param has_active_stream: Whether that group is playing something of its own.
2229 """
2230 client = MagicMock()
2231 client.group.clients = [MagicMock() for _ in range(group_members)]
2232 client.group.has_active_stream = has_active_stream
2233
2234 async def _quiesce() -> str:
2235 client.group = MagicMock(clients=[client], has_active_stream=False)
2236 return "group-1"
2237
2238 # a real group id: leaving a shared group is what earns a re-join, and None
2239 # (a solo group, which leaving simply stops) must stay distinguishable
2240 client.quiesce_to_solo_stopped = AsyncMock(side_effect=_quiesce)
2241 return client
2242
2243
2244async def test_leaving_a_shared_group_lines_up_a_rejoin() -> None:
2245 """
2246 A bridge taken out of a shared group is given an attempt to come back.
2247
2248 The group it left is captured before quiescing, because that is what moves
2249 the client into a solo group of its own.
2250 """
2251 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2252 client = _make_grouped_client()
2253 left_group = client.group
2254 bridge._sendspin_client = client
2255
2256 with patch.object(bridge, "_rejoin_attempts", MagicMock()) as rejoin:
2257 await bridge._leave_sendspin_session()
2258
2259 rejoin.assert_called_once_with(left_group)
2260
2261
2262async def test_leaving_a_solo_group_has_nothing_to_rejoin() -> None:
2263 """
2264 A solo bridge is stopped by leaving, so there is no group to return to.
2265
2266 Quiescing reports that by returning no previous group; scheduling a re-join
2267 against the group it is already alone in would put it back on PLAYING.
2268 """
2269 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2270 client = _make_grouped_client()
2271 client.quiesce_to_solo_stopped = AsyncMock(return_value=None)
2272 bridge._sendspin_client = client
2273
2274 with patch.object(bridge, "_rejoin_attempts", MagicMock()) as rejoin:
2275 await bridge._leave_sendspin_session()
2276
2277 rejoin.assert_not_called()
2278
2279
2280async def test_a_speaker_that_fails_again_right_after_a_rejoin_stays_out() -> None:
2281 """
2282 A speaker that keeps dropping out cannot cycle in and out of its group.
2283
2284 Re-joining re-runs the stream start that just failed, and a device that
2285 accepts a START before dying would otherwise earn a fresh attempt every
2286 time round, churning CLI processes and group membership indefinitely.
2287 """
2288 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2289 bridge._sendspin_client = _make_grouped_client()
2290 bridge._last_rejoin = 100.0
2291
2292 with (
2293 patch(
2294 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2295 return_value=100.0 + BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS - 1,
2296 ),
2297 patch.object(bridge, "_rejoin_attempts", MagicMock()) as rejoin,
2298 ):
2299 await bridge._leave_sendspin_session()
2300
2301 rejoin.assert_not_called()
2302
2303
2304async def test_a_speaker_that_held_its_place_earns_another_rejoin() -> None:
2305 """A device that played on for a while before failing is worth bringing back again."""
2306 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2307 bridge._sendspin_client = _make_grouped_client()
2308 bridge._last_rejoin = 100.0
2309
2310 with (
2311 patch(
2312 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2313 return_value=100.0 + BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS + 1,
2314 ),
2315 patch.object(bridge, "_rejoin_attempts", MagicMock()) as rejoin,
2316 ):
2317 await bridge._leave_sendspin_session()
2318
2319 rejoin.assert_called_once()
2320
2321
2322async def test_the_rejoin_window_is_measured_from_the_actual_rejoin() -> None:
2323 """
2324 The guard is stamped where the speaker rejoins, not where the attempt was scheduled.
2325
2326 Stamping at schedule time would tie the guard to the backoff: longer delays
2327 would put the stamp far enough in the past for the window to have expired by
2328 the time the re-joined speaker fails, letting the cycle run again.
2329 """
2330 bridge, _, group = _make_rejoin_bridge()
2331
2332 with (
2333 patch(_NO_REJOIN_DELAYS, (0,)),
2334 patch(
2335 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2336 return_value=1234.0,
2337 ),
2338 ):
2339 await bridge._rejoin_attempts(group)
2340
2341 group.add_client.assert_awaited_once()
2342 assert bridge._last_rejoin == 1234.0
2343
2344
2345async def test_a_failed_rejoin_never_stamps_the_window() -> None:
2346 """A speaker that never made it back has not held a place to be judged on."""
2347 bridge, _, group = _make_rejoin_bridge()
2348 group.add_client = AsyncMock(side_effect=OSError("group is gone"))
2349
2350 with patch(_NO_REJOIN_DELAYS, (0,)):
2351 await bridge._rejoin_attempts(group)
2352
2353 assert bridge._last_rejoin is None
2354
2355
2356async def test_a_give_up_inside_the_window_drops_a_pending_rejoin() -> None:
2357 """
2358 Leaving the speaker out means dropping the attempt that would put it back.
2359
2360 A schedule left running would contradict the decision this give-up just
2361 made, and re-add a speaker that was meant to stay out.
2362 """
2363 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2364 bridge._sendspin_client = _make_grouped_client()
2365 bridge._last_rejoin = 100.0
2366 pending = MagicMock()
2367 pending.done.return_value = False
2368 bridge._rejoin_task = pending
2369
2370 with patch(
2371 "music_assistant.providers.airplay.sendspin_bridge.time.monotonic",
2372 return_value=100.0 + BRIDGE_TRANSPORT_RECOVERY_GUARD_SECONDS - 1,
2373 ):
2374 await bridge._leave_sendspin_session()
2375
2376 assert bridge._rejoin_task is None
2377 pending.cancel.assert_called_once_with() # type: ignore[unreachable]
2378
2379
2380def _make_rejoin_bridge(
2381 *, group_members: int = 1, has_active_stream: bool = False
2382) -> tuple[SendspinAirPlayBridge, MagicMock, MagicMock]:
2383 """
2384 Build a bridge in the state a give-up leaves behind, with its client and lost group.
2385
2386 The AirPlay stream is cleared explicitly: a give-up tears it down, and a
2387 bridge still pointing at one reads as a speaker streaming outside the
2388 bridge, which is itself a reason not to re-join.
2389
2390 :param group_members: Members of the group the client sits in now.
2391 :param has_active_stream: Whether that group is playing something of its own.
2392 """
2393 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2394 cast("MagicMock", bridge.airplay_player).stream = None
2395 client = _make_grouped_client(group_members=group_members, has_active_stream=has_active_stream)
2396 bridge._sendspin_client = client
2397 group = MagicMock()
2398 group.clients = [MagicMock()]
2399 group.add_client = AsyncMock()
2400 return bridge, client, group
2401
2402
2403async def test_a_rejoin_puts_the_bridge_back_into_the_group_it_left() -> None:
2404 """The bridge re-joins through the ordinary group-add, which re-runs the stream start."""
2405 bridge, client, group = _make_rejoin_bridge()
2406
2407 with patch(_NO_REJOIN_DELAYS, (0,)):
2408 await bridge._rejoin_attempts(group)
2409
2410 group.add_client.assert_awaited_once_with(client)
2411
2412
2413async def test_a_rejoin_leaves_a_regrouped_speaker_alone() -> None:
2414 """
2415 A speaker grouped again meanwhile is never pulled out of where it was put.
2416
2417 The re-join answers a failure, not the user; landing anywhere other than the
2418 solo group the give-up left means someone else has since decided otherwise.
2419 """
2420 bridge, _, group = _make_rejoin_bridge(group_members=2)
2421
2422 with patch(_NO_REJOIN_DELAYS, (0,)):
2423 await bridge._rejoin_attempts(group)
2424
2425 group.add_client.assert_not_awaited()
2426
2427
2428async def test_a_rejoin_leaves_a_speaker_playing_on_its_own_alone() -> None:
2429 """
2430 A speaker started on its own meanwhile keeps that playback.
2431
2432 Its solo group has one member, so membership alone cannot tell it apart from
2433 the group the give-up left it in -- but adding a client to another group
2434 stops the group it came from, which here is the user's own playback.
2435 """
2436 bridge, _, group = _make_rejoin_bridge(has_active_stream=True)
2437
2438 with patch(_NO_REJOIN_DELAYS, (0,)):
2439 await bridge._rejoin_attempts(group)
2440
2441 group.add_client.assert_not_awaited()
2442
2443
2444async def test_a_rejoin_leaves_a_natively_streaming_speaker_alone() -> None:
2445 """
2446 A speaker taken over by native AirPlay is not dragged back into Sendspin.
2447
2448 Re-joining restarts the bridge transport, which would tear down a session
2449 the bridge does not own.
2450 """
2451 bridge, _, group = _make_rejoin_bridge()
2452 cast("MagicMock", bridge.airplay_player).stream = MagicMock()
2453
2454 with patch(_NO_REJOIN_DELAYS, (0,)):
2455 await bridge._rejoin_attempts(group)
2456
2457 group.add_client.assert_not_awaited()
2458
2459
2460async def test_an_offline_speaker_is_looked_for_again_before_giving_up() -> None:
2461 """
2462 A speaker missing from discovery is never re-joined, but is looked for again.
2463
2464 A rebooting device is absent from discovery for a while after it starts
2465 answering, so abandoning on the first look would spend the whole re-join
2466 budget inside the window where such a device is always missing. Running out
2467 of attempts, rather than returning on the first one, is what shows the later
2468 look happened.
2469 """
2470 bridge, _, group = _make_rejoin_bridge()
2471 cast("MagicMock", bridge.airplay_player).available = False
2472 logger = MagicMock()
2473 bridge.logger = logger
2474
2475 with patch(_NO_REJOIN_DELAYS, (0, 0)):
2476 await bridge._rejoin_attempts(group)
2477
2478 group.add_client.assert_not_awaited()
2479 assert logger.debug.call_count == 2
2480 # the give-up is only reached once the attempts run out
2481 logger.warning.assert_called_once()
2482
2483
2484async def test_a_rejoin_is_abandoned_when_the_group_is_gone() -> None:
2485 """
2486 A group everyone else has left is not a group to return to.
2487
2488 Its object outlives the members holding it, so adding the bridge back would
2489 strand it alone in a group nothing streams to.
2490 """
2491 bridge, _, group = _make_rejoin_bridge()
2492 group.clients = []
2493
2494 with patch(_NO_REJOIN_DELAYS, (0,)):
2495 await bridge._rejoin_attempts(group)
2496
2497 group.add_client.assert_not_awaited()
2498
2499
2500async def test_a_rejoin_that_keeps_failing_gives_up() -> None:
2501 """Every attempt is tried, and a speaker that never returns leaves the player idle."""
2502 bridge, _, group = _make_rejoin_bridge()
2503 group.add_client = AsyncMock(side_effect=OSError("group is gone"))
2504
2505 with patch(_NO_REJOIN_DELAYS, (0, 0)):
2506 await bridge._rejoin_attempts(group)
2507
2508 assert group.add_client.await_count == 2
2509
2510
2511async def test_a_new_stream_supersedes_a_pending_rejoin() -> None:
2512 """Joining a session by any means makes the pending re-join stale."""
2513 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2514 pending = MagicMock()
2515 pending.done.return_value = False
2516 bridge._rejoin_task = pending
2517
2518 bridge._on_stream_start(MagicMock())
2519
2520 assert bridge._rejoin_task is None
2521 pending.cancel.assert_called_once_with() # type: ignore[unreachable]
2522
2523
2524async def test_a_rejoin_never_cancels_itself() -> None:
2525 """
2526 The re-join survives the stream start it causes.
2527
2528 Adding the bridge back to the group runs the stream-start path that clears
2529 stale schedules, and that path cannot be allowed to kill the attempt making
2530 the call.
2531 """
2532 bridge, client, group = _make_rejoin_bridge()
2533
2534 async def _add_client(_client: MagicMock) -> None:
2535 bridge._rejoin_task = asyncio.current_task()
2536 bridge._on_stream_start(MagicMock())
2537
2538 group.add_client = AsyncMock(side_effect=_add_client)
2539
2540 with patch(_NO_REJOIN_DELAYS, (0,)):
2541 await bridge._rejoin_attempts(group)
2542
2543 group.add_client.assert_awaited_once_with(client)
2544
2545
2546async def test_stopping_the_bridge_drops_a_pending_rejoin() -> None:
2547 """An unloaded bridge has no group to return to."""
2548 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2549 pending = MagicMock()
2550 pending.done.return_value = False
2551 bridge._rejoin_task = pending
2552
2553 await bridge.stop()
2554
2555 assert bridge._rejoin_task is None
2556 pending.cancel.assert_called_once_with() # type: ignore[unreachable]
2557
2558
2559async def test_leaving_the_session_quiesces_the_bridge_client() -> None:
2560 """The bridge leaves a shared group (or stops a solo one) but stays registered."""
2561 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2562 client = _make_grouped_client()
2563 bridge._sendspin_client = client
2564
2565 await bridge._leave_sendspin_session()
2566
2567 client.quiesce_to_solo_stopped.assert_awaited_once_with()
2568 # staying registered is what keeps the player around for the next stream
2569 cast("MagicMock", bridge.sendspin_server).remove_client.assert_not_called()
2570
2571
2572async def test_leaving_the_session_without_a_client_is_a_noop() -> None:
2573 """
2574 Giving up before registration completed has no session to leave.
2575
2576 The call has to return without touching anything: swallowing an error from
2577 an absent client would look identical from the outside, so the absence of a
2578 complaint is what distinguishes the two.
2579 """
2580 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2581 logger = MagicMock()
2582 bridge.logger = logger
2583 assert bridge._sendspin_client is None
2584
2585 await bridge._leave_sendspin_session()
2586
2587 logger.warning.assert_not_called()
2588
2589
2590# --- An explicit stop: end playback now, without lining up a return ------------
2591
2592
2593def _bridge_manager_for(bridge: SendspinAirPlayBridge) -> SendspinBridgeManager:
2594 """Return a bridge manager holding the given bridge under its player id."""
2595 manager = SendspinBridgeManager(cast("MagicMock", bridge.provider))
2596 manager._bridges[bridge.airplay_player.player_id] = bridge
2597 return manager
2598
2599
2600async def test_an_explicit_stop_tears_the_transport_down_at_once() -> None:
2601 """
2602 A stop the user asked for stops the speaker now, not after the grace window.
2603
2604 A Sendspin stream ending defers the teardown so the next track can ride the
2605 warm binary; nothing follows a stop, and the device holds seconds of audio,
2606 so deferring there just plays out what the user asked to end.
2607 """
2608 bridge, stream = _make_anchored_bridge(running=True)
2609 writer_task = MagicMock()
2610 bridge._writer_task = writer_task
2611 start_task = bridge._airplay_stream_start_task
2612 manager = _bridge_manager_for(bridge)
2613
2614 with (
2615 patch.object(bridge, "_cleanup_old_stream", MagicMock()) as cleanup,
2616 patch.object(bridge, "_leave_sendspin_session", MagicMock()),
2617 ):
2618 assert manager.stop_streaming(bridge.airplay_player.player_id) is True
2619
2620 assert cleanup.call_args.args[:3] == (stream, writer_task, start_task)
2621 assert bridge._is_streaming is False
2622 assert bridge._airplay_stream is None
2623 # no grace window is armed: that is what the teardown would have waited out
2624 cast("MagicMock", bridge.mass).call_later.assert_not_called()
2625
2626
2627async def test_an_explicit_stop_leaves_the_session_without_a_return() -> None:
2628 """
2629 Stopping takes the speaker out of the session, and it stays out.
2630
2631 Sendspin reports playback from the group's state, so a stopped bridge that
2632 stayed in would hold the visible player on PLAYING. The re-join exists to
2633 recover a speaker that dropped out by itself; a user who stopped one has not
2634 asked for it back.
2635 """
2636 bridge, _ = _make_anchored_bridge(running=True)
2637 manager = _bridge_manager_for(bridge)
2638
2639 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
2640 manager.stop_streaming(bridge.airplay_player.player_id)
2641
2642 leave.assert_called_once_with(rejoin=False)
2643 # scheduled, not merely constructed: an unscheduled coroutine never leaves
2644 scheduled = [call.args[0] for call in cast("MagicMock", bridge.mass).create_task.call_args_list]
2645 assert leave.return_value in scheduled
2646
2647
2648async def test_a_stop_of_an_idle_bridge_keeps_its_place_in_the_group() -> None:
2649 """
2650 A bridge with nothing playing has no session to leave.
2651
2652 Its group is not reporting playback through this speaker, so quiescing it out
2653 would only cost a grouped-but-idle player its membership on a stop command.
2654 """
2655 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2656 bridge._is_streaming = False
2657 manager = _bridge_manager_for(bridge)
2658
2659 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
2660 assert manager.stop_streaming(bridge.airplay_player.player_id) is True
2661
2662 leave.assert_not_called()
2663
2664
2665async def test_a_stop_never_reaches_a_player_without_a_bridge() -> None:
2666 """An unbridged player is left to the caller's own stop path."""
2667 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2668
2669 assert _bridge_manager_for(bridge).stop_streaming("apother") is False
2670
2671
2672# --- One shared-PTP decision per Sendspin group --------------------------------
2673
2674
2675def _group_bridges(*bridges: SendspinAirPlayBridge, daemon_ready: bool) -> SendspinBridgeManager:
2676 """
2677 Put the given bridges in one Sendspin group behind a shared bridge manager.
2678
2679 :param bridges: Bridges to place in the group.
2680 :param daemon_ready: What the shared PTP daemon answers a fresh resolve.
2681 """
2682 provider = MagicMock()
2683 provider.ptp_daemon_ready = daemon_ready
2684 manager = SendspinBridgeManager(provider)
2685 provider.bridge_manager = manager
2686 group = MagicMock()
2687 for index, bridge in enumerate(bridges):
2688 bridge.provider = provider
2689 bridge._sendspin_client = MagicMock()
2690 bridge._sendspin_client.group = group
2691 manager._bridges[f"player{index}"] = bridge
2692 return manager
2693
2694
2695def test_the_first_group_member_asks_the_daemon() -> None:
2696 """With no live decision in the group, the daemon's readiness decides."""
2697 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2698 _group_bridges(bridge, daemon_ready=True)
2699
2700 assert bridge._resolve_shared_ptp() is True
2701
2702
2703@pytest.mark.parametrize(("live_decision", "daemon_ready"), [(True, False), (False, True)])
2704def test_a_later_member_adopts_the_groups_live_decision(
2705 live_decision: bool, daemon_ready: bool
2706) -> None:
2707 """
2708 A member starting later joins on the clock the group is already running.
2709
2710 Bridges in one group can start minutes apart, so what the daemon answers at
2711 the second start says nothing about the source the first member's process
2712 was spawned against. Parametrised both ways so the decision is proven to
2713 follow the sibling rather than the daemon.
2714 """
2715 playing = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2716 joiner = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2717 _group_bridges(playing, joiner, daemon_ready=daemon_ready)
2718 playing._use_shared_ptp = live_decision
2719
2720 assert joiner._resolve_shared_ptp() is live_decision
2721
2722
2723def test_a_warm_member_still_speaks_for_the_group() -> None:
2724 """
2725 A process kept for a warm reuse keeps deciding for its group.
2726
2727 Its Sendspin stream ended, but the next one rides that same cli process with
2728 the flag it was spawned with, so a sibling cold-starting alongside it has to
2729 match that flag rather than resolve against the daemon.
2730 """
2731 warm = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2732 joiner = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2733 _group_bridges(warm, joiner, daemon_ready=False)
2734 warm._use_shared_ptp = True
2735 warm._airplay_stream = _make_kept_stream()
2736 warm._started = True
2737 warm._is_streaming = False
2738
2739 assert warm.active_shared_ptp is True
2740 assert joiner._resolve_shared_ptp() is True
2741
2742
2743def test_an_idle_member_does_not_decide() -> None:
2744 """A bridge with no cli process left leaves the group to resolve fresh."""
2745 idle = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2746 starter = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2747 _group_bridges(idle, starter, daemon_ready=True)
2748 idle._use_shared_ptp = False
2749 idle._is_streaming = False
2750
2751 assert idle.active_shared_ptp is None
2752 assert starter._resolve_shared_ptp() is True
2753
2754
2755def test_another_groups_decision_is_not_adopted() -> None:
2756 """Only members of the same Sendspin group share one timing source."""
2757 stranger = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2758 starter = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2759 _group_bridges(stranger, starter, daemon_ready=False)
2760 stranger._use_shared_ptp = True
2761 # the stranger moved on to a group of its own
2762 stranger_client = MagicMock()
2763 stranger_client.group = MagicMock()
2764 stranger._sendspin_client = stranger_client
2765
2766 assert starter._resolve_shared_ptp() is False
2767
2768
2769def test_a_raop_member_carries_no_decision() -> None:
2770 """A legacy RAOP process has no shared-clock flag to hand its group."""
2771 raop = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US, protocol=StreamingProtocol.RAOP)
2772 _group_bridges(raop, daemon_ready=True)
2773
2774 assert raop._resolve_shared_ptp() is None
2775
2776
2777async def test_a_cold_start_spawns_the_cli_with_the_groups_decision() -> None:
2778 """The adopted decision reaches the cli process and is recorded on the bridge."""
2779 playing = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2780 joiner = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2781 _group_bridges(playing, joiner, daemon_ready=False)
2782 playing._use_shared_ptp = True
2783 joiner._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
2784 joiner._airplay_stream_start_task = asyncio.current_task()
2785 stream = _make_anchor_stream(ack=UNIX_NOW_MS + COLD_LEAD_MS)
2786
2787 with patch(
2788 "music_assistant.providers.airplay.sendspin_bridge.time.time",
2789 return_value=UNIX_NOW_S,
2790 ):
2791 assert await joiner._start_cold_stream(stream) is True
2792
2793 stream.connect.assert_awaited_once_with(True)
2794 assert joiner.active_shared_ptp is True
2795
2796
2797async def test_a_daemon_lost_mid_start_cannot_split_the_group() -> None:
2798 """
2799 Members starting together agree even when the daemon goes away between them.
2800
2801 The first member records its decision before it awaits its connect, so the
2802 second one finds it however the daemon answers by the time it resolves.
2803 """
2804 first = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2805 second = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2806 manager = _group_bridges(first, second, daemon_ready=True)
2807 first_stream = _make_anchor_stream(ack=UNIX_NOW_MS + COLD_LEAD_MS)
2808 second_stream = _make_anchor_stream(ack=UNIX_NOW_MS + COLD_LEAD_MS)
2809
2810 async def connect(_use_shared_ptp: bool | None) -> None:
2811 # the daemon dies while the first member is still connecting
2812 cast("MagicMock", manager.provider).ptp_daemon_ready = False
2813 await asyncio.sleep(0)
2814
2815 first_stream.connect = AsyncMock(side_effect=connect)
2816
2817 async def cold_start(bridge: SendspinAirPlayBridge, stream: MagicMock) -> None:
2818 bridge._drop_until_us = SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000
2819 bridge._airplay_stream_start_task = asyncio.current_task()
2820 await bridge._start_cold_stream(stream)
2821
2822 with patch(
2823 "music_assistant.providers.airplay.sendspin_bridge.time.time",
2824 return_value=UNIX_NOW_S,
2825 ):
2826 await asyncio.gather(cold_start(first, first_stream), cold_start(second, second_stream))
2827
2828 assert first.active_shared_ptp is True
2829 assert second.active_shared_ptp is True
2830 second_stream.connect.assert_awaited_once_with(True)
2831
2832
2833async def test_a_torn_down_bridge_stops_deciding() -> None:
2834 """
2835 The decision dies with the cli process it was spawned for.
2836
2837 A new Sendspin stream arms the bridge before it resolves, so a decision left
2838 behind by the torn-down process would be handed to the group on its behalf.
2839 """
2840 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2841 bridge._use_shared_ptp = True
2842
2843 await bridge._stop_streaming()
2844 bridge._on_stream_start(MagicMock())
2845
2846 assert bridge._is_streaming is True
2847 assert bridge.active_shared_ptp is None
2848
2849
2850def _make_warm_bridge(
2851 *,
2852 use_shared_ptp: bool | None,
2853 protocol: StreamingProtocol = StreamingProtocol.AIRPLAY2,
2854) -> SendspinAirPlayBridge:
2855 """
2856 Build a bridge holding a connected, anchored cli process on the given flag.
2857
2858 :param use_shared_ptp: The shared-PTP flag its process was spawned with,
2859 None for a process that carries no such decision.
2860 :param protocol: The streaming protocol the bridged player speaks.
2861 """
2862 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US, protocol=protocol)
2863 bridge._airplay_stream = _make_kept_stream()
2864 bridge._started = True
2865 bridge._use_shared_ptp = use_shared_ptp
2866 return bridge
2867
2868
2869def test_a_regrouped_warm_process_is_not_reused() -> None:
2870 """
2871 A process whose flag no longer matches its group has to be respawned.
2872
2873 The flag is baked in at spawn, so reusing the process would keep the bridge
2874 on the clock its old group ran on. Its start lead has to report the cold
2875 figure too, or the respawn lands past the audio Sendspin already scheduled.
2876 """
2877 regrouped = _make_warm_bridge(use_shared_ptp=False)
2878 playing = _make_warm_bridge(use_shared_ptp=True)
2879 _group_bridges(regrouped, playing, daemon_ready=True)
2880 regrouped._bridge_role = MagicMock()
2881
2882 assert regrouped._stream_is_warm_eligible() is True
2883 assert regrouped._can_reuse_stream_warm() is False
2884
2885 regrouped._refresh_bridge_timing()
2886
2887 regrouped._bridge_role.set_timing.assert_called_once_with(
2888 required_lead_time_ms=BRIDGE_COLD_START_LEAD_MS, min_buffer_ms=BRIDGE_MIN_BUFFER_MS
2889 )
2890
2891
2892def test_a_warm_process_matching_its_group_is_reused() -> None:
2893 """A group already on the process's flag costs it no respawn."""
2894 warm = _make_warm_bridge(use_shared_ptp=True)
2895 playing = _make_warm_bridge(use_shared_ptp=True)
2896 _group_bridges(warm, playing, daemon_ready=False)
2897
2898 assert warm._can_reuse_stream_warm() is True
2899
2900
2901def test_a_group_without_a_live_decision_reuses_the_warm_process() -> None:
2902 """
2903 A bridge whose group has no other live decision keeps its process.
2904
2905 Its own process is the group's decision, so a daemon that changed state
2906 since must not churn the transport on every track change.
2907 """
2908 solo = _make_warm_bridge(use_shared_ptp=True)
2909 _group_bridges(solo, daemon_ready=False)
2910
2911 assert solo._can_reuse_stream_warm() is True
2912
2913
2914def test_a_raop_member_keeps_its_warm_process_beside_an_ap2_member() -> None:
2915 """
2916 A RAOP process is never respawned over a group's shared-clock decision.
2917
2918 It carries no such decision of its own, and no respawn could give it one, so
2919 comparing it against an AirPlay 2 sibling's would cost the group a cold
2920 reconnect (and its longer start lead) on every track change for nothing.
2921 """
2922 raop = _make_warm_bridge(use_shared_ptp=None, protocol=StreamingProtocol.RAOP)
2923 ap2 = _make_warm_bridge(use_shared_ptp=True)
2924 _group_bridges(raop, ap2, daemon_ready=True)
2925 raop._bridge_role = MagicMock()
2926
2927 assert raop._can_reuse_stream_warm() is True
2928
2929 raop._refresh_bridge_timing()
2930
2931 raop._bridge_role.set_timing.assert_called_once_with(
2932 required_lead_time_ms=BRIDGE_WARM_START_LEAD_MS, min_buffer_ms=BRIDGE_MIN_BUFFER_MS
2933 )
2934
2935
2936async def test_the_real_chunk_path_records_the_decision_it_spawns_with() -> None:
2937 """
2938 Driving the bridge the way Sendspin does still records what the CLI got.
2939
2940 The start path tells whether it still owns the bridge by comparing itself
2941 against the task handle the chunk handler publishes, so the start task must
2942 not run before that handle is set. Started eagerly it would read None on its
2943 very first check and give up as if a newer start had claimed the bridge.
2944 """
2945 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
2946 _group_bridges(bridge, daemon_ready=True)
2947 stream = _make_anchor_stream(ack=UNIX_NOW_MS + COLD_LEAD_MS)
2948 started: list[asyncio.Task[None]] = []
2949
2950 async def connect(_use_shared_ptp: bool | None) -> None:
2951 # a real connect does I/O, so the task suspends here
2952 await asyncio.sleep(0)
2953
2954 stream.connect = AsyncMock(side_effect=connect)
2955
2956 loop = asyncio.get_running_loop()
2957
2958 def create_task(
2959 coro: Coroutine[None, None, None], *, eager_start: bool = True, **_kwargs: object
2960 ) -> asyncio.Task[None]:
2961 # mirrors mass.create_task, whose default eager start would run the
2962 # coroutine to its first await before this returns
2963 task = asyncio.Task(coro, loop=loop, eager_start=eager_start)
2964 started.append(task)
2965 return task
2966
2967 cast("MagicMock", bridge.mass).create_task = create_task
2968
2969 with (
2970 patch(
2971 "music_assistant.providers.airplay.sendspin_bridge.AirPlayStream",
2972 return_value=stream,
2973 ),
2974 patch(
2975 "music_assistant.providers.airplay.sendspin_bridge.time.time",
2976 return_value=UNIX_NOW_S,
2977 ),
2978 ):
2979 bridge._on_audio_chunk(_pcm_chunk(SENDSPIN_EPOCH_US + COLD_LEAD_MS * 1_000))
2980 await asyncio.gather(*started)
2981
2982 stream.connect.assert_awaited_once_with(True)
2983 assert bridge.active_shared_ptp is True
2984
2985
2986@pytest.mark.parametrize("arm", ["sendspin_stream_start", "transport_restart"])
2987def test_a_released_process_stops_deciding_for_its_group(arm: str) -> None:
2988 """
2989 A process the bridge is about to tear down no longer speaks for its group.
2990
2991 Arming the bridge for its next stream happens well before that stream
2992 resolves, so a decision left over from the released process would be handed
2993 to a sibling resolving in between - and after a regroup it is the wrong one.
2994 """
2995 regrouped = _make_warm_bridge(use_shared_ptp=False)
2996 playing = _make_warm_bridge(use_shared_ptp=True)
2997 _group_bridges(regrouped, playing, daemon_ready=True)
2998
2999 if arm == "sendspin_stream_start":
3000 regrouped._on_stream_start(MagicMock())
3001 else:
3002 regrouped._on_bridge_stream_start()
3003
3004 assert regrouped._is_streaming is True
3005 assert regrouped.active_shared_ptp is None
3006 # the sibling still holding a live process keeps deciding for the group
3007 assert playing.active_shared_ptp is True
3008
3009
3010def test_an_abandoned_process_stops_deciding_for_its_group() -> None:
3011 """Giving up on a transport takes its decision out of the group with it."""
3012 abandoned = _make_warm_bridge(use_shared_ptp=True)
3013 _group_bridges(abandoned, daemon_ready=True)
3014
3015 abandoned._abandon_streaming()
3016
3017 assert abandoned._use_shared_ptp is None
3018 assert abandoned.active_shared_ptp is None
3019
3020
3021# --- Volume/mute moved on the AirPlay side, fed back into the bridge role ------
3022
3023
3024def _make_bridge_with_role(
3025 volume: int | None = 40, muted: bool = False
3026) -> tuple[SendspinAirPlayBridge, BridgePlayerRole]:
3027 """Build a bridge whose (real) role is wired to its mocked AirPlay player."""
3028 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
3029 role = BridgePlayerRole(client=MagicMock())
3030 role.set_callbacks(
3031 on_audio_chunk=bridge._on_audio_chunk,
3032 on_volume_change=bridge._on_volume_change,
3033 on_mute_change=bridge._on_mute_change,
3034 on_stream_start=bridge._on_bridge_stream_start,
3035 on_stream_end=bridge._on_bridge_stream_end,
3036 initial_volume=volume or 25,
3037 initial_muted=muted,
3038 )
3039 bridge._bridge_role = role
3040 player = cast("MagicMock", bridge.airplay_player)
3041 player.volume_level = volume
3042 player.volume_muted = muted
3043 return bridge, role
3044
3045
3046async def test_registration_seeds_the_role_with_the_state_the_speaker_is_in() -> None:
3047 """
3048 A bridge (re)registering adopts the volume and mute the speaker is already at.
3049
3050 A bridge is torn down and rebuilt on a config change, so starting from a
3051 fixed unmuted default would re-create the divergence on every rebuild.
3052 """
3053 bridge = _make_bridge(clock_now_us=SENDSPIN_EPOCH_US)
3054 player = cast("MagicMock", bridge.airplay_player)
3055 player.volume_level = 35
3056 player.volume_muted = True
3057 role = MagicMock()
3058 server = cast("MagicMock", bridge.sendspin_server)
3059 server.register_external_player.return_value.roles_by_family.return_value = [role]
3060
3061 await bridge.start()
3062
3063 assert role.set_callbacks.call_args.kwargs["initial_volume"] == 35
3064 assert role.set_callbacks.call_args.kwargs["initial_muted"] is True
3065
3066
3067def test_device_volume_feedback_reaches_the_visible_player() -> None:
3068 """
3069 A volume the device reports itself is adopted by the role, not sent back to it.
3070
3071 The role is what the parent's volume resolves to while the bridge streams, so
3072 it has to follow the speaker; forwarding it would hand the device back the
3073 value it just reported.
3074 """
3075 bridge, role = _make_bridge_with_role(volume=40)
3076 player = cast("MagicMock", bridge.airplay_player)
3077 player.volume_level = 55
3078
3079 bridge.sync_role_volume_state()
3080
3081 assert role.get_player_volume() == 55
3082 player.volume_set.assert_not_called()
3083
3084
3085def test_a_mute_latched_on_the_airplay_side_shows_through_the_bridge() -> None:
3086 """
3087 A mute applied on the AirPlay side is visible on the Sendspin player.
3088
3089 While it is latched the AirPlay player swallows every volume command, so a
3090 bridge still reporting the speaker as unmuted leaves nothing to explain why
3091 it is silent - or to unmute it with.
3092 """
3093 bridge, role = _make_bridge_with_role(muted=False)
3094 player = cast("MagicMock", bridge.airplay_player)
3095 player.volume_muted = True
3096
3097 bridge.sync_role_volume_state()
3098
3099 assert role.get_player_muted() is True
3100 player.volume_mute.assert_not_called()
3101
3102
3103def test_an_airplay_volume_change_is_routed_to_that_player_s_bridge() -> None:
3104 """A state update carrying a volume or mute change lands on the right bridge."""
3105 bridge, role = _make_bridge_with_role(volume=40)
3106 manager = _bridge_manager_for(bridge)
3107 player = cast("MagicMock", bridge.airplay_player)
3108 player.volume_level = 55
3109 player.volume_muted = True
3110
3111 manager._on_player_state_updated(
3112 player, {"volume_level": (40, 55), "volume_muted": (False, True)}
3113 )
3114
3115 assert role.get_player_volume() == 55
3116 assert role.get_player_muted() is True
3117
3118
3119def test_state_updates_without_a_volume_change_leave_the_role_alone() -> None:
3120 """
3121 Every player's state update passes through, so unrelated ones do no work.
3122
3123 The callback runs for the whole player graph on every tick, including the
3124 position updates of a playing queue.
3125 """
3126 bridge, role = _make_bridge_with_role(volume=40)
3127 manager = _bridge_manager_for(bridge)
3128 player = cast("MagicMock", bridge.airplay_player)
3129 player.volume_level = 55
3130
3131 manager._on_player_state_updated(player, {"playback_state": ("idle", "playing")})
3132
3133 assert role.get_player_volume() == 40
3134
3135
3136def test_a_player_without_a_bridge_is_ignored() -> None:
3137 """A volume change on a player this manager knows nothing about leaves bridges alone."""
3138 bridge, role = _make_bridge_with_role(volume=40)
3139 manager = _bridge_manager_for(bridge)
3140 other_player = MagicMock()
3141 other_player.player_id = "ap0011223344ff"
3142 other_player.volume_level = 55
3143
3144 manager._on_player_state_updated(other_player, {"volume_level": (40, 55)})
3145
3146 assert role.get_player_volume() == 40
3147
3148
3149def test_the_manager_listens_for_the_state_updates_it_routes() -> None:
3150 """
3151 The bridged AirPlay players are watched for the whole life of the manager.
3152
3153 A protocol player emits no PLAYER_UPDATED event, so the controller's internal
3154 state-update subscription is the only way their volume changes are seen.
3155 """
3156 manager = SendspinBridgeManager(MagicMock())
3157
3158 cast("MagicMock", manager.mass).players.subscribe_player_state_update.assert_called_once_with(
3159 manager._on_player_state_updated
3160 )
3161
3162
3163def test_a_volume_set_through_the_bridge_settles_in_one_pass() -> None:
3164 """
3165 A volume coming down from Sendspin is not announced again on its way back.
3166
3167 The player ends up holding what the role handed it, so reading that state back
3168 has to compare equal - otherwise every command would bounce between the two.
3169 """
3170 bridge, role = _make_bridge_with_role(volume=40)
3171 player = cast("MagicMock", bridge.airplay_player)
3172 client = cast("MagicMock", role._client)
3173
3174 role.set_player_volume(70)
3175 player.volume_set.assert_called_once_with(70)
3176 # the AirPlay player records the level it was handed
3177 player.volume_level = 70
3178 client.reset_mock()
3179
3180 bridge.sync_role_volume_state()
3181
3182 assert role.get_player_volume() == 70
3183 client._signal_event.assert_not_called()
3184
3185
3186def test_a_volume_from_the_role_is_not_resolved_a_second_time() -> None:
3187 """
3188 A volume the role delivers reaches the speaker as-is, not via the controller.
3189
3190 The level arrives on the device's own scale, already resolved and scaled where
3191 it came from Music Assistant, so sending it back through the controller would
3192 resolve it a second time and return over the same route.
3193 """
3194 bridge, role = _make_bridge_with_role(volume=40)
3195 player = cast("MagicMock", bridge.airplay_player)
3196 players_ctrl = cast("MagicMock", bridge.mass).players
3197
3198 role.set_player_volume(60)
3199 role.set_player_mute(True)
3200
3201 player.volume_set.assert_called_once_with(60)
3202 player.volume_mute.assert_called_once_with(True)
3203 players_ctrl.cmd_volume_set.assert_not_called()
3204 players_ctrl.cmd_volume_mute.assert_not_called()
3205