diff --git a/Libraries/LibMedia/Audio/PlaybackStreamPulseAudio.cpp b/Libraries/LibMedia/Audio/PlaybackStreamPulseAudio.cpp index 48d2002d74..a3ef3e68f6 100644 --- a/Libraries/LibMedia/Audio/PlaybackStreamPulseAudio.cpp +++ b/Libraries/LibMedia/Audio/PlaybackStreamPulseAudio.cpp @@ -12,21 +12,18 @@ namespace Audio { -#define TRY_OR_REJECT_AND_EXIT(expression) \ - ({ \ - auto&& __temporary_result = (expression); \ - if (__temporary_result.is_error()) [[unlikely]] { \ - warnln("Failure in PulseAudio control thread: {}", __temporary_result.error().string_literal()); \ - auto event_loop = main_thread_event_loop->take(); \ - if (!event_loop) \ - return 1; \ - event_loop->deferred_invoke([promise = move(promise), error = __temporary_result.release_error()] mutable { \ - promise->reject(move(error)); \ - }); \ - internal_state->exit(); \ - return 1; \ - } \ - __temporary_result.release_value(); \ +#define TRY_OR_REJECT_AND_EXIT(expression) \ + ({ \ + auto&& __temporary_result = (expression); \ + if (__temporary_result.is_error()) [[unlikely]] { \ + warnln("Failure in PulseAudio control thread: {}", __temporary_result.error().string_literal()); \ + main_thread_event_loop.deferred_invoke([promise = move(promise), error = __temporary_result.release_error()] mutable { \ + promise->reject(move(error)); \ + }); \ + internal_state->exit(); \ + return 1; \ + } \ + __temporary_result.release_value(); \ }) NonnullRefPtr PlaybackStream::create(OutputState initial_output_state, u32 target_latency_ms, AudioDataRequestCallback&& data_request_callback) @@ -43,9 +40,10 @@ NonnullRefPtr PlaybackStreamPulseAudio::create(Ou // Create an internal state for the control thread to hold on to. auto internal_state = MUST(adopt_nonnull_ref_or_enomem(new (nothrow) InternalState())); auto playback_stream = MUST(adopt_nonnull_ref_or_enomem(new (nothrow) PlaybackStreamPulseAudio(internal_state))); + auto& main_thread_event_loop = Core::EventLoop::current(); // Create the control thread and start it. - auto thread = MUST(Threading::Thread::try_create("Audio Control"sv, [=, main_thread_event_loop = Core::EventLoop::current_weak(), data_request_callback = move(data_request_callback)]() mutable { + auto thread = MUST(Threading::Thread::try_create("Audio Control"sv, [=, &main_thread_event_loop, data_request_callback = move(data_request_callback)]() mutable { auto context = TRY_OR_REJECT_AND_EXIT(PulseAudioContext::the()); internal_state->set_stream(TRY_OR_REJECT_AND_EXIT(context->create_stream(initial_state, target_latency_ms, [data_request_callback = move(data_request_callback)](PulseAudioStream&, Span buffer) { return data_request_callback(buffer); @@ -56,10 +54,7 @@ NonnullRefPtr PlaybackStreamPulseAudio::create(Ou TRY_OR_REJECT_AND_EXIT(internal_state->stream()->set_volume(1.0)); { - auto event_loop = main_thread_event_loop->take(); - if (!event_loop) - return 1; - event_loop->deferred_invoke([promise = move(promise), playback_stream = move(playback_stream)] mutable { + main_thread_event_loop.deferred_invoke([promise = move(promise), playback_stream = move(playback_stream)] mutable { promise->resolve(move(playback_stream)); }); } diff --git a/Libraries/LibMedia/PlaybackManager.cpp b/Libraries/LibMedia/PlaybackManager.cpp index 4367518d62..2399c8507b 100644 --- a/Libraries/LibMedia/PlaybackManager.cpp +++ b/Libraries/LibMedia/PlaybackManager.cpp @@ -29,7 +29,7 @@ DecoderErrorOr> PlaybackManager::create_demuxer_for_strea return FFmpeg::FFmpegDemuxer::from_stream(stream); } -DecoderErrorOr PlaybackManager::prepare_playback_from_demuxer(WeakPlaybackManager const& self, NonnullRefPtr const& demuxer, NonnullRefPtr const& main_thread_event_loop_reference) +DecoderErrorOr PlaybackManager::prepare_playback_from_demuxer(WeakPlaybackManager const& self, NonnullRefPtr const& demuxer, Core::EventLoop& main_thread_event_loop) { // Create the video tracks and their producers. auto all_video_tracks = TRY(demuxer->get_tracks_for_type(TrackType::Video)); @@ -39,7 +39,7 @@ DecoderErrorOr PlaybackManager::prepare_playback_from_demuxer(WeakPlayback supported_video_tracks.ensure_capacity(all_video_tracks.size()); supported_video_track_datas.ensure_capacity(all_video_tracks.size()); for (auto const& track : all_video_tracks) { - auto video_producer_result = DecodedVideoProducer::try_create(main_thread_event_loop_reference, demuxer, track); + auto video_producer_result = DecodedVideoProducer::try_create(main_thread_event_loop, demuxer, track); if (video_producer_result.is_error()) continue; supported_video_tracks.append(track); @@ -56,7 +56,7 @@ DecoderErrorOr PlaybackManager::prepare_playback_from_demuxer(WeakPlayback supported_audio_tracks.ensure_capacity(all_audio_tracks.size()); supported_audio_track_datas.ensure_capacity(all_audio_tracks.size()); for (auto const& track : all_audio_tracks) { - auto audio_producer_result = DecodedAudioProducer::try_create(main_thread_event_loop_reference, demuxer, track); + auto audio_producer_result = DecodedAudioProducer::try_create(main_thread_event_loop, demuxer, track); if (audio_producer_result.is_error()) continue; auto audio_producer = audio_producer_result.release_value(); @@ -79,8 +79,7 @@ DecoderErrorOr PlaybackManager::prepare_playback_from_demuxer(WeakPlayback auto duration = demuxer->total_duration().value_or(AK::Duration::zero()); auto start_time_realtime = demuxer->start_time_realtime(); - auto main_thread_event_loop = main_thread_event_loop_reference->take(); - main_thread_event_loop->deferred_invoke([self, video_tracks = move(supported_video_tracks), video_track_datas = move(supported_video_track_datas), preferred_video_track, audio_tracks = move(supported_audio_tracks), audio_track_datas = move(supported_audio_track_datas), preferred_audio_track, duration, start_time_realtime] mutable { + main_thread_event_loop.deferred_invoke([self, video_tracks = move(supported_video_tracks), video_track_datas = move(supported_video_track_datas), preferred_video_track, audio_tracks = move(supported_audio_tracks), audio_track_datas = move(supported_audio_track_datas), preferred_audio_track, duration, start_time_realtime] mutable { if (!self) return; @@ -167,12 +166,9 @@ PlaybackManager::~PlaybackManager() m_weak_link->revoke({}); } -static void handle_media_init_error(WeakPlaybackManager self, NonnullRefPtr const& main_thread_event_loop_reference, DecoderError error) +static void handle_media_init_error(WeakPlaybackManager self, Core::EventLoop& main_thread_event_loop, DecoderError error) { - auto main_thread_event_loop = main_thread_event_loop_reference->take(); - if (!main_thread_event_loop) - return; - main_thread_event_loop->deferred_invoke([self = move(self), error = move(error)] mutable { + main_thread_event_loop.deferred_invoke([self = move(self), error = move(error)] mutable { if (!self) return; if (self->on_unsupported_format_error) @@ -183,30 +179,30 @@ static void handle_media_init_error(WeakPlaybackManager self, NonnullRefPtr const& stream) { auto self = weak(); - auto main_thread_event_loop_reference = Core::EventLoop::current_weak(); + auto& main_thread_event_loop = Core::EventLoop::current(); - Threading::ThreadPool::the().submit([self = move(self), stream, main_thread_event_loop_reference = move(main_thread_event_loop_reference)] mutable { + Threading::ThreadPool::the().submit([self = move(self), stream, &main_thread_event_loop] mutable { auto demuxer_or_error = create_demuxer_for_stream(stream); if (demuxer_or_error.is_error()) { - handle_media_init_error(move(self), move(main_thread_event_loop_reference), demuxer_or_error.release_error()); + handle_media_init_error(move(self), main_thread_event_loop, demuxer_or_error.release_error()); return; } - auto maybe_error = prepare_playback_from_demuxer(self, demuxer_or_error.release_value(), main_thread_event_loop_reference); + auto maybe_error = prepare_playback_from_demuxer(self, demuxer_or_error.release_value(), main_thread_event_loop); if (maybe_error.is_error()) - handle_media_init_error(move(self), move(main_thread_event_loop_reference), maybe_error.release_error()); + handle_media_init_error(move(self), main_thread_event_loop, maybe_error.release_error()); }); } void PlaybackManager::add_media_source(NonnullRefPtr const& demuxer) { auto self = weak(); - auto main_thread_event_loop_reference = Core::EventLoop::current_weak(); + auto& main_thread_event_loop = Core::EventLoop::current(); - Threading::ThreadPool::the().submit([self = move(self), demuxer, main_thread_event_loop_reference = move(main_thread_event_loop_reference)] mutable { - auto maybe_error = prepare_playback_from_demuxer(self, demuxer, main_thread_event_loop_reference); + Threading::ThreadPool::the().submit([self = move(self), demuxer, &main_thread_event_loop] mutable { + auto maybe_error = prepare_playback_from_demuxer(self, demuxer, main_thread_event_loop); if (maybe_error.is_error()) - handle_media_init_error(move(self), move(main_thread_event_loop_reference), maybe_error.release_error()); + handle_media_init_error(move(self), main_thread_event_loop, maybe_error.release_error()); }); } diff --git a/Libraries/LibMedia/PlaybackManager.h b/Libraries/LibMedia/PlaybackManager.h index 0c39d178d6..929d3ff770 100644 --- a/Libraries/LibMedia/PlaybackManager.h +++ b/Libraries/LibMedia/PlaybackManager.h @@ -151,7 +151,7 @@ private: } static DecoderErrorOr> create_demuxer_for_stream(NonnullRefPtr const&); - static DecoderErrorOr prepare_playback_from_demuxer(WeakPlaybackManager const&, NonnullRefPtr const&, NonnullRefPtr const&); + static DecoderErrorOr prepare_playback_from_demuxer(WeakPlaybackManager const&, NonnullRefPtr const&, Core::EventLoop&); template void replace_state_handler(Args&&... args); diff --git a/Libraries/LibMedia/Producers/DecodedAudioProducer.cpp b/Libraries/LibMedia/Producers/DecodedAudioProducer.cpp index 77a5afde42..b7730936f2 100644 --- a/Libraries/LibMedia/Producers/DecodedAudioProducer.cpp +++ b/Libraries/LibMedia/Producers/DecodedAudioProducer.cpp @@ -20,7 +20,7 @@ namespace Media { static constexpr int AUTO_SUSPEND_IDLE_TIMEOUT_MS = 10000; -DecoderErrorOr> DecodedAudioProducer::try_create(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track) +DecoderErrorOr> DecodedAudioProducer::try_create(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track) { auto converter = DECODER_TRY_ALLOC(FFmpeg::FFmpegAudioConverter::try_create()); @@ -96,7 +96,7 @@ TimeRanges DecodedAudioProducer::buffered_time_ranges() const return m_thread_data->buffered_time_ranges(); } -DecodedAudioProducer::ThreadData::ThreadData(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track, AK::Duration duration, NonnullOwnPtr&& converter) +DecodedAudioProducer::ThreadData::ThreadData(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track, AK::Duration duration, NonnullOwnPtr&& converter) : m_main_thread_event_loop(main_thread_event_loop) , m_demuxer(demuxer) , m_track(track) @@ -336,10 +336,7 @@ void DecodedAudioProducer::ThreadData::invoke_on_main_thread_while_locked(Invoke { if (m_requested_state == RequestedState::Exit) return; - auto event_loop = m_main_thread_event_loop->take(); - if (!event_loop.is_alive()) - return; - event_loop->deferred_invoke([self = NonnullRefPtr(*this), invokee = move(invokee)] mutable { + m_main_thread_event_loop.deferred_invoke([self = NonnullRefPtr(*this), invokee = move(invokee)] mutable { invokee(self); }); } diff --git a/Libraries/LibMedia/Producers/DecodedAudioProducer.h b/Libraries/LibMedia/Producers/DecodedAudioProducer.h index 5a4707f720..ff16c1373e 100644 --- a/Libraries/LibMedia/Producers/DecodedAudioProducer.h +++ b/Libraries/LibMedia/Producers/DecodedAudioProducer.h @@ -38,7 +38,7 @@ public: using ErrorHandler = Function; using BlockEndTimeHandler = Function; - static DecoderErrorOr> try_create(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track); + static DecoderErrorOr> try_create(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track); DecodedAudioProducer(NonnullRefPtr const&); ~DecodedAudioProducer(); @@ -59,7 +59,7 @@ public: private: class ThreadData final : public AtomicRefCounted { public: - ThreadData(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const&, Track const&, AK::Duration, NonnullOwnPtr&&); + ThreadData(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const&, Track const&, AK::Duration, NonnullOwnPtr&&); ~ThreadData(); void set_error_handler(ErrorHandler&&); @@ -115,7 +115,7 @@ private: void note_consumer_activity_while_locked() const; void wait_for_queue_space_or_auto_suspend_while_locked(); - NonnullRefPtr m_main_thread_event_loop; + Core::EventLoop& m_main_thread_event_loop; mutable Sync::Mutex m_mutex; mutable Sync::ConditionVariable m_wait_condition { m_mutex }; diff --git a/Libraries/LibMedia/Producers/DecodedVideoProducer.cpp b/Libraries/LibMedia/Producers/DecodedVideoProducer.cpp index e92ecca815..26d08b8488 100644 --- a/Libraries/LibMedia/Producers/DecodedVideoProducer.cpp +++ b/Libraries/LibMedia/Producers/DecodedVideoProducer.cpp @@ -18,7 +18,7 @@ namespace Media { static constexpr int AUTO_SUSPEND_IDLE_TIMEOUT_MS = 10000; -DecoderErrorOr> DecodedVideoProducer::try_create(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track) +DecoderErrorOr> DecodedVideoProducer::try_create(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track) { TRY(demuxer->create_context_for_track(track)); auto duration = TRY(demuxer->duration_of_track(track)); @@ -169,7 +169,7 @@ void DecodedVideoProducer::seek(AK::Duration timestamp) m_thread_data->seek(timestamp); } -DecodedVideoProducer::ThreadData::ThreadData(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track, AK::Duration duration) +DecodedVideoProducer::ThreadData::ThreadData(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const& demuxer, Track const& track, AK::Duration duration) : m_main_thread_event_loop(main_thread_event_loop) , m_demuxer(demuxer) , m_track(track) @@ -302,10 +302,7 @@ void DecodedVideoProducer::ThreadData::invoke_on_main_thread_while_locked(Invoke { if (m_requested_state == RequestedState::Exit) return; - auto event_loop = m_main_thread_event_loop->take(); - if (!event_loop.is_alive()) - return; - event_loop->deferred_invoke([self = NonnullRefPtr(*this), invokee = move(invokee)] mutable { + m_main_thread_event_loop.deferred_invoke([self = NonnullRefPtr(*this), invokee = move(invokee)] mutable { invokee(self); }); } diff --git a/Libraries/LibMedia/Producers/DecodedVideoProducer.h b/Libraries/LibMedia/Producers/DecodedVideoProducer.h index 4188fdde1c..76b54b8a18 100644 --- a/Libraries/LibMedia/Producers/DecodedVideoProducer.h +++ b/Libraries/LibMedia/Producers/DecodedVideoProducer.h @@ -38,7 +38,7 @@ public: using ErrorHandler = Function; using FrameEndTimeHandler = Function; - static DecoderErrorOr> try_create(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const&, Track const&); + static DecoderErrorOr> try_create(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const&, Track const&); DecodedVideoProducer(NonnullRefPtr const&); ~DecodedVideoProducer(); @@ -60,7 +60,7 @@ public: private: class ThreadData final : public AtomicRefCounted { public: - ThreadData(NonnullRefPtr const& main_thread_event_loop, NonnullRefPtr const&, Track const&, AK::Duration); + ThreadData(Core::EventLoop& main_thread_event_loop, NonnullRefPtr const&, Track const&, AK::Duration); ~ThreadData(); void set_error_handler(ErrorHandler&&); @@ -114,7 +114,7 @@ private: void note_consumer_activity_while_locked() const; void wait_for_queue_space_or_auto_suspend_while_locked(); - NonnullRefPtr m_main_thread_event_loop; + Core::EventLoop& m_main_thread_event_loop; mutable Sync::Mutex m_mutex; mutable Sync::ConditionVariable m_wait_condition { m_mutex }; diff --git a/Libraries/LibMedia/Sinks/AudioPlaybackSink.cpp b/Libraries/LibMedia/Sinks/AudioPlaybackSink.cpp index 3b9174a816..97ab37cf9f 100644 --- a/Libraries/LibMedia/Sinks/AudioPlaybackSink.cpp +++ b/Libraries/LibMedia/Sinks/AudioPlaybackSink.cpp @@ -26,7 +26,7 @@ static constexpr size_t OUTPUT_BLOCK_QUEUE_CAPACITY = 4; class AudioPlaybackSink::OutputThreadData : public AtomicRefCounted { public: OutputThreadData(PipelineStateChangeHandler on_state_changed) - : m_main_thread_event_loop(Core::EventLoop::current_weak()) + : m_main_thread_event_loop(Core::EventLoop::current()) , m_on_state_changed(move(on_state_changed)) { } @@ -36,7 +36,7 @@ public: RefPtr m_playback_stream; RefPtr m_input; - NonnullRefPtr m_main_thread_event_loop; + Core::EventLoop& m_main_thread_event_loop; mutable Sync::Mutex m_output_mutex; mutable Sync::ConditionVariable m_output_condition { m_output_mutex }; @@ -332,15 +332,13 @@ void AudioPlaybackSink::OutputThreadData::dispatch_state_if_changed(PipelineStat if (status == m_last_dispatched_status) return; m_last_dispatched_status = status; - if (auto event_loop = m_main_thread_event_loop->take(); event_loop.is_alive()) { - event_loop->deferred_invoke([self = NonnullRefPtr(*this), status, seek_id] { - Sync::MutexLocker locker { self->m_output_mutex }; - if (self->m_seek_id != seek_id) - return; - if (self->m_on_state_changed) - self->m_on_state_changed(status); - }); - } + m_main_thread_event_loop.deferred_invoke([self = NonnullRefPtr(*this), status, seek_id] { + Sync::MutexLocker locker { self->m_output_mutex }; + if (self->m_seek_id != seek_id) + return; + if (self->m_on_state_changed) + self->m_on_state_changed(status); + }); } AK::Duration AudioPlaybackSink::current_time() const diff --git a/Tests/LibMedia/TestDataProducers.cpp b/Tests/LibMedia/TestDataProducers.cpp index 25e9ede6e1..a89b7a6fe5 100644 --- a/Tests/LibMedia/TestDataProducers.cpp +++ b/Tests/LibMedia/TestDataProducers.cpp @@ -40,7 +40,7 @@ TEST_CASE(audio_producer_underspecified_5_1_channel_map) auto tracks = TRY_OR_FAIL(demuxer->get_tracks_for_type(Media::TrackType::Audio)); VERIFY(!tracks.is_empty()); - auto producer = TRY_OR_FAIL(Media::DecodedAudioProducer::try_create(Core::EventLoop::current_weak(), demuxer, tracks[0])); + auto producer = TRY_OR_FAIL(Media::DecodedAudioProducer::try_create(loop, demuxer, tracks[0])); producer->start(); diff --git a/Tests/LibMedia/TestFFmpegAudioNormalization.cpp b/Tests/LibMedia/TestFFmpegAudioNormalization.cpp index ee93ad6942..a6ee4f8400 100644 --- a/Tests/LibMedia/TestFFmpegAudioNormalization.cpp +++ b/Tests/LibMedia/TestFFmpegAudioNormalization.cpp @@ -79,7 +79,7 @@ static void decode_and_expect() auto tracks = TRY_OR_FAIL(demuxer->get_tracks_for_type(Media::TrackType::Audio)); VERIFY(!tracks.is_empty()); - auto producer = TRY_OR_FAIL(Media::DecodedAudioProducer::try_create(Core::EventLoop::current_weak(), demuxer, tracks[0])); + auto producer = TRY_OR_FAIL(Media::DecodedAudioProducer::try_create(loop, demuxer, tracks[0])); producer->set_error_handler([&](Media::DecoderError&&) { FAIL("An error occurred while decoding generated WAV data."); diff --git a/Tests/LibMedia/TestMediaCommon.h b/Tests/LibMedia/TestMediaCommon.h index a799c86545..34759cd62d 100644 --- a/Tests/LibMedia/TestMediaCommon.h +++ b/Tests/LibMedia/TestMediaCommon.h @@ -83,7 +83,7 @@ static inline void decode_audio(StringView path, u32 sample_rate, u8 channel_cou }()); auto tracks = TRY_OR_FAIL(demuxer->get_tracks_for_type(Media::TrackType::Audio)); VERIFY(!tracks.is_empty()); - auto producer = TRY_OR_FAIL(Media::DecodedAudioProducer::try_create(Core::EventLoop::current_weak(), demuxer, tracks[0])); + auto producer = TRY_OR_FAIL(Media::DecodedAudioProducer::try_create(loop, demuxer, tracks[0])); producer->set_error_handler([&](Media::DecoderError&&) { FAIL("An error occurred while decoding.");