Findings from the post-implementation review of design/audio-quality-and-latency.md. **The bandwidth gap (highest).** Tier `High` (256 kbps) and the redundant `0xD2` plane were added separately, each costed as "~1 % of the video budget", and nobody added them together: 256 kbps sent twice is 512 kbps — ~2.5 % of a 20 Mbps session but ~10 % of a 5 Mbps one. Audio rides QUIC datagrams, OUTSIDE the ABR loop, so ABR could neither see that nor reclaim it; a constrained link quietly handed a tenth of its bandwidth to audio while ABR carefully managed the rest. `plan_audio_budget` now makes tier and redundancy ONE decision against the session's resolved video bitrate, ordered by preference rather than cost — transparent audio beats redundant audio, since the field report was about quality and redundancy only pays under loss, so `High` alone outranks `Standard`+redundancy even though they cost the same. It can lower what the operator asked for, never raise it, and never goes below `Low`: a stream with unintelligible audio is worse than one spending a few percent more. **The Linux host kept the exact defect fixed on Windows.** `let _ = tx.try_send(samples)` — silent, uncounted data loss, where the encoder concatenates across the hole, so every drop is a click AND a permanent shift of everything after it. WP0.2 turned out to be Windows-only and had not said so. Linux now shares `capture_policy::CaptureStats`: drops counted and warned, plus per-window peak/RMS/delivered%. A Linux audio report was until now exactly as un-triageable as the Windows one was on 2026-08-03. **Apple's WP0.3 was half-done** — `bufferedMS` was added and wired to nothing. The drain thread now logs buffer/target/underruns/sheds like the other three, from one locked snapshot so the numbers in a line describe the same instant. Also: the Linux "audio format negotiated" line now says WHICH mode produced it, because that changes what it is worth — in stream-sink mode the host owns the sink so the mix cannot have been narrowed upstream, but in legacy monitor mode a 16 kHz Bluetooth sink would still be reported as a clean 48 kHz through PipeWire's resampler, the same way WASAPI's autoconvert hid it on Windows. Reading the monitored node's own rate needs a registry lookup this stream does not do; recorded as an open gap rather than implied to be covered. Two stale docs: `audio_wasapi.rs` cited `clients/windows/src/audio.rs` (deleted) and still described the pre-shared-policy "prime to ~3 quanta" behaviour. And the Apple ring's `prefill:` parameter, dead since the depth moved into the ring, is gone. Verified: clippy --all-targets -D warnings on Linux (docker) AND Windows (runner .133, forced clean rebuild of punktfunk-host + pf-client-core); core 167 tests; host 57 audio tests on Windows; Android clippy count identical to pristine (6, all documented arm64 artifacts); Apple ring re-simulated. The host suite's `gamestream::stream::tests::sender_delivers_batches` fails under qemu — the recorded environmental flake, unrelated to audio, green on the earlier less-loaded run. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
824 lines
41 KiB
Swift
824 lines
41 KiB
Swift
// Session audio, both directions:
|
||
//
|
||
// host → speaker: a drain thread pulls Opus packets (nextAudio, its own plane in the
|
||
// core), decodes via OpusDecoder, and writes PCM into a jitter ring; an
|
||
// AVAudioSourceNode pulls from the ring (silence on underrun with re-priming, so a
|
||
// network gap costs one dip, not permanent crackle).
|
||
//
|
||
// mic → host: a tap on the input node folds the capture to one mono bus (the chosen channel
|
||
// of a multi-channel interface, or a sum of all channels), resamples to 48 kHz mono, slices
|
||
// 10 ms chunks, Opus-encodes, and sendMic()s each packet — the host feeds them into a
|
||
// virtual PipeWire source.
|
||
//
|
||
// Engine topology. With the mic enabled and echo cancellation on (both defaults), BOTH
|
||
// directions run on ONE AVAudioEngine with the system voice processor engaged
|
||
// (`setVoiceProcessingEnabled`) — AEC needs render and capture on the same unit so it can
|
||
// subtract what the speaker is playing from what the mic hears; without it, a loudspeaker
|
||
// client feeds the host's own game audio straight back to it (the primary reported echo
|
||
// source). The voice processor can only follow the system DEFAULT devices, so explicit
|
||
// endpoint choices fall back to the old two-engine topology — see `wantsCombined` for the
|
||
// exact decision, and `startCapture` for why two engines handle arbitrary device pairs.
|
||
//
|
||
// Devices are chosen by UID ("" = system default: the engine is then never pinned to a
|
||
// concrete device and follows default-device changes).
|
||
|
||
import AVFoundation
|
||
import os
|
||
|
||
private let log = Logger(subsystem: "io.unom.punktfunk", category: "audio")
|
||
|
||
/// Render-block-owned scratch storage: freed exactly when the closure (and thus the
|
||
/// last possible render call) is released — never racing CoreAudio.
|
||
private final class ScratchBuffer {
|
||
// 8192 frames × up to 8 channels (7.1) — the render block caps `frames` at 8192.
|
||
let ptr = UnsafeMutablePointer<Float>.allocate(capacity: 8192 * 8)
|
||
deinit { ptr.deallocate() }
|
||
}
|
||
|
||
public final class SessionAudio {
|
||
private let connection: PunktfunkConnection
|
||
private let flag = StopFlag()
|
||
private let drainDone = DispatchSemaphore(value: 0)
|
||
/// Owns the engine handles + drainStarted, paired with `flag`: stop() sets the flag
|
||
/// BEFORE taking the engines, every publisher re-checks the flag under this lock
|
||
/// after publishing-side work — so a startCapture racing stop() (the mic-permission
|
||
/// callback arrives whenever the user clicks the prompt) can never leave a hot
|
||
/// microphone with no owner.
|
||
private let stateLock = NSLock()
|
||
private var playbackEngine: AVAudioEngine?
|
||
private var captureEngine: AVAudioEngine?
|
||
/// The one engine running BOTH directions when the voice processor is engaged;
|
||
/// `playbackEngine`/`captureEngine` stay nil while this is set.
|
||
private var combinedEngine: AVAudioEngine?
|
||
private var drainStarted = false
|
||
/// The mute LATCH: the effective mute the owner last asked for (see `setMicMuted`). Held
|
||
/// because the uplink engine can appear LATER than the request — the mic permission prompt
|
||
/// is answered at the user's leisure, and the engine is built on the grant — so a mute set
|
||
/// in the meantime must be waiting for it. Applied by whichever start path wins the race.
|
||
/// Guarded by `stateLock`, like the engines it applies to.
|
||
private var micMuted = false
|
||
/// The playback jitter ring — created by whichever engine starts playback first and KEPT
|
||
/// across an engine rebuild (the permission-grant upgrade in `startEngines` swaps engines,
|
||
/// not the ring, so the drain thread never has to be re-pointed). Main-thread confined,
|
||
/// like every start path.
|
||
private var ring: AudioRing?
|
||
#if !os(macOS)
|
||
/// AVAudioSession `setCategory`/`setActive` are synchronous and block on the audio server, so
|
||
/// they must not run on the main thread (UI stall — AVFoundation warns about it). PROCESS-WIDE
|
||
/// (static) so every SessionAudio shares one serial queue: the AVAudioSession is a process
|
||
/// singleton, and across a reconnect the old session's deactivate must be ordered before the
|
||
/// new session's activate (a per-instance queue would let them race and leave the new session's
|
||
/// audio deactivated). stop() enqueues its deactivate promptly so it lands before the next
|
||
/// session's activate.
|
||
private static let sessionQueue = DispatchQueue(label: "io.unom.punktfunk.audio.session")
|
||
#endif
|
||
|
||
public init(connection: PunktfunkConnection) {
|
||
self.connection = connection
|
||
}
|
||
|
||
/// Backstop for an owner dropping us without stop() — unblocks the drain thread
|
||
/// (which captures the connection strongly, NOT self) within one poll timeout.
|
||
/// Engine teardown still belongs to stop().
|
||
deinit {
|
||
flag.stop()
|
||
}
|
||
|
||
/// Start playback (and, if enabled+authorized, the mic uplink). Empty UIDs = system default
|
||
/// device; on iOS the UIDs are ignored entirely (routes are AVAudioSession-managed). On macOS
|
||
/// the engines start synchronously on the caller's (main) thread. On iOS/tvOS start() is
|
||
/// ASYNCHRONOUS: it activates the AVAudioSession off the main thread, then starts the engines on
|
||
/// a later main-queue hop (gated by `!flag.isStopped`) — so playback is live shortly after, not
|
||
/// on return. The mic may start later still if the permission prompt is pending.
|
||
/// `echoCancel` picks the engine topology — see the header note and `wantsCombined`.
|
||
public func start(
|
||
speakerUID: String, micUID: String, micChannel: Int, micEnabled: Bool, echoCancel: Bool
|
||
) {
|
||
#if os(macOS)
|
||
// No AVAudioSession on macOS — start the engines directly (caller's thread, as before).
|
||
startEngines(
|
||
speakerUID: speakerUID, micUID: micUID, micChannel: micChannel,
|
||
micEnabled: micEnabled, echoCancel: echoCancel)
|
||
#else
|
||
// Configure + activate the session OFF the main thread (it blocks on the audio server),
|
||
// then start the engines back on the main thread once it's active — engine routing/format
|
||
// depend on the active session. A stop() racing in between is caught by the flag guard.
|
||
Self.sessionQueue.async { [weak self] in
|
||
guard let self else { return }
|
||
self.activateAudioSession(micEnabled: micEnabled)
|
||
DispatchQueue.main.async { [weak self] in
|
||
guard let self, !self.flag.isStopped else { return }
|
||
self.startEngines(
|
||
speakerUID: speakerUID, micUID: micUID, micChannel: micChannel,
|
||
micEnabled: micEnabled, echoCancel: echoCancel)
|
||
}
|
||
}
|
||
#endif
|
||
}
|
||
|
||
#if !os(macOS)
|
||
/// Route + policy live in the session, not per-engine: stereo playback, mic capture when
|
||
/// enabled, Bluetooth allowed. Failure is non-fatal (defaults). Runs on `sessionQueue`.
|
||
private func activateAudioSession(micEnabled: Bool) {
|
||
let session = AVAudioSession.sharedInstance()
|
||
do {
|
||
#if os(iOS)
|
||
if micEnabled {
|
||
// .defaultToSpeaker: .playAndRecord otherwise routes to the iPhone EARPIECE; only
|
||
// affects the built-in route (headphones/BT still win).
|
||
try session.setCategory(
|
||
.playAndRecord, mode: .default,
|
||
options: [.allowBluetoothA2DP, .defaultToSpeaker])
|
||
// Uplink latency: ask for 5 ms IO quanta at the wire rate (the default ~10-23 ms
|
||
// quantum is most of the mic path's burst latency). Best-effort — the hardware
|
||
// has the final word (a Bluetooth route will ignore both), and whatever quantum
|
||
// is actually granted, the capture tap handles the buffers it gets.
|
||
try? session.setPreferredIOBufferDuration(0.005)
|
||
try? session.setPreferredSampleRate(48_000)
|
||
} else {
|
||
try session.setCategory(.playback, mode: .default)
|
||
}
|
||
#else // tvOS — no app-accessible mic
|
||
try session.setCategory(.playback, mode: .default)
|
||
#endif
|
||
try session.setActive(true)
|
||
} catch {
|
||
log.warning("AVAudioSession setup failed: \(error.localizedDescription)")
|
||
}
|
||
}
|
||
#endif
|
||
|
||
/// Build + start the engines — combined (voice-processed) or split, per `wantsCombined` —
|
||
/// with the mic uplink only when enabled + authorized. Main thread (engine setup); on
|
||
/// iOS/tvOS the session is already active by the time this runs.
|
||
private func startEngines(
|
||
speakerUID: String, micUID: String, micChannel: Int, micEnabled: Bool, echoCancel: Bool
|
||
) {
|
||
#if os(tvOS)
|
||
// No app-accessible microphone input on tvOS — playback only.
|
||
startPlayback(speakerUID: speakerUID)
|
||
#else
|
||
guard micEnabled else {
|
||
startPlayback(speakerUID: speakerUID)
|
||
return
|
||
}
|
||
let combined = wantsCombined(
|
||
speakerUID: speakerUID, micUID: micUID, micChannel: micChannel,
|
||
echoCancel: echoCancel)
|
||
switch AVCaptureDevice.authorizationStatus(for: .audio) {
|
||
case .authorized:
|
||
if combined {
|
||
startCombined(speakerUID: speakerUID, micUID: micUID, micChannel: micChannel)
|
||
} else {
|
||
startPlayback(speakerUID: speakerUID)
|
||
startCapture(micUID: micUID, micChannel: micChannel)
|
||
}
|
||
case .notDetermined:
|
||
// Playback must not wait out the permission prompt (the user answers at their
|
||
// leisure) — start it now, and on a grant either bolt the capture engine on
|
||
// (split) or swap the playback engine for the combined one (the ring and its
|
||
// drain thread carry over — see `makePlaybackChain`).
|
||
startPlayback(speakerUID: speakerUID)
|
||
AVCaptureDevice.requestAccess(for: .audio) { [weak self] granted in
|
||
DispatchQueue.main.async {
|
||
guard let self, granted, !self.flag.isStopped else { return }
|
||
if combined {
|
||
self.stateLock.lock()
|
||
let playback = self.playbackEngine
|
||
self.playbackEngine = nil
|
||
self.stateLock.unlock()
|
||
playback?.stop()
|
||
self.startCombined(
|
||
speakerUID: speakerUID, micUID: micUID, micChannel: micChannel)
|
||
} else {
|
||
self.startCapture(micUID: micUID, micChannel: micChannel)
|
||
}
|
||
}
|
||
}
|
||
default:
|
||
startPlayback(speakerUID: speakerUID)
|
||
log.warning("microphone access denied — mic uplink disabled (System Settings → Privacy)")
|
||
}
|
||
#endif
|
||
}
|
||
|
||
#if !os(tvOS)
|
||
/// One engine or two: the voice processor requires render + capture on one unit, and that
|
||
/// unit can only follow the system DEFAULT devices — so echo cancellation gets the combined
|
||
/// engine only while nothing is explicitly pinned. On macOS a chosen speaker/mic UID or a
|
||
/// picked input channel (the voice processor's capture side is its own mono mix — a
|
||
/// per-channel pick can't survive it) keeps today's two-engine path, AEC-less but honoring
|
||
/// the exact endpoints the user named. On iOS routes are session-managed and the UIDs are
|
||
/// ignored, so the toggle alone decides.
|
||
private func wantsCombined(
|
||
speakerUID: String, micUID: String, micChannel: Int, echoCancel: Bool
|
||
) -> Bool {
|
||
guard echoCancel else { return false }
|
||
#if os(macOS)
|
||
return speakerUID.isEmpty && micUID.isEmpty && micChannel == 0
|
||
#else
|
||
return true
|
||
#endif
|
||
}
|
||
#endif
|
||
|
||
/// Stop both directions. Safe from any thread; waits the drain thread out (≤ its
|
||
/// poll timeout) so the caller can close the connection right after.
|
||
public func stop() {
|
||
flag.stop() // before taking the engines — see stateLock's comment
|
||
stateLock.lock()
|
||
let capture = captureEngine
|
||
captureEngine = nil
|
||
let playback = playbackEngine
|
||
playbackEngine = nil
|
||
let combined = combinedEngine
|
||
combinedEngine = nil
|
||
let wasDraining = drainStarted
|
||
drainStarted = false
|
||
stateLock.unlock()
|
||
if let capture {
|
||
capture.inputNode.removeTap(onBus: 0)
|
||
capture.stop()
|
||
}
|
||
playback?.stop()
|
||
if let combined {
|
||
combined.inputNode.removeTap(onBus: 0)
|
||
combined.stop()
|
||
}
|
||
#if !os(macOS)
|
||
// Release the session so audio we interrupted (Music, podcasts) gets its resume cue. Like
|
||
// activation, setActive is synchronous/blocking — run it on the shared serial session queue
|
||
// (off the main thread). Enqueued HERE — engines already stopped, and BEFORE the drain wait
|
||
// below — so across a reconnect it lands ahead of the next session's activate on the shared
|
||
// queue (otherwise a deferred deactivate could deactivate the new session). Fire-and-forget.
|
||
Self.sessionQueue.async {
|
||
do {
|
||
try AVAudioSession.sharedInstance().setActive(
|
||
false, options: .notifyOthersOnDeactivation)
|
||
} catch {
|
||
log.warning("AVAudioSession deactivation failed: \(error.localizedDescription)")
|
||
}
|
||
}
|
||
#endif
|
||
if wasDraining {
|
||
_ = drainDone.wait(timeout: .now() + .milliseconds(400))
|
||
}
|
||
}
|
||
|
||
/// Silence the mic uplink (no room audio leaves the device) or restore it. THE one muting
|
||
/// mechanism: the owner composes its reasons — the user's in-stream mute and the background
|
||
/// keep-alive's privacy mute — into one effective state and passes that here, so neither can
|
||
/// clear the other (see `SessionModel.applyMicMute`).
|
||
///
|
||
/// Two-engine sessions pause/resume the capture engine; a combined session instead mutes the
|
||
/// voice processor's input (playback shares that engine and must keep running, so the engine
|
||
/// itself never pauses — the mute zeroes the mic at the IO unit, and the tap encodes silence).
|
||
/// Local and instant either way: nothing is negotiated with the host, and the packets that do
|
||
/// leave carry silence. A no-op when there's no uplink (playback-only / tvOS / mic disabled),
|
||
/// except that the state is LATCHED for an uplink that starts later. The audio SESSION stays
|
||
/// active for background playback, so iOS may keep showing the recording indicator until a
|
||
/// full reconfigure — either path stops room audio leaving the device, which is the
|
||
/// privacy-relevant part. Main thread.
|
||
public func setMicMuted(_ muted: Bool) {
|
||
stateLock.lock()
|
||
micMuted = muted
|
||
let capture = captureEngine
|
||
let combined = combinedEngine
|
||
stateLock.unlock()
|
||
apply(micMuted: muted, capture: capture, combined: combined)
|
||
}
|
||
|
||
/// Push the latched mute onto whichever engine carries the uplink. Split out from
|
||
/// `setMicMuted` because the start paths call it too, with the engine they just started —
|
||
/// that's how a mute requested before the permission grant lands on the engine the grant
|
||
/// creates. Never resumes a stopped session's engine.
|
||
private func apply(micMuted muted: Bool, capture: AVAudioEngine?, combined: AVAudioEngine?) {
|
||
if let combined {
|
||
combined.inputNode.isVoiceProcessingInputMuted = muted
|
||
return
|
||
}
|
||
guard let capture else { return }
|
||
if muted {
|
||
capture.pause()
|
||
} else if !flag.isStopped {
|
||
try? capture.start()
|
||
}
|
||
}
|
||
|
||
// MARK: - Playback (host → speaker)
|
||
|
||
/// The playback jitter ring + the source node draining it — shared by the plain playback
|
||
/// engine and the combined voice-processing engine, and REUSED across an engine rebuild
|
||
/// (same session, same ring: the drain thread keeps writing right through the swap). nil
|
||
/// when the host's channel layout can't be expressed (already logged). Main thread.
|
||
private func makePlaybackChain()
|
||
-> (ring: AudioRing, source: AVAudioSourceNode, format: AVAudioFormat)?
|
||
{
|
||
// Build the playback layout from the host-RESOLVED channel count (never the request):
|
||
// 2 = stereo / 6 = 5.1 / 8 = 7.1, canonical wire order FL FR FC LFE RL RR SL SR.
|
||
let channels = Int(connection.resolvedAudioChannels)
|
||
// 1 s interleaved capacity, scaled by the channel count. The de-jitter depth itself is
|
||
// the ring's own business now (`AudioRing.targetMS`, mirroring `JitterTuning::COREAUDIO`)
|
||
// rather than a prefill passed in here.
|
||
let ring = self.ring ?? AudioRing(capacity: 48_000 * channels, channels: channels)
|
||
self.ring = ring
|
||
|
||
// Engine-native deinterleaved float; the render block deinterleaves from the ring. Surround
|
||
// uses an explicit wire-order channel layout; the mixer downmixes to the output device when
|
||
// it has fewer speakers (e.g. an iPhone's stereo built-ins). (Explicit if/else rather than
|
||
// map/flatMap so it's correct whether the channelLayout initializer is failable or not.)
|
||
var format: AVAudioFormat?
|
||
if channels == 2 {
|
||
format = AVAudioFormat(standardFormatWithSampleRate: 48_000, channels: 2)
|
||
} else if let layout = wireChannelLayout(channels: channels) {
|
||
format = AVAudioFormat(standardFormatWithSampleRate: 48_000, channelLayout: layout)
|
||
}
|
||
guard let format else {
|
||
log.error("could not build \(channels)-channel audio format — audio disabled")
|
||
return nil
|
||
}
|
||
let scratch = ScratchBuffer() // block-owned; freed with the closure
|
||
let source = AVAudioSourceNode(format: format) { _, _, frameCount, abl -> OSStatus in
|
||
let frames = Int(frameCount)
|
||
guard frames <= 8192 else { return kAudioUnitErr_TooManyFramesToProcess }
|
||
ring.read(into: scratch.ptr, count: frames * channels)
|
||
let buffers = UnsafeMutableAudioBufferListPointer(abl)
|
||
// Deinterleave the wire-order interleaved ring into the engine's per-channel buses.
|
||
if buffers.count >= channels {
|
||
for ch in 0..<channels {
|
||
if let dst = buffers[ch].mData?.assumingMemoryBound(to: Float.self) {
|
||
for f in 0..<frames { dst[f] = scratch.ptr[f * channels + ch] }
|
||
}
|
||
}
|
||
}
|
||
return noErr
|
||
}
|
||
return (ring, source, format)
|
||
}
|
||
|
||
private func startPlayback(speakerUID: String) {
|
||
guard let (ring, source, format) = makePlaybackChain() else { return }
|
||
let engine = AVAudioEngine()
|
||
#if os(macOS)
|
||
if !speakerUID.isEmpty {
|
||
if let dev = AudioDevices.deviceID(forUID: speakerUID),
|
||
let unit = engine.outputNode.audioUnit {
|
||
if !Self.setDevice(dev, on: unit) {
|
||
log.error("could not select speaker \(speakerUID) — using default")
|
||
}
|
||
} else {
|
||
log.warning("speaker \(speakerUID) not present — using default")
|
||
}
|
||
}
|
||
#endif
|
||
engine.attach(source)
|
||
engine.connect(source, to: engine.mainMixerNode, format: format)
|
||
engine.prepare()
|
||
do {
|
||
try engine.start()
|
||
} catch {
|
||
log.error("playback engine failed to start: \(error.localizedDescription)")
|
||
return
|
||
}
|
||
stateLock.lock()
|
||
if flag.isStopped {
|
||
stateLock.unlock()
|
||
engine.stop() // stop() already ran — don't strand a started engine
|
||
return
|
||
}
|
||
playbackEngine = engine
|
||
stateLock.unlock()
|
||
startDrain(into: ring)
|
||
}
|
||
|
||
/// Idempotent — the permission-grant engine swap reaches here a second time with the
|
||
/// drain thread already feeding the (carried-over) ring.
|
||
private func startDrain(into ring: AudioRing) {
|
||
stateLock.lock()
|
||
if drainStarted {
|
||
stateLock.unlock()
|
||
return
|
||
}
|
||
drainStarted = true
|
||
stateLock.unlock()
|
||
let thread = Thread { [connection, flag, drainDone] in
|
||
defer { drainDone.signal() }
|
||
var drained = 0
|
||
// Decode happens IN-CORE (libopus multistream) — AudioToolbox's Opus path is
|
||
// stereo-only — and is handed back as interleaved f32 PCM in wire channel order.
|
||
// Per-iteration autorelease pool: no runloop on this thread (see Stage2Pipeline).
|
||
var alive = true
|
||
while alive, !flag.isStopped {
|
||
alive = autoreleasepool { () -> Bool in
|
||
let pcm: PunktfunkConnection.AudioPCM?
|
||
do {
|
||
pcm = try connection.nextAudioPcm(timeoutMs: 100)
|
||
} catch {
|
||
return false // session closed
|
||
}
|
||
guard let pcm, pcm.frameCount > 0 else { return true }
|
||
pcm.samples.withUnsafeBufferPointer { p in
|
||
if let base = p.baseAddress {
|
||
ring.write(base, count: pcm.frameCount * pcm.channels)
|
||
}
|
||
}
|
||
// Periodic vitals (~10 s at the protocol's 5 ms frames). The other three clients
|
||
// log buffer depth and underruns; without this an Apple audio report — latency or
|
||
// dropout — arrives with no numbers at all, which is the position every platform
|
||
// was in before the 2026-08 audio work.
|
||
drained += 1
|
||
if drained % 2_000 == 0 {
|
||
let s = ring.stats
|
||
log.info(
|
||
"audio: buffer_ms=\(s.bufferedMS) target_ms=\(s.targetMS) underruns=\(s.underruns) drift_sheds=\(s.sheds)"
|
||
)
|
||
}
|
||
return true
|
||
}
|
||
}
|
||
}
|
||
thread.name = "punktfunk-audio"
|
||
thread.qualityOfService = .userInteractive
|
||
thread.start()
|
||
}
|
||
|
||
// MARK: - Mic (mic → host)
|
||
|
||
#if !os(tvOS)
|
||
/// One engine, both directions: engage the system voice processor on the shared IO unit
|
||
/// (AEC + noise suppression + AGC), hang the playback source off its render side and the
|
||
/// mic tap off its capture side. Every failure falls back to a WORKING configuration —
|
||
/// the split path (no AEC) when the voice processor won't engage, plain playback when the
|
||
/// mic chain can't be built — a session never loses audio to the echo-cancel feature.
|
||
private func startCombined(speakerUID: String, micUID: String, micChannel: Int) {
|
||
let engine = AVAudioEngine()
|
||
let input = engine.inputNode
|
||
do {
|
||
// Before anything reads the input's format: the voice processor changes it (often
|
||
// to its own mono mix, sometimes at a lower rate) — installMicTap reads the format
|
||
// AFTER this, so the converter chain adapts to whatever the processor emits.
|
||
try input.setVoiceProcessingEnabled(true)
|
||
} catch {
|
||
log.warning("""
|
||
voice processing unavailable (\(error.localizedDescription)) — separate \
|
||
engines, no echo cancellation
|
||
""")
|
||
startPlayback(speakerUID: speakerUID)
|
||
startCapture(micUID: micUID, micChannel: micChannel)
|
||
return
|
||
}
|
||
// Symmetric enable for the render side; with both directions on one engine the
|
||
// input-node enable already covers it, so a refusal here is not a failure.
|
||
try? engine.outputNode.setVoiceProcessingEnabled(true)
|
||
// This is a game stream, not a call: never duck the host's audio under the outgoing
|
||
// voice. .min is the closest to "off" the API offers, and advanced (selective)
|
||
// ducking stays off with it.
|
||
input.voiceProcessingOtherAudioDuckingConfiguration = .init(
|
||
enableAdvancedDucking: false, duckingLevel: .min)
|
||
|
||
guard let (ring, source, format) = makePlaybackChain() else {
|
||
// Playback impossible (logged) — keep the uplink alive, as the split path would.
|
||
startCapture(micUID: micUID, micChannel: micChannel)
|
||
return
|
||
}
|
||
engine.attach(source)
|
||
engine.connect(source, to: engine.mainMixerNode, format: format)
|
||
guard installMicTap(on: input, micUID: micUID, micChannel: micChannel) else {
|
||
// Mic chain unavailable (logged) — keep the session audible on the plain playback
|
||
// engine rather than playing through an idle voice processor.
|
||
startPlayback(speakerUID: speakerUID)
|
||
return
|
||
}
|
||
engine.prepare()
|
||
do {
|
||
try engine.start()
|
||
} catch {
|
||
log.error("combined engine failed to start: \(error.localizedDescription)")
|
||
input.removeTap(onBus: 0)
|
||
startPlayback(speakerUID: speakerUID) // no echo cancellation beats no audio
|
||
return
|
||
}
|
||
stateLock.lock()
|
||
if flag.isStopped {
|
||
stateLock.unlock()
|
||
input.removeTap(onBus: 0)
|
||
engine.stop() // stop() already ran — don't strand a started engine (or a hot mic)
|
||
return
|
||
}
|
||
combinedEngine = engine
|
||
let muted = micMuted // latched before this engine existed (a mute during the prompt)
|
||
stateLock.unlock()
|
||
apply(micMuted: muted, capture: nil, combined: engine)
|
||
startDrain(into: ring)
|
||
log.info("audio engines joined — voice processing (echo cancellation) active")
|
||
}
|
||
|
||
/// The split path: capture on its OWN engine, playback on another — the pre-echo-cancel
|
||
/// topology, kept verbatim. Two engines, not one — a single AVAudioEngine ties
|
||
/// input+output to one aggregate clock, separate engines keep arbitrary mic/speaker
|
||
/// combinations trivial. That freedom is exactly why the voice processor can't ride this
|
||
/// path (AEC needs both directions on one unit) and why explicitly pinned endpoints land
|
||
/// here — see `wantsCombined`.
|
||
private func startCapture(micUID: String, micChannel: Int) {
|
||
let engine = AVAudioEngine()
|
||
let input = engine.inputNode
|
||
#if os(macOS)
|
||
if !micUID.isEmpty {
|
||
if let dev = AudioDevices.deviceID(forUID: micUID), let unit = input.audioUnit {
|
||
if !Self.setDevice(dev, on: unit) {
|
||
log.error("could not select microphone \(micUID) — using default")
|
||
}
|
||
} else {
|
||
log.warning("microphone \(micUID) not present — using default")
|
||
}
|
||
}
|
||
#endif
|
||
guard installMicTap(on: input, micUID: micUID, micChannel: micChannel) else { return }
|
||
engine.prepare()
|
||
do {
|
||
try engine.start()
|
||
} catch {
|
||
log.error("capture engine failed to start: \(error.localizedDescription)")
|
||
input.removeTap(onBus: 0)
|
||
return
|
||
}
|
||
stateLock.lock()
|
||
if flag.isStopped {
|
||
// stop() ran while we were starting (the permission prompt resolves at the
|
||
// user's leisure) — tear the engine down ourselves, nobody else owns it now.
|
||
stateLock.unlock()
|
||
input.removeTap(onBus: 0)
|
||
engine.stop()
|
||
return
|
||
}
|
||
captureEngine = engine
|
||
let muted = micMuted // latched before this engine existed (a mute during the prompt)
|
||
stateLock.unlock()
|
||
apply(micMuted: muted, capture: engine, combined: nil)
|
||
log.info("mic uplink started (\(micUID.isEmpty ? "default input" : micUID))")
|
||
}
|
||
|
||
/// Resolve the input's live format + fold plan, build the mono→Opus chain, and install the
|
||
/// capture tap on `input` — everything mic except engine ownership, shared verbatim by the
|
||
/// combined and split topologies. Reads `input.outputFormat(forBus:)` at call time, so the
|
||
/// chain follows whatever the node emits: the raw device format, or the voice processor's
|
||
/// own mix when that's enabled. False (logged) when no input is usable or the encoder
|
||
/// can't be built; the tap is installed on true.
|
||
private func installMicTap(
|
||
on input: AVAudioInputNode, micUID: String, micChannel: Int
|
||
) -> Bool {
|
||
let inFormat = input.outputFormat(forBus: 0)
|
||
guard inFormat.sampleRate > 0, inFormat.channelCount > 0 else {
|
||
log.error("no usable input device — mic uplink disabled")
|
||
return false
|
||
}
|
||
|
||
// Multi-channel-interface handling. A pro interface exposes N discrete inputs with the mic
|
||
// on ONE of them, but AVAudioConverter's N→stereo downmix takes channels 0/1 — dead
|
||
// silence when the mic sits higher up (the classic "host receives zeros"). So we fold the
|
||
// input to a single mono bus OURSELVES and resample that. micChannel: 0 = Auto (sum every
|
||
// channel — a lone hot mic passes at full level), n≥1 pins 1-based input channel n.
|
||
let inChannels = Int(inFormat.channelCount)
|
||
let pinnedChannel: Int? = {
|
||
guard micChannel >= 1 else { return nil }
|
||
let idx = micChannel - 1
|
||
guard idx < inChannels else {
|
||
log.warning(
|
||
"mic channel \(micChannel) out of range (device has \(inChannels)) — mixing all")
|
||
return nil
|
||
}
|
||
return idx
|
||
}()
|
||
let channelPlan = pinnedChannel.map { "channel \($0 + 1)/\(inChannels)" }
|
||
?? (inChannels > 1 ? "mix \(inChannels)ch→mono" : "mono")
|
||
|
||
// Name the device we're ACTUALLY recording from + its format + how we fold it, once per
|
||
// session. This single line localizes the whole class of "host receives silence" failures
|
||
// that otherwise need a host-side tone injection to pin down: a UID that silently fell back
|
||
// to the default, the wrong device being live, or the wrong channel picked.
|
||
#if os(macOS)
|
||
if let unit = input.audioUnit, let live = Self.currentDevice(of: unit),
|
||
let dev = AudioDevices.describe(live) {
|
||
if !micUID.isEmpty, dev.uid != micUID {
|
||
log.warning("""
|
||
mic selection not honored — requested \(micUID) but capturing from \
|
||
\(dev.name) [\(dev.uid)]; the device's UID likely changed (replug) — \
|
||
reselect it in Settings
|
||
""")
|
||
}
|
||
log.info("""
|
||
mic capture: \(dev.name) [\(dev.uid)] — \(Int(inFormat.sampleRate)) Hz, \
|
||
\(inChannels) ch, \(channelPlan)
|
||
""")
|
||
} else {
|
||
log.info("""
|
||
mic capture: <device unavailable> — \(Int(inFormat.sampleRate)) Hz, \
|
||
\(inChannels) ch, \(channelPlan)
|
||
""")
|
||
}
|
||
#else
|
||
log.info(
|
||
"mic capture: \(Int(inFormat.sampleRate)) Hz, \(inChannels) ch, \(channelPlan)")
|
||
#endif
|
||
|
||
// Encode a single mono bus (folded from `inFormat` in the tap): the resampler goes
|
||
// mono@inputSR → the encoder's 48 kHz mono, so it handles the rate change and the
|
||
// wrong-channel downmix never happens. Mono end to end — the host's decoder upmixes,
|
||
// so the old duplicate-into-stereo step only cost bits and cycles.
|
||
//
|
||
// `mono`/`staging` are the per-callback scratch buffers, preallocated HERE (grown only
|
||
// if a larger-than-expected device quantum ever arrives) — the steady-state tap path
|
||
// allocates nothing.
|
||
let scratchFrames: AVAudioFrameCount = 8192
|
||
let stagingCapacity = { (frames: AVAudioFrameCount) -> AVAudioFrameCount in
|
||
AVAudioFrameCount(
|
||
(Double(frames) * 48_000 / inFormat.sampleRate).rounded(.up)) + 64
|
||
}
|
||
guard let monoFormat = AVAudioFormat(
|
||
commonFormat: .pcmFormatFloat32, sampleRate: inFormat.sampleRate,
|
||
channels: 1, interleaved: false),
|
||
let encoder = try? OpusEncoder(),
|
||
let resampler = AVAudioConverter(from: monoFormat, to: encoder.pcmFormat),
|
||
let chunk = AVAudioPCMBuffer(
|
||
pcmFormat: encoder.pcmFormat, frameCapacity: encoder.framesPerPacket),
|
||
let monoScratch = AVAudioPCMBuffer(
|
||
pcmFormat: monoFormat, frameCapacity: scratchFrames),
|
||
let stagingScratch = AVAudioPCMBuffer(
|
||
pcmFormat: encoder.pcmFormat, frameCapacity: stagingCapacity(scratchFrames))
|
||
else {
|
||
log.error("Opus encoder unavailable — mic uplink disabled")
|
||
return false
|
||
}
|
||
|
||
// Tap-thread-confined state: fold into `mono`, resample into `staging`, accumulate in
|
||
// `fifo`, slice `framesPerPacket` (10 ms) chunks for the encoder.
|
||
var mono = monoScratch
|
||
var staging = stagingScratch
|
||
var fifo: [Float] = []
|
||
fifo.reserveCapacity(48_000)
|
||
var seq: UInt32 = 0
|
||
let connection = connection
|
||
let flag = flag
|
||
|
||
// Silence tripwire (tap-confined): a "recording" app can be handed pure digital zeros —
|
||
// a zeroed input-volume slider, a stale TCC grant, a muted device, OR the wrong channel
|
||
// picked — and everything downstream looks alive while the host gets silence. Track the
|
||
// peak of the EXTRACTED mono bus over the first ~10 s (not the raw device — a mic present
|
||
// on a channel we didn't grab must still read as silence) and emit exactly ONE verdict.
|
||
// This is the log line whose absence made the last occurrence take a host-side tone.
|
||
let silenceWindow = Int(inFormat.sampleRate * 10)
|
||
let deviceLabel = micUID.isEmpty ? "default input" : micUID
|
||
var framesInspected = 0
|
||
var inputPeak: Float = 0
|
||
var levelReported = false
|
||
|
||
// 480 frames = 10 ms, matching the packet duration. Advisory — CoreAudio delivers the
|
||
// device quantum whatever we ask (the old 2048 request came back as 42.7 ms bursts, most
|
||
// of the uplink's latency) — but where the system honors it, the tap fires per-packet.
|
||
input.installTap(onBus: 0, bufferSize: 480, format: inFormat) { buffer, _ in
|
||
if flag.isStopped { return }
|
||
let frames = Int(buffer.frameLength)
|
||
guard frames > 0, let src = buffer.floatChannelData else { return }
|
||
if frames > Int(mono.frameCapacity) {
|
||
// A quantum larger than the scratch (bufferSize is advisory both ways) — regrow
|
||
// once to the new high-water mark; the steady state stays allocation-free.
|
||
guard let biggerMono = AVAudioPCMBuffer(
|
||
pcmFormat: monoFormat, frameCapacity: buffer.frameLength),
|
||
let biggerStaging = AVAudioPCMBuffer(
|
||
pcmFormat: encoder.pcmFormat,
|
||
frameCapacity: stagingCapacity(buffer.frameLength))
|
||
else { return }
|
||
mono = biggerMono
|
||
staging = biggerStaging
|
||
}
|
||
guard let dst = mono.floatChannelData?[0] else { return }
|
||
mono.frameLength = buffer.frameLength
|
||
|
||
// Fold the multi-channel input down to the one mono bus we encode.
|
||
Self.foldToMono(
|
||
input: src, frames: frames, channels: Int(buffer.format.channelCount),
|
||
interleaved: buffer.format.isInterleaved, pinned: pinnedChannel, out: dst)
|
||
|
||
if !levelReported {
|
||
var localPeak: Float = 0
|
||
for i in 0..<frames where abs(dst[i]) > localPeak { localPeak = abs(dst[i]) }
|
||
if localPeak > inputPeak { inputPeak = localPeak }
|
||
framesInspected += frames
|
||
if framesInspected >= silenceWindow {
|
||
levelReported = true
|
||
if inputPeak == 0 {
|
||
log.warning("""
|
||
mic uplink has been pure digital SILENCE for 10 s (\(deviceLabel), \
|
||
\(channelPlan)) — check the input level (System Settings → Sound → \
|
||
Input), Privacy & Security → Microphone, and the Microphone channel in \
|
||
Settings; the host is receiving zeros
|
||
""")
|
||
} else {
|
||
let dbfs = 20 * log10(inputPeak)
|
||
log.info("""
|
||
mic uplink OK — peak \(String(format: "%.1f", dbfs)) dBFS over first \
|
||
10 s (\(deviceLabel), \(channelPlan))
|
||
""")
|
||
}
|
||
}
|
||
}
|
||
|
||
var fed = false
|
||
var convError: NSError?
|
||
let status = resampler.convert(to: staging, error: &convError) { _, outStatus in
|
||
if fed {
|
||
outStatus.pointee = .noDataNow
|
||
return nil
|
||
}
|
||
fed = true
|
||
outStatus.pointee = .haveData
|
||
return mono
|
||
}
|
||
guard status != .error, let p = staging.floatChannelData?[0] else { return }
|
||
fifo.append(contentsOf: UnsafeBufferPointer(
|
||
start: p, count: Int(staging.frameLength)))
|
||
|
||
// Consume whole chunks through a head index, then drop the eaten prefix in ONE
|
||
// move of the sub-chunk remainder. The old per-chunk removeFirst memmoved the
|
||
// entire backlog for every packet — O(n) on the render-adjacent tap thread.
|
||
let samplesPerChunk = Int(encoder.framesPerPacket)
|
||
var head = 0
|
||
while fifo.count - head >= samplesPerChunk {
|
||
chunk.frameLength = encoder.framesPerPacket
|
||
fifo.withUnsafeBufferPointer { src in
|
||
chunk.floatChannelData![0].update(
|
||
from: src.baseAddress! + head, count: samplesPerChunk)
|
||
}
|
||
head += samplesPerChunk
|
||
guard let packets = try? encoder.encode(chunk) else { continue }
|
||
for packet in packets {
|
||
connection.sendMic(
|
||
packet, seq: seq, ptsNs: DispatchTime.now().uptimeNanoseconds)
|
||
seq &+= 1
|
||
}
|
||
}
|
||
if head > 0 { fifo.removeFirst(head) } // keeps capacity — no realloc
|
||
}
|
||
return true
|
||
}
|
||
|
||
/// Fold `channels` of input (`floatChannelData` layout: `interleaved` → one buffer strided by
|
||
/// channel count; else one buffer per channel) down to a single mono bus in `out` (`frames`
|
||
/// long). `pinned` (0-based, must be `< channels`) copies exactly that channel — the fix for a
|
||
/// mic on one input of a multi-channel interface; `nil` sums every channel, clamped to
|
||
/// [-1, 1], so a lone hot channel still passes at full level instead of the silent 0/1 the
|
||
/// default N→stereo downmix would grab. Pure + `internal` for unit testing the index math.
|
||
static func foldToMono(
|
||
input: UnsafePointer<UnsafeMutablePointer<Float>>, frames: Int, channels: Int,
|
||
interleaved: Bool, pinned: Int?, out: UnsafeMutablePointer<Float>
|
||
) {
|
||
if let ch = pinned, ch < channels {
|
||
if interleaved {
|
||
let d = input[0]
|
||
for i in 0..<frames { out[i] = d[i * channels + ch] }
|
||
} else {
|
||
let d = input[ch]
|
||
for i in 0..<frames { out[i] = d[i] }
|
||
}
|
||
} else if interleaved {
|
||
let d = input[0]
|
||
for i in 0..<frames {
|
||
var s: Float = 0
|
||
for c in 0..<channels { s += d[i * channels + c] }
|
||
out[i] = max(-1, min(1, s))
|
||
}
|
||
} else {
|
||
let d0 = input[0]
|
||
for i in 0..<frames { out[i] = d0[i] }
|
||
for c in 1..<channels {
|
||
let d = input[c]
|
||
for i in 0..<frames { out[i] += d[i] }
|
||
}
|
||
if channels > 1 { for i in 0..<frames { out[i] = max(-1, min(1, out[i])) } }
|
||
}
|
||
}
|
||
#endif
|
||
|
||
#if os(macOS)
|
||
private static func setDevice(_ id: AudioDeviceID, on unit: AudioUnit) -> Bool {
|
||
var dev = id
|
||
return AudioUnitSetProperty(
|
||
unit, kAudioOutputUnitProperty_CurrentDevice, kAudioUnitScope_Global, 0,
|
||
&dev, UInt32(MemoryLayout<AudioDeviceID>.size)) == noErr
|
||
}
|
||
|
||
/// Read back the AUHAL's live device — the definitive "what are we actually capturing
|
||
/// from", which catches a selection that succeeded on paper but silently fell back to
|
||
/// the system default (a stale/changed UID, a device that vanished between resolve and
|
||
/// start). 0 / an error means we couldn't tell.
|
||
private static func currentDevice(of unit: AudioUnit) -> AudioDeviceID? {
|
||
var dev = AudioDeviceID(0)
|
||
var size = UInt32(MemoryLayout<AudioDeviceID>.size)
|
||
let status = AudioUnitGetProperty(
|
||
unit, kAudioOutputUnitProperty_CurrentDevice, kAudioUnitScope_Global, 0, &dev, &size)
|
||
guard status == noErr, dev != 0 else { return nil }
|
||
return dev
|
||
}
|
||
#endif
|
||
}
|