Fix iOS voice capture and screen sharing
This commit is contained in:
@@ -66,8 +66,25 @@ public sealed class AdaptivePcmBuffer
|
||||
}
|
||||
if (available <= 0) { primed = false; phase = 0; destination.Clear(); return 0; }
|
||||
|
||||
double correction = Math.Clamp((available - target) / (SampleRate * 2.0), -MaximumCorrection, MaximumCorrection);
|
||||
// Reads observe a packetized producer. Occupancy naturally moves by roughly a callback
|
||||
// or codec block even when both clocks are exactly 48 kHz; interpreting that sawtooth as
|
||||
// clock drift biases the resampler fast until it periodically drains the ring. Only
|
||||
// correct sustained movement outside two codec blocks of quantization headroom.
|
||||
double error = available - target, positiveDeadband = requestedFrames * 2.0;
|
||||
// Low occupancy is real underflow pressure and must be corrected immediately. The
|
||||
// packetization sawtooth is above the target, so only the positive side needs a deadband.
|
||||
double drift = error < 0 ? error : Math.Max(0, error - positiveDeadband);
|
||||
double correction = Math.Clamp(drift / (SampleRate * 2.0), -MaximumCorrection, MaximumCorrection);
|
||||
double step = 1.0 + correction;
|
||||
// A caller encoding fixed-size frames cannot use a partial read. Limit the correction
|
||||
// to what the current occupancy can supply for the whole destination; if even the
|
||||
// maximum slow-down cannot do that, preserve every queued sample and re-prime.
|
||||
double maximumWholeReadStep = (available - phase) / requestedFrames;
|
||||
if (maximumWholeReadStep < 1.0 - MaximumCorrection)
|
||||
{
|
||||
primed = false; phase = 0; destination.Clear(); return 0;
|
||||
}
|
||||
step = Math.Min(step, maximumWholeReadStep);
|
||||
int produced = 0, read = readFrame;
|
||||
for (int frame = 0; frame < requestedFrames; frame++)
|
||||
{
|
||||
|
||||
@@ -4,6 +4,9 @@ using Voicecat.V1;
|
||||
|
||||
namespace VoiceCat.Audio;
|
||||
|
||||
public readonly record struct LocalAudioDiagnostics(long Cycles, long StarvedCycles, long EncodedPackets,
|
||||
long RejectedPackets, int BufferedFrames);
|
||||
|
||||
public sealed class AudioEngine : IDisposable
|
||||
{
|
||||
private readonly object gate = new();
|
||||
@@ -42,7 +45,7 @@ public sealed class AudioEngine : IDisposable
|
||||
maintenance = MaintainAsync();
|
||||
if (startWorker)
|
||||
{
|
||||
worker = new Thread(Work) { IsBackground = true, Name = "VoiceCat managed audio" };
|
||||
worker = new Thread(Work) { IsBackground = true, Name = "VoiceCat managed audio", Priority = ThreadPriority.Highest };
|
||||
worker.Start();
|
||||
}
|
||||
}
|
||||
@@ -103,6 +106,12 @@ public sealed class AudioEngine : IDisposable
|
||||
foreach (LocalStream stream in Volatile.Read(ref routes).Local) if (stream.Info.StreamId == streamId) return (stream.Level, stream.Talking);
|
||||
return default;
|
||||
}
|
||||
public LocalAudioDiagnostics GetLocalDiagnostics(uint streamId)
|
||||
{
|
||||
foreach (LocalStream stream in Volatile.Read(ref routes).Local)
|
||||
if (stream.Info.StreamId == streamId) return stream.Diagnostics;
|
||||
return default;
|
||||
}
|
||||
public void SetLocalGain(uint streamId, float gain)
|
||||
{
|
||||
if (!float.IsFinite(gain) || gain < 0 || gain > 4) throw new ArgumentOutOfRangeException(nameof(gain));
|
||||
@@ -153,13 +162,28 @@ public sealed class AudioEngine : IDisposable
|
||||
{
|
||||
ProcessCycle();
|
||||
deadline += Stopwatch.Frequency / 50;
|
||||
double remaining = (deadline - Stopwatch.GetTimestamp()) * 1000.0 / Stopwatch.Frequency;
|
||||
if (remaining > 0) Thread.Sleep((int)Math.Ceiling(remaining));
|
||||
else if (remaining < -100) deadline = Stopwatch.GetTimestamp();
|
||||
WaitUntil(deadline);
|
||||
if ((deadline - Stopwatch.GetTimestamp()) * 1000.0 / Stopwatch.Frequency < -100)
|
||||
deadline = Stopwatch.GetTimestamp();
|
||||
}
|
||||
}
|
||||
catch (Exception exception) { Failure = exception; stop.Cancel(); }
|
||||
}
|
||||
|
||||
// Millisecond-rounded sleeps periodically overshoot Core Audio's hardware clock enough to
|
||||
// empty its small handoff buffer. Sleep for the coarse portion, then hold the audio worker
|
||||
// to the Stopwatch deadline for the final sub-millisecond interval.
|
||||
private void WaitUntil(long deadline)
|
||||
{
|
||||
while (!stop.IsCancellationRequested)
|
||||
{
|
||||
long remaining = deadline - Stopwatch.GetTimestamp();
|
||||
if (remaining <= 0) return;
|
||||
double milliseconds = remaining * 1000.0 / Stopwatch.Frequency;
|
||||
if (milliseconds > 2) Thread.Sleep(Math.Max(1, (int)milliseconds - 1));
|
||||
else Thread.SpinWait(64);
|
||||
}
|
||||
}
|
||||
private async Task MaintainAsync()
|
||||
{
|
||||
try
|
||||
|
||||
@@ -28,8 +28,11 @@ internal sealed class LocalStream : IDisposable
|
||||
private int feeding;
|
||||
private int buffered;
|
||||
private int starvedSamples;
|
||||
private long cycles, starvedCycles, encodedPackets, rejectedPackets;
|
||||
private uint timestamp;
|
||||
private bool wasTransmitting, marker;
|
||||
internal LocalAudioDiagnostics Diagnostics => new(Volatile.Read(ref cycles), Volatile.Read(ref starvedCycles),
|
||||
Volatile.Read(ref encodedPackets), Volatile.Read(ref rejectedPackets), Input.CountFrames);
|
||||
|
||||
internal bool Feed(ReadOnlySpan<short> pcm, int channels)
|
||||
{
|
||||
@@ -67,9 +70,11 @@ internal sealed class LocalStream : IDisposable
|
||||
|
||||
internal void Process(AudioEngine engine, EncodedVoiceSender sender)
|
||||
{
|
||||
Interlocked.Increment(ref cycles);
|
||||
var input = capture.AsSpan(0, 960 * CaptureChannels);
|
||||
if (Input.Read(input) != input.Length)
|
||||
{
|
||||
Interlocked.Increment(ref starvedCycles);
|
||||
Level = 0; Talking = false; starvedSamples += 960;
|
||||
if (starvedSamples >= 9600) { buffered = 0; wasTransmitting = false; }
|
||||
return;
|
||||
@@ -114,7 +119,9 @@ internal sealed class LocalStream : IDisposable
|
||||
{
|
||||
int length = encoder.Encode(wire.AsSpan(0, frame), packet);
|
||||
VoiceFrameFlags flags = (Info.Audio.Fec ? VoiceFrameFlags.FecPresent : VoiceFrameFlags.None) | (marker ? VoiceFrameFlags.Marker : VoiceFrameFlags.None);
|
||||
sender(Info.Ssrc, timestamp, packet.AsSpan(0, length), flags); marker = false;
|
||||
if (sender(Info.Ssrc, timestamp, packet.AsSpan(0, length), flags)) Interlocked.Increment(ref encodedPackets);
|
||||
else Interlocked.Increment(ref rejectedPackets);
|
||||
marker = false;
|
||||
timestamp = unchecked(timestamp + (uint)encoder.Options.SamplesPerChannel);
|
||||
buffered -= frame;
|
||||
wire.AsSpan(frame, buffered).CopyTo(wire);
|
||||
|
||||
@@ -16,7 +16,7 @@ internal sealed class ClientMediaTransport : IAsyncDisposable
|
||||
private readonly byte[] binding = new byte[VoiceFrameHeader.Size + 16];
|
||||
private readonly byte[] keepalive = new byte[VoiceFrameHeader.Size];
|
||||
private readonly PacketQueue packets = new();
|
||||
private readonly Task sending;
|
||||
private readonly Thread sending;
|
||||
private readonly Task receiving;
|
||||
private readonly TaskCompletionSource bound = new(TaskCreationOptions.RunContinuationsAsynchronously);
|
||||
internal event EncodedVoiceHandler? Received;
|
||||
@@ -34,12 +34,13 @@ internal sealed class ClientMediaTransport : IAsyncDisposable
|
||||
token.CopyTo(binding.AsSpan(VoiceFrameHeader.Size));
|
||||
new VoiceFrameHeader(MediaFrameType.Keepalive, 0, 0, 0, 0, 0).Write(keepalive);
|
||||
receiving = ReceiveAsync();
|
||||
sending = SendAsync();
|
||||
sending = new Thread(Send) { IsBackground = true, Name = "VoiceCat UDP sender", Priority = ThreadPriority.AboveNormal };
|
||||
sending.Start();
|
||||
}
|
||||
|
||||
internal bool TrySend(VoiceFrameHeader header, ReadOnlySpan<byte> payload) => packets.TryWrite(header, payload);
|
||||
|
||||
private async Task SendAsync()
|
||||
private void Send()
|
||||
{
|
||||
byte[] plain = new byte[1275], packet = new byte[1275 + VoiceFrameHeader.Size + MediaEncryptor.TagSize];
|
||||
long nextKeepalive = 0;
|
||||
@@ -49,19 +50,19 @@ internal sealed class ClientMediaTransport : IAsyncDisposable
|
||||
{
|
||||
if (Environment.TickCount64 >= nextKeepalive)
|
||||
{
|
||||
if (!bound.Task.IsCompleted) await socket.SendAsync(binding, SocketFlags.None, stop.Token).ConfigureAwait(false);
|
||||
await socket.SendAsync(keepalive, SocketFlags.None, stop.Token).ConfigureAwait(false);
|
||||
if (!bound.Task.IsCompleted) socket.Send(binding, SocketFlags.None);
|
||||
socket.Send(keepalive, SocketFlags.None);
|
||||
nextKeepalive = Environment.TickCount64 + (bound.Task.IsCompleted ? 5000 : 250);
|
||||
}
|
||||
while (packets.TryRead(plain, out VoiceFrameHeader header, out int length))
|
||||
{
|
||||
int size = crypto.Encryptor.Encrypt(header, plain.AsSpan(0, length), packet);
|
||||
await socket.SendAsync(packet.AsMemory(0, size), SocketFlags.None, stop.Token).ConfigureAwait(false);
|
||||
socket.Send(packet.AsSpan(0, size), SocketFlags.None);
|
||||
}
|
||||
await Task.Delay(5, stop.Token).ConfigureAwait(false);
|
||||
Thread.Sleep(1);
|
||||
}
|
||||
}
|
||||
catch (Exception exception) when (exception is OperationCanceledException or SocketException or ObjectDisposedException)
|
||||
catch (Exception exception) when (exception is SocketException or ObjectDisposedException)
|
||||
{ if (!stop.IsCancellationRequested) bound.TrySetException(exception); }
|
||||
finally { stop.Cancel(); }
|
||||
}
|
||||
@@ -91,7 +92,7 @@ internal sealed class ClientMediaTransport : IAsyncDisposable
|
||||
public async ValueTask DisposeAsync()
|
||||
{
|
||||
stop.Cancel(); socket.Dispose();
|
||||
try { await Task.WhenAll(sending, receiving).ConfigureAwait(false); }
|
||||
try { sending.Join(); await receiving.ConfigureAwait(false); }
|
||||
finally { System.Security.Cryptography.CryptographicOperations.ZeroMemory(binding); stop.Dispose(); }
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user