Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ BaseAudioContext::BaseAudioContext(
: state_(ContextState::SUSPENDED),
sampleRate_(sampleRate),
audioEventHandlerRegistry_(audioEventHandlerRegistry),
audioEventProducer_(audioEventHandlerRegistry->createAudioEventProducer()),
stateChangeEvent_(audioEventHandlerRegistry),
pendingPromisesOffloader_(
std::make_unique<task_offloader::TaskOffloader<
Expand All @@ -38,7 +39,7 @@ BaseAudioContext::BaseAudioContext(
disposer_(
std::make_unique<utils::DisposerImpl<DISPOSER_PAYLOAD_SIZE>>(AUDIO_SCHEDULER_CAPACITY)),
graph_(std::make_shared<utils::graph::Graph>(AUDIO_SCHEDULER_CAPACITY, disposer_.get())),
deferredEvents_(audioEventHandlerRegistry) {}
deferredEvents_(audioEventHandlerRegistry, audioEventProducer_) {}

void BaseAudioContext::initialize(const AudioDestinationNode *destination) {
destination_ = destination;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,13 @@ class BaseAudioContext : public std::enable_shared_from_this<BaseAudioContext> {
return deferredEvents_;
}

/// @brief This context's dispatch lane, for nodes that emit events from the render thread.
/// Every node of a context shares it because they all render on that one thread; a node of
/// another context must never use it.
[[nodiscard]] std::shared_ptr<AudioEventProducer> getAudioEventProducer() const {
return audioEventProducer_;
}

template <typename F>
bool scheduleAudioEvent(F &&event) noexcept { // NOLINT(cppcoreguidelines-missing-std-forward)
std::scoped_lock lock(driverMutex_);
Expand Down Expand Up @@ -176,6 +183,8 @@ class BaseAudioContext : public std::enable_shared_from_this<BaseAudioContext> {
private:
std::atomic<float> sampleRate_;
std::shared_ptr<IAudioEventHandlerRegistry> audioEventHandlerRegistry_;
/// context's own lane into the registry's dispatch queue, shared with every node it owns.
std::shared_ptr<AudioEventProducer> audioEventProducer_;

EventCaller<AudioEvent::STATE_CHANGE> stateChangeEvent_;
/// Ledger backing dispatchStateChange()'s dedupe; contexts start suspended.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ AudioBufferBaseSourceNode::AudioBufferBaseSourceNode(
detuneParam_)),
positionChanged_(
context->getAudioEventHandlerRegistry(),
context->getAudioEventProducer(),
static_cast<int>(context->getSampleRate())) {
setOnPositionChangedInterval(options.onPositionChangedInterval);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@ AudioBufferQueueSourceNode::AudioBufferQueueSourceNode(
const std::shared_ptr<BaseAudioContext> &context,
const BaseAudioBufferSourceOptions &options)
: AudioBufferBaseSourceNode(context, options),
onBufferEndedEvent_(context->getAudioEventHandlerRegistry()) {
onBufferEndedEvent_(
context->getAudioEventHandlerRegistry(),
context->getAudioEventProducer()) {
if (options.pitchCorrection) {
// If pitch correction is enabled, add extra frames at the end
// to compensate for processing latency.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ AudioBufferSourceNode::AudioBufferSourceNode(
loopSkip_(options.loopSkip),
loopStart_(options.loopStart),
loopEnd_(options.loopEnd),
onLoopEndedEvent_(context->getAudioEventHandlerRegistry()) {
onLoopEndedEvent_(context->getAudioEventHandlerRegistry(), context->getAudioEventProducer()) {
auto onLoopEnded = [this]() {
sendOnLoopEndedEvent();
};
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,12 @@ AudioFileSourceNode::AudioFileSourceNode(
targetPlaybackRate_(options.playbackRate),
positionChanged_(
context->getAudioEventHandlerRegistry(),
context->getAudioEventProducer(),
static_cast<int>(context->getSampleRate() * ON_POSITION_CHANGED_INTERVAL),
true),
bufferingStateDispatcher_(
context->getAudioEventHandlerRegistry(),
context->getAudioEventProducer(),
static_cast<int>(context->getSampleRate() * ON_BUFFERING_STATE_DEBOUNCE_INTERVAL)) {
decoderState_->playbackRate.store(options.playbackRate, std::memory_order_release);
decoderState_->preservesPitch.store(options.preservesPitch, std::memory_order_release);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ AudioScheduledSourceNode::AudioScheduledSourceNode(
startTime_(-1.0),
stopTime_(-1.0),
playbackState_(PlaybackState::UNSCHEDULED),
onEndedEvent_(context->getAudioEventHandlerRegistry()) {}
onEndedEvent_(context->getAudioEventHandlerRegistry(), context->getAudioEventProducer()) {}

void AudioScheduledSourceNode::start(double when) {
playbackState_ = PlaybackState::SCHEDULED;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,7 @@ namespace audioapi {
AudioEventHandlerRegistry::AudioEventHandlerRegistry(
jsi::Runtime *runtime,
const std::shared_ptr<react::CallInvoker> &callInvoker)
: callInvoker_(callInvoker),
runtime_(runtime),
dispatchQueue_(kDispatchCapacity),
audioProducerToken_(dispatchQueue_) {
: callInvoker_(callInvoker), runtime_(runtime), dispatchQueue_(kDispatchCapacity) {
// Dispatch worker: wait for an item, dequeue it, then hop to the JS thread.
workerThread_ = std::thread([this]() {
while (true) {
Expand Down Expand Up @@ -86,15 +83,20 @@ bool AudioEventHandlerRegistry::dispatchEvent(
return true;
}

std::shared_ptr<AudioEventProducer> AudioEventHandlerRegistry::createAudioEventProducer() {
return std::make_shared<AudioEventProducer>(dispatchQueue_);
}

bool AudioEventHandlerRegistry::dispatchEventFromAudioThread(
AudioEventProducer &producer,
AudioEvent eventName,
uint64_t listenerId,
AudioEventPayload &&payload) noexcept {
if (runtime_ == nullptr) {
return false;
}
if (!dispatchQueue_.try_enqueue(
audioProducerToken_,
producer.token(),
DispatchEvent{
.event = eventName, .listenerId = listenerId, .payload = std::move(payload)})) {
return false;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
#include <ReactCommon/CallInvoker.h>
#include <audioapi/events/AudioEvent.h>
#include <audioapi/events/AudioEventPayload.h>
#include <audioapi/events/AudioEventProducer.h>
#include <audioapi/events/IAudioEventHandlerRegistry.h>
#include <audioapi/libs/concurrentqueue/concurrentqueue.h>
#include <audioapi/libs/concurrentqueue/lightweightsemaphore.h>
Expand All @@ -27,8 +28,8 @@ using namespace facebook;
/// Both entry points enqueue a DispatchEvent and signal itemsAvailable_:
///
/// - dispatchEventFromAudioThread() — real-time audio thread (e.g. `processNode()`).
/// Uses a pre-created moodycamel ProducerToken so try_enqueue() is wait-free and never
/// allocates. Drops the event if the queue is full.
/// Enqueues through the caller's own AudioEventProducer so try_enqueue() is wait-free and
/// never allocates. Drops the event if the queue is full.
///
/// - dispatchEvent() — any other (non-RT) thread (worker, platform/JNI callbacks, recorder
/// cleanup, AudioAPIModule, etc.). Uses the implicit (multi-producer) enqueue path.
Expand Down Expand Up @@ -57,9 +58,12 @@ class AudioEventHandlerRegistry : public IAudioEventHandlerRegistry,
uint64_t listenerId,
AudioEventPayload &&payload) noexcept override;

std::shared_ptr<AudioEventProducer> createAudioEventProducer() override;

/// @brief Enqueue an event from the real-time audio thread.
/// Wait-free and allocation-free on the calling thread; drops when the queue is full.
bool dispatchEventFromAudioThread(
AudioEventProducer &producer,
AudioEvent eventName,
uint64_t listenerId,
AudioEventPayload &&payload) noexcept override;
Expand All @@ -84,12 +88,9 @@ class AudioEventHandlerRegistry : public IAudioEventHandlerRegistry,
std::unordered_map<AudioEvent, std::unordered_map<uint64_t, std::shared_ptr<jsi::Function>>>
eventHandlers_;

// Single producer-to-consumer channel for every thread. Declared before
// audioProducerToken_ so the token can bind to it during construction.
// Single producer-to-consumer channel for every thread. Audio threads bind their own
// AudioEventProducer to it; every other thread uses the implicit-producer path.
moodycamel::ConcurrentQueue<DispatchEvent> dispatchQueue_;
// Dedicated token for the audio thread; lets it enqueue without the implicit-producer
// lookup/allocation that the first enqueue from a new thread would otherwise trigger.
moodycamel::ProducerToken audioProducerToken_;
// Counts queued items; workerThread_ waits on it instead of busy-spinning.
moodycamel::LightweightSemaphore itemsAvailable_;
std::atomic<bool> isExiting_{false};
Expand Down
Comment thread
maciejmakowski2003 marked this conversation as resolved.
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
#pragma once

#include <audioapi/libs/concurrentqueue/concurrentqueue.h>
#include <audioapi/utils/Macros.h>

namespace audioapi {

/// @brief One audio thread's private lane into the event registry's dispatch queue.
///
/// A moodycamel ProducerToken is single-producer by construction: do not share tokens between threads.
///
/// every context owns a producer and hands it to
/// IAudioEventHandlerRegistry::dispatchEventFromAudioThread.
///
/// @note A producer may move between threads over time — an OfflineAudioContext spawns a
/// fresh render thread on every resume — but never concurrently
class AudioEventProducer {
public:
/// @param queue The registry's dispatch queue; the token binds to it for its whole life.
template <typename TQueue>
explicit AudioEventProducer(TQueue &queue) : token_(queue) {}

~AudioEventProducer() = default;

DELETE_COPY_AND_MOVE(AudioEventProducer);

[[nodiscard]] moodycamel::ProducerToken &token() noexcept {
return token_;
}

private:
moodycamel::ProducerToken token_;
};

} // namespace audioapi
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,12 @@ namespace audioapi {
/// `scheduleAudioEvent` path) — no lock, so no other thread may touch it.
class DeferredEventQueue {
public:
explicit DeferredEventQueue(std::shared_ptr<IAudioEventHandlerRegistry> registry)
: registry_(std::move(registry)) {}
/// @param audioEventProducer The owning context's dispatch lane — `dispatchDue` always runs
/// on that context's render thread.
DeferredEventQueue(
std::shared_ptr<IAudioEventHandlerRegistry> registry,
std::shared_ptr<AudioEventProducer> audioEventProducer)
: registry_(std::move(registry)), audioEventProducer_(std::move(audioEventProducer)) {}

/// @brief Queues @p event for dispatch once the clock reaches @p dueTime.
/// @return False when there is nothing to queue (@p callbackId unset) or no
Expand Down Expand Up @@ -74,16 +78,20 @@ class DeferredEventQueue {
};

void dispatch(const DeferredEvent &deferred) const {
if (registry_ == nullptr) {
if (registry_ == nullptr || audioEventProducer_ == nullptr) {
return;
}

registry_->dispatchEventFromAudioThread(
deferred.event, deferred.callbackId, AudioEventPayload{EmptyPayload{}});
*audioEventProducer_,
deferred.event,
deferred.callbackId,
AudioEventPayload{EmptyPayload{}});
}

BoundedPriorityQueue<DeferredEvent, MAX_PENDING_EVENTS, ByDueTime> pending_;
std::shared_ptr<IAudioEventHandlerRegistry> registry_;
std::shared_ptr<AudioEventProducer> audioEventProducer_;
Comment thread
maciejmakowski2003 marked this conversation as resolved.
};

} // namespace audioapi
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,15 @@ class EventCaller {
explicit EventCaller(const std::shared_ptr<IAudioEventHandlerRegistry> &audioEventHandlerRegistry)
: eventHandlerRegistry_(audioEventHandlerRegistry) {}

/// @param audioEventProducer The owning context's dispatch lane. Required to dispatch from
/// the audio thread; without it `dispatchFromAudioThread` reports failure instead of
/// enqueueing through a lane that belongs to another thread.
EventCaller(
const std::shared_ptr<IAudioEventHandlerRegistry> &audioEventHandlerRegistry,
std::shared_ptr<AudioEventProducer> audioEventProducer)
: eventHandlerRegistry_(audioEventHandlerRegistry),
audioEventProducer_(std::move(audioEventProducer)) {}

~EventCaller() {
unregisterCallback();
}
Expand Down Expand Up @@ -91,16 +100,18 @@ class EventCaller {
requires EventPayloadFor<Event, Payload>
bool dispatchFromAudioThread(Payload &&payload) const noexcept {
const auto callbackId = getCallbackId();
if (eventHandlerRegistry_ == nullptr || callbackId == 0) {
if (eventHandlerRegistry_ == nullptr || audioEventProducer_ == nullptr || callbackId == 0) {
return false;
}

return eventHandlerRegistry_->dispatchEventFromAudioThread(
Event, callbackId, AudioEventPayload{std::forward<Payload>(payload)});
*audioEventProducer_, Event, callbackId, AudioEventPayload{std::forward<Payload>(payload)});
}

private:
std::shared_ptr<IAudioEventHandlerRegistry> eventHandlerRegistry_;
/// Null for events only ever dispatched NOT on the audio thread
std::shared_ptr<AudioEventProducer> audioEventProducer_;
std::atomic<uint64_t> callbackId_{0};
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@

namespace audioapi {

class AudioEventProducer;

class IAudioEventHandlerRegistry {
public:
IAudioEventHandlerRegistry() = default;
Expand All @@ -31,7 +33,14 @@ class IAudioEventHandlerRegistry {
uint64_t listenerId,
AudioEventPayload &&payload) noexcept = 0;

/// @brief Creates a dispatch lane for one audio thread. Every owner of an audio
/// thread needs its own — see AudioEventProducer for why sharing one corrupts the queue.
virtual std::shared_ptr<AudioEventProducer> createAudioEventProducer() = 0;

/// @param producer The calling audio thread's own producer, never one shared with
/// another thread.
virtual bool dispatchEventFromAudioThread(
AudioEventProducer &producer,
AudioEvent eventName,
uint64_t listenerId,
AudioEventPayload &&payload) noexcept = 0;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,13 +1,15 @@
#include <audioapi/utils/events/BufferingStateDispatcher.h>

#include <memory>
#include <utility>

namespace audioapi {

BufferingStateDispatcher::BufferingStateDispatcher(
const std::shared_ptr<IAudioEventHandlerRegistry> &audioEventHandlerRegistry,
std::shared_ptr<AudioEventProducer> audioEventProducer,
int startThresholdFrames)
: bufferingStateChangeEvent_(audioEventHandlerRegistry),
: bufferingStateChangeEvent_(audioEventHandlerRegistry, std::move(audioEventProducer)),
startThresholdFrames_(startThresholdFrames) {}

void BufferingStateDispatcher::assignCallbackId(uint64_t callbackId) noexcept {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ class BufferingStateDispatcher {
public:
BufferingStateDispatcher(
const std::shared_ptr<IAudioEventHandlerRegistry> &audioEventHandlerRegistry,
std::shared_ptr<AudioEventProducer> audioEventProducer,
int startThresholdFrames);

void assignCallbackId(uint64_t callbackId) noexcept;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,16 @@

#include <algorithm>
#include <memory>
#include <utility>

namespace audioapi {

PositionChangedDispatcher::PositionChangedDispatcher(
const std::shared_ptr<IAudioEventHandlerRegistry> &audioEventHandlerRegistry,
std::shared_ptr<AudioEventProducer> audioEventProducer,
int intervalInFrames,
bool shouldFlush)
: positionChangedEvent_(audioEventHandlerRegistry),
: positionChangedEvent_(audioEventHandlerRegistry, std::move(audioEventProducer)),
shouldFlush_(shouldFlush),
intervalInFrames_(intervalInFrames) {}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ class PositionChangedDispatcher {
public:
PositionChangedDispatcher(
const std::shared_ptr<IAudioEventHandlerRegistry> &audioEventHandlerRegistry,
std::shared_ptr<AudioEventProducer> audioEventProducer,
int intervalInFrames,
bool shouldFlush = false);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

#include <audioapi/events/AudioEvent.h>
#include <audioapi/events/AudioEventPayload.h>
#include <audioapi/events/AudioEventProducer.h>
#include <audioapi/events/IAudioEventHandlerRegistry.h>
#include <gmock/gmock.h>
#include <memory>
Expand All @@ -24,6 +25,19 @@ class MockAudioEventHandlerRegistry : public IAudioEventHandlerRegistry {
MOCK_METHOD(
bool,
dispatchEventFromAudioThread,
(AudioEvent eventName, uint64_t listenerId, AudioEventPayload &&payload),
(AudioEventProducer & producer,
AudioEvent eventName,
uint64_t listenerId,
AudioEventPayload &&payload),
(noexcept, override));

/// Real producers, not mocked: the expectations above never touch the token, but callers
/// refuse to dispatch from the audio thread without one.
std::shared_ptr<AudioEventProducer> createAudioEventProducer() override {
return std::make_shared<AudioEventProducer>(producerQueue_);
}

private:
/// Only ever a binding target for the producers handed out above; nothing is enqueued.
moodycamel::ConcurrentQueue<int> producerQueue_;
};
Loading
Loading