| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806 |
- # -*- coding: utf-8 -*-
- """
- High-Level Satellite API für Plugins.
- Bietet eine einfache Schnittstelle für die Interaktion mit Satellites.
- """
- from __future__ import annotations
- import asyncio
- import uuid
- from dataclasses import dataclass, field
- from datetime import datetime
- from enum import Enum, auto
- from typing import TYPE_CHECKING, Any, Callable, Coroutine, Iterator
- if TYPE_CHECKING:
- from trixy_core.satellite.satellite import Satellite
- from trixy_core.satellite.satellite_manager import SatelliteManager
- from trixy_core.application import IApplication
- class StreamState(Enum):
- """Zustand eines Audio-Streams."""
- IDLE = auto()
- PLAYING = auto()
- PAUSED = auto()
- STOPPED = auto()
- @dataclass
- class TTSRequest:
- """
- Anfrage für Text-to-Speech.
- """
- id: str
- """Eindeutige Anfrage-ID."""
- text: str
- """Zu sprechender Text."""
- satellite_ids: list[str]
- """Ziel-Satellites."""
- voice: str = "default"
- """Stimme/Profil."""
- speed: float = 1.0
- """Sprechgeschwindigkeit."""
- volume: float = 1.0
- """Lautstärke (0.0-1.0)."""
- priority: int = 0
- """Priorität (höher = wichtiger)."""
- created_at: datetime = field(default_factory=datetime.now)
- """Erstellungszeitpunkt."""
- @dataclass
- class StreamSession:
- """
- Eine Audio-Streaming-Session.
- Verwaltet das Streaming von Audio zu mehreren Satellites.
- """
- id: str
- """Session-ID."""
- state: StreamState = StreamState.IDLE
- """Aktueller Zustand."""
- satellite_ids: set[str] = field(default_factory=set)
- """Verbundene Satellites."""
- source: str = ""
- """Aktuelle Audio-Quelle."""
- volume: float = 1.0
- """Lautstärke."""
- created_at: datetime = field(default_factory=datetime.now)
- """Erstellungszeitpunkt."""
- metadata: dict[str, Any] = field(default_factory=dict)
- """Zusätzliche Metadaten."""
- class SatelliteProxy:
- """
- Proxy für einzelnen Satellite mit High-Level-API.
- Ermöglicht einfache Interaktion mit einem Satellite.
- Example:
- satellite = api.satellites[0]
- await satellite.say("Hallo Welt")
- await satellite.play("http://musik.mp3")
- await satellite.stop()
- """
- def __init__(
- self,
- satellite: "Satellite",
- api: "SatelliteAPI",
- ) -> None:
- """
- Initialisiert den Proxy.
- Args:
- satellite: Der zugrundeliegende Satellite.
- api: Referenz zur API.
- """
- self._satellite = satellite
- self._api = api
- @property
- def id(self) -> str:
- """Satellite-ID."""
- return self._satellite.id
- @property
- def alias(self) -> str:
- """Alias/Name."""
- return self._satellite.alias
- @property
- def room(self) -> str:
- """Raum."""
- return self._satellite.room_id
- @property
- def is_connected(self) -> bool:
- """Verbindungsstatus."""
- return self._satellite.is_connected
- @property
- def raw(self) -> "Satellite":
- """Zugriff auf das rohe Satellite-Objekt."""
- return self._satellite
- async def say(
- self,
- text: str,
- voice: str = "default",
- speed: float = 1.0,
- volume: float = 1.0,
- wait: bool = True,
- ) -> bool:
- """
- Spricht Text über TTS.
- Args:
- text: Der zu sprechende Text.
- voice: Stimme/Profil.
- speed: Sprechgeschwindigkeit.
- volume: Lautstärke.
- wait: Ob auf Abschluss gewartet wird.
- Returns:
- True bei Erfolg.
- """
- return await self._api.say(
- text,
- satellites=[self._satellite.id],
- voice=voice,
- speed=speed,
- volume=volume,
- wait=wait,
- )
- async def say_raw(self, audio_data: bytes) -> bool:
- """
- Sendet rohe Audio-Daten.
- Args:
- audio_data: Audio-Bytes (16KHz, 16-bit, mono).
- Returns:
- True bei Erfolg.
- """
- return await self._satellite.say(audio_data)
- async def play(
- self,
- source: str,
- volume: float = 1.0,
- ) -> StreamSession | None:
- """
- Startet Audio/Musik-Wiedergabe.
- Args:
- source: Audio-Quelle (URL oder Dateipfad).
- volume: Lautstärke.
- Returns:
- StreamSession oder None bei Fehler.
- """
- return await self._api.play(
- source,
- satellites=[self._satellite.id],
- volume=volume,
- )
- async def stop(self) -> bool:
- """
- Stoppt alle Wiedergabe.
- Returns:
- True bei Erfolg.
- """
- return await self._api.stop_satellite(self._satellite.id)
- async def set_volume(self, volume: float) -> bool:
- """
- Setzt die Lautstärke.
- Args:
- volume: Lautstärke (0.0-1.0).
- Returns:
- True bei Erfolg.
- """
- return await self._api.set_volume(self._satellite.id, volume)
- def __repr__(self) -> str:
- return f"SatelliteProxy({self.alias!r}, room={self.room!r})"
- class SatelliteCollection:
- """
- Sammlung von Satellites mit Index-Zugriff.
- Ermöglicht verschiedene Zugriffsmethoden auf Satellites.
- Example:
- # Nach Index
- satellite = api.satellites[0]
- # Nach Alias
- satellite = api.satellites["wohnzimmer"]
- # Nach Raum
- satellites = api.satellites.room("küche")
- # Alle verbundenen
- for sat in api.satellites.connected:
- await sat.say("Hallo")
- """
- def __init__(
- self,
- manager: "SatelliteManager",
- api: "SatelliteAPI",
- ) -> None:
- """
- Initialisiert die Collection.
- Args:
- manager: Der SatelliteManager.
- api: Referenz zur API.
- """
- self._manager = manager
- self._api = api
- def __getitem__(self, key: int | str) -> SatelliteProxy | list[SatelliteProxy]:
- """
- Zugriff nach Index, ID, Alias oder Selektor.
- Args:
- key: Index (int), ID, Alias oder Selektor ("room:name").
- Returns:
- SatelliteProxy oder Liste von Proxies.
- """
- if isinstance(key, int):
- satellite = self._manager[key]
- return SatelliteProxy(satellite, self._api)
- # Versuche als ID
- satellite = self._manager.get(key)
- if satellite:
- return SatelliteProxy(satellite, self._api)
- # Versuche als Alias
- for sat in self._manager:
- if sat.alias.lower() == key.lower():
- return SatelliteProxy(sat, self._api)
- # Selektor-Syntax
- if ":" in key:
- result = self._manager[key]
- if isinstance(result, list):
- return [SatelliteProxy(s, self._api) for s in result]
- return SatelliteProxy(result, self._api)
- raise KeyError(f"Satellite nicht gefunden: {key}")
- def __iter__(self) -> Iterator[SatelliteProxy]:
- """Iterator über alle Satellites."""
- for satellite in self._manager:
- yield SatelliteProxy(satellite, self._api)
- def __len__(self) -> int:
- """Anzahl der Satellites."""
- return len(self._manager)
- def __contains__(self, key: str) -> bool:
- """Prüft ob Satellite existiert."""
- return key in self._manager
- @property
- def all(self) -> list[SatelliteProxy]:
- """Alle Satellites."""
- return [SatelliteProxy(s, self._api) for s in self._manager.get_all()]
- @property
- def connected(self) -> list[SatelliteProxy]:
- """Alle verbundenen Satellites."""
- return [SatelliteProxy(s, self._api) for s in self._manager.get_connected()]
- @property
- def disconnected(self) -> list[SatelliteProxy]:
- """Alle getrennten Satellites."""
- return [SatelliteProxy(s, self._api) for s in self._manager.get_disconnected()]
- def room(self, room_id: str) -> list[SatelliteProxy]:
- """
- Satellites in einem Raum.
- Args:
- room_id: Raum-ID.
- Returns:
- Liste der Satellites im Raum.
- """
- return [SatelliteProxy(s, self._api) for s in self._manager.get_by_room(room_id)]
- def find(
- self,
- room: str | None = None,
- alias: str | None = None,
- connected: bool | None = None,
- ) -> list[SatelliteProxy]:
- """
- Sucht Satellites nach Kriterien.
- Args:
- room: Filter nach Raum.
- alias: Filter nach Alias (Teilübereinstimmung).
- connected: Filter nach Verbindungsstatus.
- Returns:
- Liste passender Satellites.
- """
- return [
- SatelliteProxy(s, self._api)
- for s in self._manager.find(room=room, alias=alias, connected=connected)
- ]
- class SatelliteAPI:
- """
- High-Level API für Satellite-Interaktion.
- Bietet eine einfache Schnittstelle für Plugins.
- Example:
- api = SatelliteAPI(application)
- # Einzelner Satellite
- await api.satellites[0].say("Hallo")
- # Alle Satellites
- await api.say("Guten Morgen")
- # Bestimmter Raum
- await api.say("Willkommen", room="wohnzimmer")
- # Musik-Streaming
- stream = await api.play("http://radio.mp3")
- stream = await api.add_to_stream(stream.id, "küche")
- """
- def __init__(self, application: "IApplication") -> None:
- """
- Initialisiert die API.
- Args:
- application: Die Anwendungsinstanz.
- """
- self._application = application
- self._streams: dict[str, StreamSession] = {}
- self._tts_queue: asyncio.Queue[TTSRequest] = asyncio.Queue()
- @property
- def satellites(self) -> SatelliteCollection:
- """Zugriff auf Satellites."""
- manager = self._get_satellite_manager()
- return SatelliteCollection(manager, self)
- @property
- def streams(self) -> dict[str, StreamSession]:
- """Aktive Streams."""
- return dict(self._streams)
- def _get_satellite_manager(self) -> "SatelliteManager":
- """Holt den SatelliteManager."""
- # Annahme: SatelliteManager ist als Service registriert
- manager = getattr(self._application, "satellite_manager", None)
- if manager is None:
- # Fallback: Aus Service-Container
- container = getattr(self._application, "services", None)
- if container:
- manager = container.get("SatelliteManager")
- if manager is None:
- raise RuntimeError("SatelliteManager nicht verfügbar")
- return manager
- def _get_event_manager(self):
- """Holt den EventManager."""
- return getattr(self._application, "event_manager", None)
- async def say(
- self,
- text: str,
- satellites: list[str] | None = None,
- room: str | None = None,
- voice: str = "default",
- speed: float = 1.0,
- volume: float = 1.0,
- wait: bool = True,
- priority: int = 0,
- ) -> bool:
- """
- Spricht Text auf Satellites.
- Args:
- text: Der zu sprechende Text.
- satellites: Liste von Satellite-IDs (None = alle).
- room: Raum-Filter.
- voice: Stimme/Profil.
- speed: Sprechgeschwindigkeit.
- volume: Lautstärke.
- wait: Ob auf Abschluss gewartet wird.
- priority: Priorität.
- Returns:
- True bei Erfolg.
- """
- # Template-Platzhalter aufloesen
- from trixy_core.utils.template_formatter import format_template
- text = format_template(text, application=self._application)
- # Ziel-Satellites ermitteln
- target_ids = self._resolve_targets(satellites, room)
- if not target_ids:
- return False
- # TTS-Request erstellen
- request = TTSRequest(
- id=str(uuid.uuid4()),
- text=text,
- satellite_ids=target_ids,
- voice=voice,
- speed=speed,
- volume=volume,
- priority=priority,
- )
- # Event auslösen für TTS-Plugin
- event_manager = self._get_event_manager()
- if event_manager:
- from trixy_core.events.event_data.basic import TTSRequest as TTSRequestEvent
- event_data = TTSRequestEvent(
- request_id=request.id,
- satellite_id=",".join(target_ids),
- text=text,
- voice=voice,
- speed=speed,
- volume=volume,
- source="satellite_api.say"
- )
- if wait:
- # Warte-Logik mit Future
- completion_future: asyncio.Future[bool] = asyncio.Future()
- async def on_tts_completed(event_name: str, data) -> None:
- """Handler für tts_completed Event."""
- if getattr(data, "request_id", None) == request.id:
- completion_future.set_result(getattr(data, "success", True))
- # Temporären Handler registrieren
- event_manager.register("tts_completed", on_tts_completed)
- try:
- # Request auslösen
- await event_manager.trigger("tts_request", event_data)
- if event_data.is_cancelled():
- return False
- # Auf Abschluss warten (mit Timeout)
- try:
- await asyncio.wait_for(completion_future, timeout=60.0)
- except asyncio.TimeoutError:
- return False
- finally:
- # Handler wieder entfernen
- event_manager.unregister("tts_completed", on_tts_completed)
- return completion_future.result() if completion_future.done() else False
- else:
- # Ohne Warten
- await event_manager.trigger("tts_request", event_data)
- return not event_data.is_cancelled()
- return False
- async def say_raw(
- self,
- audio_data: bytes,
- satellites: list[str] | None = None,
- room: str | None = None,
- ) -> int:
- """
- Sendet rohe Audio-Daten.
- Args:
- audio_data: Audio-Bytes (16KHz, 16-bit, mono).
- satellites: Liste von Satellite-IDs.
- room: Raum-Filter.
- Returns:
- Anzahl erfolgreich gesendeter Satellites.
- """
- target_ids = self._resolve_targets(satellites, room)
- manager = self._get_satellite_manager()
- count = 0
- for sat_id in target_ids:
- satellite = manager.get(sat_id)
- if satellite and satellite.is_connected:
- if await satellite.say(audio_data):
- count += 1
- return count
- async def play(
- self,
- source: str,
- satellites: list[str] | None = None,
- room: str | None = None,
- volume: float = 1.0,
- ) -> StreamSession | None:
- """
- Startet Audio-Streaming.
- Args:
- source: Audio-Quelle (URL, Dateipfad).
- satellites: Ziel-Satellites.
- room: Raum-Filter.
- volume: Lautstärke.
- Returns:
- StreamSession oder None.
- """
- target_ids = self._resolve_targets(satellites, room)
- if not target_ids:
- return None
- session = StreamSession(
- id=str(uuid.uuid4()),
- state=StreamState.PLAYING,
- satellite_ids=set(target_ids),
- source=source,
- volume=volume,
- )
- self._streams[session.id] = session
- # Event für Streaming-Start
- event_manager = self._get_event_manager()
- if event_manager:
- from trixy_core.events.event_data.base import EventData
- @dataclass
- class StreamStartEvent(EventData):
- session_id: str = session.id
- source: str = source
- satellite_ids: list[str] = field(default_factory=lambda: target_ids)
- volume: float = volume
- await event_manager.trigger("stream_start", StreamStartEvent())
- return session
- async def add_to_stream(
- self,
- stream_id: str,
- satellite: str | list[str],
- ) -> bool:
- """
- Fügt Satellite(s) zu einem laufenden Stream hinzu.
- Args:
- stream_id: Stream-ID.
- satellite: Satellite-ID oder Liste.
- Returns:
- True bei Erfolg.
- """
- session = self._streams.get(stream_id)
- if not session or session.state != StreamState.PLAYING:
- return False
- if isinstance(satellite, str):
- satellite = [satellite]
- for sat_id in satellite:
- session.satellite_ids.add(sat_id)
- # Event für Stream-Update
- event_manager = self._get_event_manager()
- if event_manager:
- from trixy_core.events.event_data.base import EventData
- @dataclass
- class StreamUpdateEvent(EventData):
- session_id: str = stream_id
- added_satellites: list[str] = field(default_factory=lambda: satellite)
- await event_manager.trigger("stream_update", StreamUpdateEvent())
- return True
- async def remove_from_stream(
- self,
- stream_id: str,
- satellite: str | list[str],
- ) -> bool:
- """
- Entfernt Satellite(s) von einem Stream.
- Args:
- stream_id: Stream-ID.
- satellite: Satellite-ID oder Liste.
- Returns:
- True bei Erfolg.
- """
- session = self._streams.get(stream_id)
- if not session:
- return False
- if isinstance(satellite, str):
- satellite = [satellite]
- for sat_id in satellite:
- session.satellite_ids.discard(sat_id)
- # Stream beenden wenn keine Satellites mehr
- if not session.satellite_ids:
- return await self.stop_stream(stream_id)
- return True
- async def stop_stream(self, stream_id: str) -> bool:
- """
- Stoppt einen Stream.
- Args:
- stream_id: Stream-ID.
- Returns:
- True bei Erfolg.
- """
- session = self._streams.pop(stream_id, None)
- if not session:
- return False
- session.state = StreamState.STOPPED
- # Event für Stream-Stop
- event_manager = self._get_event_manager()
- if event_manager:
- from trixy_core.events.event_data.base import EventData
- @dataclass
- class StreamStopEvent(EventData):
- session_id: str = stream_id
- await event_manager.trigger("stream_stop", StreamStopEvent())
- return True
- async def stop_satellite(self, satellite_id: str) -> bool:
- """
- Stoppt alle Streams für einen Satellite.
- Args:
- satellite_id: Satellite-ID.
- Returns:
- True wenn mindestens ein Stream gestoppt.
- """
- stopped = False
- for session in list(self._streams.values()):
- if satellite_id in session.satellite_ids:
- await self.remove_from_stream(session.id, satellite_id)
- stopped = True
- return stopped
- async def stop_all(self) -> int:
- """
- Stoppt alle Streams.
- Returns:
- Anzahl gestoppter Streams.
- """
- count = len(self._streams)
- for stream_id in list(self._streams.keys()):
- await self.stop_stream(stream_id)
- return count
- async def set_volume(
- self,
- satellite: str | None = None,
- volume: float = 1.0,
- ) -> bool:
- """
- Setzt die Lautstärke.
- Args:
- satellite: Satellite-ID (None = alle).
- volume: Lautstärke (0.0-1.0).
- Returns:
- True bei Erfolg.
- """
- # Event für Volume-Change
- event_manager = self._get_event_manager()
- if event_manager:
- from trixy_core.events.event_data.base import EventData
- @dataclass
- class VolumeChangeEvent(EventData):
- satellite_id: str | None = satellite
- volume: float = volume
- await event_manager.trigger("volume_change", VolumeChangeEvent())
- return True
- return False
- def _resolve_targets(
- self,
- satellites: list[str] | None,
- room: str | None,
- ) -> list[str]:
- """
- Ermittelt die Ziel-Satellite-IDs.
- Args:
- satellites: Explizite Liste.
- room: Raum-Filter.
- Returns:
- Liste von Satellite-IDs.
- """
- manager = self._get_satellite_manager()
- if satellites:
- return satellites
- if room:
- return [s.id for s in manager.get_by_room(room) if s.is_connected]
- # Alle verbundenen
- return [s.id for s in manager.get_connected()]
- # Convenience-Funktion für schnellen Zugriff
- def get_satellite_api(application: "IApplication") -> SatelliteAPI:
- """
- Erstellt oder holt eine SatelliteAPI-Instanz.
- Args:
- application: Die Anwendungsinstanz.
- Returns:
- SatelliteAPI-Instanz.
- """
- # Cache auf Application-Ebene
- if not hasattr(application, "_satellite_api"):
- application._satellite_api = SatelliteAPI(application)
- return application._satellite_api
|