2#include <ossia/detail/config.hpp>
4#if defined(OSSIA_ENABLE_PIPEWIRE)
5#if __has_include(<libremidi/backends/linux/pipewire/context.hpp>) \
6 && __has_include(<pipewire/filter.h>) \
7 && __has_include(<spa/param/latency-utils.h>)
8#define OSSIA_AUDIO_PIPEWIRE 1
10#include <ossia/audio/audio_engine.hpp>
11#include <ossia/audio/pipewire_quantum.hpp>
13#include <ossia/detail/pod_vector.hpp>
14#include <ossia/detail/thread.hpp>
16#include <libremidi/backends/linux/pipewire/context.hpp>
17#include <libremidi/backends/linux/pipewire/filter.hpp>
18#include <libremidi/backends/linux/pipewire/loader.hpp>
19#include <libremidi/backends/linux/pipewire/subscription.hpp>
20#include <libremidi/backends/linux/pipewire/types.hpp>
22#include <pipewire/filter.h>
23#include <pipewire/keys.h>
24#include <pipewire/properties.h>
25#include <spa/param/latency-utils.h>
26#include <spa/utils/result.h>
28#include <fmt/format.h>
49 std::vector<std::string> inputs;
50 std::vector<std::string> outputs;
75class pipewire_audio_protocol :
public audio_engine
82 std::shared_ptr<libremidi::pipewire::context> loop;
84 std::vector<pw_proxy*> links;
86 std::vector<port*> input_ports;
87 std::vector<port*> output_ports;
91 explicit pipewire_audio_protocol(
92 std::shared_ptr<libremidi::pipewire::context> ctx,
93 const audio_setup& setup)
94 : loop{std::move(ctx)}
96 if (!loop || !loop->ok())
99 auto& pw = libremidi::pipewire::load();
100 if (!pw.filter_available)
103 if (setup.buffer_size <= 0 || setup.rate <= 0)
104 throw std::runtime_error(
"PipeWire: invalid buffer size or sample rate");
109 this->effective_buffer_size = setup.buffer_size;
110 this->effective_sample_rate = setup.rate;
111 this->effective_inputs = setup.inputs.size();
112 this->effective_outputs = setup.outputs.size();
113 m_quantum.expected = setup.buffer_size;
114 m_rate.expected = setup.rate;
115 m_silence.assign(setup.buffer_size, 0.f);
116 m_scratch.assign(setup.buffer_size, 0.f);
117 m_cycle_in.resize(setup.inputs.size());
118 m_cycle_out.resize(setup.outputs.size());
119 m_chunk_in.resize(setup.inputs.size());
120 m_chunk_out.resize(setup.outputs.size());
121 m_buf_in.resize(setup.inputs.size());
122 m_buf_out.resize(setup.outputs.size());
125 static constexpr const struct pw_filter_events filter_events = {
126 .version = PW_VERSION_FILTER_EVENTS,
133 .process = &on_process,
135#if PW_VERSION_CORE > 3
140 std::string default_sink_name = loop->default_audio_sink_name();
142 "PipeWire filter: default sink name = '{}'", default_sink_name);
144 bool created =
false;
145 loop->with_lock([&] {
146 auto* filter_props = pw.properties_new(
147 PW_KEY_MEDIA_TYPE,
"Audio",
148 PW_KEY_MEDIA_CATEGORY,
"Duplex",
149 PW_KEY_MEDIA_ROLE,
"DSP",
150 PW_KEY_MEDIA_NAME, setup.name.c_str(),
151 PW_KEY_NODE_NAME, setup.name.c_str(),
152 PW_KEY_NODE_GROUP,
"group.dsp.0",
153 PW_KEY_NODE_DESCRIPTION,
"ossia score",
157 fmt::format(
"{}/{}", setup.buffer_size, setup.rate).c_str(),
158 PW_KEY_NODE_FORCE_QUANTUM,
159 fmt::format(
"{}", setup.buffer_size).c_str(),
160 PW_KEY_NODE_FORCE_RATE, fmt::format(
"{}", setup.rate).c_str(),
167 PW_KEY_NODE_LOCK_RATE,
"true",
168 PW_KEY_NODE_TRANSPORT_SYNC,
"true",
169 PW_KEY_NODE_ALWAYS_PROCESS,
"true",
170 PW_KEY_NODE_PAUSE_ON_IDLE,
"false",
171 PW_KEY_NODE_SUSPEND_ON_IDLE,
"false",
175 if (!default_sink_name.empty())
178 filter_props, PW_KEY_TARGET_OBJECT, default_sink_name.c_str());
181 this->filter = pw.filter_new_simple(
182 loop->bare_loop(), setup.name.c_str(), filter_props,
183 &filter_events,
this);
188 for (
const auto& name : setup.inputs)
190 auto* p =
static_cast<port*
>(pw.filter_add_port(
191 this->filter, PW_DIRECTION_INPUT,
192 PW_FILTER_PORT_FLAG_MAP_BUFFERS,
sizeof(
struct port),
194 PW_KEY_FORMAT_DSP,
"32 bit float mono audio",
195 PW_KEY_PORT_NAME, name.c_str(),
nullptr),
197 input_ports.push_back(p);
200 for (
const auto& name : setup.outputs)
202 auto* p =
static_cast<port*
>(pw.filter_add_port(
203 this->filter, PW_DIRECTION_OUTPUT,
204 PW_FILTER_PORT_FLAG_MAP_BUFFERS,
sizeof(
struct port),
206 PW_KEY_FORMAT_DSP,
"32 bit float mono audio",
207 PW_KEY_PORT_NAME, name.c_str(),
nullptr),
209 output_ports.push_back(p);
212 created = (pw.filter_connect(
213 this->filter, PW_FILTER_FLAG_RT_PROCESS,
nullptr, 0)
218 throw std::runtime_error(
"PipeWire: could not create filter instance");
222 loop->with_lock([&] {
223 pw.filter_destroy(this->filter);
224 this->filter =
nullptr;
226 throw std::runtime_error(
"PipeWire: cannot connect");
229 if (!loop->synchronize())
232 "PipeWire: synchronize() failed after filter_connect — engine inactive");
237 auto node_id = filter_node_id();
238 while (node_id == 0xFFFFFFFFu)
240 if (!loop->synchronize())
243 "PipeWire: synchronize() failed while waiting for node id");
246 node_id = filter_node_id();
252 const auto num_in = input_ports.size();
253 const auto num_out = output_ports.size();
254 bool have_ports =
false;
255 for (
int j = 0; j < 200; ++j)
257 auto snap = loop->snapshot();
258 if (
const auto* self = snap.find_by_id(node_id))
260 if (self->inputs.size() >= num_in
261 && self->outputs.size() >= num_out)
267 if (!loop->synchronize())
270 "PipeWire: synchronize() failed while waiting for ports");
277 "PipeWire: ports never appeared in graph — engine inactive");
287 std::uint32_t seen = 0;
288 for (
int j = 0; j < 100; ++j)
290 seen = m_observed_rate.load(std::memory_order_relaxed);
293 std::this_thread::sleep_for(std::chrono::milliseconds(10));
295 if (seen != 0 && seen !=
static_cast<std::uint32_t
>(setup.rate))
298 "PipeWire: the graph runs at {} Hz, not the requested {} Hz; "
299 "the engine will use {} Hz",
300 seen, setup.rate, seen);
301 this->effective_sample_rate = seen;
314 m_watchdog = std::thread{[
this] { watchdog_main(); }};
317 std::uint32_t filter_node_id() const noexcept
321 auto& pw = libremidi::pipewire::load();
322 if (!pw.filter_get_node_id)
324 return pw.filter_get_node_id(this->filter);
332 const auto our_node = filter_node_id();
333 if (our_node == 0xFFFFFFFFu)
336 const std::string default_sink = loop->default_audio_sink_name();
337 const std::string default_source = loop->default_audio_source_name();
339 "PipeWire autoconnect: defaults src='{}' sink='{}'",
340 default_source, default_sink);
342 std::vector<std::uint32_t> source_outputs;
343 std::vector<std::uint32_t> sink_inputs;
344 std::vector<std::uint32_t> self_in_ids, self_out_ids;
345 bool have_self =
false;
347 for (
int attempt = 0; attempt < 50; ++attempt)
349 auto snap = loop->snapshot();
352 self_out_ids.clear();
354 if (
const auto* self_node = snap.find_by_id(our_node))
357 for (
const auto& p : self_node->inputs)
358 self_in_ids.push_back(p.id);
359 for (
const auto& p : self_node->outputs)
360 self_out_ids.push_back(p.id);
363 source_outputs.clear();
365 if (!default_source.empty())
367 if (
const auto* n = snap.find_by_name(default_source))
369 for (
const auto& p : n->outputs)
370 source_outputs.push_back(p.id);
373 if (!default_sink.empty())
375 if (
const auto* n = snap.find_by_name(default_sink))
377 for (
const auto& p : n->inputs)
378 sink_inputs.push_back(p.id);
384 if (have_self && !source_outputs.empty() && !sink_inputs.empty())
386 if (!loop->synchronize())
394 "PipeWire autoconnect: src_ports={}, sink_ports={}, "
395 "self_in={}, self_out={}",
396 source_outputs.size(), sink_inputs.size(),
397 self_in_ids.size(), self_out_ids.size());
400 auto snap = loop->snapshot();
402 "PipeWire autoconnect: snapshot has {} nodes", snap.nodes.size());
403 for (
const auto& n : snap.nodes)
406 "PipeWire autoconnect: snap id={} name='{}' class='{}' "
407 "inputs={} outputs={}",
408 n.id, n.name, n.media_class_str, n.inputs.size(),
413 for (std::size_t i = 0;
414 i < self_in_ids.size() && i < source_outputs.size(); ++i)
416 if (
auto* link = libremidi::pipewire::link_ports(
417 *loop, source_outputs[i], self_in_ids[i]))
418 links.push_back(link);
420 for (std::size_t i = 0;
421 i < self_out_ids.size() && i < sink_inputs.size(); ++i)
423 if (
auto* link = libremidi::pipewire::link_ports(
424 *loop, self_out_ids[i], sink_inputs[i]))
425 links.push_back(link);
429 void wait(
int ms)
override
432 std::this_thread::sleep_for(std::chrono::milliseconds(ms));
435 bool running()
const override {
return loop && activated; }
442 audio_engine::stop();
446 m_watchdog_quit.store(
true, std::memory_order_release);
447 if (m_watchdog.joinable())
457 auto& pw = libremidi::pipewire::load();
463 loop->with_lock([&] {
464 if (
int res = pw.filter_disconnect(this->filter); res < 0)
467 "PipeWire: filter_disconnect failed: {}", spa_strerror(res));
470 (void)loop->synchronize();
473 for (
auto* link : this->links)
474 libremidi::pipewire::unlink_ports(*loop, link);
479 loop->with_lock([&] {
480 pw.filter_destroy(this->filter);
481 this->filter =
nullptr;
485 (void)loop->synchronize();
489 ~pipewire_audio_protocol()
override { stop(); }
503 static pw_buffer* dequeue_cycle_buffer(
504 const auto& pw,
void* port,
float*& data, std::uint32_t& safe)
noexcept
506 pw_buffer* b = pw.filter_dequeue_buffer(port);
507 if (!b || !b->buffer || b->buffer->n_datas < 1
508 || !b->buffer->datas[0].data)
513 auto& d = b->buffer->datas[0];
514 data =
static_cast<float*
>(d.data);
515 const std::uint32_t cap = d.maxsize /
sizeof(float);
521 static void requeue_cycle_buffer(
522 const auto& pw,
void* port, pw_buffer* b,
bool output,
523 std::uint32_t frames)
noexcept
527 if (output && b->buffer && b->buffer->n_datas >= 1)
529 if (
auto* chunk = b->buffer->datas[0].chunk)
532 chunk->size = frames *
sizeof(float);
533 chunk->stride =
sizeof(float);
537 pw.filter_queue_buffer(port, b);
540 static bool can_dequeue(
const auto& pw)
noexcept
542 return pw.filter_dequeue_buffer && pw.filter_queue_buffer;
547 std::uint32_t fetch_cycle_buffers(
const auto& pw, std::uint32_t nframes)
549 const auto inputs = input_ports.size();
550 const auto outputs = output_ports.size();
552 if (!can_dequeue(pw))
556 for (std::size_t i = 0; i < inputs; i++)
557 m_cycle_in[i] =
static_cast<float*
>(
558 pw.filter_get_dsp_buffer(input_ports[i], nframes));
559 for (std::size_t i = 0; i < outputs; i++)
560 m_cycle_out[i] =
static_cast<float*
>(
561 pw.filter_get_dsp_buffer(output_ports[i], nframes));
565 std::uint32_t safe = nframes;
566 for (std::size_t i = 0; i < inputs; i++)
567 m_buf_in[i] = dequeue_cycle_buffer(pw, input_ports[i], m_cycle_in[i], safe);
568 for (std::size_t i = 0; i < outputs; i++)
569 m_buf_out[i] = dequeue_cycle_buffer(pw, output_ports[i], m_cycle_out[i], safe);
571 for (std::size_t i = 0; i < inputs; i++)
572 requeue_cycle_buffer(pw, input_ports[i], m_buf_in[i],
false, safe);
573 for (std::size_t i = 0; i < outputs; i++)
574 requeue_cycle_buffer(pw, output_ports[i], m_buf_out[i],
true, safe);
579 clear_buffers(pipewire_audio_protocol& self, std::uint32_t nframes,
582 auto& pw = libremidi::pipewire::load();
583 if (!can_dequeue(pw))
585 for (std::size_t i = 0; i < outputs; i++)
587 auto* chan =
static_cast<float*
>(
588 pw.filter_get_dsp_buffer(self.output_ports[i], nframes));
590 for (std::size_t j = 0; j < nframes; j++)
596 for (std::size_t i = 0; i < outputs; i++)
599 std::uint32_t safe = nframes;
600 auto* b = dequeue_cycle_buffer(pw, self.output_ports[i], data, safe);
602 for (std::size_t j = 0; j < safe; j++)
604 requeue_cycle_buffer(pw, self.output_ports[i], b,
true, data ? safe : 0);
608 void do_process(std::uint32_t nframes,
double secs,
double rate)
610 auto& pw = libremidi::pipewire::load();
614 const auto inputs = input_ports.size();
615 const auto outputs = output_ports.size();
619 clear_buffers(*
this, nframes, outputs);
623 const std::uint32_t frames = fetch_cycle_buffers(pw, nframes);
625 bool missing_input =
false;
626 for (std::size_t i = 0; i < inputs; i++)
627 missing_input |= !m_cycle_in[i];
629 std::memset(m_silence.data(), 0, m_silence.size() *
sizeof(
float));
636 const auto block =
static_cast<std::uint32_t
>(effective_buffer_size);
637 ossia::pipewire::for_each_chunk(
638 frames, block, [&](std::uint32_t offset, std::uint32_t n) {
639 ossia::pipewire::assign_chunk_pointers(
640 m_cycle_in.data(), m_chunk_in.data(), inputs, offset,
642 ossia::pipewire::assign_chunk_pointers(
643 m_cycle_out.data(), m_chunk_out.data(), outputs, offset,
646 ossia::audio_tick_state ts{
647 m_chunk_in.data(), m_chunk_out.data(), (int)inputs, (
int)outputs,
648 n, secs + offset / rate};
654 static void on_process(
void* userdata,
struct spa_io_position* position)
656 [[maybe_unused]]
static const thread_local auto _ = [] {
657 ossia::set_thread_name(
"ossia audio 0");
658 ossia::set_thread_pinned(thread_type::Audio, 0);
662 if (!userdata || !position)
665 auto& self = *
static_cast<pipewire_audio_protocol*
>(userdata);
666 self.m_cycles.fetch_add(1, std::memory_order_relaxed);
667 const std::uint32_t nframes = position->clock.duration;
668 const std::uint32_t rate = position->clock.rate.denom;
669 const double current_time = position->clock.nsec * 1e-9;
678 using event = ossia::pipewire::quantum_tracker::event;
681 switch (self.m_quantum.observe(nframes))
683 case event::mismatch:
684 if (may_log(self.m_quantum_logs))
686 "PipeWire: graph quantum is {} but {} was requested; "
687 "adapting by processing in chunks",
688 nframes, self.effective_buffer_size);
690 case event::recovered:
691 if (may_log(self.m_quantum_logs))
693 "PipeWire: graph quantum restored to {}", nframes);
701 switch (self.m_rate.observe(rate))
703 case event::mismatch:
704 if (may_log(self.m_rate_logs))
706 "PipeWire: graph sample rate is {} but {} was requested; "
707 "audio will play at the wrong speed until it is restored",
708 rate, self.m_rate.expected);
710 case event::recovered:
711 if (may_log(self.m_rate_logs))
713 "PipeWire: graph sample rate restored to {}", rate);
721 self.m_observed_rate.store(rate, std::memory_order_relaxed);
729 nframes, current_time, rate != 0 ? rate : self.m_rate.expected);
735 std::atomic<int> stall_timeout_ms{3000};
736 std::atomic<std::uint32_t> stalls_detected{};
737 std::atomic<std::uint32_t> recover_attempts{};
738 static constexpr std::uint32_t max_recover_attempts = 5;
743 ossia::set_thread_name(
"ossia pw wdog");
744 auto& pw = libremidi::pipewire::load();
746 std::uint64_t last = m_cycles.load(std::memory_order_relaxed);
747 auto last_progress = std::chrono::steady_clock::now();
748 while (!m_watchdog_quit.load(std::memory_order_acquire))
750 std::this_thread::sleep_for(std::chrono::milliseconds(100));
751 const auto now = std::chrono::steady_clock::now();
752 const auto cur = m_cycles.load(std::memory_order_relaxed);
759 recover_attempts.store(0, std::memory_order_relaxed);
762 const auto stalled_ms
763 = std::chrono::duration_cast<std::chrono::milliseconds>(
766 if (stalled_ms < stall_timeout_ms.load(std::memory_order_relaxed))
769 stalls_detected.fetch_add(1, std::memory_order_relaxed);
771 = recover_attempts.fetch_add(1, std::memory_order_relaxed) + 1;
772 if (attempt > max_recover_attempts)
775 "PipeWire: still no process cycles after {} reconnections; "
776 "giving up. The daemon is likely not scheduling any audio "
777 "(no usable driver); restarting PipeWire may help",
778 max_recover_attempts);
783 "PipeWire: no process cycles for {} ms; re-exporting the node "
784 "(attempt {}/{}). This happens when the graph has no usable "
785 "driver or the node was left unscheduled",
786 stalled_ms, attempt, max_recover_attempts);
788 bool connect_failed =
false;
789 loop->with_lock([&] {
792 if (
int res = pw.filter_disconnect(this->filter); res < 0)
794 "PipeWire: watchdog filter_disconnect failed: {}",
796 if (
int res = pw.filter_connect(
797 this->filter, PW_FILTER_FLAG_RT_PROCESS,
nullptr, 0);
803 connect_failed =
true;
805 "PipeWire: watchdog filter_connect failed: {}; automatic "
806 "recovery is not possible, restart the audio engine",
812 (void)loop->synchronize();
815 last = m_cycles.load(std::memory_order_relaxed);
816 last_progress = std::chrono::steady_clock::now();
820 std::atomic<std::uint64_t> m_cycles{};
821 std::thread m_watchdog;
822 std::atomic_bool m_watchdog_quit{};
824 ossia::pipewire::quantum_tracker m_quantum{};
825 ossia::pipewire::quantum_tracker m_rate{};
829 std::atomic<std::uint32_t> m_observed_rate{};
833 std::uint32_t m_quantum_logs{};
834 std::uint32_t m_rate_logs{};
836 static bool may_log(std::uint32_t& n)
noexcept
842 "PipeWire: the graph configuration keeps changing; further "
843 "changes will not be logged");
851 ossia::pod_vector<float*> m_cycle_in, m_cycle_out;
852 ossia::pod_vector<float*> m_chunk_in, m_chunk_out;
853 ossia::pod_vector<pw_buffer*> m_buf_in, m_buf_out;
858 ossia::pod_vector<float> m_silence;
859 ossia::pod_vector<float> m_scratch;
spdlog::logger & logger() noexcept
Where the errors will be logged. Default is stderr.
Definition context.cpp:120