/
/
/
1"""State synchronization handler."""
2
3from __future__ import annotations
4
5import asyncio
6import contextlib
7import time
8from functools import partial
9from typing import TYPE_CHECKING, Any
10
11from music_assistant_models.enums import PlaybackState
12from music_assistant_models.errors import PlayerCommandFailed
13from music_assistant_models.player import DeviceInfo
14from pywam.lib.api_call import ApiCall
15from pywam.lib.exceptions import PywamError
16from pywam.speaker import Speaker
17
18from music_assistant.constants import ATTR_ANNOUNCEMENT_IN_PROGRESS
19from music_assistant.providers.samsung_wam.consts import MANUFACTURER_NAME
20from music_assistant.providers.samsung_wam.features.base import (
21 WamPlayerFeatureBase,
22 handle_pywam_errors,
23 retry_command,
24)
25from music_assistant.providers.samsung_wam.features.playback.models import WamSource
26
27from .consts import HEALTH_CHECK_TIMEOUT
28from .mapper import StateSyncMapper
29
30if TYPE_CHECKING:
31 from .models import WamSpeakerAttributes
32
33
34def get_speaker_status() -> ApiCall:
35 """(UIC) Get speaker status."""
36 return ApiCall(api_type="UIC", method="GetSpeakerStatus", expected_response="SpeakerStatus")
37
38
39def get_current_play_time() -> ApiCall:
40 """(UIC) Get current play time."""
41 return ApiCall(api_type="UIC", method="GetCurrentPlayTime", expected_response="MusicPlayTime")
42
43
44class StateSyncHandler(WamPlayerFeatureBase):
45 """Encapsulates polling, state updates, and connection recovery."""
46
47 _terminating_wifi_stream: bool = False
48
49 def apply_initial_state(self, attrs: WamSpeakerAttributes) -> None:
50 """
51 Apply an initial state snapshot to the player.
52
53 :param attrs: The speaker attributes to apply.
54 """
55 self.player._attr_device_info = DeviceInfo(
56 model=attrs.model or "Unknown",
57 manufacturer=MANUFACTURER_NAME,
58 software_version=attrs.software_version,
59 identifiers=self.player._attr_device_info.identifiers,
60 )
61 self.player._attr_name = attrs.name
62 self.suppress_speaker_status_events(self.speaker)
63 self.player._attr_available = True
64
65 def subscribe_speaker_events(self) -> None:
66 """Register the speaker event subscriber."""
67 if self.speaker:
68 self.speaker.events.register_subscriber(partial(self.on_speaker_event), info_level=0)
69
70 async def poll(self) -> None:
71 """Poll the player for state updates and handle connection recovery."""
72 try:
73 async with self.player.connection_lock:
74 if not self.player.connected:
75 self.logger.debug("Poller found disconnected speaker; attempting reconnect")
76 self._mark_player_unavailable()
77 await self._reconnect_speaker()
78
79 await self.check_status()
80
81 if self.player.playback_state == PlaybackState.PLAYING:
82 await self.update_play_time()
83
84 if not self.player.available:
85 self.logger.info("Player %s is back online", self.player.log_name)
86 self.player._attr_available = True
87
88 self.player.update_state()
89
90 except (ConnectionError, PywamError, TimeoutError, PlayerCommandFailed) as err:
91 self.logger.debug(
92 "Poll failed for %s (%s); dropping connection",
93 self.player.log_name,
94 err.__class__.__name__,
95 )
96 await self.disconnect_speaker()
97 self._mark_player_unavailable()
98
99 def _mark_player_unavailable(self) -> None:
100 """Mark the player as unavailable and reset playback state."""
101 if not self.player.available:
102 return
103 self.player._attr_available = False
104 self.player._attr_playback_state = PlaybackState.IDLE
105 self.player.stream_active = False
106 self.player.update_state()
107
108 @retry_command()
109 @handle_pywam_errors
110 async def check_status(self) -> None:
111 """Check whether the speaker is reachable."""
112 if not self.player.connected:
113 raise ConnectionError("Speaker is not connected")
114 async with asyncio.timeout(HEALTH_CHECK_TIMEOUT):
115 await self.speaker.client.request(get_speaker_status())
116
117 @retry_command()
118 @handle_pywam_errors
119 async def update_play_time(self) -> None:
120 """Fetch and apply the current play time from the speaker."""
121 if not self.player.connected:
122 return
123
124 async with asyncio.timeout(HEALTH_CHECK_TIMEOUT):
125 response = await self.speaker.client.request(get_current_play_time())
126
127 if not response or not hasattr(response, "data") or "playtime" not in response.data:
128 return
129
130 try:
131 playtime = float(response.data["playtime"])
132 self.player._attr_elapsed_time = playtime
133 self.player._attr_elapsed_time_last_updated = time.time()
134 except ValueError, TypeError:
135 self.logger.debug("Failed to parse playtime from response: %s", response.data)
136
137 def on_speaker_event(self, event: Any = None) -> None:
138 """
139 Handle a state update event broadcast by the speaker.
140
141 :param event: The payload data emitted from the speaker.
142 """
143 # Skip during announcements: playtime events would corrupt elapsed time tracking
144 # and cause the resume after the announcement to seek to the wrong position
145 if (
146 event
147 and hasattr(event, "data")
148 and isinstance(event.data, dict)
149 and "playtime" in event.data
150 and not self.player.extra_data.get(ATTR_ANNOUNCEMENT_IN_PROGRESS)
151 ):
152 try:
153 playtime = float(event.data["playtime"])
154 self.player._attr_elapsed_time = playtime
155 self.player._attr_elapsed_time_last_updated = time.time()
156 except ValueError, TypeError:
157 pass
158
159 self.refresh_state(notify_provider=True)
160
161 def refresh_state(self, notify_provider: bool = False) -> None:
162 """
163 Re-apply current speaker state to the player.
164
165 :param notify_provider: Trigger an event signal if True.
166 """
167 if not self.player.connected:
168 return
169
170 speaker_attrs = StateSyncMapper.create_speaker_attributes(self.speaker)
171 group_children = self.player.prov.groups.states.get(self.player.player_id, set())
172
173 # If the speaker has externally switched away from the Wi-Fi stream (e.g. a phone
174 # connecting via Bluetooth), clear the stream flag and capture the new source
175 external_source = None
176 if (
177 self.player.stream_active
178 and speaker_attrs.source
179 and speaker_attrs.source not in (WamSource.WIFI, "Unknown")
180 ):
181 self.player.stream_active = False
182 external_source = speaker_attrs.source
183
184 queue_id = None
185 if queue := self.mass.player_queues.get(self.player.player_id):
186 queue_id = queue.queue_id
187
188 StateSyncMapper.apply_attributes_to_player(
189 player=self.player,
190 speaker_attrs=speaker_attrs,
191 group_children=group_children,
192 stream_active=self.player.stream_active,
193 )
194
195 self.player.update_state()
196
197 # Must be set after the state update to prevent the source display from briefly reverting
198 if external_source:
199 self.player.set_active_mass_source(external_source)
200 # Stopping the queue alone won't terminate the stream as the speaker keeps
201 # its HTTP connection alive in the background
202 if queue_id:
203 self.mass.create_task(self.mass.player_queues.stop(queue_id))
204 self.mass.create_task(self._terminate_wifi_stream(external_source))
205
206 if notify_provider:
207 self.player.signal_state_update_event()
208 self.player.prov.groups.on_player_state_changed(self.player)
209
210 async def ensure_speaker_connected(self) -> None:
211 """Ensure the speaker is connected, reconnecting if necessary."""
212 async with self.player.connection_lock:
213 if not self.player.connected:
214 await self._reconnect_speaker()
215
216 async def _terminate_wifi_stream(self, target_source: str) -> None:
217 """
218 Terminate the audio stream after an external source switch.
219
220 :param target_source: The source the speaker should end up on (e.g. 'Bluetooth').
221 """
222 # The speaker keeps its HTTP connection to the stream URL open in the background
223 # even when on an external source â this sequence forces it to close
224 if self._terminating_wifi_stream:
225 return
226 if not self.player.connected or self.player.synced_to:
227 return
228 if self.mass.players.get_player(self.player.player_id) is not self.player:
229 return
230
231 self._terminating_wifi_stream = True
232 try:
233 try:
234 await self.speaker.select_source(str(WamSource.WIFI))
235 await self.player.await_state_change(
236 lambda: self.speaker.attribute.source == str(WamSource.WIFI),
237 2.0,
238 )
239 await self.speaker.cmd_pause()
240 except (PywamError, ConnectionError, TimeoutError) as err:
241 self.logger.debug(
242 "Stream termination (Wi-Fi phase) failed for %s: %s",
243 self.player.log_name,
244 err,
245 )
246
247 try:
248 await self.speaker.select_source(target_source)
249 except (PywamError, ConnectionError, TimeoutError) as err:
250 self.logger.debug(
251 "Stream termination (source restore) failed for %s: %s",
252 self.player.log_name,
253 err,
254 )
255 finally:
256 # Restore the target source to cancel any timer that fired while the
257 # speaker was briefly back on Wi-Fi
258 self.player.set_active_mass_source(target_source)
259 finally:
260 self._terminating_wifi_stream = False
261
262 @staticmethod
263 def suppress_speaker_status_events(speaker: Speaker) -> None:
264 """
265 Replace the SpeakerStatus event handler with a no-op to prevent log noise.
266
267 :param speaker: The Speaker instance to patch.
268 """
269
270 def _handle_speaker_status_event(_event: Any) -> bool:
271 return False
272
273 speaker.events.event_SpeakerStatus = _handle_speaker_status_event
274
275 async def _reconnect_speaker(self) -> None:
276 """Reconnect to the speaker using a fresh client."""
277 await self.disconnect_speaker()
278
279 new_speaker = Speaker(self.player.ip_address)
280 await new_speaker.connect()
281
282 try:
283 await new_speaker.update()
284
285 attrs = StateSyncMapper.create_speaker_attributes(new_speaker)
286 if not attrs.mac or attrs.mac != self.player.player_id:
287 raise ConnectionError("MAC address verification failed after reconnect")
288
289 # Guard against the case where this player instance has been replaced (e.g. after
290 # a provider restart). Continuing would leave a dangling pywam connection open
291 if self.mass.players.get_player(self.player.player_id) is not self.player:
292 raise ConnectionError("Player instance superseded; aborting reconnect")
293
294 self.suppress_speaker_status_events(new_speaker)
295 self.player.speaker = new_speaker
296 self.subscribe_speaker_events()
297
298 except Exception:
299 await new_speaker.disconnect()
300 raise
301
302 async def unload(self) -> None:
303 """Handle cleanup when the player is being permanently unloaded."""
304 self.player.prov.groups.unregister_player(self.player)
305 await self.disconnect_speaker()
306
307 async def disconnect_speaker(self) -> None:
308 """Safely disconnect the underlying pywam speaker."""
309 if self.speaker:
310 with contextlib.suppress(Exception):
311 await self.speaker.disconnect()
312