/
/
1"""Runtime execution mixin for AI Radio."""
2# mypy: disable-error-code="attr-defined"
3
4from __future__ import annotations
5
6import asyncio
7import datetime
8import logging
9import random
10import time
11from collections import defaultdict
12from copy import deepcopy
13from pathlib import Path
14from typing import TYPE_CHECKING, Any, cast
15from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
16
17from aiohttp import ClientTimeout
18from music_assistant_models.enums import (
19 EventType,
20 ImageType,
21 MediaType,
22 PlaybackState,
23)
24from music_assistant_models.errors import MusicAssistantError
25from music_assistant_models.media_items import (
26 MediaItemImage,
27 ProviderMapping,
28 SoundEffect,
29 UniqueList,
30)
31
32from music_assistant.controllers.player_queues.helpers import build_queue_item
33from music_assistant.helpers.datetime import now, utc
34from music_assistant.helpers.json import json_loads
35from music_assistant.helpers.plugin_engines import resolve_ai_engine, resolve_tts_engine
36from music_assistant.helpers.uri import create_uri
37
38from .constants import (
39 AI_QUERY_TIMEOUT_SECONDS,
40 ATTR_HOST_ID,
41 ATTR_MAX_CHARS,
42 ATTR_PROMPT,
43 ATTR_SESSION_ID,
44 ATTR_STATION_ID,
45 ATTR_WEATHER_REQUIRED,
46 ATTR_WEB_SEARCH_MODE,
47 CONF_AI_ENGINE,
48 CONF_TIMEZONE,
49 CONF_TTS_ENGINE,
50 CONF_WEATHER_CITY,
51 CONF_WEATHER_COUNTRY,
52 CONF_WEATHER_PROVIDER,
53 CONF_WEATHER_TIMEOUT,
54 DEFAULT_LLM_INSTRUCTIONS,
55 DEFAULT_WEATHER_PROVIDER,
56 DEFAULT_WEATHER_TIMEOUT_SECONDS,
57 DEFERRED_PLACEHOLDERS,
58 FAHRENHEIT_COUNTRY_CODES,
59 SHOW_START_TIMEOUT_SECONDS,
60 TTS_PRONUNCIATION_INSTRUCTIONS,
61 VALID_WEB_SEARCH_MODES,
62 WEATHER_PLACEHOLDER_TOKENS,
63 WEB_SEARCH_MODE_RANK,
64)
65from .helpers import (
66 build_slots,
67 coerce_float,
68 coerce_int,
69 is_empty_section,
70 pick_weighted_choice,
71 slugify,
72 track_songinfo,
73 utc_now_iso,
74)
75from .models import (
76 PlannedSection,
77 SessionState,
78 Slot,
79)
80
81if TYPE_CHECKING:
82 from music_assistant_models.config_entries import ConfigValueType, ProviderConfig
83 from music_assistant_models.event import MassEvent
84 from music_assistant_models.media_items import PlayableMediaItemType
85 from music_assistant_models.queue_item import QueueItem
86
87 from music_assistant.mass import MusicAssistant
88 from music_assistant.models.plugin import AIEngine, TTSEngine
89
90
91# the sticky queue DJ re-plans on every queue change, so an uncached forecast lookup would
92# add two HTTP round trips to each one. Weather does not move meaningfully within this window
93WEATHER_TOKENS_CACHE_SECONDS = 300
94
95
96class AIRadioRuntimeMixin:
97 """Mixin with all runtime logic for AI Radio runs."""
98
99 # (fetched_at, tokens) of the last weather lookup, shared by the show and DJ paths
100 _weather_tokens_cache: tuple[float, dict[str, str]] | None = None
101
102 if TYPE_CHECKING:
103 mass: MusicAssistant
104 config: ProviderConfig
105 logger: logging.Logger
106 _sessions: dict[str, SessionState]
107
108 def get_setup_value(self, key: str, default: ConfigValueType = None) -> ConfigValueType:
109 """Return a value collected by this provider's setup flow."""
110
111 def _schedule_replan(self, queue_id: str) -> None:
112 """Request a replan pass for the given queue."""
113
114 async def set_queue_dj(self, queue_id: str, host_id: str | None) -> dict[str, str]:
115 """Enable, switch or disable the sticky AI DJ on a queue."""
116
117 def _set_session_progress(
118 self,
119 session: SessionState,
120 phase: str,
121 **details: Any,
122 ) -> None:
123 """Set progress payload with a stable phase key."""
124 session.progress = {
125 "phase": phase,
126 # Keep legacy key for compatibility with older UI code.
127 "step": phase,
128 **details,
129 }
130
131 def _build_program(self, station: dict[str, Any], host: dict[str, Any]) -> dict[str, Any]:
132 """Merge a station and its host into the dict the planner consumes."""
133 sections, missing = self._materialize_sections(list(host.get("section_ids", [])))
134 if missing:
135 raise MusicAssistantError(
136 f"Host references unknown sections: {', '.join(sorted(set(missing)))}"
137 )
138 return {
139 **deepcopy(station),
140 "host_id": str(host.get("id", "")),
141 "instructions": str(host.get("instructions", "")),
142 "tts_engine": str(host.get("tts_engine", "")),
143 "language": str(host.get("language", "")),
144 "options": deepcopy(host.get("options", {})),
145 "sections": sections,
146 "section_order": deepcopy(host.get("section_order", [])),
147 "merge_section_id": str(host.get("merge_section_id", "")),
148 }
149
150 async def _run_session(self, session_id: str, program: dict[str, Any]) -> None:
151 """Run one session in the background."""
152 session = self._sessions[session_id]
153 session.started_at = utc_now_iso()
154 self.logger.info(
155 "AI Radio run started: session=%s station=%s",
156 session.session_id,
157 session.station_id,
158 )
159 try:
160 result = await self._run_show(session, program)
161 session.result = result
162 queue_stopped = result.get("ended_reason") == "queue_stopped"
163 session.status = "stopped" if queue_stopped else "completed"
164 self.logger.info(
165 "AI Radio run %s: session=%s station=%s",
166 session.status,
167 session.session_id,
168 session.station_id,
169 )
170 except asyncio.CancelledError:
171 session.status = "stopped"
172 self.logger.info(
173 "AI Radio run cancelled: session=%s station=%s",
174 session.session_id,
175 session.station_id,
176 )
177 raise
178 except Exception as err:
179 session.status = "failed"
180 session.error = str(err).strip() or err.__class__.__name__
181 self.logger.exception("AI Radio session failed")
182 finally:
183 session.ended_at = utc_now_iso()
184 # a show session blocks queue DJ replans while it runs, so ending it must
185 # re-arm the DJ itself instead of waiting on the next queue change
186 if session.queue_id:
187 self._schedule_replan(session.queue_id)
188
189 async def _run_show(
190 self,
191 session: SessionState,
192 program: dict[str, Any],
193 ) -> dict[str, Any]:
194 """Plan and queue the whole show in one pass, then start playback."""
195 program = deepcopy(program)
196 self.logger.debug(
197 "Show starting for station '%s' (%s)",
198 program.get("name", "AI Radio"),
199 program.get("id", ""),
200 )
201 self._set_session_progress(session, "fetch_source_tracks")
202 # runtime_tokens only feeds the require_placeholders_present guards below; its
203 # resolved text is discarded here and re-fetched fresh when each clip renders
204 runtime_tokens = await self._prepare_runtime_tokens(program)
205 player_id = str(program.get("default_player_id") or "").strip()
206 if not player_id:
207 raise MusicAssistantError("AI Radio requires a target player")
208 if not self.mass.players.get_player(player_id):
209 raise MusicAssistantError(f"Unknown target player: {player_id}")
210
211 tracks, playlist_name = await self._fetch_source_tracks(program)
212 tracks = self._apply_source_shuffle(tracks, program)
213 tracks = self._apply_track_duration_limit(tracks, program)
214 if not tracks:
215 raise MusicAssistantError("No source tracks available after applying station limits")
216
217 # a grouped player plays from the group leader's queue, so resolve the
218 # active queue up front and target that one for queueing and polling
219 queue_id = player_id
220 active_queue = self.mass.player_queues.get_active_queue(player_id)
221 if active_queue is not None:
222 queue_id = str(active_queue.queue_id)
223 # a queue runs one host at a time; the show is now that host, so any sticky
224 # DJ assignment on the queue is cleared before the show takes it over
225 await self.set_queue_dj(queue_id, None)
226 self.mass.player_queues.clear(queue_id)
227 session.queue_id = queue_id
228
229 # a shuffled queue reorders each batch, scattering sections away from their tracks
230 await self.mass.player_queues.set_shuffle(queue_id, False)
231
232 cumulative_minutes = [0.0]
233 for track in tracks:
234 duration = track.get("duration")
235 seconds = (
236 float(duration) if isinstance(duration, (int, float)) and duration > 0 else 210.0
237 )
238 cumulative_minutes.append(cumulative_minutes[-1] + (seconds / 60.0))
239
240 self._set_session_progress(session, "planning_sections", total_tracks=len(tracks))
241 planned_sections, _history = self._plan_sections(
242 session_id=session.session_id,
243 tracks=tracks,
244 program=program,
245 track_index_offset=0,
246 minute_offset=0.0,
247 history_state={},
248 allowed_slot_when=None,
249 runtime_tokens=runtime_tokens,
250 )
251 queue_items = self._compose_queue_items(
252 queue_id=queue_id,
253 session=session,
254 program=program,
255 tracks=tracks,
256 sections=planned_sections,
257 )
258 if not queue_items:
259 raise MusicAssistantError("No queue entries were generated")
260
261 self._set_session_progress(
262 session,
263 "initializing_queue",
264 total_tracks=len(tracks),
265 queue_entries=len(queue_items),
266 queue_id=queue_id,
267 )
268 # load() stages the items without starting playback, so every clip already carries its
269 # prompt by the time anything can ask for its audio
270 await self.mass.player_queues.load(
271 queue_id,
272 queue_items=queue_items,
273 keep_remaining=False,
274 keep_played=False,
275 shuffle=False,
276 )
277 await self.mass.player_queues.play_index(queue_id, 0)
278 self._set_session_progress(
279 session,
280 "running",
281 total_tracks=len(tracks),
282 queue_entries=len(queue_items),
283 queue_id=queue_id,
284 )
285 has_clips = any(ATTR_SESSION_ID in item.extra_attributes for item in queue_items)
286 ended_reason = await self._await_show_end(
287 session, queue_id, len(queue_items) - 1, has_clips=has_clips
288 )
289 return {
290 "ended_reason": ended_reason,
291 "source_playlist_name": playlist_name,
292 "source_tracks": len(tracks),
293 "queue_id": queue_id,
294 "queue_entries": len(queue_items),
295 "planned_sections": len(planned_sections),
296 "skipped_sections": session.skipped_sections,
297 }
298
299 async def _await_show_end(
300 self, session: SessionState, queue_id: str, last_index: int, *, has_clips: bool
301 ) -> str:
302 """
303 Block until this session's show is over and report why it ended.
304
305 :param session: The session whose clips are in the queue.
306 :param queue_id: The queue playing the show.
307 :param last_index: Queue index of the final entry this session enqueued.
308 :param has_clips: Whether this run enqueued any AI Radio clips at all. A clip-free
309 show (every section was skipped by its rules) must not be mistaken for one whose
310 clips were cleared out from under it, so that rule is skipped entirely here.
311 :return: ``"source_exhausted"`` when the show played out, ``"queue_stopped"`` when the
312 queue was stopped or taken over before reaching the end.
313 :raises MusicAssistantError: if playback never starts within
314 :data:`SHOW_START_TIMEOUT_SECONDS`.
315 """
316 finished = asyncio.Event()
317 playback_started = asyncio.Event()
318 # a queue that has not started yet must never be mistaken for a stopped one
319 playback_seen = False
320 ended_reason = "queue_stopped"
321
322 def _check_show_state() -> None:
323 nonlocal playback_seen, ended_reason
324 queue = self.mass.player_queues.get(queue_id)
325 if queue is None:
326 finished.set()
327 return
328 if queue.state in (PlaybackState.PLAYING, PlaybackState.PAUSED):
329 playback_seen = True
330 playback_started.set()
331 if has_clips and not self._session_has_clips(queue_id, session.session_id):
332 finished.set()
333 return
334 if not playback_seen or queue.state != PlaybackState.IDLE:
335 return
336 # playing out and being stopped both end IDLE, so position is the discriminator
337 current_index = queue.current_index
338 if current_index is not None and current_index >= last_index:
339 ended_reason = "source_exhausted"
340 self.logger.info(
341 "Queue %s went idle at index %s of %s, ending show (%s)",
342 queue_id,
343 current_index,
344 last_index,
345 ended_reason,
346 )
347 finished.set()
348
349 def _on_queue_event(_event: MassEvent) -> None:
350 _check_show_state()
351
352 unsubscribe = self.mass.subscribe(
353 _on_queue_event,
354 (EventType.QUEUE_UPDATED, EventType.QUEUE_ITEMS_UPDATED, EventType.PLAYER_REMOVED),
355 id_filter=queue_id,
356 )
357 try:
358 # the queue may already have gone away, or (for a show with clips) already lost
359 # them, by the time this subscribes; IDLE-after-playout still needs a fresh event,
360 # since playback_seen is not latched yet
361 _check_show_state()
362 await self._await_playback_start(playback_started, finished)
363 await finished.wait()
364 finally:
365 unsubscribe()
366 return ended_reason
367
368 async def _await_playback_start(
369 self, playback_started: asyncio.Event, finished: asyncio.Event
370 ) -> None:
371 """
372 Wait for the show to either start playing or end before it ever did.
373
374 :param playback_started: Set once the queue is first observed playing or paused.
375 :param finished: Set once the show is over, however that came about.
376 :raises MusicAssistantError: if neither happens within
377 :data:`SHOW_START_TIMEOUT_SECONDS`.
378 """
379 if playback_started.is_set() or finished.is_set():
380 return
381 # a player that never comes online (or whose clips all fail) must not pin this
382 # session's "running" status, and its max-concurrent-runs slot, forever
383 wait_tasks = (
384 asyncio.ensure_future(playback_started.wait()),
385 asyncio.ensure_future(finished.wait()),
386 )
387 try:
388 done, _pending = await asyncio.wait(
389 wait_tasks,
390 timeout=SHOW_START_TIMEOUT_SECONDS,
391 return_when=asyncio.FIRST_COMPLETED,
392 )
393 finally:
394 for task in wait_tasks:
395 if not task.done():
396 task.cancel()
397 if not done:
398 raise MusicAssistantError(
399 f"Playback did not start within {SHOW_START_TIMEOUT_SECONDS}s"
400 )
401
402 def _session_has_clips(self, queue_id: str, session_id: str) -> bool:
403 """Return whether any queue item still belongs to the given session."""
404 page_size = 500
405 offset = 0
406 while True:
407 page = self.mass.player_queues.items(queue_id, limit=page_size, offset=offset)
408 if not page:
409 return False
410 if any(item.extra_attributes.get(ATTR_SESSION_ID) == session_id for item in page):
411 return True
412 if len(page) < page_size:
413 return False
414 offset += page_size
415
416 async def _fetch_source_tracks(
417 self, station: dict[str, Any]
418 ) -> tuple[list[dict[str, Any]], str]:
419 """Load and normalize source playlist tracks."""
420 playlist_id = str(station.get("source_playlist_id", "")).strip()
421 provider = str(station.get("source_playlist_provider", "library")).strip() or "library"
422 if not playlist_id:
423 raise MusicAssistantError("Station is missing source_playlist_id")
424
425 playlist = await self.mass.music.playlists.get(playlist_id, provider)
426 playlist_name = playlist.name
427 tracks = [track async for track in self.mass.music.playlists.tracks(playlist_id, provider)]
428 normalized: list[dict[str, Any]] = []
429 for track in tracks:
430 artist = ""
431 track_artists = getattr(track, "artists", None)
432 if isinstance(track_artists, list) and track_artists:
433 artist = str(track_artists[0].name)
434 uri = await self._track_to_uri(track)
435 if not uri:
436 self.logger.warning(
437 "Skipping source track with no resolvable uri: %s - %s (item_id=%s)",
438 artist,
439 track.name,
440 track.item_id,
441 )
442 continue
443 normalized.append(
444 {
445 "index": len(normalized),
446 "item_id": track.item_id,
447 "name": track.name,
448 "artist": artist,
449 "songinfo": f"{artist} - {track.name}".strip(" -"),
450 "duration": track.duration,
451 "uri": uri,
452 "media_item": track,
453 }
454 )
455 return normalized, playlist_name
456
457 async def _track_to_uri(self, track: PlayableMediaItemType) -> str:
458 """Resolve a stable URI for a source track."""
459 if track.uri:
460 return track.uri
461 ordered_mappings = sorted(
462 track.provider_mappings,
463 key=lambda mapping: mapping.quality,
464 reverse=True,
465 )
466 for mapping in ordered_mappings:
467 if not mapping.available:
468 continue
469 return create_uri(MediaType.TRACK, mapping.provider_instance, mapping.item_id)
470 return ""
471
472 def _apply_source_shuffle(
473 self, tracks: list[dict[str, Any]], station: dict[str, Any]
474 ) -> list[dict[str, Any]]:
475 """Return the source tracks in random order when the station asks for it."""
476 if not station.get("shuffle_source_tracks", True) or not tracks:
477 return tracks
478 indices = list(range(len(tracks)))
479 random.Random().shuffle(indices)
480 result: list[dict[str, Any]] = []
481 for new_index, old_index in enumerate(indices):
482 updated = deepcopy(tracks[old_index])
483 updated["index"] = new_index
484 updated["source_index"] = old_index
485 result.append(updated)
486 self.logger.info("Shuffled %d source tracks", len(result))
487 return result
488
489 def _apply_track_duration_limit(
490 self, tracks: list[dict[str, Any]], station: dict[str, Any]
491 ) -> list[dict[str, Any]]:
492 """Truncate the given tracks to the configured playtime cap, preserving their order."""
493 max_duration = float(station.get("max_duration_minutes", 0) or 0)
494 if max_duration <= 0 or not tracks:
495 return tracks
496 chosen: list[int] = []
497 total_minutes = 0.0
498 for index, track in enumerate(tracks):
499 duration = track.get("duration")
500 seconds = (
501 float(duration) if isinstance(duration, (int, float)) and duration > 0 else 210.0
502 )
503 chosen.append(index)
504 total_minutes += seconds / 60.0
505 if total_minutes > max_duration:
506 break
507 result: list[dict[str, Any]] = []
508 for new_index, old_index in enumerate(chosen):
509 updated = deepcopy(tracks[old_index])
510 updated["index"] = new_index
511 updated["source_index"] = old_index
512 result.append(updated)
513 self.logger.info(
514 "Applied source playtime cap: %.1f min requested, %d -> %d tracks selected",
515 max_duration,
516 len(tracks),
517 len(result),
518 )
519 return result
520
521 def _plan_sections( # noqa: PLR0915
522 self,
523 session_id: str,
524 tracks: list[dict[str, Any]],
525 program: dict[str, Any],
526 track_index_offset: int,
527 minute_offset: float,
528 history_state: dict[str, list[tuple[int, float]]],
529 allowed_slot_when: list[str] | None,
530 runtime_tokens: dict[str, str],
531 decided_next_item_ids: set[str] | None = None,
532 ) -> tuple[list[PlannedSection], dict[str, list[tuple[int, float]]]]:
533 """Evaluate section rules and produce planning entries."""
534 sections = program.get("sections", [])
535 section_order = program.get("section_order", [])
536 if not isinstance(sections, list) or not sections:
537 raise MusicAssistantError("Station has no sections configured")
538 if not isinstance(section_order, list) or not section_order:
539 raise MusicAssistantError("Station has no section_order configured")
540
541 section_by_id = {
542 str(section.get("id", "")).strip(): section
543 for section in sections
544 if str(section.get("id", "")).strip()
545 }
546 slots = build_slots(tracks)
547 history = {section_id: list(events) for section_id, events in history_state.items()}
548 selected: list[tuple[str, Slot, dict[str, str]]] = []
549 rng = random.Random()
550
551 def slot_event(slot: Slot) -> tuple[int, float]:
552 song_local = slot.next_index if slot.next_index is not None else len(tracks)
553 return track_index_offset + song_local, minute_offset + slot.minute_mark
554
555 def register_event(section_id: str, slot: Slot) -> None:
556 if is_empty_section(section_id):
557 return
558 history.setdefault(section_id, []).append(slot_event(slot))
559
560 for slot in slots:
561 if allowed_slot_when and slot.when not in allowed_slot_when:
562 continue
563 if (
564 decided_next_item_ids
565 and slot.when == "between_songs"
566 and slot.next_index is not None
567 and str(tracks[slot.next_index].get("item_id", "")) in decided_next_item_ids
568 ):
569 # the caller settled this slot in an earlier run: re-evaluating it would
570 # consume a chance roll and register its event a second time
571 continue
572 matching_rules = [
573 rule for rule in section_order if str(rule.get("when", "")).strip() == slot.when
574 ]
575 if not matching_rules:
576 continue
577 static, deferred = self._resolve_placeholders(
578 program=program,
579 tracks=tracks,
580 slot=slot,
581 runtime_tokens=runtime_tokens,
582 )
583 # guards may require a deferred token to be present, so they see the merged view;
584 # only the static half is substituted into the stored prompt
585 guard_values = {**deferred, **static}
586 for rule in matching_rules:
587 flow = rule.get("flow", [])
588 if not isinstance(flow, list):
589 continue
590 for flow_item in flow:
591 if not isinstance(flow_item, dict):
592 continue
593 if "MUST" in flow_item:
594 section_id = str(flow_item["MUST"]).strip()
595 if not section_id:
596 continue
597 if is_empty_section(section_id):
598 continue
599 selected.append((section_id, slot, static))
600 register_event(section_id, slot)
601 continue
602 if "ALTERNATIVE" in flow_item:
603 alternative = flow_item["ALTERNATIVE"]
604 if not isinstance(alternative, dict):
605 continue
606 section_id = pick_weighted_choice(alternative.get("choices", []), rng)
607 if is_empty_section(section_id):
608 continue
609 selected.append((section_id, slot, static))
610 register_event(section_id, slot)
611 continue
612 if "OPTIONAL" in flow_item:
613 optional = flow_item["OPTIONAL"]
614 if not isinstance(optional, dict):
615 continue
616 section_id = str(optional.get("section", "")).strip()
617 if not section_id:
618 continue
619 chance_raw = coerce_float(optional.get("chance"), 0.0)
620 chance = chance_raw / 100.0 if chance_raw > 1 else chance_raw
621 if rng.random() > chance:
622 continue
623 guards = optional.get("guards", {}) if isinstance(optional, dict) else {}
624 if not self._passes_optional_guards(
625 section_id=section_id,
626 guards=guards if isinstance(guards, dict) else {},
627 history=history,
628 slot=slot,
629 tracks=tracks,
630 placeholders=guard_values,
631 track_index_offset=track_index_offset,
632 minute_offset=minute_offset,
633 ):
634 continue
635 if is_empty_section(section_id):
636 continue
637 selected.append((section_id, slot, static))
638 register_event(section_id, slot)
639
640 merge_section_id = str(program.get("merge_section_id", "")).strip()
641 meta_section = section_by_id.get(merge_section_id) if merge_section_id else None
642 grouped: dict[str, list[tuple[str, Slot, dict[str, str]]]] = defaultdict(list)
643 for item in selected:
644 section_id, slot, placeholders = item
645 key = f"{slot.when}:{slot.at_index}"
646 grouped[key].append((section_id, slot, placeholders))
647
648 weather_guarded_ids = self._weather_guarded_section_ids(program)
649 planned: list[PlannedSection] = []
650 order_index = 0
651 processed_keys: set[str] = set()
652 for section_id, slot, placeholders in selected:
653 key = f"{slot.when}:{slot.at_index}"
654 grouped_items = grouped[key]
655 if (
656 len(grouped_items) > 1
657 and slot.when == "between_songs"
658 and meta_section
659 and key not in processed_keys
660 ):
661 processed_keys.add(key)
662 merged = self._build_meta_section_plan(
663 grouped_items=grouped_items,
664 meta_section=meta_section,
665 placeholders=placeholders,
666 order=order_index,
667 section_by_id=section_by_id,
668 session_id=session_id,
669 history_events=[(item[0], slot_event(item[1])) for item in grouped_items],
670 weather_guarded_ids=weather_guarded_ids,
671 )
672 planned.append(merged)
673 order_index += 1
674 continue
675 if key in processed_keys:
676 continue
677 section = section_by_id.get(section_id)
678 if not section:
679 continue
680 if str(section.get("type", "ai_text")).strip().lower() != "ai_text":
681 continue
682 prompt = self._apply_placeholders(str(section.get("prompt", "")), placeholders)
683 weather_required = section_id in weather_guarded_ids
684 max_chars = int((section.get("constraints") or {}).get("max_chars", 0) or 0)
685 if max_chars > 0:
686 prompt += (
687 f"\n\nTarget length: around {max_chars} characters. It may exceed by up to "
688 "15% if needed to finish naturally. Never stop mid-sentence."
689 )
690 planned.append(
691 PlannedSection(
692 order=order_index,
693 clip_id=f"{session_id}_{order_index:03d}",
694 section_id=section_id,
695 section_name=self._resolve_section_name(section, section_id),
696 when=slot.when,
697 insert_at_index=slot.at_index,
698 prompt=prompt,
699 max_chars=max_chars,
700 web_search_mode=self._resolve_web_search_mode(section, section_id),
701 weather_required=weather_required,
702 history_events=[(section_id, slot_event(slot))],
703 )
704 )
705 order_index += 1
706
707 return planned, history
708
709 def _passes_optional_guards(
710 self,
711 section_id: str,
712 guards: dict[str, Any],
713 history: dict[str, list[tuple[int, float]]],
714 slot: Slot,
715 tracks: list[dict[str, Any]],
716 placeholders: dict[str, str],
717 track_index_offset: int,
718 minute_offset: float,
719 ) -> bool:
720 """Evaluate OPTIONAL section guards."""
721 min_gap_songs = coerce_int(guards.get("min_gap_songs"), 0)
722 max_per_60min = coerce_int(guards.get("max_per_60min"), 0)
723 required_placeholders = guards.get("require_placeholders_present", [])
724 events = history.get(section_id, [])
725 song_local = slot.next_index if slot.next_index is not None else len(tracks)
726 song_global = track_index_offset + song_local
727 minute_global = minute_offset + slot.minute_mark
728
729 if min_gap_songs > 0 and events:
730 if song_global - events[-1][0] < min_gap_songs:
731 return False
732 if max_per_60min > 0:
733 in_window = [event for event in events if (minute_global - event[1]) <= 60.0]
734 if len(in_window) >= max_per_60min:
735 return False
736 if isinstance(required_placeholders, list):
737 for token in required_placeholders:
738 if not placeholders.get(str(token), "").strip():
739 return False
740 return True
741
742 def _build_meta_section_plan(
743 self,
744 grouped_items: list[tuple[str, Slot, dict[str, str]]],
745 meta_section: dict[str, Any],
746 placeholders: dict[str, str],
747 order: int,
748 section_by_id: dict[str, dict[str, Any]],
749 session_id: str,
750 history_events: list[tuple[str, tuple[int, float]]],
751 weather_guarded_ids: set[str],
752 ) -> PlannedSection:
753 """Build a merged ai_meta section for one slot."""
754 section_ids = [item[0] for item in grouped_items]
755 slot = grouped_items[0][1]
756 prompt_lines: list[str] = []
757 total_max_chars = 0
758 max_web_mode = "disabled"
759 merged_names: list[str] = []
760 # a weather+news merge must still air the news half, so only all-guarded merges require it
761 all_weather_required = all(section_id in weather_guarded_ids for section_id in section_ids)
762 for index, section_id in enumerate(section_ids, start=1):
763 section = section_by_id.get(section_id, {})
764 section_name = self._resolve_section_name(section, section_id)
765 merged_names.append(section_name)
766 prompt_base = self._apply_placeholders(str(section.get("prompt", "")), placeholders)
767 max_chars = int((section.get("constraints") or {}).get("max_chars", 0) or 0)
768 total_max_chars += max_chars
769 prompt_lines.append(f"{index}. [{section_id}] {prompt_base}")
770 mode = self._resolve_web_search_mode(section, section_id)
771 if WEB_SEARCH_MODE_RANK[mode] > WEB_SEARCH_MODE_RANK[max_web_mode]:
772 max_web_mode = mode
773
774 meta_prompt = self._apply_placeholders(str(meta_section.get("prompt", "")), placeholders)
775 prompt_block = "\n".join(prompt_lines)
776 if "<section_drafts>" in meta_prompt:
777 meta_prompt = meta_prompt.replace("<section_drafts>", prompt_block)
778 else:
779 meta_prompt = f"{meta_prompt}\n\nSection prompts:\n{prompt_block}\n"
780 meta_prompt += (
781 "\n\nCreate one single moderator script that naturally combines all requested parts. "
782 "Return plain text only."
783 )
784 if total_max_chars > 0:
785 meta_prompt += (
786 f"\n\nTarget length: around {total_max_chars} characters total. It may exceed "
787 "by up to 15% if needed to finish naturally. Never stop mid-sentence."
788 )
789 section_id = f"multi_{'_'.join(slugify(item) for item in section_ids)}"
790 section_name = " + ".join(dict.fromkeys(merged_names))
791 return PlannedSection(
792 order=order,
793 clip_id=f"{session_id}_{order:03d}",
794 section_id=section_id,
795 section_name=section_name,
796 when=slot.when,
797 insert_at_index=slot.at_index,
798 prompt=meta_prompt,
799 max_chars=total_max_chars,
800 web_search_mode=max_web_mode,
801 weather_required=all_weather_required,
802 history_events=history_events,
803 )
804
805 def _compose_queue_items(
806 self,
807 queue_id: str,
808 session: SessionState,
809 program: dict[str, Any],
810 tracks: list[dict[str, Any]],
811 sections: list[PlannedSection],
812 ) -> list[QueueItem]:
813 """
814 Build the queue items for a whole show.
815
816 Clips carry their render state in ``extra_attributes`` from the moment they are built, so
817 a clip is renderable as soon as the queue holds it.
818
819 :param queue_id: The queue the items are built for.
820 :param session: The session that owns the show.
821 :param program: The station+host program being played.
822 :param tracks: The normalized source tracks, in play order.
823 :param sections: The planned sections to interleave between them.
824 """
825 sections_by_index: dict[int, list[PlannedSection]] = defaultdict(list)
826 for item in sections:
827 sections_by_index[item.insert_at_index].append(item)
828 items: list[QueueItem] = []
829 for index in range(len(tracks) + 1):
830 for section in sorted(sections_by_index.get(index, []), key=lambda item: item.order):
831 items.append(
832 self._section_to_clip_item(queue_id, session.session_id, program, section)
833 )
834 if index < len(tracks) and (media_item := tracks[index].get("media_item")) is not None:
835 items.append(build_queue_item(queue_id, media_item))
836 return items
837
838 def _section_to_clip_item(
839 self,
840 queue_id: str,
841 session_id: str,
842 program: dict[str, Any],
843 section: PlannedSection,
844 ) -> QueueItem:
845 """Build the queue item for a not-yet-rendered clip."""
846 clip = SoundEffect(
847 item_id=section.clip_id,
848 provider=self.instance_id,
849 name=section.section_name,
850 provider_mappings={
851 ProviderMapping(
852 item_id=section.clip_id,
853 provider_domain=self.domain,
854 provider_instance=self.instance_id,
855 )
856 },
857 )
858 clip.metadata.images = UniqueList(
859 [
860 MediaItemImage(
861 type=ImageType.THUMB,
862 path=self._ai_radio_cover_image_path(),
863 provider="builtin",
864 remotely_accessible=False,
865 )
866 ]
867 )
868 queue_item = build_queue_item(queue_id, clip)
869 # the section name already travels as the item's own name, so it is not duplicated here
870 queue_item.extra_attributes.update(
871 {
872 ATTR_SESSION_ID: session_id,
873 ATTR_STATION_ID: str(program.get("id") or ""),
874 ATTR_HOST_ID: str(program.get("host_id") or ""),
875 ATTR_PROMPT: section.prompt,
876 ATTR_MAX_CHARS: section.max_chars,
877 ATTR_WEB_SEARCH_MODE: section.web_search_mode,
878 ATTR_WEATHER_REQUIRED: section.weather_required,
879 }
880 )
881 return queue_item
882
883 @staticmethod
884 def _ai_radio_cover_image_path() -> str:
885 """Return the explicit AI Radio playlist cover image path."""
886 return str(Path(__file__).with_name("air.png"))
887
888 async def _prepare_runtime_tokens(self, program: dict[str, Any]) -> dict[str, str]:
889 """Prepare runtime tokens (including weather placeholders) for one run."""
890 if not self._program_uses_weather_placeholders(program):
891 return {}
892 return await self._prepare_weather_tokens()
893
894 async def _prepare_weather_tokens(self) -> dict[str, str]:
895 """Return the weather placeholder tokens, fetching them at most once per cache window."""
896 cached = self._weather_tokens_cache
897 if cached is not None and (time.monotonic() - cached[0]) < WEATHER_TOKENS_CACHE_SECONDS:
898 return dict(cached[1])
899 tokens = await self._fetch_weather_tokens()
900 # failed and disabled lookups are cached too, so a broken forecast source cannot
901 # put its timeout in front of every replan pass
902 self._weather_tokens_cache = (time.monotonic(), tokens)
903 return dict(tokens)
904
905 async def _fetch_weather_tokens(self) -> dict[str, str]:
906 """Fetch and format weather placeholder tokens from the configured provider."""
907 runtime_tokens: dict[str, str] = {}
908
909 weather_provider = (
910 str(self.config.get_value(CONF_WEATHER_PROVIDER) or DEFAULT_WEATHER_PROVIDER)
911 .strip()
912 .lower()
913 )
914 if weather_provider in {"", "none", "disabled", "off"}:
915 return runtime_tokens
916 if weather_provider != "open_meteo":
917 self.logger.warning(
918 "Unsupported weather provider '%s' for AI Radio station",
919 weather_provider,
920 )
921 return runtime_tokens
922
923 city, country = self._extract_location()
924 if not city or not country:
925 self.logger.warning(
926 "Weather placeholders used but no location configured "
927 "(set the weather_city/weather_country provider options)"
928 )
929 return runtime_tokens
930
931 configured_timeout = self.config.get_value(CONF_WEATHER_TIMEOUT)
932 timeout_seconds = max(5, coerce_int(configured_timeout, DEFAULT_WEATHER_TIMEOUT_SECONDS))
933 try:
934 weather_hourly, weather_daily = await self._fetch_open_meteo_weather(
935 city=city,
936 country=country,
937 timeout_seconds=timeout_seconds,
938 )
939 except Exception as err:
940 self.logger.warning(
941 "Weather lookup failed for '%s, %s': %s",
942 city,
943 country,
944 err,
945 )
946 return runtime_tokens
947
948 if weather_hourly:
949 runtime_tokens["<weather_hourly>"] = weather_hourly
950 if weather_daily:
951 runtime_tokens["<weather_daily>"] = weather_daily
952 return runtime_tokens
953
954 def _program_uses_weather_placeholders(self, program: dict[str, Any]) -> bool:
955 """Return whether the program references weather placeholders."""
956 for section in program.get("sections", []):
957 prompt = str(section.get("prompt", ""))
958 if any(token in prompt for token in WEATHER_PLACEHOLDER_TOKENS):
959 return True
960
961 for rule in program.get("section_order", []):
962 flow = rule.get("flow", [])
963 for item in flow:
964 optional = item.get("OPTIONAL")
965 if not optional:
966 continue
967 guards = optional.get("guards", {})
968 required = guards.get("require_placeholders_present", [])
969 if any(str(token) in WEATHER_PLACEHOLDER_TOKENS for token in required):
970 return True
971 return False
972
973 def _weather_guarded_section_ids(self, program: dict[str, Any]) -> set[str]:
974 """Return OPTIONAL section ids whose guards require a weather placeholder."""
975 guarded: set[str] = set()
976 for rule in program.get("section_order", []):
977 flow = rule.get("flow", [])
978 for item in flow:
979 optional = item.get("OPTIONAL")
980 if not optional:
981 continue
982 section_id = str(optional.get("section", "")).strip()
983 guards = optional.get("guards", {})
984 required = guards.get("require_placeholders_present", [])
985 if section_id and any(str(t) in WEATHER_PLACEHOLDER_TOKENS for t in required):
986 guarded.add(section_id)
987 return guarded
988
989 def _extract_location(self) -> tuple[str, str]:
990 """Extract weather location (city/country) from the provider config."""
991 city = str(self.config.get_value(CONF_WEATHER_CITY) or "").strip()
992 country = str(self.config.get_value(CONF_WEATHER_COUNTRY) or "").strip()
993 return city, country
994
995 def _configured_now(self) -> datetime.datetime:
996 """Return the current time in the configured timezone, falling back to host local time."""
997 tz_name = str(self.config.get_value(CONF_TIMEZONE) or "").strip()
998 if tz_name:
999 try:
1000 return utc().astimezone(ZoneInfo(tz_name))
1001 except ZoneInfoNotFoundError, ValueError:
1002 # a typo must not take the run down, but it should not pass unnoticed either
1003 self.logger.warning(
1004 "Ignoring invalid timezone %r, falling back to the host timezone", tz_name
1005 )
1006 return now()
1007
1008 async def _fetch_open_meteo_weather(
1009 self,
1010 city: str,
1011 country: str,
1012 timeout_seconds: int,
1013 ) -> tuple[str, str]:
1014 """Fetch weather strings from Open-Meteo for weather placeholders."""
1015 use_fahrenheit = country.upper() in FAHRENHEIT_COUNTRY_CODES
1016 geocode_params: dict[str, str | int] = {
1017 "name": city,
1018 "count": 10,
1019 "language": "en",
1020 "format": "json",
1021 }
1022 country_code = country.upper() if len(country) == 2 and country.isalpha() else ""
1023 if country_code:
1024 geocode_params["countryCode"] = country_code
1025 geocode = await self._open_meteo_get_json(
1026 "https://geocoding-api.open-meteo.com/v1/search",
1027 geocode_params,
1028 timeout_seconds,
1029 )
1030 results = geocode.get("results", [])
1031 if not isinstance(results, list) or not results:
1032 raise MusicAssistantError(f"No geocoding result for {city}, {country}")
1033
1034 selected: dict[str, Any] | None = None
1035 country_lc = country.lower()
1036 for candidate in results:
1037 if not isinstance(candidate, dict):
1038 continue
1039 candidate_country = str(candidate.get("country", "")).strip().lower()
1040 candidate_country_code = str(candidate.get("country_code", "")).strip().upper()
1041 if candidate_country and candidate_country == country_lc:
1042 selected = candidate
1043 break
1044 if country_code and candidate_country_code == country_code:
1045 selected = candidate
1046 break
1047
1048 if selected is None:
1049 if country:
1050 # a same-named city in another country is worse than no forecast at all
1051 raise MusicAssistantError(
1052 f"No geocoding result for {city} matched configured country {country}"
1053 )
1054 first = results[0]
1055 selected = first if isinstance(first, dict) else None
1056
1057 if not isinstance(selected, dict):
1058 raise MusicAssistantError(f"No valid geocoding result for {city}, {country}")
1059 latitude_value: object = selected.get("latitude")
1060 longitude_value: object = selected.get("longitude")
1061 if not isinstance(latitude_value, (int, float, str)) or not isinstance(
1062 longitude_value, (int, float, str)
1063 ):
1064 raise MusicAssistantError(
1065 f"Geocoding result for {city}, {country} has invalid coordinates"
1066 )
1067 try:
1068 lat = float(latitude_value)
1069 lon = float(longitude_value)
1070 except ValueError as err:
1071 raise MusicAssistantError(
1072 f"Geocoding result for {city}, {country} has invalid coordinates"
1073 ) from err
1074 timezone_name = str(selected.get("timezone") or "UTC")
1075 forecast_params: dict[str, str | int | float] = {
1076 "latitude": lat,
1077 "longitude": lon,
1078 "current": "temperature_2m,apparent_temperature,weather_code",
1079 "hourly": "temperature_2m,precipitation_probability,weather_code",
1080 "daily": (
1081 "temperature_2m_max,temperature_2m_min,precipitation_probability_max,weather_code"
1082 ),
1083 "forecast_days": 3,
1084 "timezone": timezone_name,
1085 }
1086 if use_fahrenheit:
1087 forecast_params["temperature_unit"] = "fahrenheit"
1088 forecast = await self._open_meteo_get_json(
1089 "https://api.open-meteo.com/v1/forecast",
1090 forecast_params,
1091 timeout_seconds,
1092 )
1093 return self._format_weather_strings(forecast, unit_suffix="F" if use_fahrenheit else "C")
1094
1095 async def _open_meteo_get_json(
1096 self,
1097 base_url: str,
1098 params: dict[str, Any],
1099 timeout_seconds: int,
1100 ) -> dict[str, Any]:
1101 """Perform one Open-Meteo GET request."""
1102 async with self.mass.http_session.get(
1103 base_url,
1104 params=params,
1105 timeout=ClientTimeout(total=timeout_seconds),
1106 ) as response:
1107 payload = await response.read()
1108 if response.status >= 400:
1109 raise MusicAssistantError(
1110 f"Open-Meteo request failed ({response.status}): "
1111 f"{payload.decode(errors='ignore')}"
1112 )
1113 data = json_loads(payload)
1114 if not isinstance(data, dict):
1115 raise MusicAssistantError("Open-Meteo response is not a JSON object")
1116 return data
1117
1118 def _format_weather_strings(
1119 self, payload: dict[str, Any], unit_suffix: str = "C"
1120 ) -> tuple[str, str]:
1121 """Format Open-Meteo payload into weather placeholder strings."""
1122 hourly = payload.get("hourly", {})
1123 daily = payload.get("daily", {})
1124 current = payload.get("current", {})
1125 if not isinstance(hourly, dict):
1126 hourly = {}
1127 if not isinstance(daily, dict):
1128 daily = {}
1129 if not isinstance(current, dict):
1130 current = {}
1131
1132 hourly_times = hourly.get("time", [])
1133 hourly_temp = hourly.get("temperature_2m", [])
1134 hourly_prec = hourly.get("precipitation_probability", [])
1135 if not isinstance(hourly_times, list):
1136 hourly_times = []
1137 if not isinstance(hourly_temp, list):
1138 hourly_temp = []
1139 if not isinstance(hourly_prec, list):
1140 hourly_prec = []
1141
1142 current_time = str(current.get("time") or "").strip()
1143 start_index = 0
1144 if current_time:
1145 # current.time sits on a 15-minute grid while hourly.time is on whole hours;
1146 # the summary starts at the first hour that is not in the past
1147 for index, hour_time in enumerate(hourly_times):
1148 if str(hour_time) >= current_time:
1149 start_index = index
1150 break
1151
1152 max_items = min(len(hourly_times), len(hourly_temp), len(hourly_prec))
1153 hourly_parts: list[str] = []
1154 for index in range(start_index, min(start_index + 6, max_items)):
1155 ts = str(hourly_times[index]).replace("T", " ")
1156 hourly_parts.append(
1157 f"{ts}: {self._format_number(hourly_temp[index])}{unit_suffix}, "
1158 f"rain {self._format_number(hourly_prec[index])}%"
1159 )
1160 current_text = ""
1161 if current:
1162 current_text = (
1163 f"now {self._format_number(current.get('temperature_2m'))}{unit_suffix} "
1164 f"(feels {self._format_number(current.get('apparent_temperature'))}{unit_suffix})"
1165 )
1166 weather_hourly = "; ".join(([current_text] if current_text else []) + hourly_parts)
1167
1168 daily_times = daily.get("time", [])
1169 max_t = daily.get("temperature_2m_max", [])
1170 min_t = daily.get("temperature_2m_min", [])
1171 max_prec = daily.get("precipitation_probability_max", [])
1172 if not isinstance(daily_times, list):
1173 daily_times = []
1174 if not isinstance(max_t, list):
1175 max_t = []
1176 if not isinstance(min_t, list):
1177 min_t = []
1178 if not isinstance(max_prec, list):
1179 max_prec = []
1180 daily_parts: list[str] = []
1181 for index in range(min(len(daily_times), len(max_t), len(min_t), len(max_prec))):
1182 daily_parts.append(
1183 f"{daily_times[index]}: "
1184 f"{self._format_number(min_t[index])}-{self._format_number(max_t[index])}"
1185 f"{unit_suffix}, rain {self._format_number(max_prec[index])}%"
1186 )
1187 weather_daily = "; ".join(daily_parts)
1188 return weather_hourly, weather_daily
1189
1190 def _format_number(self, value: Any) -> str:
1191 """Format weather numeric values compactly for prompts."""
1192 try:
1193 # the host reads these out loud, where a decimal place only clutters the line
1194 numeric = round(float(value))
1195 except Exception:
1196 return str(value)
1197 return str(numeric)
1198
1199 def _resolve_placeholders(
1200 self,
1201 program: dict[str, Any],
1202 tracks: list[dict[str, Any]],
1203 slot: Slot,
1204 runtime_tokens: dict[str, str],
1205 ) -> tuple[dict[str, str], dict[str, str]]:
1206 """
1207 Resolve placeholders for one slot, split by when they are substituted.
1208
1209 :param program: The station+host program being planned.
1210 :param tracks: The track list the slot indexes into.
1211 :param slot: The insertion slot being filled.
1212 :param runtime_tokens: Weather tokens fetched for this run.
1213 :return: ``(static, deferred)`` â static values are fixed by the track order and are
1214 substituted at plan time; deferred values describe the moment of airing and are
1215 substituted at render time.
1216 """
1217 prev_track = tracks[slot.prev_index] if slot.prev_index is not None else None
1218 next_track = tracks[slot.next_index] if slot.next_index is not None else None
1219 very_next_track = tracks[slot.very_next_index] if slot.very_next_index is not None else None
1220 static = {
1221 "<prev_songinfo>": track_songinfo(prev_track),
1222 "<next_songinfo>": track_songinfo(next_track),
1223 "<very_next_songinfo>": track_songinfo(very_next_track),
1224 }
1225 deferred = dict.fromkeys(DEFERRED_PLACEHOLDERS, "")
1226 deferred["<timestamp>"] = self._configured_now().strftime("%Y-%m-%d %H:%M %Z")
1227 for key, value in runtime_tokens.items():
1228 if str(key) in DEFERRED_PLACEHOLDERS:
1229 deferred[str(key)] = str(value)
1230 else:
1231 static[str(key)] = str(value)
1232 return static, deferred
1233
1234 def _apply_placeholders(self, prompt: str, values: dict[str, str]) -> str:
1235 """Apply placeholder replacements in a prompt."""
1236 text = prompt
1237 for key, value in values.items():
1238 text = text.replace(key, value)
1239 return text
1240
1241 def _resolve_section_name(self, section: dict[str, Any], fallback_id: str) -> str:
1242 """Resolve section display name."""
1243 name = str(section.get("name", "")).strip()
1244 return name or fallback_id.replace("_", " ")
1245
1246 def _resolve_web_search_mode(self, section: dict[str, Any], section_id: str) -> str:
1247 """Resolve and validate section web search mode."""
1248 mode = str(section.get("web_search", "disabled")).strip().lower()
1249 if mode not in VALID_WEB_SEARCH_MODES:
1250 raise MusicAssistantError(
1251 f"Invalid web_search mode '{mode}' in section '{section_id}'. "
1252 f"Allowed: {sorted(VALID_WEB_SEARCH_MODES)}"
1253 )
1254 return mode
1255
1256 async def _generate_text(
1257 self, instructions: str, prompt: str, web_mode: str, language: str | None = None
1258 ) -> str:
1259 """Generate one section text using the configured AI engine."""
1260 instructions = instructions.strip() or DEFAULT_LLM_INSTRUCTIONS
1261 query_parts: list[str] = []
1262 if instructions:
1263 query_parts.append(f"Program instructions:\n{instructions}")
1264 query_parts.append(f"Pronunciation rules:\n{TTS_PRONUNCIATION_INSTRUCTIONS}")
1265 # stated as a default so a station can still ask for another language in its instructions
1266 query_parts.append(
1267 "Unless the program instructions ask for another language, write the output "
1268 f"in the language matching the locale '{language or self.mass.metadata.locale}'."
1269 )
1270 if web_mode == "force":
1271 query_parts.append(
1272 "Web mode: force. Use current up-to-date information where relevant."
1273 )
1274 elif web_mode == "allow":
1275 query_parts.append(
1276 "Web mode: allow. Use current information if it improves the answer."
1277 )
1278 query_parts.append(
1279 f"Task: Write one concise spoken radio section.\n\n{prompt}\n\nReturn plain text only."
1280 )
1281 query = "\n\n".join(query_parts)
1282 engine = await self._get_ai_engine()
1283 self.logger.debug(
1284 "AI query prepared: engine=%s web_mode=%s query_chars=%d",
1285 engine.uid,
1286 web_mode,
1287 len(query),
1288 )
1289 try:
1290 async with asyncio.timeout(AI_QUERY_TIMEOUT_SECONDS) as query_timeout:
1291 response = await engine.provider.ai_query(query, engine_id=engine.id)
1292 except Exception as err:
1293 # expired() tells our own cap apart from a timeout raised inside the engine
1294 if isinstance(err, TimeoutError) and query_timeout.expired():
1295 raise MusicAssistantError(
1296 f"AI engine '{engine.uid}' did not respond within {AI_QUERY_TIMEOUT_SECONDS}s"
1297 ) from err
1298 error_name = err.__class__.__name__
1299 error_text = str(err).strip()
1300 if error_name == "NotConnected":
1301 raise MusicAssistantError(
1302 "AI engine "
1303 f"'{engine.uid}' is not connected. Reconnect the provider "
1304 "(for example Home Assistant) and retry."
1305 ) from err
1306 details = error_text or error_name
1307 raise MusicAssistantError(f"AI engine '{engine.uid}' query failed: {details}") from err
1308 if not response or not str(response).strip():
1309 raise MusicAssistantError(
1310 f"AI engine '{engine.uid}' returned an empty response for section text"
1311 )
1312 text = str(response).strip()
1313 self.logger.debug(
1314 "AI query response received: engine=%s chars=%d",
1315 engine.uid,
1316 len(text),
1317 )
1318 return text
1319
1320 async def _get_ai_engine(self) -> AIEngine:
1321 """Return the engine used for AI_QUERY tasks, honouring the configured selection."""
1322 selected = cast("str | None", self.get_setup_value(CONF_AI_ENGINE))
1323 if engine := await resolve_ai_engine(self.mass, selected):
1324 return engine
1325 raise MusicAssistantError(
1326 "No AI engine available. Set up a plugin that provides AI (for example Home "
1327 "Assistant with an ai_task entity) and select it in the AI Radio settings."
1328 )
1329
1330 async def _stop_session_queue(self, session: SessionState) -> None:
1331 """Stop playback of the queue a run was playing on."""
1332 queue_id = session.queue_id
1333 if queue_id is None:
1334 return
1335 queue = self.mass.player_queues.get(queue_id)
1336 if queue is None or getattr(queue, "state", None) == PlaybackState.IDLE:
1337 return
1338 try:
1339 await self.mass.player_queues.stop(queue_id)
1340 except MusicAssistantError as err:
1341 self.logger.debug("Could not stop queue %s: %s", queue_id, err)
1342
1343 async def _get_tts_engine(self, engine_uid: str | None = None) -> TTSEngine:
1344 """Return the engine used for TTS tasks, preferring a host-specific engine_uid."""
1345 if engine_uid:
1346 if engine := await resolve_tts_engine(self.mass, engine_uid):
1347 return engine
1348 self.logger.warning(
1349 "Host TTS engine %s is unavailable, falling back to the provider default",
1350 engine_uid,
1351 )
1352 selected = cast("str | None", self.get_setup_value(CONF_TTS_ENGINE))
1353 if engine := await resolve_tts_engine(self.mass, selected):
1354 return engine
1355 raise MusicAssistantError(
1356 "No text-to-speech engine available. Set up a plugin that provides text-to-speech "
1357 "(for example Home Assistant with a TTS entity) and select it in the AI Radio "
1358 "settings."
1359 )
1360