using System.IO.MemoryMappedFiles; using VoiceCat.Core; using Voicecat.V1; namespace VoiceCat.iOS; internal sealed class BroadcastAudioPump : IAsyncDisposable { private const uint Magic = 0x56434252, Version = 1; private const int Header = 64, Capacity = 96_000, Frame = 960; private readonly CancellationTokenSource stop = new(); private readonly short[] scratch = new short[Frame * 2]; private Thread? worker; private VoiceCatClient? client; private MemoryMappedFile? map; private MemoryMappedViewAccessor? view; private string? path; private uint streamId; private bool active; private int generation; private ulong observedWrite; private ulong idleWrite; private bool idleWriteSeen; internal event Action? Changed; internal bool IsActive => active; #if DEBUG internal bool IsMappingOpen => view is not null; internal long EncodedPackets => streamId == 0 ? 0 : client?.Audio.GetLocalDiagnostics(streamId).EncodedPackets ?? 0; #endif internal void Start(VoiceCatClient owner) { client = owner; worker = new Thread(Run) { IsBackground = true, Name = "VoiceCat iOS screen audio", Priority = ThreadPriority.AboveNormal }; worker.Start(); } private void Run() { CancellationToken token = stop.Token; try { while (!token.IsCancellationRequested) { try { DrainAsync(token).GetAwaiter().GetResult(); } catch (OperationCanceledException) when (token.IsCancellationRequested) { break; } catch (Exception exception) { // A rejected screen stream is a session error, not an unhandled exception on // this background Thread. Keep the app and microphone call alive. Console.Error.WriteLine($"VoiceCat screen audio pump failed: {exception}"); CloseMapping(); StopStream(); SetActive(false); if (!token.IsCancellationRequested) Thread.Sleep(1000); } // The shared file must not remain mapped while the app is idle or suspending. if (!token.IsCancellationRequested) Thread.Sleep(view is null ? 100 : 5); } } finally { CloseMapping(); } } private async Task DrainAsync(CancellationToken token) { if (!OpenMapping()) return; MemoryMappedViewAccessor pages = view!; bool active = pages.ReadUInt32(16) != 0; if (!active) { StopStream(); SetActive(false); CloseMapping(); return; } VoiceCatClient owner = client ?? throw new IOException("Client disconnected."); if (streamId == 0) { // A ring left marked active by a killed extension must not start a phantom stream. // Wait for a fresh producer write after this connection opened the mapping. ulong currentWrite = pages.ReadUInt64(24); if (currentWrite == observedWrite) return; observedWrite = currentWrite; int startGeneration = Volatile.Read(ref generation); VoiceSubscriptionResult subscription = await owner.SubscribeVoiceAsync(cancellationToken: token).ConfigureAwait(false); if (!subscription.Ok) throw new InvalidOperationException(subscription.Error); StreamInfo stream = await owner.StartStreamAsync(StreamKind.StreamScreenAudio, "Screen audio", 2, token).ConfigureAwait(false); if (startGeneration != Volatile.Read(ref generation) || pages.ReadUInt32(16) == 0 || !ReferenceEquals(client, owner)) { try { owner.StopStream(stream.StreamId); } catch (Exception exception) when (exception is IOException or InvalidOperationException) { } return; } // Keep the newest 120 ms captured during stream negotiation. This covers the // largest codec frame plus the input buffer target; dropping through the latest // write loses short sounds before the first drain can feed them. ulong latestWrite = pages.ReadUInt64(24); ulong earliest = latestWrite > (ulong)(Frame * 12) ? latestWrite - (ulong)(Frame * 12) : 0; pages.Write(32, Math.Max(pages.ReadUInt64(32), earliest)); streamId = stream.StreamId; SetActive(true); } ulong write = pages.ReadUInt64(24), read = pages.ReadUInt64(32); if (write - read > Capacity) read = write - Capacity; while (write - read >= Frame * 2) { for (int sample = 0; sample < scratch.Length; sample++) { ulong index = (read + (ulong)sample) % Capacity; scratch[sample] = pages.ReadInt16(Header + checked((long)index * sizeof(short))); } if (!owner.Audio.FeedPcm(streamId, scratch, 2)) break; read += (ulong)scratch.Length; pages.Write(32, read); } } // Keep a mapping only while a broadcast is active. An inactive mapping holds an open handle // into the shared App Group container through suspension, which iOS may terminate as a shared // file lock (0xdead10cc). Poll the small header with a short-lived read when idle, then map // once for the active broadcast's high-rate drain. private bool OpenMapping() { path ??= ResolvePath(); if (path is null) return false; bool present = File.Exists(path); if (view is not null) { if (present) return true; CloseMapping(); return false; } if (!present) { idleWriteSeen = false; return false; } using (var file = new FileStream(path, FileMode.Open, FileAccess.Read, FileShare.ReadWrite | FileShare.Delete)) { if (file.Length < Header) { idleWriteSeen = false; return false; } Span header = stackalloc byte[32]; file.ReadExactly(header); if (System.Buffers.Binary.BinaryPrimitives.ReadUInt32LittleEndian(header) != Magic || System.Buffers.Binary.BinaryPrimitives.ReadUInt32LittleEndian(header[4..]) != Version || System.Buffers.Binary.BinaryPrimitives.ReadUInt32LittleEndian(header[16..]) == 0) { idleWriteSeen = false; return false; } ulong currentWrite = System.Buffers.Binary.BinaryPrimitives.ReadUInt64LittleEndian(header[24..]); if (!idleWriteSeen) { idleWrite = currentWrite; idleWriteSeen = true; return false; } if (currentWrite == idleWrite) return false; } MemoryMappedFile candidate = MemoryMappedFile.CreateFromFile(path, FileMode.Open, null, Header + Capacity * sizeof(short), MemoryMappedFileAccess.ReadWrite); try { MemoryMappedViewAccessor pages = candidate.CreateViewAccessor(0, Header + Capacity * sizeof(short), MemoryMappedFileAccess.ReadWrite); if (pages.ReadUInt32(0) != Magic || pages.ReadUInt32(4) != Version) { pages.Dispose(); throw new InvalidDataException("Unsupported broadcast ring."); } map = candidate; view = pages; observedWrite = idleWrite; return true; } catch { candidate.Dispose(); throw; } } private static string? ResolvePath() { NSUrl? root = NSFileManager.DefaultManager.GetContainerUrl(IosConstants.AppGroup); return root?.Path is null ? null : Path.Combine(root.Path, "voicecat", "broadcast_audio.ring"); } private void CloseMapping() { view?.Dispose(); view = null; map?.Dispose(); map = null; observedWrite = 0; idleWriteSeen = false; } private void StopStream() { uint id = streamId; streamId = 0; if (id == 0 || client?.State != ClientConnectionState.Connected) return; try { client.StopStream(id); } catch (Exception exception) when (exception is IOException or InvalidOperationException or ObjectDisposedException) { } } internal void RequestStop() { Interlocked.Increment(ref generation); SetActive(false); } private void SetActive(bool value) { if (active == value) return; active = value; Changed?.Invoke(); } public ValueTask DisposeAsync() { stop.Cancel(); Interlocked.Increment(ref generation); worker?.Join(); worker = null; StopStream(); SetActive(false); stop.Dispose(); return ValueTask.CompletedTask; } }