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