c43dd3ca4a
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.
269 lines
9.5 KiB
Swift
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) }
|
|
}
|
|
}
|