api.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806
  1. # -*- coding: utf-8 -*-
  2. """
  3. High-Level Satellite API für Plugins.
  4. Bietet eine einfache Schnittstelle für die Interaktion mit Satellites.
  5. """
  6. from __future__ import annotations
  7. import asyncio
  8. import uuid
  9. from dataclasses import dataclass, field
  10. from datetime import datetime
  11. from enum import Enum, auto
  12. from typing import TYPE_CHECKING, Any, Callable, Coroutine, Iterator
  13. if TYPE_CHECKING:
  14. from trixy_core.satellite.satellite import Satellite
  15. from trixy_core.satellite.satellite_manager import SatelliteManager
  16. from trixy_core.application import IApplication
  17. class StreamState(Enum):
  18. """Zustand eines Audio-Streams."""
  19. IDLE = auto()
  20. PLAYING = auto()
  21. PAUSED = auto()
  22. STOPPED = auto()
  23. @dataclass
  24. class TTSRequest:
  25. """
  26. Anfrage für Text-to-Speech.
  27. """
  28. id: str
  29. """Eindeutige Anfrage-ID."""
  30. text: str
  31. """Zu sprechender Text."""
  32. satellite_ids: list[str]
  33. """Ziel-Satellites."""
  34. voice: str = "default"
  35. """Stimme/Profil."""
  36. speed: float = 1.0
  37. """Sprechgeschwindigkeit."""
  38. volume: float = 1.0
  39. """Lautstärke (0.0-1.0)."""
  40. priority: int = 0
  41. """Priorität (höher = wichtiger)."""
  42. created_at: datetime = field(default_factory=datetime.now)
  43. """Erstellungszeitpunkt."""
  44. @dataclass
  45. class StreamSession:
  46. """
  47. Eine Audio-Streaming-Session.
  48. Verwaltet das Streaming von Audio zu mehreren Satellites.
  49. """
  50. id: str
  51. """Session-ID."""
  52. state: StreamState = StreamState.IDLE
  53. """Aktueller Zustand."""
  54. satellite_ids: set[str] = field(default_factory=set)
  55. """Verbundene Satellites."""
  56. source: str = ""
  57. """Aktuelle Audio-Quelle."""
  58. volume: float = 1.0
  59. """Lautstärke."""
  60. created_at: datetime = field(default_factory=datetime.now)
  61. """Erstellungszeitpunkt."""
  62. metadata: dict[str, Any] = field(default_factory=dict)
  63. """Zusätzliche Metadaten."""
  64. class SatelliteProxy:
  65. """
  66. Proxy für einzelnen Satellite mit High-Level-API.
  67. Ermöglicht einfache Interaktion mit einem Satellite.
  68. Example:
  69. satellite = api.satellites[0]
  70. await satellite.say("Hallo Welt")
  71. await satellite.play("http://musik.mp3")
  72. await satellite.stop()
  73. """
  74. def __init__(
  75. self,
  76. satellite: "Satellite",
  77. api: "SatelliteAPI",
  78. ) -> None:
  79. """
  80. Initialisiert den Proxy.
  81. Args:
  82. satellite: Der zugrundeliegende Satellite.
  83. api: Referenz zur API.
  84. """
  85. self._satellite = satellite
  86. self._api = api
  87. @property
  88. def id(self) -> str:
  89. """Satellite-ID."""
  90. return self._satellite.id
  91. @property
  92. def alias(self) -> str:
  93. """Alias/Name."""
  94. return self._satellite.alias
  95. @property
  96. def room(self) -> str:
  97. """Raum."""
  98. return self._satellite.room_id
  99. @property
  100. def is_connected(self) -> bool:
  101. """Verbindungsstatus."""
  102. return self._satellite.is_connected
  103. @property
  104. def raw(self) -> "Satellite":
  105. """Zugriff auf das rohe Satellite-Objekt."""
  106. return self._satellite
  107. async def say(
  108. self,
  109. text: str,
  110. voice: str = "default",
  111. speed: float = 1.0,
  112. volume: float = 1.0,
  113. wait: bool = True,
  114. ) -> bool:
  115. """
  116. Spricht Text über TTS.
  117. Args:
  118. text: Der zu sprechende Text.
  119. voice: Stimme/Profil.
  120. speed: Sprechgeschwindigkeit.
  121. volume: Lautstärke.
  122. wait: Ob auf Abschluss gewartet wird.
  123. Returns:
  124. True bei Erfolg.
  125. """
  126. return await self._api.say(
  127. text,
  128. satellites=[self._satellite.id],
  129. voice=voice,
  130. speed=speed,
  131. volume=volume,
  132. wait=wait,
  133. )
  134. async def say_raw(self, audio_data: bytes) -> bool:
  135. """
  136. Sendet rohe Audio-Daten.
  137. Args:
  138. audio_data: Audio-Bytes (16KHz, 16-bit, mono).
  139. Returns:
  140. True bei Erfolg.
  141. """
  142. return await self._satellite.say(audio_data)
  143. async def play(
  144. self,
  145. source: str,
  146. volume: float = 1.0,
  147. ) -> StreamSession | None:
  148. """
  149. Startet Audio/Musik-Wiedergabe.
  150. Args:
  151. source: Audio-Quelle (URL oder Dateipfad).
  152. volume: Lautstärke.
  153. Returns:
  154. StreamSession oder None bei Fehler.
  155. """
  156. return await self._api.play(
  157. source,
  158. satellites=[self._satellite.id],
  159. volume=volume,
  160. )
  161. async def stop(self) -> bool:
  162. """
  163. Stoppt alle Wiedergabe.
  164. Returns:
  165. True bei Erfolg.
  166. """
  167. return await self._api.stop_satellite(self._satellite.id)
  168. async def set_volume(self, volume: float) -> bool:
  169. """
  170. Setzt die Lautstärke.
  171. Args:
  172. volume: Lautstärke (0.0-1.0).
  173. Returns:
  174. True bei Erfolg.
  175. """
  176. return await self._api.set_volume(self._satellite.id, volume)
  177. def __repr__(self) -> str:
  178. return f"SatelliteProxy({self.alias!r}, room={self.room!r})"
  179. class SatelliteCollection:
  180. """
  181. Sammlung von Satellites mit Index-Zugriff.
  182. Ermöglicht verschiedene Zugriffsmethoden auf Satellites.
  183. Example:
  184. # Nach Index
  185. satellite = api.satellites[0]
  186. # Nach Alias
  187. satellite = api.satellites["wohnzimmer"]
  188. # Nach Raum
  189. satellites = api.satellites.room("küche")
  190. # Alle verbundenen
  191. for sat in api.satellites.connected:
  192. await sat.say("Hallo")
  193. """
  194. def __init__(
  195. self,
  196. manager: "SatelliteManager",
  197. api: "SatelliteAPI",
  198. ) -> None:
  199. """
  200. Initialisiert die Collection.
  201. Args:
  202. manager: Der SatelliteManager.
  203. api: Referenz zur API.
  204. """
  205. self._manager = manager
  206. self._api = api
  207. def __getitem__(self, key: int | str) -> SatelliteProxy | list[SatelliteProxy]:
  208. """
  209. Zugriff nach Index, ID, Alias oder Selektor.
  210. Args:
  211. key: Index (int), ID, Alias oder Selektor ("room:name").
  212. Returns:
  213. SatelliteProxy oder Liste von Proxies.
  214. """
  215. if isinstance(key, int):
  216. satellite = self._manager[key]
  217. return SatelliteProxy(satellite, self._api)
  218. # Versuche als ID
  219. satellite = self._manager.get(key)
  220. if satellite:
  221. return SatelliteProxy(satellite, self._api)
  222. # Versuche als Alias
  223. for sat in self._manager:
  224. if sat.alias.lower() == key.lower():
  225. return SatelliteProxy(sat, self._api)
  226. # Selektor-Syntax
  227. if ":" in key:
  228. result = self._manager[key]
  229. if isinstance(result, list):
  230. return [SatelliteProxy(s, self._api) for s in result]
  231. return SatelliteProxy(result, self._api)
  232. raise KeyError(f"Satellite nicht gefunden: {key}")
  233. def __iter__(self) -> Iterator[SatelliteProxy]:
  234. """Iterator über alle Satellites."""
  235. for satellite in self._manager:
  236. yield SatelliteProxy(satellite, self._api)
  237. def __len__(self) -> int:
  238. """Anzahl der Satellites."""
  239. return len(self._manager)
  240. def __contains__(self, key: str) -> bool:
  241. """Prüft ob Satellite existiert."""
  242. return key in self._manager
  243. @property
  244. def all(self) -> list[SatelliteProxy]:
  245. """Alle Satellites."""
  246. return [SatelliteProxy(s, self._api) for s in self._manager.get_all()]
  247. @property
  248. def connected(self) -> list[SatelliteProxy]:
  249. """Alle verbundenen Satellites."""
  250. return [SatelliteProxy(s, self._api) for s in self._manager.get_connected()]
  251. @property
  252. def disconnected(self) -> list[SatelliteProxy]:
  253. """Alle getrennten Satellites."""
  254. return [SatelliteProxy(s, self._api) for s in self._manager.get_disconnected()]
  255. def room(self, room_id: str) -> list[SatelliteProxy]:
  256. """
  257. Satellites in einem Raum.
  258. Args:
  259. room_id: Raum-ID.
  260. Returns:
  261. Liste der Satellites im Raum.
  262. """
  263. return [SatelliteProxy(s, self._api) for s in self._manager.get_by_room(room_id)]
  264. def find(
  265. self,
  266. room: str | None = None,
  267. alias: str | None = None,
  268. connected: bool | None = None,
  269. ) -> list[SatelliteProxy]:
  270. """
  271. Sucht Satellites nach Kriterien.
  272. Args:
  273. room: Filter nach Raum.
  274. alias: Filter nach Alias (Teilübereinstimmung).
  275. connected: Filter nach Verbindungsstatus.
  276. Returns:
  277. Liste passender Satellites.
  278. """
  279. return [
  280. SatelliteProxy(s, self._api)
  281. for s in self._manager.find(room=room, alias=alias, connected=connected)
  282. ]
  283. class SatelliteAPI:
  284. """
  285. High-Level API für Satellite-Interaktion.
  286. Bietet eine einfache Schnittstelle für Plugins.
  287. Example:
  288. api = SatelliteAPI(application)
  289. # Einzelner Satellite
  290. await api.satellites[0].say("Hallo")
  291. # Alle Satellites
  292. await api.say("Guten Morgen")
  293. # Bestimmter Raum
  294. await api.say("Willkommen", room="wohnzimmer")
  295. # Musik-Streaming
  296. stream = await api.play("http://radio.mp3")
  297. stream = await api.add_to_stream(stream.id, "küche")
  298. """
  299. def __init__(self, application: "IApplication") -> None:
  300. """
  301. Initialisiert die API.
  302. Args:
  303. application: Die Anwendungsinstanz.
  304. """
  305. self._application = application
  306. self._streams: dict[str, StreamSession] = {}
  307. self._tts_queue: asyncio.Queue[TTSRequest] = asyncio.Queue()
  308. @property
  309. def satellites(self) -> SatelliteCollection:
  310. """Zugriff auf Satellites."""
  311. manager = self._get_satellite_manager()
  312. return SatelliteCollection(manager, self)
  313. @property
  314. def streams(self) -> dict[str, StreamSession]:
  315. """Aktive Streams."""
  316. return dict(self._streams)
  317. def _get_satellite_manager(self) -> "SatelliteManager":
  318. """Holt den SatelliteManager."""
  319. # Annahme: SatelliteManager ist als Service registriert
  320. manager = getattr(self._application, "satellite_manager", None)
  321. if manager is None:
  322. # Fallback: Aus Service-Container
  323. container = getattr(self._application, "services", None)
  324. if container:
  325. manager = container.get("SatelliteManager")
  326. if manager is None:
  327. raise RuntimeError("SatelliteManager nicht verfügbar")
  328. return manager
  329. def _get_event_manager(self):
  330. """Holt den EventManager."""
  331. return getattr(self._application, "event_manager", None)
  332. async def say(
  333. self,
  334. text: str,
  335. satellites: list[str] | None = None,
  336. room: str | None = None,
  337. voice: str = "default",
  338. speed: float = 1.0,
  339. volume: float = 1.0,
  340. wait: bool = True,
  341. priority: int = 0,
  342. ) -> bool:
  343. """
  344. Spricht Text auf Satellites.
  345. Args:
  346. text: Der zu sprechende Text.
  347. satellites: Liste von Satellite-IDs (None = alle).
  348. room: Raum-Filter.
  349. voice: Stimme/Profil.
  350. speed: Sprechgeschwindigkeit.
  351. volume: Lautstärke.
  352. wait: Ob auf Abschluss gewartet wird.
  353. priority: Priorität.
  354. Returns:
  355. True bei Erfolg.
  356. """
  357. # Template-Platzhalter aufloesen
  358. from trixy_core.utils.template_formatter import format_template
  359. text = format_template(text, application=self._application)
  360. # Ziel-Satellites ermitteln
  361. target_ids = self._resolve_targets(satellites, room)
  362. if not target_ids:
  363. return False
  364. # TTS-Request erstellen
  365. request = TTSRequest(
  366. id=str(uuid.uuid4()),
  367. text=text,
  368. satellite_ids=target_ids,
  369. voice=voice,
  370. speed=speed,
  371. volume=volume,
  372. priority=priority,
  373. )
  374. # Event auslösen für TTS-Plugin
  375. event_manager = self._get_event_manager()
  376. if event_manager:
  377. from trixy_core.events.event_data.basic import TTSRequest as TTSRequestEvent
  378. event_data = TTSRequestEvent(
  379. request_id=request.id,
  380. satellite_id=",".join(target_ids),
  381. text=text,
  382. voice=voice,
  383. speed=speed,
  384. volume=volume,
  385. source="satellite_api.say"
  386. )
  387. if wait:
  388. # Warte-Logik mit Future
  389. completion_future: asyncio.Future[bool] = asyncio.Future()
  390. async def on_tts_completed(event_name: str, data) -> None:
  391. """Handler für tts_completed Event."""
  392. if getattr(data, "request_id", None) == request.id:
  393. completion_future.set_result(getattr(data, "success", True))
  394. # Temporären Handler registrieren
  395. event_manager.register("tts_completed", on_tts_completed)
  396. try:
  397. # Request auslösen
  398. await event_manager.trigger("tts_request", event_data)
  399. if event_data.is_cancelled():
  400. return False
  401. # Auf Abschluss warten (mit Timeout)
  402. try:
  403. await asyncio.wait_for(completion_future, timeout=60.0)
  404. except asyncio.TimeoutError:
  405. return False
  406. finally:
  407. # Handler wieder entfernen
  408. event_manager.unregister("tts_completed", on_tts_completed)
  409. return completion_future.result() if completion_future.done() else False
  410. else:
  411. # Ohne Warten
  412. await event_manager.trigger("tts_request", event_data)
  413. return not event_data.is_cancelled()
  414. return False
  415. async def say_raw(
  416. self,
  417. audio_data: bytes,
  418. satellites: list[str] | None = None,
  419. room: str | None = None,
  420. ) -> int:
  421. """
  422. Sendet rohe Audio-Daten.
  423. Args:
  424. audio_data: Audio-Bytes (16KHz, 16-bit, mono).
  425. satellites: Liste von Satellite-IDs.
  426. room: Raum-Filter.
  427. Returns:
  428. Anzahl erfolgreich gesendeter Satellites.
  429. """
  430. target_ids = self._resolve_targets(satellites, room)
  431. manager = self._get_satellite_manager()
  432. count = 0
  433. for sat_id in target_ids:
  434. satellite = manager.get(sat_id)
  435. if satellite and satellite.is_connected:
  436. if await satellite.say(audio_data):
  437. count += 1
  438. return count
  439. async def play(
  440. self,
  441. source: str,
  442. satellites: list[str] | None = None,
  443. room: str | None = None,
  444. volume: float = 1.0,
  445. ) -> StreamSession | None:
  446. """
  447. Startet Audio-Streaming.
  448. Args:
  449. source: Audio-Quelle (URL, Dateipfad).
  450. satellites: Ziel-Satellites.
  451. room: Raum-Filter.
  452. volume: Lautstärke.
  453. Returns:
  454. StreamSession oder None.
  455. """
  456. target_ids = self._resolve_targets(satellites, room)
  457. if not target_ids:
  458. return None
  459. session = StreamSession(
  460. id=str(uuid.uuid4()),
  461. state=StreamState.PLAYING,
  462. satellite_ids=set(target_ids),
  463. source=source,
  464. volume=volume,
  465. )
  466. self._streams[session.id] = session
  467. # Event für Streaming-Start
  468. event_manager = self._get_event_manager()
  469. if event_manager:
  470. from trixy_core.events.event_data.base import EventData
  471. @dataclass
  472. class StreamStartEvent(EventData):
  473. session_id: str = session.id
  474. source: str = source
  475. satellite_ids: list[str] = field(default_factory=lambda: target_ids)
  476. volume: float = volume
  477. await event_manager.trigger("stream_start", StreamStartEvent())
  478. return session
  479. async def add_to_stream(
  480. self,
  481. stream_id: str,
  482. satellite: str | list[str],
  483. ) -> bool:
  484. """
  485. Fügt Satellite(s) zu einem laufenden Stream hinzu.
  486. Args:
  487. stream_id: Stream-ID.
  488. satellite: Satellite-ID oder Liste.
  489. Returns:
  490. True bei Erfolg.
  491. """
  492. session = self._streams.get(stream_id)
  493. if not session or session.state != StreamState.PLAYING:
  494. return False
  495. if isinstance(satellite, str):
  496. satellite = [satellite]
  497. for sat_id in satellite:
  498. session.satellite_ids.add(sat_id)
  499. # Event für Stream-Update
  500. event_manager = self._get_event_manager()
  501. if event_manager:
  502. from trixy_core.events.event_data.base import EventData
  503. @dataclass
  504. class StreamUpdateEvent(EventData):
  505. session_id: str = stream_id
  506. added_satellites: list[str] = field(default_factory=lambda: satellite)
  507. await event_manager.trigger("stream_update", StreamUpdateEvent())
  508. return True
  509. async def remove_from_stream(
  510. self,
  511. stream_id: str,
  512. satellite: str | list[str],
  513. ) -> bool:
  514. """
  515. Entfernt Satellite(s) von einem Stream.
  516. Args:
  517. stream_id: Stream-ID.
  518. satellite: Satellite-ID oder Liste.
  519. Returns:
  520. True bei Erfolg.
  521. """
  522. session = self._streams.get(stream_id)
  523. if not session:
  524. return False
  525. if isinstance(satellite, str):
  526. satellite = [satellite]
  527. for sat_id in satellite:
  528. session.satellite_ids.discard(sat_id)
  529. # Stream beenden wenn keine Satellites mehr
  530. if not session.satellite_ids:
  531. return await self.stop_stream(stream_id)
  532. return True
  533. async def stop_stream(self, stream_id: str) -> bool:
  534. """
  535. Stoppt einen Stream.
  536. Args:
  537. stream_id: Stream-ID.
  538. Returns:
  539. True bei Erfolg.
  540. """
  541. session = self._streams.pop(stream_id, None)
  542. if not session:
  543. return False
  544. session.state = StreamState.STOPPED
  545. # Event für Stream-Stop
  546. event_manager = self._get_event_manager()
  547. if event_manager:
  548. from trixy_core.events.event_data.base import EventData
  549. @dataclass
  550. class StreamStopEvent(EventData):
  551. session_id: str = stream_id
  552. await event_manager.trigger("stream_stop", StreamStopEvent())
  553. return True
  554. async def stop_satellite(self, satellite_id: str) -> bool:
  555. """
  556. Stoppt alle Streams für einen Satellite.
  557. Args:
  558. satellite_id: Satellite-ID.
  559. Returns:
  560. True wenn mindestens ein Stream gestoppt.
  561. """
  562. stopped = False
  563. for session in list(self._streams.values()):
  564. if satellite_id in session.satellite_ids:
  565. await self.remove_from_stream(session.id, satellite_id)
  566. stopped = True
  567. return stopped
  568. async def stop_all(self) -> int:
  569. """
  570. Stoppt alle Streams.
  571. Returns:
  572. Anzahl gestoppter Streams.
  573. """
  574. count = len(self._streams)
  575. for stream_id in list(self._streams.keys()):
  576. await self.stop_stream(stream_id)
  577. return count
  578. async def set_volume(
  579. self,
  580. satellite: str | None = None,
  581. volume: float = 1.0,
  582. ) -> bool:
  583. """
  584. Setzt die Lautstärke.
  585. Args:
  586. satellite: Satellite-ID (None = alle).
  587. volume: Lautstärke (0.0-1.0).
  588. Returns:
  589. True bei Erfolg.
  590. """
  591. # Event für Volume-Change
  592. event_manager = self._get_event_manager()
  593. if event_manager:
  594. from trixy_core.events.event_data.base import EventData
  595. @dataclass
  596. class VolumeChangeEvent(EventData):
  597. satellite_id: str | None = satellite
  598. volume: float = volume
  599. await event_manager.trigger("volume_change", VolumeChangeEvent())
  600. return True
  601. return False
  602. def _resolve_targets(
  603. self,
  604. satellites: list[str] | None,
  605. room: str | None,
  606. ) -> list[str]:
  607. """
  608. Ermittelt die Ziel-Satellite-IDs.
  609. Args:
  610. satellites: Explizite Liste.
  611. room: Raum-Filter.
  612. Returns:
  613. Liste von Satellite-IDs.
  614. """
  615. manager = self._get_satellite_manager()
  616. if satellites:
  617. return satellites
  618. if room:
  619. return [s.id for s in manager.get_by_room(room) if s.is_connected]
  620. # Alle verbundenen
  621. return [s.id for s in manager.get_connected()]
  622. # Convenience-Funktion für schnellen Zugriff
  623. def get_satellite_api(application: "IApplication") -> SatelliteAPI:
  624. """
  625. Erstellt oder holt eine SatelliteAPI-Instanz.
  626. Args:
  627. application: Die Anwendungsinstanz.
  628. Returns:
  629. SatelliteAPI-Instanz.
  630. """
  631. # Cache auf Application-Ebene
  632. if not hasattr(application, "_satellite_api"):
  633. application._satellite_api = SatelliteAPI(application)
  634. return application._satellite_api