Files
fegger c43dd3ca4a feat(ios): native receiver app (Swift, min iOS 17)
A native iOS receiver so an iPhone can act as the second receiver, speaking the existing signaling + RTP protocol (no C++ changes) and mirroring the Android receiver (Phase 8) source-to-source.

- RTP core (header/packet, jitter buffer, H.264 depacketizer) ported from the Android receiver
- BSD-socket signaling server (dual-stack, most-recent-peer, never-throwing sends) + NSBonjourServices
- VideoToolbox H.264 decode (in-band SPS/PPS, real-time, rebuilds on size change) -> AVSampleBufferDisplayLayer
- PLI keyframe recovery (500 ms) + pendingOffer for late surface attach
- XcodeGen project + bootstrap.sh; XCTest port of the Android suite + new coverage
- .gitignore for generated artifacts; CHANGELOG; PHASES + MEMORY updated

Status: authored; on-device validation pending a Mac + Xcode 26 + iPhone 16.
2026-09-10 21:31:53 +02:00

269 lines
9.5 KiB
Swift

import Foundation
import CoreGraphics
/// The receiver pipeline, mirroring the C++ ReceiverPipeline and the Kotlin
/// receiver:
///
/// signaling server (offer answer)
/// UDP RTP jitter buffer depacketize VideoToolbox render sink
///
/// Recovery matches the C++ receiver: a damaged frame is dropped and a PLI
/// (rate-limited to one per 500 ms) asks the sender for a keyframe.
///
/// Concurrency: all shared state and every decode call run on a single serial
/// queue; the reader thread only polls/receives UDP and forwards datagrams to
/// that queue. UI callbacks are marshalled to the main thread.
final class ReceiverPipeline {
static let desiredUdpPort: UInt16 = 5004
static let desiredSignalingPort: UInt16 = 5005
static let pliMinIntervalMs: Int = 500
static let receiveBufferSize: Int32 = 4 * 1024 * 1024
private let localIP: String
private let displaySize: () -> CGSize
private let onStatus: (String) -> Void
private let onFirstFrame: () -> Void
private let onVideoSize: (CGSize) -> Void
private let queue = DispatchQueue(label: "sc.receiver.pipeline")
private var running = false
private var udp: UdpTransport?
private var signaling: SignalingServer?
private var readerThread: Thread?
private let jitter = JitterBuffer()
private var depacketizer = H264Depacketizer()
private var decoder: H264VideoToolboxDecoder?
private var renderSink: RenderSink?
private var pendingOffer: SignalingMessage?
private var activeSession = ""
private var firstFrameSeen = false
private var currentVideoSize = CGSize.zero
private var lastPliAtMs: Int64 = 0
init(localIP: String,
displaySize: @escaping () -> CGSize,
onStatus: @escaping (String) -> Void,
onFirstFrame: @escaping () -> Void,
onVideoSize: @escaping (CGSize) -> Void) {
self.localIP = localIP
self.displaySize = displaySize
self.onStatus = onStatus
self.onFirstFrame = onFirstFrame
self.onVideoSize = onVideoSize
}
/// The current decoded resolution (zero until the first keyframe).
func videoSize() -> CGSize {
return queue.sync { currentVideoSize }
}
// MARK: - Lifecycle
/// Binds the ports, starts the signaling server, and starts reading RTP.
func start() {
queue.async { [weak self] in
guard let self, !self.running else { return }
let udp = UdpTransport()
guard udp.bind(preferredPort: Self.desiredUdpPort, receiveBufferSize: Self.receiveBufferSize) else {
self.postStatus("Failed to bind the media port")
return
}
let signaling = SignalingServer(
onOffer: { [weak self] offer in self?.queue.async { self?.handleOffer(offer) } },
onPli: { _ in })
guard signaling.start(Self.desiredSignalingPort) else {
udp.close()
self.postStatus("Failed to start signaling")
return
}
self.running = true
self.udp = udp
self.signaling = signaling
let thread = Thread { [weak self] in self?.readLoop(udp: udp) }
thread.name = "rtp-reader"
self.readerThread = thread
thread.start()
let ip = self.localIP
let mediaPort = udp.port
let sigPort = signaling.port
self.postStatus("Listening on \(ip) (media :\(mediaPort), signaling :\(sigPort))\nWaiting for a sender… (fall back to: screencast --send --peer \(ip):\(sigPort))")
}
}
/// Points the (current or future) decoder at a render sink. A pending offer
/// (accepted while no sink existed) configures its decoder now and requests
/// a keyframe, since the sender only emits one when asked.
func attachSink(_ sink: RenderSink) {
queue.async { [weak self] in
guard let self else { return }
self.renderSink = sink
sink.attach()
guard let offer = self.pendingOffer else { return }
self.pendingOffer = nil
// The decoder configures on the first keyframe; the sink is already
// attached, so creation cannot fail. A late-configured decoder needs
// a keyframe (the sender only emits one when asked).
self.decoder = H264VideoToolboxDecoder(renderSink: sink)
self.requestPli()
}
}
/// Forgets a destroyed render sink so a later offer cannot render into it.
func detachSink() {
queue.async { [weak self] in
guard let self else { return }
self.decoder?.release()
self.decoder = nil
self.renderSink?.detach()
self.renderSink = nil
}
}
/// Stops listening; the pipeline can be started again.
func stop() {
queue.async { [weak self] in
guard let self, self.running else { return }
self.running = false
self.activeSession = ""
self.udp?.close()
self.udp = nil
self.readerThread?.join()
self.readerThread = nil
self.signaling?.close()
self.signaling = nil
self.decoder?.release()
self.decoder = nil
self.pendingOffer = nil
self.firstFrameSeen = false
self.currentVideoSize = .zero
self.postStatus("Stopped")
}
}
// MARK: - Media path
private func readLoop(udp: UdpTransport) {
while true {
switch udp.poll(timeoutMs: 100) {
case 1:
if let data = udp.receiveDatagram() {
queue.async { [weak self] in self?.processDatagram(data) }
}
case 0:
continue // timeout: re-check on the next poll
default:
return // closed/errored: the socket was closed by stop()
}
}
}
private func processDatagram(_ data: [UInt8]) {
guard let packet = RtpPacket.parse(data) else { return }
for released in jitter.push(packet) {
handleDepacketized(released)
}
}
private func handleDepacketized(_ packet: RtpPacket) {
let result = depacketizer.depacketize(packet)
if let accessUnit = result.accessUnit {
let presentation = Int(packet.header.timestamp & 0xFFFFFFFF)
guard let decoder = decoder else { return }
if !decoder.decode(accessUnit: accessUnit, rtpTimestamp: presentation, isKeyFrame: result.isKeyFrame) {
// Input/decode trouble: the dropped frame corrupts the GOP
// until the next keyframe ask for one.
requestPli()
}
if !firstFrameSeen {
firstFrameSeen = true
postFirstFrame()
}
postVideoSizeIfChanged()
}
if result.frameDropped {
requestPli()
}
}
private func postVideoSizeIfChanged() {
guard let size = decoder?.outputSize(), size != .zero else { return }
if size != currentVideoSize {
currentVideoSize = size
postVideoSize(size)
}
}
// MARK: - Signaling
private func handleOffer(_ offer: SignalingMessage) {
guard case let .offer(sessionId, codec, _, _, _, _, _, _) = offer else { return }
if codec != "h264" {
postStatus("Unsupported codec: \(codec)")
return
}
// New session: pristine decoder, reassembly state, and session.
decoder?.release()
decoder = nil
pendingOffer = nil
if renderSink != nil {
decoder = H264VideoToolboxDecoder(renderSink: renderSink!)
} else {
// No surface yet: park the offer; attachSink configures later.
pendingOffer = offer
}
depacketizer = H264Depacketizer()
jitter.clear()
firstFrameSeen = false
currentVideoSize = .zero
activeSession = sessionId
let size = displaySize()
let answer = SignalingMessage.answer(
sessionId: sessionId,
rtpAddress: "", // the sender targets the address of its own signaling connection
rtpPort: Int(udp?.port ?? 0),
displayWidth: Int(size.width),
displayHeight: Int(size.height))
signaling?.send(answer)
if decoder != nil {
postStatus("Session \(sessionId) negotiated — waiting for the first frame…")
} else {
postStatus("Session \(sessionId) negotiated — waiting for the display surface…")
}
}
/// Rate-limited keyframe request, callable from any thread (it always runs
/// on the pipeline queue in practice).
private func requestPli() {
let session = activeSession
guard !session.isEmpty else { return }
let now = Int64(Date().timeIntervalSince1970 * 1000)
if now - lastPliAtMs < Int64(Self.pliMinIntervalMs) { return }
lastPliAtMs = now
signaling?.send(.pli(sessionId: session))
}
// MARK: - UI callbacks (main thread)
private func postStatus(_ text: String) {
DispatchQueue.main.async { [onStatus] in onStatus(text) }
}
private func postFirstFrame() {
DispatchQueue.main.async { [onFirstFrame] in onFirstFrame() }
}
private func postVideoSize(_ size: CGSize) {
DispatchQueue.main.async { [onVideoSize] in onVideoSize(size) }
}
}