415 lines
16 KiB
Python
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()
|