diff --git a/src/libs/audio/impl/utils/PcmDecodeStreamer.cpp b/src/libs/audio/impl/utils/PcmDecodeStreamer.cpp index ef25633f..c7d80565 100644 --- a/src/libs/audio/impl/utils/PcmDecodeStreamer.cpp +++ b/src/libs/audio/impl/utils/PcmDecodeStreamer.cpp @@ -21,26 +21,25 @@ #include -#include "audio/Exception.hpp" #include "core/ILogger.hpp" +#include "audio/Exception.hpp" #include "audio/IAudioOutput.hpp" #include "audio/IPcmDecoder.hpp" namespace lms::audio::utils { - std::shared_ptr createPcmDecodeStreamer(boost::asio::io_context& ioContext, const PcmDecodeStreamerParameters& parameters) + std::unique_ptr createPcmDecodeStreamer(boost::asio::io_context& ioContext, const PcmDecodeStreamerParameters& parameters) { - return std::make_shared(ioContext, parameters); + return std::make_unique(ioContext, parameters); } PcmDecodeStreamer::PcmDecodeStreamer(boost::asio::io_context& ioContext, const PcmDecodeStreamerParameters& parameters) : _ioContext{ ioContext } , _strand{ _ioContext } , _outputStream{ parameters.outputStream } - , _pcmDecoder{ audio::createPcmDecoder(parameters.file, parameters.offset, parameters.pcmParameters) } { - prepareBuffers(); + prepareBuffers(parameters.bufferCount, parameters.bufferDuration); } PcmDecodeStreamer::~PcmDecodeStreamer() @@ -48,42 +47,51 @@ namespace lms::audio::utils assert(!isWritePending()); } - void PcmDecodeStreamer::start(DecodeCompleteCallback cb) + void PcmDecodeStreamer::start(const std::filesystem::path& path, std::chrono::microseconds offset, DecodeCompleteCallback cb) { assert(cb); - assert(!_decodeCompleteCallback); + assert(isComplete()); // previous job must be finished or cancelled + + _pcmDecoder = audio::createPcmDecoder(path, offset, getPcmParameters()); // may throw + _aborted = false; + _eofReached = false; _ioContext.get_executor().on_work_started(); _decodeCompleteCallback = std::move(cb); - boost::asio::post(_strand, [self = shared_from_this()] { - self->decodeSome(); + boost::asio::post(_strand, [this] { + decodeSome(); }); } void PcmDecodeStreamer::abort() { - boost::asio::post(_strand, [self = shared_from_this()] { - LMS_LOG(AUDIO, DEBUG, "Processing abort"); - self->_aborted = true; - self->_outputStream.flush(); + if (isComplete()) + return; - if (!self->isWritePending()) - self->notifyDecodeComplete(); + boost::asio::post(_strand, [this] { + LMS_LOG(AUDIO, DEBUG, "Processing abort"); + _aborted = true; + _outputStream.flush(); + + if (!isWritePending()) + notifyDecodeComplete(); }); } + bool PcmDecodeStreamer::isComplete() const + { + return !_pcmDecoder; + } + const audio::PcmParameters& PcmDecodeStreamer::getPcmParameters() const { - return _pcmDecoder->getParameters(); + return _outputStream.getParameters(); } - void PcmDecodeStreamer::prepareBuffers() + void PcmDecodeStreamer::prepareBuffers(std::size_t bufferCount, std::chrono::microseconds bufferDuration) { - constexpr std::chrono::milliseconds bufferDuration{ 100 }; - - const audio::PcmParameters& pcmParams{ getPcmParameters() }; - std::size_t sampleCountPerBuffer{ static_cast(std::chrono::duration_cast(bufferDuration).count() * pcmParams.sampleRate / std::chrono::microseconds::period::den) }; + const std::size_t sampleCountPerBuffer{ static_cast(std::chrono::duration_cast(bufferDuration).count() * getPcmParameters().sampleRate / std::chrono::microseconds::period::den) }; const std::size_t bufferSize{ sampleCountToByteCount(sampleCountPerBuffer) }; _buffers.resize(bufferCount); @@ -127,11 +135,11 @@ namespace lms::audio::utils bufferDesc.isWritePending = true; buffer = { buffer.data(), sampleCountToByteCount(sampleCount) }; - _outputStream.asyncWrite(buffer, [self = shared_from_this(), bufferIndex] { - boost::asio::post(self->_strand, [self, bufferIndex] { self->onBufferWriteComplete(bufferIndex); }); + _outputStream.asyncWrite(buffer, [this, bufferIndex] { + boost::asio::post(_strand, [this, bufferIndex] { onBufferWriteComplete(bufferIndex); }); }); - boost::asio::post(_strand, [self = shared_from_this()] { self->decodeSome(); }); + boost::asio::post(_strand, [this] { decodeSome(); }); } } @@ -185,10 +193,12 @@ namespace lms::audio::utils void PcmDecodeStreamer::notifyDecodeComplete() { - boost::asio::post(_ioContext, [self = shared_from_this(), cb = std::move(_decodeCompleteCallback)] { - LMS_LOG(AUDIO, DEBUG, "Decode complete notification"); - cb(self->_aborted); - self->_ioContext.get_executor().on_work_finished(); + LMS_LOG(AUDIO, DEBUG, "Decode complete notification"); + + boost::asio::post(_ioContext, [this, cb = std::move(_decodeCompleteCallback)] { + _pcmDecoder.reset(); + cb(_aborted); + _ioContext.get_executor().on_work_finished(); }); } diff --git a/src/libs/audio/impl/utils/PcmDecodeStreamer.hpp b/src/libs/audio/impl/utils/PcmDecodeStreamer.hpp index bfaa2da6..8e017c91 100644 --- a/src/libs/audio/impl/utils/PcmDecodeStreamer.hpp +++ b/src/libs/audio/impl/utils/PcmDecodeStreamer.hpp @@ -19,14 +19,12 @@ #pragma once -#include -#include +#include #include #include #include -#include "audio/PcmTypes.hpp" #include "audio/utils/IPcmDecodeStreamer.hpp" namespace lms::audio @@ -37,7 +35,7 @@ namespace lms::audio namespace lms::audio::utils { - class PcmDecodeStreamer : public IPcmDecodeStreamer, public std::enable_shared_from_this + class PcmDecodeStreamer : public IPcmDecodeStreamer { public: PcmDecodeStreamer(boost::asio::io_context& ioContext, const PcmDecodeStreamerParameters& parameters); @@ -47,12 +45,13 @@ namespace lms::audio::utils PcmDecodeStreamer& operator=(const PcmDecodeStreamer&) = delete; private: - void start(DecodeCompleteCallback cb) override; - void abort() override; // will call DecodeCompleteCallback once done + void start(const std::filesystem::path& path, std::chrono::microseconds offset, DecodeCompleteCallback cb) override; + void abort() override; + bool isComplete() const override; const audio::PcmParameters& getPcmParameters() const; - void prepareBuffers(); + void prepareBuffers(std::size_t bufferCount, std::chrono::microseconds bufferDuration); bool isWritePending() const; void decodeSome(); std::size_t readSamples(std::span buffer); @@ -72,7 +71,6 @@ namespace lms::audio::utils audio::IAudioOutputStream& _outputStream; std::unique_ptr _pcmDecoder; - static constexpr std::size_t bufferCount{ 4 }; std::vector _buffers; std::size_t _nextBufferIndex{}; bool _eofReached{}; diff --git a/src/libs/audio/include/audio/utils/IPcmDecodeStreamer.hpp b/src/libs/audio/include/audio/utils/IPcmDecodeStreamer.hpp index d69e8265..5a831f7a 100644 --- a/src/libs/audio/include/audio/utils/IPcmDecodeStreamer.hpp +++ b/src/libs/audio/include/audio/utils/IPcmDecodeStreamer.hpp @@ -33,22 +33,26 @@ namespace lms::audio namespace lms::audio::utils { + // helper class to decode files to PCM samples, fed into the provided output stream class IPcmDecodeStreamer { public: virtual ~IPcmDecodeStreamer() = default; using DecodeCompleteCallback = std::function; - virtual void start(DecodeCompleteCallback cb) = 0; - virtual void abort() = 0; // will call DecodeCompleteCallback once done + + // DecodeCompleteCallback is fired once the file has been fully decoded (but still buffered in output) + // You can start a new file only if the previous one is finished (i.e. once the callback is fired) + virtual void start(const std::filesystem::path& path, std::chrono::microseconds offset, DecodeCompleteCallback cb) = 0; + virtual void abort() = 0; // will call DecodeCompleteCallback once aborted + virtual bool isComplete() const = 0; }; struct PcmDecodeStreamerParameters { audio::IAudioOutputStream& outputStream; - std::filesystem::path file; - std::chrono::microseconds offset; - audio::PcmParameters pcmParameters; + std::size_t bufferCount; + std::chrono::milliseconds bufferDuration; }; - std::shared_ptr createPcmDecodeStreamer(boost::asio::io_context& ioContext, const PcmDecodeStreamerParameters& params); + std::unique_ptr createPcmDecodeStreamer(boost::asio::io_context& ioContext, const PcmDecodeStreamerParameters& parameters); } // namespace lms::audio::utils \ No newline at end of file diff --git a/src/libs/services/jukebox/impl/JukeboxService.cpp b/src/libs/services/jukebox/impl/JukeboxService.cpp index c2d919ba..b7a19410 100644 --- a/src/libs/services/jukebox/impl/JukeboxService.cpp +++ b/src/libs/services/jukebox/impl/JukeboxService.cpp @@ -20,7 +20,9 @@ #include "JukeboxService.hpp" #include +#include #include +#include #include #include @@ -59,9 +61,6 @@ namespace lms::jukebox if (_decoder) _decoder->abort(); - if (_outputStream) - _outputStream->flush(); - _ioContextRunner.wait(); LMS_LOG(JUKEBOX, INFO, "Service stopped!"); } @@ -72,7 +71,6 @@ namespace lms::jukebox std::unique_lock lock{ _mutex }; - deleteDecoder(); if (trackIndex >= _tracks.size()) { LMS_LOG(JUKEBOX, INFO, "Requested track index out of bound: stopping"); @@ -80,12 +78,15 @@ namespace lms::jukebox return; } - if (createDecoder(trackIndex, offset)) + if (!_outputStream) + return; + + abortDecoder(); + if (startDecoder(trackIndex, offset)) { _currentTrackIndex = trackIndex; _currentTrackPlaybackTimeOffset = _outputStream->getPlaybackTime(); _currentTrackStartTimeOffset = offset; - startDecoder(); _outputStream->resume(); } // TODO if failure, switch to the next song? @@ -202,10 +203,20 @@ namespace lms::jukebox void JukeboxService::onStreamReady() { + audio::utils::PcmDecodeStreamerParameters params{ + .outputStream = *_outputStream, + .bufferCount = 50, + .bufferDuration = std::chrono::milliseconds{ 1000 }, + }; + + _decoder = audio::utils::createPcmDecodeStreamer(_ioContext, params); } - bool JukeboxService::createDecoder(std::size_t trackIndex, std::chrono::microseconds offset) + bool JukeboxService::startDecoder(std::size_t trackIndex, std::chrono::microseconds offset) { + if (!_decoder) + return false; + std::filesystem::path trackPath; { auto& session{ _db.getTLSSession() }; @@ -223,48 +234,42 @@ namespace lms::jukebox try { - audio::utils::PcmDecodeStreamerParameters params{ - .outputStream = *_outputStream, - .file = trackPath, - .offset = offset, - .pcmParameters = _pcmParams, - }; - - _decoder = audio::utils::createPcmDecodeStreamer(_ioContext, params); + _decoder->start(trackPath, offset, [this](bool aborted) { + onDecodeFinished(aborted); + }); } catch (const audio::Exception& e) { - LMS_LOG(JUKEBOX, ERROR, "Failed to create PCM decoder for track " << trackPath); + LMS_LOG(JUKEBOX, ERROR, "Failed to start PCM decoder for track " << trackPath); return false; } return true; } - void JukeboxService::deleteDecoder() + void JukeboxService::abortDecoder() { + // Must not be called from within owned io_context if (_decoder) { _decoder->abort(); - _decoder.reset(); - } - } - void JukeboxService::startDecoder() - { - _decoder->start([this](bool aborted) { - onDecodeFinished(aborted); - }); + // Should be hopefully short since flushing/aborting + while (!_decoder->isComplete()) + std::this_thread::yield(); // TODO: execute some io_context stuff? + } } void JukeboxService::onDecodeFinished(bool aborted) { if (aborted) - return; // already setup to play next song + return; // already setup to play next song, if needed std::unique_lock lock{ _mutex }; - _decoder.reset(); + _currentTrackPlaybackTimeOffset = {}; + _currentTrackStartTimeOffset = {}; + if (!_currentTrackIndex) { _outputStream->pause(); @@ -274,15 +279,19 @@ namespace lms::jukebox if (++(*_currentTrackIndex) >= _tracks.size()) { _currentTrackIndex.reset(); + // let the output stream run out of data return; } - if (createDecoder(*_currentTrackIndex)) + if (startDecoder(*_currentTrackIndex)) { - startDecoder(); _currentTrackPlaybackTimeOffset = _outputStream->getPlaybackTime() + _outputStream->getLatency(); _currentTrackStartTimeOffset = {}; } + else + { + _currentTrackIndex.reset(); + } // TODO if failure, switch to the next song? } } // namespace lms::jukebox \ No newline at end of file diff --git a/src/libs/services/jukebox/impl/JukeboxService.hpp b/src/libs/services/jukebox/impl/JukeboxService.hpp index 75253464..cd5deed9 100644 --- a/src/libs/services/jukebox/impl/JukeboxService.hpp +++ b/src/libs/services/jukebox/impl/JukeboxService.hpp @@ -64,9 +64,8 @@ namespace lms::jukebox void onContextReady(); void onStreamReady(); - bool createDecoder(std::size_t trackIndex, std::chrono::microseconds offset = {}); - void deleteDecoder(); - void startDecoder(); + bool startDecoder(std::size_t trackIndex, std::chrono::microseconds offset = {}); + void abortDecoder(); void onDecodeFinished(bool aborted); // TODO: make configurable or use detected output params @@ -91,6 +90,6 @@ namespace lms::jukebox std::unique_ptr _outputContext; std::unique_ptr _outputStream; - std::shared_ptr _decoder; + std::unique_ptr _decoder; }; } // namespace lms::jukebox \ No newline at end of file diff --git a/src/tools/audioplay/LmsAudioPlay.cpp b/src/tools/audioplay/LmsAudioPlay.cpp index f2220c08..cd5f27a6 100644 --- a/src/tools/audioplay/LmsAudioPlay.cpp +++ b/src/tools/audioplay/LmsAudioPlay.cpp @@ -72,13 +72,12 @@ namespace lms { audio::utils::PcmDecodeStreamerParameters params{ .outputStream = *_outputStream, - .file = _filePath, - .offset = _offset, - .pcmParameters = _pcmParams, + .bufferCount = 2, + .bufferDuration = std::chrono::milliseconds{ 100 }, }; _fileStreamer = audio::utils::createPcmDecodeStreamer(_ioContext, params); - _fileStreamer->start([this](bool aborted) { + _fileStreamer->start(_filePath, _offset, [this](bool aborted) { if (aborted) std::cerr << "Playback aborted!" << std::endl; @@ -98,7 +97,7 @@ namespace lms }); // Gives some time for the buffer to fill in - _playTimer.expires_from_now(std::chrono::milliseconds{ 50 }); + _playTimer.expires_after(std::chrono::milliseconds{ 50 }); _playTimer.async_wait([this](const boost::system::error_code& ec) { if (ec) return;