Files
2026-08-22 20:34:42 +02:00

415 lines
16 KiB
Python

"""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()