"""Der Miracast-Handshake als Quelle (Source). Der Wi-Fi-Display-Standard schreibt die Schritte M1 bis M7 fest. Wichtig: Die Steuerverbindung baut **der Fernseher zu uns** auf, nicht umgekehrt - genauso macht es Android intern (RemoteDisplay.listen). Wir hören deshalb auf Port 7236 und warten, dass er anruft. Die RTSP-Rollen bleiben davon unberührt: Wir stellen M1, M3, M4 und M5, der Fernseher M2, M6 und M7. """ import logging import socket import threading from dataclasses import dataclass, field from typing import Callable, Dict, Optional from . import formats log = logging.getLogger("wfd.rtsp") DEFAULT_PORTS = (7236, 8554) @dataclass class RtspMessage: is_request: bool method: str = "" uri: str = "" status_code: int = 0 status_text: str = "" headers: Dict[str, str] = field(default_factory=dict) body: str = "" @property def cseq(self) -> int: try: return int(self.headers.get("cseq", "0").strip()) except ValueError: return 0 @property def session(self) -> Optional[str]: value = self.headers.get("session") return value.split(";")[0].strip() if value else None def param(self, name: str) -> Optional[str]: """Sucht einen wfd-Parameter im Nachrichtenrumpf.""" for line in self.body.splitlines(): if line.lower().startswith(name.lower() + ":"): return line.split(":", 1)[1].strip() return None def read_message(stream) -> Optional[RtspMessage]: """Liest genau eine Nachricht; None, wenn die Gegenseite auflegt.""" lines = [] while True: raw = stream.readline() if not raw: return None text = raw.decode("utf-8", "replace").rstrip("\r\n") if text == "": if not lines: continue break lines.append(text) if len(lines) > 80: return None headers = {} for line in lines[1:]: if ":" in line: key, value = line.split(":", 1) headers[key.strip().lower()] = value.strip() length = int(headers.get("content-length", "0") or 0) body = stream.read(length).decode("utf-8", "replace") if length else "" first = lines[0] if first.startswith("RTSP/"): parts = first.split(" ", 2) return RtspMessage( is_request=False, status_code=int(parts[1]) if len(parts) > 1 and parts[1].isdigit() else 0, status_text=parts[2] if len(parts) > 2 else "", headers=headers, body=body, ) parts = first.split(" ") return RtspMessage( is_request=True, method=parts[0] if parts else "", uri=parts[1] if len(parts) > 1 else "", headers=headers, body=body, ) def build_request(method: str, uri: str, cseq: int, headers=None, body: str = "") -> bytes: out = [f"{method} {uri} RTSP/1.0", f"CSeq: {cseq}"] for key, value in (headers or {}).items(): out.append(f"{key}: {value}") if body: out.append("Content-Type: text/parameters") out.append(f"Content-Length: {len(body.encode())}") return ("\r\n".join(out) + "\r\n\r\n" + body).encode() def build_response(cseq: int, headers=None, body: str = "", status: str = "200 OK") -> bytes: out = [f"RTSP/1.0 {status}", f"CSeq: {cseq}"] for key, value in (headers or {}).items(): out.append(f"{key}: {value}") if body: out.append("Content-Type: text/parameters") out.append(f"Content-Length: {len(body.encode())}") return ("\r\n".join(out) + "\r\n\r\n" + body).encode() class SourceSession: """Führt den Handshake über eine bestehende Verbindung.""" def __init__(self, sock: socket.socket, local_address: str, max_height: int, on_play: Callable[[str, int, formats.VideoFormat], None], on_stopped: Callable[[str], None], on_status: Callable[[str], None] = lambda text: None): self.sock = sock self.local_address = local_address self.max_height = max_height self.on_play = on_play self.on_stopped = on_stopped self.on_status = on_status self.stream = sock.makefile("rwb") self.peer = sock.getpeername()[0] self.running = True self._cseq = 1 self._pending = {} # CSeq -> Schrittname self._options_answered = False self._peer_options_seen = False self._capabilities_requested = False self._format: Optional[formats.VideoFormat] = None self._sink_rtp_port = 0 self._session_id = "1" self._lock = threading.Lock() self.local_rtp_port = 0 # ------------------------------------------------------------- Ablauf def run(self): try: # M1: Wir melden uns und fragen, was die Gegenseite kann. self._send_request("OPTIONS", "*", {"Require": "org.wfa.wfd1.0"}, step="M1") while self.running: message = read_message(self.stream) if message is None: break if message.is_request: self._handle_request(message) else: self._handle_response(message) self.on_stopped("Verbindung beendet") except Exception as error: # noqa: BLE001 - alles melden, nichts verschlucken if self.running: self.on_stopped(f"Fehler: {error}") finally: self.running = False try: self.sock.close() except OSError: pass def stop(self): self.running = False try: self.sock.close() except OSError: pass def keep_alive(self): """Lebenszeichen, sonst trennt der Fernseher nach etwa einer Minute.""" if self.running: self._send_request("GET_PARAMETER", f"rtsp://{self.local_address}/wfd1.0") # --------------------------------------------- Anfragen des Fernsehers def _handle_request(self, message: RtspMessage): method = message.method.upper() log.info("<- %s %s", method, message.uri) self.on_status(f"Fernseher: {method}") if method == "OPTIONS": # M2: Er fragt seinerseits, was wir können. self._write(build_response(message.cseq, { "Public": "org.wfa.wfd1.0, GET_PARAMETER, SET_PARAMETER, SETUP, PLAY, PAUSE, TEARDOWN" })) self._peer_options_seen = True self._maybe_request_capabilities() elif method == "SETUP": # M6: Sitzung einrichten; hier nennt er den Port für das Bild. transport = message.headers.get("transport", "") for token in transport.split(";"): if token.startswith("client_port="): value = token.split("=", 1)[1].split("-")[0] if value.isdigit(): self._sink_rtp_port = int(value) self.on_status(f"Fernseher erwartet das Bild auf Port {self._sink_rtp_port}") reply_transport = f"RTP/AVP/UDP;unicast;client_port={self._sink_rtp_port}" if self.local_rtp_port: reply_transport += f";server_port={self.local_rtp_port}" self._write(build_response(message.cseq, { "Session": f"{self._session_id};timeout=60", "Transport": reply_transport, })) elif method == "PLAY": # M7: Los geht's. self._write(build_response(message.cseq, {"Session": f"{self._session_id};timeout=60"})) if self._format and self._sink_rtp_port: self.on_status("Wiedergabe angefordert - Bild geht raus") self.on_play(self.peer, self._sink_rtp_port, self._format) else: self.on_status("PLAY kam, aber Format oder Port fehlen noch") elif method == "TEARDOWN": self._write(build_response(message.cseq, {"Session": self._session_id})) self.running = False self.on_stopped("Der Fernseher hat die Übertragung beendet") elif method in ("GET_PARAMETER", "SET_PARAMETER", "PAUSE"): self._write(build_response(message.cseq)) else: self._write(build_response(message.cseq, status="501 Not Implemented")) # ------------------------------------------ Antworten auf unsere Fragen def _handle_response(self, message: RtspMessage): step = self._pending.pop(message.cseq, None) log.info("-> Antwort %s auf %s", message.status_code, step or f"CSeq {message.cseq}") if not 200 <= message.status_code < 300: self.running = False self.on_stopped(f"Der Fernseher hat abgelehnt ({message.status_code} {message.status_text})") return if step == "M1": self._options_answered = True self._maybe_request_capabilities() # Bleibt sein OPTIONS aus, fragen wir trotzdem weiter. threading.Timer(1.5, lambda: self._maybe_request_capabilities(force=True)).start() elif step == "M3": self._negotiate(message) elif step == "M4": self._trigger_setup() def _maybe_request_capabilities(self, force: bool = False): with self._lock: if self._capabilities_requested or not self._options_answered: return if not self._peer_options_seen and not force: return self._capabilities_requested = True # M3: Fähigkeiten abfragen. body = ("wfd_video_formats\r\n" "wfd_audio_codecs\r\n" "wfd_client_rtp_ports\r\n" "wfd_content_protection\r\n") self._send_request("GET_PARAMETER", f"rtsp://{self.local_address}/wfd1.0", body=body, step="M3") def _negotiate(self, message: RtspMessage): """M4: Format festlegen und zurückmelden.""" value = message.param("wfd_video_formats") if not value: self.running = False self.on_stopped("Der Fernseher hat keine Bildformate genannt") return capabilities = formats.parse_video_formats(value) chosen = formats.choose(capabilities, self.max_height) if not chosen: self.running = False self.on_stopped("Kein gemeinsames Bildformat gefunden") return self._format = chosen self.on_status(f"Format ausgehandelt: {chosen}") ports_value = message.param("wfd_client_rtp_ports") port = formats.parse_rtp_port(ports_value) if port: self._sink_rtp_port = port capability = capabilities[0] body = ( f"wfd_video_formats: {formats.build_selection(chosen, capability.profile, capability.level)}\r\n" f"wfd_presentation_URL: rtsp://{self.local_address}/wfd1.0/streamid=0 none\r\n" f"wfd_client_rtp_ports: {ports_value or f'RTP/AVP/UDP;unicast {self._sink_rtp_port} 0 mode=play'}\r\n" ) self._send_request("SET_PARAMETER", f"rtsp://{self.local_address}/wfd1.0", body=body, step="M4") def _trigger_setup(self): """M5: Den Fernseher bitten, jetzt SETUP zu schicken.""" self._send_request("SET_PARAMETER", f"rtsp://{self.local_address}/wfd1.0", body="wfd_trigger_method: SETUP\r\n", step="M5") # ------------------------------------------------------------- Technik def _send_request(self, method: str, uri: str, headers=None, body: str = "", step: str = ""): with self._lock: cseq = self._cseq self._cseq += 1 if step: self._pending[cseq] = step self._write(build_request(method, uri, cseq, headers, body)) log.info("-> %s %s (%s)", method, uri, step or f"CSeq {cseq}") def _write(self, data: bytes): try: self.stream.write(data) self.stream.flush() except OSError as error: log.warning("Senden fehlgeschlagen: %s", error) class SourceServer: """Hört auf den Fernseher und übergibt jede Verbindung an eine Sitzung.""" def __init__(self, local_address: str, max_height: int, on_play, on_stopped, on_status=lambda text: None, ports=DEFAULT_PORTS): self.local_address = local_address self.max_height = max_height self.on_play = on_play self.on_stopped = on_stopped self.on_status = on_status self.ports = ports self._servers = [] self._session: Optional[SourceSession] = None self._adopted = threading.Event() def start(self) -> list: """Öffnet die Steuerkanäle. Gibt die tatsächlich belegten Ports zurück.""" opened = [] for port in self.ports: try: server = socket.socket(socket.AF_INET, socket.SOCK_STREAM) server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) server.bind(("0.0.0.0", port)) server.listen(1) self._servers.append(server) threading.Thread(target=self._accept_loop, args=(server, port), daemon=True, name=f"rtsp-listen-{port}").start() opened.append(port) self.on_status(f"Warte auf den Fernseher (Port {port})") except OSError as error: log.warning("Port %s nicht verfügbar: %s", port, error) return opened def _accept_loop(self, server: socket.socket, port: int): try: while not self._adopted.is_set(): sock, address = server.accept() self.on_status(f"Der Fernseher hat sich verbunden ({address[0]})") if not self._adopt(sock): sock.close() except OSError: pass def _adopt(self, sock: socket.socket) -> bool: """Die erste Verbindung gewinnt.""" if self._adopted.is_set(): return False self._adopted.set() for server in self._servers: try: server.close() except OSError: pass self._servers.clear() sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) self._session = SourceSession(sock, self.local_address, self.max_height, self.on_play, self.on_stopped, self.on_status) threading.Thread(target=self._session.run, daemon=True, name="rtsp-session").start() return True def try_outgoing(self, host: str, port: int): """Gegenrichtung - manche Geräte erwarten es andersherum.""" def attempt(): if self._adopted.is_set(): return try: sock = socket.create_connection((host, port), timeout=8) except OSError as error: log.debug("Ausgehender Versuch auf %s:%s erfolglos: %s", host, port, error) return self.on_status(f"Steuerkanal zum Fernseher aufgebaut (Port {port})") if not self._adopt(sock): sock.close() threading.Thread(target=attempt, daemon=True, name="rtsp-outgoing").start() @property def session(self) -> Optional[SourceSession]: return self._session def stop(self): self._adopted.set() for server in self._servers: try: server.close() except OSError: pass self._servers.clear() if self._session: self._session.stop()