diff --git a/AudioCore/Impl/BlockingRingBuffer.cs b/AudioCore/Impl/BlockingRingBuffer.cs index 039d647..20ca2ae 100644 --- a/AudioCore/Impl/BlockingRingBuffer.cs +++ b/AudioCore/Impl/BlockingRingBuffer.cs @@ -22,7 +22,11 @@ public class BlockingRingBuffer while (written < srcLen) { - ct.ThrowIfCancellationRequested(); + if ( ct.IsCancellationRequested ) + { + Debug.WriteLine("BlockingRingBuffer: No room in the buffer to write. Timeout."); + return; + } int remaining = srcLen - written; @@ -33,12 +37,11 @@ public class BlockingRingBuffer ? _ringWrite - _ringRead : _ring.Length - _ringRead + _ringWrite; - free = _ring.Length - used - 1; // leave 1 byte to distinguish full/empty + free = _ring.Length - used - 1; } if (free <= 0) { - // No room → block until space becomes available Thread.Sleep(1); continue; } @@ -49,7 +52,6 @@ public class BlockingRingBuffer { int first = Math.Min(toWrite, _ring.Length - _ringWrite); - // Write first segment src.Slice(written, first) .CopyTo(new Span(_ring, _ringWrite, first)); @@ -58,7 +60,6 @@ public class BlockingRingBuffer int leftover = toWrite - first; if (leftover > 0) { - // Wrap-around segment src.Slice(written + first, leftover) .CopyTo(new Span(_ring, _ringWrite, leftover)); @@ -72,13 +73,19 @@ public class BlockingRingBuffer public int WaitForOutput(CancellationToken token) { - while (!token.IsCancellationRequested) + while (true) { + if (token.IsCancellationRequested) + { + Debug.WriteLine("BlockingRingBuffer: No data in the buffer to read. Timeout."); + return 0; + } + lock (_ringLock) { var available = (_ringWrite >= _ringRead) - ? _ringWrite - _ringRead - : _ring.Length - _ringRead + _ringWrite; + ? _ringWrite - _ringRead + : _ring.Length - _ringRead + _ringWrite; if (available > 0) return available; @@ -86,9 +93,6 @@ public class BlockingRingBuffer Thread.Sleep(2); } - - Debug.WriteLine("BlockingRingBuffer: Timeout waiting for output"); - return 0; } public int DrainRing(Span dest, int maxBytes) diff --git a/AudioCore/Impl/RubberBandTimeStretchEngine.cs b/AudioCore/Impl/RubberBandTimeStretchEngine.cs index dc54e63..4b06599 100644 --- a/AudioCore/Impl/RubberBandTimeStretchEngine.cs +++ b/AudioCore/Impl/RubberBandTimeStretchEngine.cs @@ -27,7 +27,7 @@ public sealed class RubberBandTimeStretchEngine : ITimeStretchEngine, IDisposabl _channels = channels; var bytesPerSecond = sampleRate * channels * sizeof(float); - _ring = new BlockingRingBuffer(1 * bytesPerSecond); + _ring = new BlockingRingBuffer(10 * bytesPerSecond); } public void Configure(PlaybackSpeedSettings settings) @@ -37,25 +37,17 @@ public sealed class RubberBandTimeStretchEngine : ITimeStretchEngine, IDisposabl _speed = settings.Speed; - if (Math.Abs(_speed - 1.0f) < 0.01f) - { - DisposeProcess(); - _ring.ResetRing(); - } - else - { - RestartProcess(); - } + DisposeProcess(); + _ring.ResetRing(); } - public Task Submit(MixedAudioBlock input) + public Task Submit(MixedAudioBlock input, CancellationToken token) { // No-stretch path: enqueue block and signal semaphore if (Math.Abs(_speed - 1.0f) < 0.01f) { - using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); var b = MemoryMarshal.AsBytes(input.Buffer.Span); - _ring.WriteToOutput(b, b.Length, cts.Token); + _ring.WriteToOutput(b, b.Length, token); return Task.CompletedTask; } @@ -72,21 +64,19 @@ public sealed class RubberBandTimeStretchEngine : ITimeStretchEngine, IDisposabl return Task.CompletedTask; } - public async Task Receive() + public async Task Receive(CancellationToken token) { - var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); - int available = 0; - while (!cts.IsCancellationRequested) + while (!token.IsCancellationRequested) { - available = _ring.WaitForOutput(cts.Token); + available = _ring.WaitForOutput(token); if (available > 0) break; await Task.Delay(2).ConfigureAwait(false); } - if (cts.IsCancellationRequested) + if (token.IsCancellationRequested) return default; var maxFloats = available / sizeof(float); @@ -132,13 +122,6 @@ public sealed class RubberBandTimeStretchEngine : ITimeStretchEngine, IDisposabl _readerThread.Start(); } - private void RestartProcess() - { - DisposeProcess(); - _ring.ResetRing(); - StartProcess(); - } - private void ReaderLoop() { var buf = new byte[4096]; diff --git a/AudioCore/Impl/StemPlaybackEngine.cs b/AudioCore/Impl/StemPlaybackEngine.cs index a851d79..86b88e8 100644 --- a/AudioCore/Impl/StemPlaybackEngine.cs +++ b/AudioCore/Impl/StemPlaybackEngine.cs @@ -1,4 +1,6 @@ -namespace AudioCore.Impl; +using System.Diagnostics; + +namespace AudioCore.Impl; public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable { @@ -238,7 +240,7 @@ public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable var decodeTask = DecodeLoopAsync(pipeline, ct); var stretchTask = StretchLoopAsync(pipeline, ct); - await Task.WhenAny(decodeTask, stretchTask); + await Task.WhenAny(decodeTask, stretchTask).ConfigureAwait(false); // When either loop ends, stop output if (pipeline.OutputStarted) @@ -250,6 +252,8 @@ public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable private async Task DecodeLoopAsync(PipelineState pipeline, CancellationToken ct) { + await Task.Yield(); + try { while (!ct.IsCancellationRequested) @@ -263,12 +267,12 @@ public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable lock (_stateLock) { - playing = _isPlaying; - mixerSnapshot = Mixer; + playing = _isPlaying; + mixerSnapshot = Mixer; decodersSnapshot = pipeline.Decoders; - loopStart = _loopStartFrames; - loopEnd = _loopEndFrames; - loopEnabled = _loopRegion.IsEnabled; + loopStart = _loopStartFrames; + loopEnd = _loopEndFrames; + loopEnabled = _loopRegion.IsEnabled; progressReporter = _progressReporter; } @@ -286,7 +290,8 @@ public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable if (!decoder.TryDecodeNextBlock(out var block)) { eof = true; - foreach (var b in _stemBlocks) b.Dispose(); + foreach (var b in _stemBlocks) + b.Dispose(); _stemBlocks.Clear(); break; } @@ -307,7 +312,7 @@ public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable b.Dispose(); _stemBlocks.Clear(); - await _timeStretchEngine.Submit(mixed); + await _timeStretchEngine.Submit(mixed, ct).ConfigureAwait(false); var progress = TimeSpan.FromSeconds( (double)_currentFramePosition / _outputDevice.SampleRate); @@ -336,11 +341,12 @@ public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable private async Task StretchLoopAsync(PipelineState pipeline, CancellationToken ct) { + await Task.Yield(); try { while (!ct.IsCancellationRequested) { - var stretched = await _timeStretchEngine.Receive(); + var stretched = await _timeStretchEngine.Receive(ct).ConfigureAwait(false); if (stretched.Buffer == null) { @@ -357,7 +363,10 @@ public sealed class StemPlaybackEngine : IStemPlaybackEngine, IDisposable _outputDevice.Write(stretched.Buffer.Span); } } - catch { } + catch(Exception ex) + { + Debug.WriteLine($"StemPlaybackEngine: Error in StretchLoopAsync: {ex.Message}"); + } } diff --git a/AudioCore/Impl/WasapiOutputDevice.cs b/AudioCore/Impl/WasapiOutputDevice.cs index beb8385..55deedd 100644 --- a/AudioCore/Impl/WasapiOutputDevice.cs +++ b/AudioCore/Impl/WasapiOutputDevice.cs @@ -101,18 +101,39 @@ public sealed class WasapiOutputDevice : IAudioOutputDevice, IDisposable } + private Lock _lock = new(); + private void Send(byte[] bytes) { - // Wait until buffer has enough free space - while (_buffer.BufferedBytes + bytes.Length > _buffer.BufferLength) - { - // Sleep a tiny amount to let WASAPI consume data - Thread.Sleep(2); - } + if (_out.PlaybackState != PlaybackState.Playing) + return; - _buffer.AddSamples(bytes, 0, bytes.Length); + lock (_lock) + { + + int offset = 0; + + while (offset < bytes.Length) + { + if (_out.PlaybackState != PlaybackState.Playing) + return; + + int free = _buffer.BufferLength - _buffer.BufferedBytes; + if (free <= 0) + { + Thread.Sleep(2); + continue; + } + + int toWrite = Math.Min(free, bytes.Length - offset); + + _buffer.AddSamples(bytes, offset, toWrite); + offset += toWrite; + } + } } + public void Dispose() { _out.Dispose(); diff --git a/AudioCore/Interfaces/ITimeStretchEngine.cs b/AudioCore/Interfaces/ITimeStretchEngine.cs index 3e299c6..8617dde 100644 --- a/AudioCore/Interfaces/ITimeStretchEngine.cs +++ b/AudioCore/Interfaces/ITimeStretchEngine.cs @@ -11,6 +11,6 @@ public interface ITimeStretchEngine void Configure(PlaybackSpeedSettings settings); // Streaming block processing - Task Submit(MixedAudioBlock input); - Task Receive(); + Task Submit(MixedAudioBlock input, CancellationToken token); + Task Receive(CancellationToken token); } diff --git a/AudioCore_Tests/StemPlaybackEngine_Tests.cs b/AudioCore_Tests/StemPlaybackEngine_Tests.cs index 7a32927..add488d 100644 --- a/AudioCore_Tests/StemPlaybackEngine_Tests.cs +++ b/AudioCore_Tests/StemPlaybackEngine_Tests.cs @@ -110,13 +110,13 @@ public sealed class StemPlaybackEngine_Tests // no-op for tests } - public Task Submit(MixedAudioBlock input) + public Task Submit(MixedAudioBlock input, CancellationToken token) { _lastInput = input; return Task.CompletedTask; } - public Task Receive() + public Task Receive(CancellationToken token) { if (_lastInput.Buffer == null) return Task.FromResult(default(TimeStretchedAudioBlock)); diff --git a/AudioCore_Tests/TimeStretchEngine_Tests.cs b/AudioCore_Tests/TimeStretchEngine_Tests.cs index 6a6ed6b..25e0f50 100644 --- a/AudioCore_Tests/TimeStretchEngine_Tests.cs +++ b/AudioCore_Tests/TimeStretchEngine_Tests.cs @@ -34,8 +34,8 @@ public sealed class TimeStretchEngine_Tests var input = MakeBlock(5000); - await engine.Submit(input); - var output = await engine.Receive(); + await engine.Submit(input, CancellationToken.None); + var output = await engine.Receive(CancellationToken.None); Assert.IsGreaterThan(0, output.Frames); Assert.AreEqual(2, output.Channels); @@ -59,12 +59,12 @@ public sealed class TimeStretchEngine_Tests var input = MakeBlock(1000); for (var i = 0; i < 25; i++) - await engine.Submit(input); + await engine.Submit(input, CancellationToken.None); var normalFrames = 0; while (true) { - using var data = await engine.Receive(); + using var data = await engine.Receive(CancellationToken.None); normalFrames += data.Frames; if (data.Buffer == null) break; @@ -75,12 +75,12 @@ public sealed class TimeStretchEngine_Tests for (var i = 0; i < 25; i++) - await engine.Submit(input); + await engine.Submit(input, CancellationToken.None); var fasterFrames = 0; while (true) { - using var data = await engine.Receive(); + using var data = await engine.Receive(CancellationToken.None); fasterFrames += data.Frames; if (data.Buffer == null) break; @@ -98,12 +98,12 @@ public sealed class TimeStretchEngine_Tests var input = MakeBlock(1000); for (var i = 0; i < 25; i++) - await engine.Submit(input); + await engine.Submit(input, CancellationToken.None); var normalFrames = 0; while (true) { - using var data = await engine.Receive(); + using var data = await engine.Receive(CancellationToken.None); normalFrames += data.Frames; if (data.Buffer == null) break; @@ -112,12 +112,12 @@ public sealed class TimeStretchEngine_Tests engine.Configure(new PlaybackSpeedSettings { Speed = 0.5f }); for (var i = 0; i < 25; i++) - await engine.Submit(input); + await engine.Submit(input, CancellationToken.None); var slowerFrames = 0; while (true) { - using var data = await engine.Receive(); + using var data = await engine.Receive(CancellationToken.None); slowerFrames += data.Frames; if (data.Buffer == null) break; @@ -135,13 +135,13 @@ public sealed class TimeStretchEngine_Tests var input = MakeBlock(100); - await engine.Submit(input); - var before = await engine.Receive(); + await engine.Submit(input, CancellationToken.None); + var before = await engine.Receive(CancellationToken.None); engine.Configure(new PlaybackSpeedSettings { Speed = 0.75f }); - await engine.Submit(input); - var after = await engine.Receive(); + await engine.Submit(input, CancellationToken.None); + var after = await engine.Receive(CancellationToken.None); Assert.AreEqual(0, after.Frames); } diff --git a/Directory.Packages.props b/Directory.Packages.props index 2b65ff4..e57eeff 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -3,19 +3,19 @@ true - - - - - + + + + + - + - + \ No newline at end of file