diff --git a/Package.resolved b/Package.resolved index b5d9e2fc..dc979f8f 100644 --- a/Package.resolved +++ b/Package.resolved @@ -63,6 +63,15 @@ "version" : "12.13.0" } }, + { + "identity" : "cxxqueue", + "kind" : "remoteSourceControl", + "location" : "https://github.com/sbooth/CXXQueue", + "state" : { + "revision" : "7e4e21a78c4bde49d6b4f4ec4220dfbf81e8e1c5", + "version" : "0.1.0" + } + }, { "identity" : "cxxtaglib", "kind" : "remoteSourceControl", diff --git a/Package.swift b/Package.swift index 8896cb02..2250b01c 100644 --- a/Package.swift +++ b/Package.swift @@ -28,6 +28,7 @@ let package = Package( .package(url: "https://github.com/sbooth/CXXAudioRingBuffer", .upToNextMinor(from: "0.1.1")), .package(url: "https://github.com/sbooth/CXXDispatchSemaphore", .upToNextMinor(from: "0.4.1")), .package(url: "https://github.com/sbooth/CXXMessageQueue", .upToNextMinor(from: "0.2.0")), + .package(url: "https://github.com/sbooth/CXXQueue", .upToNextMinor(from: "0.1.0")), .package(url: "https://github.com/sbooth/CXXUnfairLock", .upToNextMinor(from: "0.3.1")), // Standalone dependencies from source @@ -66,6 +67,7 @@ let package = Package( .product(name: "CXXAudioRingBuffer", package: "CXXAudioRingBuffer"), .product(name: "CXXDispatchSemaphore", package: "CXXDispatchSemaphore"), .product(name: "CXXMessageQueue", package: "CXXMessageQueue"), + .product(name: "CXXQueue", package: "CXXQueue"), .product(name: "CXXUnfairLock", package: "CXXUnfairLock"), // Standalone dependencies .product(name: "dumb", package: "CDUMB"), diff --git a/Sources/CSFBAudioEngine/Player/AudioPlayer.h b/Sources/CSFBAudioEngine/Player/AudioPlayer.h index dba628d6..c2f6606c 100644 --- a/Sources/CSFBAudioEngine/Player/AudioPlayer.h +++ b/Sources/CSFBAudioEngine/Player/AudioPlayer.h @@ -15,6 +15,7 @@ #import #import #import +#import #import @@ -34,6 +35,33 @@ namespace sfb { +namespace detail { + +/// A descriptor for a decoded chunk of audio. +struct DecodedChunkDescriptor final { + /// The playback generation at the time this chunk was decoded + uint64_t playbackGeneration_{0}; + /// Decoder sequence number that produced the audio. + uint64_t sequenceNumber_{0}; + /// Decoder frame position for the first audio frame in the chunk. + int64_t framePosition_{0}; + /// Number of audio frames in the chunk. + uint32_t frameLength_{0}; +}; + +/// A descriptor for a rendering chunk of audio. +struct RenderingChunkDescriptor final { + /// The decoded chunk descriptor. + DecodedChunkDescriptor descriptor_{}; + /// The number of frames consumed from `descriptor_` + uint32_t framesConsumed_{0}; + + /// Returns the number of frames remaining in this chunk + [[nodiscard]] uint32_t framesRemaining() const noexcept { return descriptor_.frameLength_ - framesConsumed_; } +}; + +} /* namespace detail */ + // MARK: - AudioPlayer /// SFBAudioPlayer implementation @@ -54,7 +82,12 @@ class AudioPlayer final { using DecoderStateVector = std::vector>; /// Ring buffer transferring audio between the decoding thread and the render block - spsc::AudioRingBuffer audioRingBuffer_; + spsc::AudioRingBuffer audioBuffer_; + /// Queue transferring audio metadata between the decoding thread and the render block + spsc::Queue audioMetadata_; + /// The current transport epoch + std::atomic playbackGeneration_{1}; + static_assert(std::atomic::is_always_lock_free, "Lock-free std::atomic required"); /// Active decoders and associated state DecoderStateVector activeDecoders_; @@ -249,6 +282,8 @@ class AudioPlayer final { /// Render block implementation OSStatus render(BOOL &isSilence, const AudioTimeStamp ×tamp, AVAudioFrameCount frameCount, AudioBufferList *_Nonnull outputData) noexcept; + /// The current rendering chunk descriptor + detail::RenderingChunkDescriptor renderingChunk_{}; // MARK: - Events @@ -311,6 +346,9 @@ class AudioPlayer final { /// Returns the first decoder state in `activeDecoders_` that has not been canceled DecoderState *_Nullable firstActiveDecoderState() const noexcept; + /// Returns the decoder state in `activeDecoders_` with the specified sequence number + DecoderState *_Nullable decoderStateWithSequenceNumber(uint64_t sequenceNumber) const noexcept; + public: // MARK: - AVAudioEngine Notification Handling diff --git a/Sources/CSFBAudioEngine/Player/AudioPlayer.mm b/Sources/CSFBAudioEngine/Player/AudioPlayer.mm index b7e97583..f5193532 100644 --- a/Sources/CSFBAudioEngine/Player/AudioPlayer.mm +++ b/Sources/CSFBAudioEngine/Player/AudioPlayer.mm @@ -25,13 +25,14 @@ #import #import #import +#import #import namespace { -/// The default ring buffer capacity -constexpr std::size_t ringBufferCapacity = 16384; -/// The minimum number of frames to write to the ring buffer +/// The default audio ring buffer capacity in frames +constexpr std::size_t audioBufferCapacity = 16384; +/// The minimum number of frames to write to the audio ring buffer constexpr AVAudioFrameCount ringBufferChunkSize = 2048; /// The number of nanoseconds in one second @@ -273,7 +274,7 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re /// Returns `true` if a seek is pending bool isSeekRequested() const noexcept; /// Performs the pending seek request - bool performSeek(NSError **error) noexcept; + std::optional performSeek(NSError **error) noexcept; }; std::atomic AudioPlayer::DecoderState::sequenceCounter_{1}; @@ -389,7 +390,7 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re return true; } - this->framesDecoded_.fetch_add(framesDecoded, std::memory_order_acq_rel); + this->framesDecoded_.fetch_add(framesDecoded, std::memory_order_release); // Only PCM to PCM conversions are performed if (![converter_ convertToBuffer:buffer fromBuffer:decodeBuffer_ error:error]) { @@ -422,7 +423,7 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re } /// Performs the pending seek request -inline bool AudioPlayer::DecoderState::performSeek(NSError **error) noexcept { +inline std::optional AudioPlayer::DecoderState::performSeek(NSError **error) noexcept { const auto requestedFrame = requestedFrame_.load(std::memory_order_acquire); #if DEBUG assert(requestedFrame != SFBUnknownFramePosition); @@ -443,28 +444,26 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re *error = seekError; } clearSeekRequest(); - return false; + return std::nullopt; } // Reset the converter to flush any buffers [converter_ reset]; - const auto framePosition = decoder_.framePosition; - if (framePosition != SFBUnknownFramePosition) { - if (framePosition != requestedFrame) { - os_log_info(log_, "Inaccurate seek to frame %lld, got %lld", requestedFrame, framePosition); - } - - // Update the frame counters accordingly - // A seek is handled in essentially the same way as initial playback - framesDecoded_.store(framePosition, std::memory_order_release); - framesRendered_.store(framePosition, std::memory_order_release); - } - // Clear the seek request clearSeekRequest(); - return framePosition != SFBUnknownFramePosition; + const auto framePosition = decoder_.framePosition; + if (framePosition == SFBUnknownFramePosition) { + os_log_error(log_, "Unknown frame position in %{public}@ after seeking to frame %lld", decoder_, + requestedFrame); + return std::nullopt; + } + if (framePosition != requestedFrame) { + os_log_info(log_, "Inaccurate seek to frame %lld, got %lld", requestedFrame, framePosition); + } + + return framePosition; } } /* namespace sfb */ @@ -482,19 +481,19 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re throw std::runtime_error("Unable to create AVAudioFormat"); } - // Allocate the audio ring buffer moving audio from the decoder queue to the render block - if (!audioRingBuffer_.allocate(*(format.streamDescription), ringBufferCapacity)) { + // Allocate the audio buffer carrying audio from the decoder thread to the render block + if (!audioBuffer_.allocate(*(format.streamDescription), audioBufferCapacity)) { os_log_error(log_, - "Unable to create audio ring buffer: spsc::AudioRingBuffer::allocate failed with format " + "Unable to create audio buffer: spsc::AudioRingBuffer::allocate failed with format " "%{public}@ and capacity %zu", - SFBASBDFormatDescription(format.streamDescription), ringBufferCapacity); + SFBASBDFormatDescription(format.streamDescription), audioBufferCapacity); throw std::runtime_error("spsc::AudioRingBuffer::allocate failed"); } // ======================================== // Event Processing Setup - // Create the dispatch queue used for event processing + // Create the dispatch queue used for asynchronous event processing auto attr = dispatch_queue_attr_make_with_qos_class(DISPATCH_QUEUE_SERIAL, QOS_CLASS_USER_INITIATED, 0); if (attr == nullptr) { os_log_error(log_, "dispatch_queue_attr_make_with_qos_class failed"); @@ -1169,10 +1168,15 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re // Whether there is a mismatch between the rendering format and the next decoder's processing format auto formatMismatch = false; + /// Sets the decoder state's error and cancellation request flag + const auto setErrorAndRequestCancel = [](DecoderState *_Nonnull decoderState, NSError *_Nonnull error) noexcept { + decoderState->error_ = error; + decoderState->setFlags(DecoderState::Flags::cancelRequested); + }; + while (!stoken.stop_requested()) { // The decoder state being processed DecoderState *decoderState = nullptr; - auto ringBufferStale = false; { std::lock_guard lock{activeDecodersMutex_}; @@ -1180,9 +1184,10 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re // Process cancellations auto signal = false; auto anyCanceled = false; + for (const auto &decoderState : activeDecoders_) { - const auto flags = decoderState->loadFlags(); - if (bits::is_set_or_is_clear(flags, DecoderState::Flags::isCanceled, + const auto decoderFlags = decoderState->loadFlags(); + if (bits::is_set_or_is_clear(decoderFlags, DecoderState::Flags::isCanceled, DecoderState::Flags::cancelRequested)) { continue; } @@ -1193,9 +1198,12 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re os_log_error(log_, "Aborting decoding for %{public}@ due to error", decoderState->decoder_); } - // Drain the ring buffer if the decoder could have contributed any stale frames - if (bits::is_set(flags, DecoderState::Flags::decodingStarted)) { - ringBufferStale = true; + if (bits::is_set(decoderFlags, DecoderState::Flags::decodingStarted)) { + // Drain the ring buffer since the decoder could have contributed stale frames + setFlags(Flags::drainRequired); + + // Increment the playback epoch to expire any inflight events + playbackGeneration_.fetch_add(1, std::memory_order_release); } decoderState->setFlags(DecoderState::Flags::isCanceled); @@ -1225,24 +1233,33 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re // Process pending seeks if (decoderState != nullptr && decoderState->isSeekRequested()) { - // Mute until the seek is complete and the ring buffer is refilled - setFlags(Flags::isMuted); - - if (NSError *seekError = nil; !decoderState->performSeek(&seekError)) { - decoderState->error_ = seekError; - decoderState->setFlags(DecoderState::Flags::cancelRequested); + NSError *seekError = nil; + const auto framePosition = decoderState->performSeek(&seekError); + if (!framePosition.has_value()) { + setErrorAndRequestCancel(decoderState, seekError); continue; } - if (const auto frame = decoderState->framesDecoded_.load(std::memory_order_acquire); - events_.enqueue(EventCommand::seek, decoderState->sequenceNumber_, frame)) { + // Mute until the seek is complete and the ring buffer is drained and refilled + setFlags(Flags::isMuted | Flags::drainRequired); + + { + // Ensure the playback epoch increment and frame counter updates occur together + std::lock_guard lock{activeDecodersMutex_}; + + // Increment the playback epoch to expire any inflight events + playbackGeneration_.fetch_add(1, std::memory_order_release); + + decoderState->framesDecoded_.store(framePosition.value(), std::memory_order_release); + decoderState->framesRendered_.store(framePosition.value(), std::memory_order_release); + } + + if (events_.enqueue(EventCommand::seek, decoderState->sequenceNumber_, framePosition.value())) { eventSemaphore_.signal(); } else { os_log_fault(log_, "Error writing decoder seek event"); } - ringBufferStale = true; - if (bits::is_set(decoderState->loadFlags(), DecoderState::Flags::decodingComplete)) { os_log_debug(log_, "Resuming decoding for %{public}@", decoderState->decoder_); @@ -1268,21 +1285,36 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re continue; } - const auto flags = nextDecoderState->loadFlags(); - if (bits::is_set(flags, DecoderState::Flags::isCanceled)) { + const auto nextDecoderFlags = nextDecoderState->loadFlags(); + if (bits::is_set(nextDecoderFlags, DecoderState::Flags::isCanceled)) { continue; } - if (bits::is_set(flags, DecoderState::Flags::decodingStarted)) { + + if (bits::is_set(nextDecoderFlags, DecoderState::Flags::decodingStarted)) { os_log_debug(log_, "Suspending decoding for %{public}@", nextDecoderState->decoder_); // TODO: Investigate a per-state buffer to mitigate frame loss if (nextDecoderState->decoder_.supportsSeeking) { nextDecoderState->requestSeekToFrame(0); - if (NSError *seekError = nil; !nextDecoderState->performSeek(&seekError)) { - nextDecoderState->error_ = seekError; - nextDecoderState->setFlags(DecoderState::Flags::cancelRequested); + + NSError *seekError = nil; + const auto framePosition = nextDecoderState->performSeek(&seekError); + if (!framePosition.has_value()) { + setErrorAndRequestCancel(nextDecoderState.get(), seekError); continue; } + + nextDecoderState->framesDecoded_.store(framePosition.value(), + std::memory_order_release); + nextDecoderState->framesRendered_.store(framePosition.value(), + std::memory_order_release); + + if (events_.enqueue(EventCommand::seek, nextDecoderState->sequenceNumber_, + framePosition.value())) { + eventSemaphore_.signal(); + } else { + os_log_fault(log_, "Error writing decoder seek event"); + } } else { os_log_error(log_, "Discarding %lld frames from %{public}@", nextDecoderState->framesDecoded_.load(std::memory_order_acquire), @@ -1302,26 +1334,17 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re } } - // Request a drain of the ring buffer during the next render cycle to prevent audible artifacts from seeking or - // cancellation - if (ringBufferStale) { - setFlags(Flags::drainRequired); - } - // Get the earliest decoder state that has not completed decoding { std::lock_guard lock{activeDecodersMutex_}; const auto iter = std::ranges::find_if(activeDecoders_, [](const auto &decoderState) noexcept { - const auto flags = decoderState->loadFlags(); - return bits::has_none(flags, DecoderState::Flags::isCanceled | DecoderState::Flags::decodingComplete); + const auto decoderFlags = decoderState->loadFlags(); + return bits::has_none(decoderFlags, + DecoderState::Flags::isCanceled | DecoderState::Flags::decodingComplete); }); - if (iter != activeDecoders_.cend()) { - decoderState = iter->get(); - } else { - decoderState = nullptr; - } + decoderState = iter != activeDecoders_.cend() ? iter->get() : nullptr; } // Dequeue the next decoder if there are no decoders that haven't completed decoding @@ -1362,8 +1385,7 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re if (!decoderState->decoder_.isOpen) { if (NSError *error = nil; ![decoderState->decoder_ openReturningError:&error]) { os_log_error(log_, "Error opening %{public}@: %{public}@", decoderState->decoder_, error); - decoderState->error_ = error; - decoderState->setFlags(DecoderState::Flags::cancelRequested); + setErrorAndRequestCancel(decoderState, error); continue; } @@ -1379,10 +1401,9 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re "Error allocating decoder state data: DecoderStateData::allocate failed with frame " "capacity %u", ringBufferChunkSize); - decoderState->error_ = [NSError errorWithDomain:SFBAudioPlayerErrorDomain - code:SFBAudioPlayerErrorCodeInternalError - userInfo:nil]; - decoderState->setFlags(DecoderState::Flags::cancelRequested); + setErrorAndRequestCancel(decoderState, [NSError errorWithDomain:SFBAudioPlayerErrorDomain + code:SFBAudioPlayerErrorCodeInternalError + userInfo:nil]); continue; } @@ -1406,10 +1427,10 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re os_log_error(log_, "Error creating AVAudioPCMBuffer with format %{public}@ and frame capacity %u", stringDescribingAVAudioFormat(renderFormat), ringBufferChunkSize); - decoderState->error_ = [NSError errorWithDomain:SFBAudioPlayerErrorDomain - code:SFBAudioPlayerErrorCodeInternalError - userInfo:nil]; - decoderState->setFlags(DecoderState::Flags::cancelRequested); + setErrorAndRequestCancel(decoderState, + [NSError errorWithDomain:SFBAudioPlayerErrorDomain + code:SFBAudioPlayerErrorCodeInternalError + userInfo:nil]); continue; } } @@ -1430,18 +1451,17 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re }(); if (okToReconfigure) { - clearFlags(Flags::drainRequired); - formatMismatch = false; - os_log_debug(log_, "Non-gapless join for %{public}@", decoderState->decoder_); auto renderFormat = decoderState->converter_.outputFormat; if (NSError *error = nil; !configureProcessingGraphAndRingBufferForFormat(renderFormat, &error)) { - decoderState->error_ = error; - decoderState->setFlags(DecoderState::Flags::cancelRequested); + setErrorAndRequestCancel(decoderState, error); continue; } + clearFlags(Flags::drainRequired); + formatMismatch = false; + // Allocate the buffer that is the intermediary between the decoder state and the ring buffer if (auto format = buffer.format; format.channelCount != renderFormat.channelCount || format.sampleRate != renderFormat.sampleRate) { @@ -1451,10 +1471,10 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re os_log_error(log_, "Error creating AVAudioPCMBuffer with format %{public}@ and frame capacity %u", stringDescribingAVAudioFormat(renderFormat), ringBufferChunkSize); - decoderState->error_ = [NSError errorWithDomain:SFBAudioPlayerErrorDomain - code:SFBAudioPlayerErrorCodeInternalError - userInfo:nil]; - decoderState->setFlags(DecoderState::Flags::cancelRequested); + setErrorAndRequestCancel(decoderState, + [NSError errorWithDomain:SFBAudioPlayerErrorDomain + code:SFBAudioPlayerErrorCodeInternalError + userInfo:nil]); continue; } } @@ -1466,12 +1486,18 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re if (decoderState != nullptr) { if (const auto flags = loadFlags(); bits::is_clear(flags, Flags::drainRequired)) { - // Decode and write chunks to the ring buffer - while (audioRingBuffer_.freeSpace() >= ringBufferChunkSize) { + // Decode and write chunks and metadata to the ring buffers + while (audioBuffer_.freeSpace() >= ringBufferChunkSize && !audioMetadata_.isFull()) { + + // The chunk descriptor for the chunk to be decoded + detail::DecodedChunkDescriptor descriptor{}; + descriptor.playbackGeneration_ = playbackGeneration_.load(std::memory_order_acquire); + descriptor.sequenceNumber_ = decoderState->sequenceNumber_; + // Decoding started - if (const auto flags = decoderState->loadFlags(); - bits::is_clear(flags, DecoderState::Flags::decodingStarted)) { - const auto suspended = bits::is_set(flags, DecoderState::Flags::decodingSuspended); + if (const auto decoderFlags = decoderState->loadFlags(); + bits::is_clear(decoderFlags, DecoderState::Flags::decodingStarted)) { + const auto suspended = bits::is_set(decoderFlags, DecoderState::Flags::decodingSuspended); if (!suspended) { os_log_debug(log_, "Decoding starting for %{public}@", decoderState->decoder_); @@ -1492,26 +1518,32 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re } } + descriptor.framePosition_ = decoderState->framesDecoded_.load(std::memory_order_acquire); + // Decode audio into the buffer, converting to the rendering format in the process if (NSError *error = nil; !decoderState->decodeAudio(buffer, &error)) { - decoderState->error_ = error; - decoderState->setFlags(DecoderState::Flags::cancelRequested); + setErrorAndRequestCancel(decoderState, error); goto next_outer_iteration; } - // Write the decoded audio to the ring buffer for rendering - const auto framesWritten = audioRingBuffer_.write(buffer.audioBufferList, buffer.frameLength); + descriptor.frameLength_ = buffer.frameLength; + + // Write the decoded chunk descriptor to the metadata buffer + if (!audioMetadata_.push(descriptor)) { + os_log_fault(log_, "Error writing chunk descriptor: spsc::Queue::push failed"); + } + + // Write the decoded audio to the audio buffer for rendering + const auto framesWritten = audioBuffer_.write(buffer.audioBufferList, buffer.frameLength); if (framesWritten != buffer.frameLength) { - os_log_fault( - log_, - "Error writing audio to ring buffer: spsc::AudioRingBuffer::write failed for %u frames", - buffer.frameLength); + os_log_fault(log_, "Error writing audio: spsc::AudioRingBuffer::write failed for %u frames", + buffer.frameLength); } // Decoding complete - if (const auto flags = decoderState->loadFlags(); - bits::is_set(flags, DecoderState::Flags::decodingComplete)) { - const auto resumed = bits::is_set(flags, DecoderState::Flags::decodingResumed); + if (const auto decoderFlags = decoderState->loadFlags(); + bits::is_set(decoderFlags, DecoderState::Flags::decodingComplete)) { + const auto resumed = bits::is_set(decoderFlags, DecoderState::Flags::decodingResumed); // Submit the decoding complete event for the first completion only if (!resumed) { @@ -1552,14 +1584,14 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re } else { // Determine timeout based on ring buffer free space // Attempt to keep the ring buffer 75% full - const auto targetMaxFreeSpace = audioRingBuffer_.capacity() / 4; - const auto freeSpace = audioRingBuffer_.freeSpace(); + const auto targetMaxFreeSpace = audioBuffer_.capacity() / 4; + const auto freeSpace = audioBuffer_.freeSpace(); if (freeSpace > targetMaxFreeSpace) { // Minimal timeout if the ring buffer has more free space than desired deltaNanos = static_cast(2.5 * NSEC_PER_MSEC); } else { - const auto duration = (targetMaxFreeSpace - freeSpace) / audioRingBuffer_.format().mSampleRate; + const auto duration = (targetMaxFreeSpace - freeSpace) / audioBuffer_.format().mSampleRate; deltaNanos = static_cast(duration * NSEC_PER_SEC); } } @@ -1588,7 +1620,9 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re // Discard any stale frames in the ring buffer from a seek or decoder cancelation if (bits::is_set(flags, Flags::drainRequired)) { - audioRingBuffer_.drain(); + audioBuffer_.drain(); + audioMetadata_.discardAll(); + renderingChunk_ = {}; clearFlags(Flags::drainRequired); zeroABL(outputData); isSilence = YES; @@ -1602,20 +1636,49 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re return noErr; } + /// Computes the event time for a given frame offset + const auto eventTimeForFrameOffset = [&](uint32_t frameOffset) noexcept -> uint64_t { + const auto deltaSeconds = frameOffset / audioBuffer_.format().mSampleRate; + const auto scaledNanos = static_cast(deltaSeconds * timestamp.mRateScalar * nanosecondsPerSecond); + return timestamp.mHostTime + host_time::fromNanoseconds(scaledNanos); + }; + // Read audio from the ring buffer - const auto framesRead = audioRingBuffer_.read(outputData, frameCount); + const auto framesRead = static_cast(audioBuffer_.read(outputData, frameCount)); + + // Read and process chunk descriptors for the rendered audio if (framesRead > 0) { - if (!events_.enqueue(EventCommand::framesRendered, timestamp.mHostTime, timestamp.mRateScalar, - static_cast(framesRead))) { - setFlags(Flags::renderEventDropped); - } + auto framesRemaining = framesRead; + do { + // Read the next chunk descriptor if needed + if (renderingChunk_.framesRemaining() == 0) { + if (!audioMetadata_.pop(renderingChunk_.descriptor_)) { + setFlags(Flags::renderEventDropped); + break; + } + renderingChunk_.framesConsumed_ = 0; + } + + const auto chunkFrames = std::min(renderingChunk_.framesRemaining(), framesRemaining); + if (chunkFrames > 0) { + const auto eventTime = eventTimeForFrameOffset(framesRead - framesRemaining); + if (!events_.enqueue(EventCommand::framesRendered, eventTime, + renderingChunk_.descriptor_.sequenceNumber_, chunkFrames, + renderingChunk_.descriptor_.playbackGeneration_)) { + setFlags(Flags::renderEventDropped); + break; + } + + renderingChunk_.framesConsumed_ += chunkFrames; + framesRemaining -= chunkFrames; + } + } while (framesRemaining > 0); } else { isSilence = YES; } if (framesRead != frameCount) { - if (!events_.enqueue(EventCommand::renderBufferUnderrun, timestamp.mHostTime, static_cast(framesRead), - static_cast(frameCount))) { + if (!events_.enqueue(EventCommand::renderBufferUnderrun, timestamp.mHostTime, framesRead, frameCount)) { setFlags(Flags::renderEventDropped); } } @@ -1710,9 +1773,8 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re { std::lock_guard lock{activeDecodersMutex_}; - if (const auto iter = std::ranges::find(activeDecoders_, sequenceNumber, &DecoderState::sequenceNumber_); - iter != activeDecoders_.cend()) { - decoder = (*iter)->decoder_; + if (const auto *decoderState = decoderStateWithSequenceNumber(sequenceNumber); decoderState != nullptr) { + decoder = decoderState->decoder_; } else { os_log_error(log_, "Decoder state with sequence number %llu missing for decoding started event", sequenceNumber); @@ -1752,9 +1814,8 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re { std::lock_guard lock{activeDecodersMutex_}; - if (const auto iter = std::ranges::find(activeDecoders_, sequenceNumber, &DecoderState::sequenceNumber_); - iter != activeDecoders_.cend()) { - decoder = (*iter)->decoder_; + if (const auto *decoderState = decoderStateWithSequenceNumber(sequenceNumber); decoderState != nullptr) { + decoder = decoderState->decoder_; } else { os_log_error(log_, "Decoder state with sequence number %llu missing for decoding complete event", sequenceNumber); @@ -1787,9 +1848,8 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re { std::lock_guard lock{activeDecodersMutex_}; - if (const auto iter = std::ranges::find(activeDecoders_, sequenceNumber, &DecoderState::sequenceNumber_); - iter != activeDecoders_.cend()) { - decoder = (*iter)->decoder_; + if (auto *decoderState = decoderStateWithSequenceNumber(sequenceNumber); decoderState != nullptr) { + decoder = decoderState->decoder_; } else { os_log_error(log_, "Decoder state with sequence number %llu missing for decoder seek event", sequenceNumber); @@ -1900,25 +1960,39 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re bool sfb::AudioPlayer::processFramesRenderedEvent() noexcept { EventCommand command; - // The host time and rate scalar from the render cycle's timestamp - uint64_t hostTime; - double rateScalar; + // The event time calculated from the render cycle's host time and rate scalar + uint64_t eventTime; + // The decoder sequence number for the decoder providing the frames + uint64_t sequenceNumber; // The number of valid frames rendered - uint32_t framesRendered; - if (!events_.dequeue(command, hostTime, rateScalar, framesRendered)) { - os_log_error(log_, "Missing timestamp or frames rendered for frames rendered event"); + uint32_t frameCount; + // The playback generation of the chunk containing the frames + uint64_t playbackGeneration; + if (!events_.dequeue(command, eventTime, sequenceNumber, frameCount, playbackGeneration)) { + os_log_error(log_, "Missing event time, decoder sequence number, frame count, or playback generation for " + "frames rendered event"); return false; } #if DEBUG assert(command == EventCommand::framesRendered); - assert(framesRendered > 0); + assert(frameCount > 0); #endif /* DEBUG */ - // Perform bookkeeping to apportion the rendered frames appropriately + // If a frames rendered event was posted it means valid frames were rendered + // during that render cycle. + // + // However, between the time the frames rendered event was queued and when it is processed + // a decoder may have been canceled or a seek may have occurred, making the event stale. + // + // This is indicated by an increment in the transport epoch/playback generation. // - // framesRendered contains the number of valid frames that were rendered - // but they could have come from multiple decoders + // NB: The generation check must happen under activeDecodersMutex_, together with the frame + // counter reads below. A seek on the decoding thread bumps playbackGeneration_ and resets a + // decoder state's framesDecoded_/framesRendered_ under the same mutex (see the seek handling + // in the decoding loop). Checking the generation before taking the lock leaves a window in + // which the seek's reset can land between the check and the reads below, making a stale + // event's frame count get applied to post-seek counters. struct RenderingEventDetails { enum class Type { @@ -1936,50 +2010,20 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re { std::lock_guard lock{activeDecodersMutex_}; - AVAudioFramePosition framesRemainingToDistribute = framesRendered; - - auto iter = activeDecoders_.cbegin(); - while (iter != activeDecoders_.cend()) { - const auto flags = (*iter)->loadFlags(); - - // Skip uninitialized decoders - if (bits::is_set(flags, DecoderState::Flags::needsInitialization)) { - ++iter; - continue; - } + // Discard stale events from previous playback generations + if (playbackGeneration != playbackGeneration_.load(std::memory_order_acquire)) { + os_log_debug(log_, "Discarding stale frames rendered event"); + return true; + } - // If a frames rendered event was posted it means valid frames were rendered - // during that render cycle. - // - // However, between the time the frames rendered event was posted and when it is processed - // - A decoder may have been canceled - // - A seek can occur - // - // Bookkeeping is handled no differently for canceled decoders but rendering notifications are suppressed - // - // In the case of a seek the frames from that event are not valid and should be discarded. - - const auto decoderFramesDecoded = (*iter)->framesDecoded_.load(std::memory_order_acquire); - const auto decoderFramesRendered = (*iter)->framesRendered_.load(std::memory_order_acquire); - const auto decoderFramesRemaining = decoderFramesDecoded - decoderFramesRendered; - - if (decoderFramesRemaining == 0) { -#if DEBUG - os_log_debug(log_, "Not accounting for %lld frames in frames rendered event", - framesRemainingToDistribute); -#endif /* DEBUG */ - break; - } + if (const auto iter = std::ranges::find(activeDecoders_, sequenceNumber, &DecoderState::sequenceNumber_); + iter != activeDecoders_.cend()) { + const auto decoderFlags = (*iter)->loadFlags(); // Rendering is starting - if (bits::has_none(flags, DecoderState::Flags::isCanceled | DecoderState::Flags::renderingStarted)) { + if (bits::is_clear(decoderFlags, DecoderState::Flags::renderingStarted)) { (*iter)->setFlags(DecoderState::Flags::renderingStarted); - const auto frameOffset = framesRendered - framesRemainingToDistribute; - const auto deltaSeconds = frameOffset / (*iter)->sampleRate(); - const auto eventTime = hostTime + host_time::fromNanoseconds(static_cast( - deltaSeconds * rateScalar * nanosecondsPerSecond)); - try { queuedEvents.push_back({RenderingEventDetails::Type::willStart, (*iter)->decoder_, eventTime}); } catch (const std::exception &e) { @@ -1988,20 +2032,15 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re } } - const auto framesFromThisDecoder = std::min(decoderFramesRemaining, framesRemainingToDistribute); - - (*iter)->framesRendered_.fetch_add(framesFromThisDecoder, std::memory_order_acq_rel); - framesRemainingToDistribute -= framesFromThisDecoder; + const auto framesDecoded = (*iter)->framesDecoded_.load(std::memory_order_acquire); + const auto framesRendered = (*iter)->framesRendered_.fetch_add(frameCount, std::memory_order_acq_rel); + const auto framesRemaining = framesDecoded - framesRendered; +#if DEBUG + assert(framesRemaining >= frameCount); +#endif /* DEBUG */ // Rendering is complete - if (bits::is_set_and_is_clear(flags, DecoderState::Flags::decodingComplete, - DecoderState::Flags::isCanceled) && - framesFromThisDecoder == decoderFramesRemaining) { - const auto frameOffset = framesRendered - framesRemainingToDistribute; - const auto deltaSeconds = frameOffset / (*iter)->sampleRate(); - const auto eventTime = hostTime + host_time::fromNanoseconds(static_cast( - deltaSeconds * rateScalar * nanosecondsPerSecond)); - + if (bits::is_set(decoderFlags, DecoderState::Flags::decodingComplete) && frameCount == framesRemaining) { try { queuedEvents.push_back({RenderingEventDetails::Type::willComplete, (*iter)->decoder_, eventTime}); } catch (const std::exception &e) { @@ -2010,15 +2049,12 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re } os_log_debug(log_, "Deleting decoder state for %{public}@", (*iter)->decoder_); - iter = activeDecoders_.erase(iter); - } else { - ++iter; - } - - // All frames processed - if (framesRemainingToDistribute == 0) { - break; + activeDecoders_.erase(iter); } + } else { + os_log_error(log_, "Decoder state with sequence number %llu missing for frames rendered event", + sequenceNumber); + return false; } } @@ -2244,15 +2280,19 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re #if DEBUG activeDecodersMutex_.assertIsOwner(); #endif /* DEBUG */ - const auto iter = std::ranges::find_if(activeDecoders_, [](const auto &decoderState) noexcept { - const auto flags = decoderState->loadFlags(); - return bits::has_none(flags, DecoderState::Flags::needsInitialization | DecoderState::Flags::isCanceled); + const auto decoderFlags = decoderState->loadFlags(); + return bits::has_none(decoderFlags, DecoderState::Flags::needsInitialization | DecoderState::Flags::isCanceled); }); - if (iter == activeDecoders_.cend()) { - return nullptr; - } - return iter->get(); + return iter != activeDecoders_.cend() ? iter->get() : nullptr; +} + +auto sfb::AudioPlayer::decoderStateWithSequenceNumber(uint64_t sequenceNumber) const noexcept -> DecoderState * { +#if DEBUG + activeDecodersMutex_.assertIsOwner(); +#endif /* DEBUG */ + const auto iter = std::ranges::find(activeDecoders_, sequenceNumber, &DecoderState::sequenceNumber_); + return iter != activeDecoders_.cend() ? iter->get() : nullptr; } // MARK: - AVAudioEngine Notification Handling @@ -2463,11 +2503,11 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re // Allocate a temporary ring buffer for the new format before touching the engine or graph spsc::AudioRingBuffer ringBuffer; - if (!ringBuffer.allocate(*(format.streamDescription), ringBufferCapacity)) { + if (!ringBuffer.allocate(*(format.streamDescription), audioBufferCapacity)) { os_log_error(log_, - "Unable to create audio ring buffer: spsc::AudioRingBuffer::allocate failed with format " + "Unable to create audio buffer: spsc::AudioRingBuffer::allocate failed with format " "%{public}@ and capacity %zu", - SFBASBDFormatDescription(format.streamDescription), ringBufferCapacity); + SFBASBDFormatDescription(format.streamDescription), audioBufferCapacity); if (error != nullptr) { *error = [NSError errorWithDomain:SFBAudioPlayerErrorDomain code:SFBAudioPlayerErrorCodeInternalError @@ -2491,9 +2531,11 @@ Flags clearFlags(Flags flags, std::memory_order order = std::memory_order_acq_re outputBus:0] firstObject]; [engine_ disconnectNodeOutput:sourceNode_]; - // Adopt the new ring buffer - // The move is not thread-safe but the engine is stopped - audioRingBuffer_ = std::move(ringBuffer); + // Adopt the new ring buffer and reset the render state + // These operations are not thread-safe but the engine is stopped + audioBuffer_ = std::move(ringBuffer); + audioMetadata_.discardAll(); + renderingChunk_ = {}; // Reconnect the source node to the next node in the processing chain // This is the mixer node in the default configuration, but additional nodes may