using ytLive.Helpers;
using ytLive.Services;
namespace ytLive.Services.Audio;
///
/// Owns the mic + desktop/game capture sources (TASK 4 ship step 4) and drives
/// the footer meters; since TASK 9 it is also the live audio path into the
/// encoder. Capture runs for the app's lifetime so both bars stay live in
/// preview: mic samples are downmixed to mono, resampled to 48 kHz, put
/// through the voice chain (bass → treble → noise gate → compressor), THEN
/// level-metered and queued; loopback samples are resampled to 48 kHz stereo
/// and queued. While live (), a 10 ms loop drains both
/// queues, applies the honest gains (mic = MicVolume, loopback = unity
/// × auto-duck when the mic is hot — the slider drives system volume via mixes them to stereo float, runs the
/// master limiter (a -1 dBFS ceiling so the sum never clips the encoder), and
/// writes the chunk to the encoder's audio pipe.
///
public sealed class AudioMixer : IDisposable
{
public const int OutputSampleRate = 48000;
private static readonly TimeSpan DefaultMixInterval = TimeSpan.FromMilliseconds(10);
private readonly IAudioSource _mic;
private readonly IAudioSource _loopback;
private readonly AudioLevelMeter _meter;
private readonly AudioLevelMeter _loopbackMeter;
private readonly Action? _log;
private readonly Func? _micGain;
private readonly Func? _loopbackGain;
private readonly TimeSpan _mixInterval;
private readonly Func? _syncOffsetMs;
private readonly AudioSyncDelay _syncDelay;
private bool _started;
// Live path (TASK 9): per-source resampling → voice chain on the mic →
// thread-safe queues → the ducker + gain/mix loop → the encoder's pipe.
private readonly AudioRingBuffer _micBuffer;
private readonly AudioRingBuffer _loopbackBuffer;
private readonly VoiceFilterChain _voiceChain;
private readonly AutoDucker _ducker;
private readonly MasterLimiter _masterLimiter = new();
private TinyResampler? _micResampler;
private TinyResampler? _loopbackResampler;
private CancellationTokenSource? _liveCts;
private NamedPipeAudioWriter _pipe = new();
private float[]? _micChunk;
private float[]? _loopbackChunk;
private float[]? _mixBuffer;
private float[]? _delayedMix;
private DateTime _lastLoopErrorLogged = DateTime.MinValue;
private DateTime _nextStatsLog = DateTime.UtcNow + TimeSpan.FromSeconds(5);
private long _micDrainedTotal;
private long _loopDrainedTotal;
private float _peakMix;
public AudioMixer(
IAudioSource mic,
IAudioSource loopback,
Action? log = null,
Func? micGain = null,
Func? loopbackGain = null,
Func? syncOffsetMs = null,
TimeSpan? mixInterval = null)
{
_mic = mic;
_loopback = loopback;
_meter = new AudioLevelMeter();
_loopbackMeter = new AudioLevelMeter();
_log = log;
_micGain = micGain;
_loopbackGain = loopbackGain;
_syncOffsetMs = syncOffsetMs;
_syncDelay = new AudioSyncDelay(OutputSampleRate);
_mixInterval = mixInterval ?? DefaultMixInterval;
var seconds = _mixInterval.TotalSeconds;
var framesPerTick = Math.Max(1, (int)Math.Round(OutputSampleRate * seconds));
_micBuffer = new AudioRingBuffer(OutputSampleRate * 2);
_loopbackBuffer = new AudioRingBuffer(OutputSampleRate * 2 * 2);
_voiceChain = new VoiceFilterChain(OutputSampleRate);
_ducker = new AutoDucker();
_micChunk = new float[framesPerTick];
_loopbackChunk = new float[framesPerTick * 2];
_mixBuffer = new float[framesPerTick * 2];
_delayedMix = new float[framesPerTick * 2];
_mic.Started += OnMicStarted;
_mic.SampleReady += OnMicSample;
_loopback.SampleReady += OnLoopbackSample;
_mic.Failed += OnMicFailed;
_loopback.Failed += OnLoopbackFailed;
}
/// Current smoothed mic level (0..1), post voice chain.
public float MicLevel => _meter.Level;
/// Current smoothed desktop/game level (0..1).
public float LoopbackLevel => _loopbackMeter.Level;
/// Raised whenever the smoothed mic level changes.
public event Action? MicLevelChanged;
/// Raised whenever the smoothed desktop/game level changes.
public event Action? LoopbackLevelChanged;
/// Raised when the mic capture comes up (the status dot goes green).
public event Action? MicConnected;
/// Raised when the mic capture fails or dies (the status dot goes yellow).
public event Action? MicFailed;
public void Start()
{
if (_started)
return;
_started = true;
_meter.Reset();
_loopback.Start();
_mic.Start();
}
/// Swaps the mic source without touching loopback — used when the
/// creator picks a different device mid-session. The level resets and the
/// new source raises or .
public void RestartMic()
{
_mic.Stop();
_meter.Reset();
_micBuffer.Clear();
_micResampler?.Reset();
_voiceChain.Reset();
MicLevelChanged?.Invoke(0);
_mic.Start();
}
public void Stop()
{
if (!_started)
return;
_started = false;
StopLive();
_mic.Stop();
_loopback.Stop();
_meter.Reset();
_loopbackMeter.Reset();
MicLevelChanged?.Invoke(0);
LoopbackLevelChanged?.Invoke(0);
}
/// Begins the live mix loop: drains the capture queues, applies the
/// honest gains + auto-duck, and streams stereo float into the named audio
/// pipe. Idempotent — safe to call once per go-live.
public void StartLive(string pipeName)
{
if (_liveCts != null)
return;
var cts = new CancellationTokenSource();
_liveCts = cts;
_pipe = new NamedPipeAudioWriter();
_pipe.Start(pipeName);
_log?.Invoke($"Audio live: started, pipe '{pipeName}'");
_ = Task.Run(() => LiveLoopAsync(_pipe, cts.Token));
}
/// Ends the live mix loop and closes the audio pipe — ffmpeg sees
/// EOF on the audio input. Safe when not live: the writer's Stop() is idempotent
/// (StopStream now calls this on every rollback path).
public void StopLive()
{
var cts = _liveCts;
_liveCts = null;
if (cts != null)
{
_log?.Invoke("Audio live: stopped");
cts.Cancel();
}
_pipe.Stop();
}
public void Dispose()
{
Stop();
_mic.Started -= OnMicStarted;
_mic.SampleReady -= OnMicSample;
_loopback.SampleReady -= OnLoopbackSample;
_mic.Failed -= OnMicFailed;
_loopback.Failed -= OnLoopbackFailed;
_mic.Dispose();
_loopback.Dispose();
}
private void OnMicStarted()
{
MicConnected?.Invoke();
}
private void OnMicSample(AudioSample sample)
{
var mono = DownmixToMono(sample);
mono = ResampleMic(mono, sample.SampleRate);
for (var i = 0; i < mono.Length; i++)
mono[i] = _voiceChain.Process(mono[i]);
_micBuffer.Write(mono);
// Push unconditionally: the ?. on the event would otherwise skip the
// argument (and the meter update) when nothing is subscribed yet. The
// level is the POST-filtered signal — the meter shows what the stream
// will carry.
var level = _meter.Push(new AudioSample(mono, OutputSampleRate, 1));
MicLevelChanged?.Invoke(level);
}
private void OnLoopbackSample(AudioSample sample)
{
var data = ResampleLoopback(sample.Samples, sample.SampleRate);
_loopbackBuffer.Write(data);
var level = _loopbackMeter.Push(new AudioSample(data, OutputSampleRate, sample.Channels));
LoopbackLevelChanged?.Invoke(level);
}
private void OnMicFailed(Exception ex)
{
_log?.Invoke($"Mic capture failed: {ex.Message}");
MicLevelChanged?.Invoke(0);
MicFailed?.Invoke(ex);
}
private void OnLoopbackFailed(Exception ex)
{
_log?.Invoke($"Desktop audio capture failed: {ex.Message}");
}
private static float[] DownmixToMono(AudioSample sample)
{
if (sample.Channels <= 1)
return sample.Samples;
var frames = sample.Samples.Length / sample.Channels;
var mono = new float[frames];
for (var i = 0; i < frames; i++)
{
var sum = 0f;
for (var c = 0; c < sample.Channels; c++)
sum += sample.Samples[i * sample.Channels + c];
mono[i] = sum / sample.Channels;
}
return mono;
}
private float[] ResampleMic(float[] mono, int inputRate)
{
if (inputRate == OutputSampleRate)
return mono;
_micResampler ??= new TinyResampler(inputRate, OutputSampleRate);
if (_micResampler.NeedsResampling)
return _micResampler.Process(mono);
return mono;
}
private float[] ResampleLoopback(float[] interleaved, int inputRate)
{
if (inputRate == OutputSampleRate)
return interleaved;
_loopbackResampler ??= new TinyResampler(inputRate, OutputSampleRate);
if (_loopbackResampler.NeedsResampling)
return _loopbackResampler.Process(interleaved);
return interleaved;
}
private async Task LiveLoopAsync(IAudioPipeWriter pipe, CancellationToken cancellationToken)
{
var mixBuffer = _mixBuffer!;
var micChunk = _micChunk!;
var loopbackChunk = _loopbackChunk!;
while (!cancellationToken.IsCancellationRequested)
{
var nextTick = DateTime.UtcNow + _mixInterval;
try
{
var (_, micDrained, loopDrained) = FillAndMix(micChunk, loopbackChunk, mixBuffer);
await pipe.WriteAsync(mixBuffer, cancellationToken).ConfigureAwait(false);
_micDrainedTotal += micDrained;
_loopDrainedTotal += loopDrained;
for (var i = 0; i < mixBuffer.Length; i++)
{
var abs = Math.Abs(mixBuffer[i]);
if (abs > _peakMix) _peakMix = abs;
}
var now = DateTime.UtcNow;
if (now >= _nextStatsLog && _log != null)
{
_nextStatsLog = now + TimeSpan.FromSeconds(5);
_log($"Audio live: pipe connected={pipe.IsConnected} dropped writes={pipe.DroppedWrites} " +
$"micLevel={_meter.Level:F3} loopLevel={_loopbackMeter.Level:F3} " +
$"drained {_micDrainedTotal:N0}/{_loopDrainedTotal:N0} samples peakMix {_peakMix:F3} (per 5s)");
_micDrainedTotal = 0;
_loopDrainedTotal = 0;
_peakMix = 0;
}
}
catch (OperationCanceledException)
{
break;
}
catch (Exception ex)
{
// Throttled + never on a dead session: the 2026-09-01 first-launch flood
// was one logged NRE every 10ms tick (see _delayedMix fix). Log at most
// once per 5s, and stop looping if the pipe writer is gone for good.
if (cancellationToken.IsCancellationRequested)
break;
var now = DateTime.UtcNow;
if (now - _lastLoopErrorLogged >= TimeSpan.FromSeconds(5))
{
_lastLoopErrorLogged = now;
AppLog.Write(ex, "Audio live loop error");
_log?.Invoke($"Audio live loop error: {ex.Message}");
}
}
var delay = nextTick - DateTime.UtcNow;
if (delay > TimeSpan.Zero)
{
try
{
await Task.Delay(delay, cancellationToken).ConfigureAwait(false);
}
catch (OperationCanceledException)
{
break;
}
}
}
}
/// Drains one tick's worth of mic + loopback, silence-fills any
/// underrun, applies the ducker and the honest gains, and mixes to stereo.
/// Returns the pre-duck mic RMS and how many samples each buffer contributed
/// (the 5s live-loop telemetry reads those to name whether capture or the
/// pipe starved a silent recording).
private (float MicRms, int MicDrained, int LoopDrained) FillAndMix(float[] micChunk, float[] loopbackChunk, float[] mix)
{
var micCount = _micBuffer.Read(micChunk, micChunk.Length);
var loopCount = _loopbackBuffer.Read(loopbackChunk, loopbackChunk.Length);
float micRms = 0;
if (micCount == 0)
{
Array.Clear(micChunk, 0, micChunk.Length);
}
else
{
double sumSquares = 0;
for (var i = 0; i < micCount; i++)
sumSquares += micChunk[i] * micChunk[i];
micRms = (float)Math.Sqrt(sumSquares / micCount);
if (micCount < micChunk.Length)
Array.Clear(micChunk, micCount, micChunk.Length - micCount);
}
if (loopCount == 0)
{
Array.Clear(loopbackChunk, 0, loopbackChunk.Length);
}
else if (loopCount < loopbackChunk.Length)
{
Array.Clear(loopbackChunk, loopCount, loopbackChunk.Length - loopCount);
}
var duck = _ducker.Update(micRms);
var micGain = (float)(_micGain?.Invoke() ?? 1.0);
var loopGain = (float)(_loopbackGain?.Invoke() ?? 1.0) * duck;
var frames = micChunk.Length;
for (var i = 0; i < frames; i++)
{
var m = micChunk[i] * micGain;
mix[i * 2] = m + loopbackChunk[i * 2] * loopGain;
mix[i * 2 + 1] = m + loopbackChunk[i * 2 + 1] * loopGain;
}
_masterLimiter.Process(mix);
_syncDelay.Configure(_syncOffsetMs?.Invoke() ?? 0);
var delayed = _delayedMix;
if (delayed is null || delayed.Length < mix.Length)
delayed = _delayedMix = new float[mix.Length];
_syncDelay.Process(mix, delayed);
Array.Copy(delayed, mix, mix.Length);
return (micRms, micCount, loopCount);
}
}