very, very simple WP with midi CC out
This commit is contained in:
@@ -0,0 +1,53 @@
|
||||
#pragma once
|
||||
#include <libremidi/config.hpp>
|
||||
|
||||
#include <cinttypes>
|
||||
#include <cstdint>
|
||||
#include <functional>
|
||||
#include <string>
|
||||
|
||||
extern "C" {
|
||||
struct pw_main_loop;
|
||||
struct pw_filter;
|
||||
struct spa_io_position;
|
||||
}
|
||||
|
||||
namespace libremidi
|
||||
{
|
||||
using pipewire_callback_function = std::function<void(spa_io_position*)>;
|
||||
struct pipewire_callback
|
||||
{
|
||||
int64_t token;
|
||||
pipewire_callback_function callback;
|
||||
};
|
||||
|
||||
struct pipewire_input_configuration
|
||||
{
|
||||
std::string client_name = "libremidi client";
|
||||
|
||||
pw_main_loop* context{};
|
||||
pw_filter* filter{};
|
||||
std::function<void(pipewire_callback)> set_process_func;
|
||||
std::function<void(int64_t)> clear_process_func;
|
||||
};
|
||||
|
||||
struct pipewire_output_configuration
|
||||
{
|
||||
std::string client_name = "libremidi client";
|
||||
|
||||
pw_main_loop* context{};
|
||||
pw_filter* filter{};
|
||||
std::function<void(pipewire_callback)> set_process_func;
|
||||
std::function<void(int64_t)> clear_process_func;
|
||||
|
||||
int64_t output_buffer_size{65536};
|
||||
};
|
||||
|
||||
struct pipewire_observer_configuration
|
||||
{
|
||||
std::string client_name = "libremidi client";
|
||||
|
||||
pw_main_loop* context{};
|
||||
};
|
||||
|
||||
}
|
||||
@@ -0,0 +1,600 @@
|
||||
#pragma once
|
||||
#include <libremidi/backends/linux/pipewire.hpp>
|
||||
#include <libremidi/detail/memory.hpp>
|
||||
|
||||
#include <pipewire/filter.h>
|
||||
#include <pipewire/pipewire.h>
|
||||
#include <spa/control/control.h>
|
||||
#include <spa/param/props.h>
|
||||
#include <spa/utils/defs.h>
|
||||
#include <spa/utils/result.h>
|
||||
|
||||
#include <algorithm>
|
||||
#include <array>
|
||||
#include <atomic>
|
||||
#include <functional>
|
||||
#include <iostream>
|
||||
#include <memory>
|
||||
#include <string>
|
||||
#include <unordered_map>
|
||||
#include <vector>
|
||||
|
||||
#include <spa/param/audio/format-utils.h>
|
||||
|
||||
#pragma GCC diagnostic push
|
||||
#pragma GCC diagnostic ignored "-Wmissing-field-initializers"
|
||||
namespace libremidi
|
||||
{
|
||||
template <typename K, typename V>
|
||||
using hash_map = std::unordered_map<K, V>;
|
||||
|
||||
struct pipewire_instance
|
||||
{
|
||||
const libpipewire& pw = libpipewire::instance();
|
||||
pipewire_instance()
|
||||
{
|
||||
/// Initialize the PipeWire main loop, context, etc.
|
||||
int argc = 0;
|
||||
char* argv[] = {NULL};
|
||||
char** aa = argv;
|
||||
pw.init(&argc, &aa);
|
||||
}
|
||||
|
||||
~pipewire_instance() { pw.deinit(); }
|
||||
};
|
||||
|
||||
struct pipewire_context
|
||||
{
|
||||
struct listened_port
|
||||
{
|
||||
uint32_t id{};
|
||||
pw_port* port{};
|
||||
std::unique_ptr<spa_hook> listener;
|
||||
};
|
||||
|
||||
struct port_info
|
||||
{
|
||||
uint32_t id{};
|
||||
|
||||
std::string format;
|
||||
std::string port_name;
|
||||
std::string port_alias;
|
||||
std::string object_path;
|
||||
std::string node_id;
|
||||
std::string port_id;
|
||||
|
||||
bool physical{};
|
||||
bool terminal{};
|
||||
bool monitor{};
|
||||
pw_direction direction{};
|
||||
};
|
||||
|
||||
struct node
|
||||
{
|
||||
std::vector<port_info> inputs;
|
||||
std::vector<port_info> outputs;
|
||||
};
|
||||
|
||||
struct graph
|
||||
{
|
||||
mutable std::mutex mtx;
|
||||
libremidi::hash_map<uint32_t, node> physical_audio;
|
||||
libremidi::hash_map<uint32_t, node> physical_midi;
|
||||
libremidi::hash_map<uint32_t, node> software_audio;
|
||||
libremidi::hash_map<uint32_t, node> software_midi;
|
||||
libremidi::hash_map<uint32_t, port_info> port_cache;
|
||||
|
||||
void for_each_port(auto func)
|
||||
{
|
||||
for (auto& map : {physical_audio, physical_midi, software_audio, software_midi})
|
||||
{
|
||||
for (auto& [id, node] : map)
|
||||
{
|
||||
for (auto& port : node.inputs)
|
||||
func(port);
|
||||
for (auto& port : node.outputs)
|
||||
func(port);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void remove_port(uint32_t id)
|
||||
{
|
||||
port_cache.erase(id);
|
||||
for (auto map : {&physical_audio, &physical_midi, &software_audio, &software_midi})
|
||||
{
|
||||
for (auto& [_, node] : *map)
|
||||
{
|
||||
std::erase_if(node.inputs, [id](const port_info& p) { return p.id == id; });
|
||||
std::erase_if(node.outputs, [id](const port_info& p) { return p.id == id; });
|
||||
}
|
||||
}
|
||||
}
|
||||
} current_graph;
|
||||
|
||||
explicit pipewire_context(pw_main_loop* inst)
|
||||
: main_loop{inst}
|
||||
, owns_main_loop{false}
|
||||
{
|
||||
assert(main_loop);
|
||||
|
||||
initialize();
|
||||
}
|
||||
|
||||
explicit pipewire_context(std::shared_ptr<pipewire_instance> inst)
|
||||
: global_instance{inst}
|
||||
, owns_main_loop{true}
|
||||
{
|
||||
this->main_loop = pw.main_loop_new(nullptr);
|
||||
if (!this->main_loop)
|
||||
{
|
||||
// libremidi::logger().error("PipeWire: main_loop_new failed!");
|
||||
return;
|
||||
}
|
||||
initialize();
|
||||
}
|
||||
|
||||
void initialize()
|
||||
{
|
||||
this->lp = pw.main_loop_get_loop(this->main_loop);
|
||||
if (!lp)
|
||||
{
|
||||
// libremidi::logger().error("PipeWire: main_loop_get_loop failed!");
|
||||
return;
|
||||
}
|
||||
|
||||
this->context = pw.context_new(lp, nullptr, 0);
|
||||
if (!this->context)
|
||||
{
|
||||
// libremidi::logger().error("PipeWire: context_new failed!");
|
||||
return;
|
||||
}
|
||||
|
||||
this->core = pw.context_connect(this->context, nullptr, 0);
|
||||
if (!this->core)
|
||||
{
|
||||
// libremidi::logger().error("PipeWire: context_connect failed!");
|
||||
return;
|
||||
}
|
||||
|
||||
this->registry = pw_core_get_registry(this->core, PW_VERSION_REGISTRY, 0);
|
||||
if (!this->registry)
|
||||
{
|
||||
// libremidi::logger().error("PipeWire: core_get_registry failed!");
|
||||
return;
|
||||
}
|
||||
|
||||
initialize_observation();
|
||||
|
||||
synchronize();
|
||||
|
||||
// Add a manual 1ms event loop iteration at the end of
|
||||
// ctor to ensure synchronous clients will still see the ports
|
||||
pw_loop_iterate(this->lp, 1);
|
||||
}
|
||||
|
||||
void initialize_observation()
|
||||
{
|
||||
// Register a listener which will listen on when ports are added / removed
|
||||
spa_zero(registry_listener);
|
||||
|
||||
static constexpr const struct pw_registry_events registry_events = {
|
||||
.version = PW_VERSION_REGISTRY_EVENTS,
|
||||
.global =
|
||||
[](void* object, uint32_t id, uint32_t /*permissions*/, const char* type,
|
||||
uint32_t /*version*/, const struct spa_dict* /*props*/) {
|
||||
pipewire_context& self = *(pipewire_context*)object;
|
||||
if (strcmp(type, PW_TYPE_INTERFACE_Port) == 0)
|
||||
self.register_port(id, type);
|
||||
},
|
||||
.global_remove =
|
||||
[](void* object, uint32_t id) {
|
||||
pipewire_context& self = *(pipewire_context*)object;
|
||||
self.unregister_port(id);
|
||||
},
|
||||
};
|
||||
|
||||
// Start listening
|
||||
pw_registry_add_listener(this->registry, &this->registry_listener, ®istry_events, this);
|
||||
}
|
||||
|
||||
void register_port(uint32_t id, const char* type)
|
||||
{
|
||||
auto port = (pw_port*)pw_registry_bind(registry, id, type, PW_VERSION_PORT, 0);
|
||||
port_listener.push_back({id, port, std::make_unique<spa_hook>()});
|
||||
auto& l = port_listener.back();
|
||||
|
||||
static constexpr const struct pw_port_events port_events = {
|
||||
.version = PW_VERSION_PORT_EVENTS,
|
||||
.info
|
||||
= [](void* object,
|
||||
const pw_port_info* info) { ((pipewire_context*)object)->update_port_info(info); },
|
||||
};
|
||||
pw_port_add_listener(l.port, l.listener.get(), &port_events, this);
|
||||
}
|
||||
|
||||
void unregister_port(uint32_t id)
|
||||
{
|
||||
// When a port is removed:
|
||||
// Notify
|
||||
std::unique_lock _{current_graph.mtx, std::defer_lock};
|
||||
if (on_port_removed)
|
||||
{
|
||||
_.lock();
|
||||
if (auto it = current_graph.port_cache.find(id); it != current_graph.port_cache.end())
|
||||
{
|
||||
auto copy = it->second;
|
||||
_.unlock();
|
||||
on_port_removed(copy);
|
||||
}
|
||||
else
|
||||
{
|
||||
_.unlock();
|
||||
}
|
||||
}
|
||||
|
||||
// Remove from the graph
|
||||
{
|
||||
_.lock();
|
||||
current_graph.remove_port(id);
|
||||
_.unlock();
|
||||
}
|
||||
|
||||
// Remove from the listeners
|
||||
auto it
|
||||
= std::find_if(port_listener.begin(), port_listener.end(), [&](const listened_port& l) {
|
||||
return l.id == id;
|
||||
});
|
||||
if (it != port_listener.end())
|
||||
{
|
||||
pw.proxy_destroy((pw_proxy*)it->port);
|
||||
port_listener.erase(it);
|
||||
}
|
||||
}
|
||||
|
||||
void synchronize()
|
||||
{
|
||||
pending = 0;
|
||||
done = 0;
|
||||
|
||||
if (!core)
|
||||
return;
|
||||
|
||||
spa_hook core_listener;
|
||||
|
||||
static constexpr struct pw_core_events core_events = {
|
||||
.version = PW_VERSION_CORE_EVENTS,
|
||||
.done =
|
||||
[](void* object, uint32_t id, int seq) {
|
||||
auto& self = *(pipewire_context*)object;
|
||||
if(id == PW_ID_CORE && seq == self.pending)
|
||||
{
|
||||
self.done = 1;
|
||||
libpipewire::instance().main_loop_quit(self.main_loop);
|
||||
}
|
||||
},
|
||||
};
|
||||
|
||||
spa_zero(core_listener);
|
||||
pw_core_add_listener(core, &core_listener, &core_events, this);
|
||||
|
||||
pending = pw_core_sync(core, PW_ID_CORE, 0);
|
||||
while (!done)
|
||||
{
|
||||
pw.main_loop_run(this->main_loop);
|
||||
}
|
||||
spa_hook_remove(&core_listener);
|
||||
}
|
||||
|
||||
[[nodiscard]] pw_proxy* link_ports(uint32_t out_port, uint32_t in_port)
|
||||
{
|
||||
auto props = pw.properties_new(
|
||||
PW_KEY_LINK_OUTPUT_PORT, std::to_string(out_port).c_str(), PW_KEY_LINK_INPUT_PORT,
|
||||
std::to_string(in_port).c_str(), nullptr);
|
||||
|
||||
auto proxy = (pw_proxy*)pw_core_create_object(
|
||||
this->core, "link-factory", PW_TYPE_INTERFACE_Link, PW_VERSION_LINK, &props->dict, 0);
|
||||
|
||||
if (!proxy)
|
||||
{
|
||||
std::cerr << "PipeWire: could not allocate link\n";
|
||||
pw.properties_free(props);
|
||||
return nullptr;
|
||||
}
|
||||
|
||||
synchronize();
|
||||
pw.properties_free(props);
|
||||
return proxy;
|
||||
}
|
||||
|
||||
void unlink_ports(pw_proxy* link) { pw.proxy_destroy(link); }
|
||||
|
||||
void update_port_info(const pw_port_info* info)
|
||||
{
|
||||
const spa_dict_item* item{};
|
||||
|
||||
port_info p;
|
||||
p.id = info->id;
|
||||
|
||||
spa_dict_for_each(item, info->props)
|
||||
{
|
||||
std::string_view k{item->key}, v{item->value};
|
||||
if (k == "format.dsp")
|
||||
p.format = v;
|
||||
else if (k == "port.name")
|
||||
p.port_name = v;
|
||||
else if (k == "port.alias")
|
||||
p.port_alias = v;
|
||||
else if (k == "object.path")
|
||||
p.object_path = v;
|
||||
else if (k == "port.id")
|
||||
p.port_id = v;
|
||||
else if (k == "node.id")
|
||||
p.node_id = v;
|
||||
else if (k == "port.physical" && v == "true")
|
||||
p.physical = true;
|
||||
else if (k == "port.terminal" && v == "true")
|
||||
p.terminal = true;
|
||||
else if (k == "port.monitor" && v == "true")
|
||||
p.monitor = true;
|
||||
else if (k == "port.direction")
|
||||
{
|
||||
if (v == "out")
|
||||
{
|
||||
p.direction = pw_direction::SPA_DIRECTION_OUTPUT;
|
||||
}
|
||||
else
|
||||
{
|
||||
p.direction = pw_direction::SPA_DIRECTION_INPUT;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (p.node_id.empty())
|
||||
return;
|
||||
|
||||
const auto nid = std::stoul(p.node_id);
|
||||
auto get_node = [&]() -> node* {
|
||||
if (p.physical)
|
||||
{
|
||||
if (p.format.find("audio") != p.format.npos)
|
||||
return &this->current_graph.physical_audio[nid];
|
||||
else if (p.format.find("midi") != p.format.npos)
|
||||
return &this->current_graph.physical_midi[nid];
|
||||
}
|
||||
else
|
||||
{
|
||||
if (p.format.find("audio") != p.format.npos)
|
||||
return &this->current_graph.software_audio[nid];
|
||||
else if (p.format.find("midi") != p.format.npos)
|
||||
return &this->current_graph.software_midi[nid];
|
||||
}
|
||||
return nullptr;
|
||||
};
|
||||
|
||||
{
|
||||
std::lock_guard _{current_graph.mtx};
|
||||
current_graph.port_cache[p.id] = p;
|
||||
if (auto node = get_node())
|
||||
{
|
||||
if (p.direction == pw_direction::SPA_DIRECTION_OUTPUT)
|
||||
node->outputs.push_back(p);
|
||||
else
|
||||
node->inputs.push_back(p);
|
||||
}
|
||||
}
|
||||
|
||||
if (on_port_added)
|
||||
on_port_added(p);
|
||||
}
|
||||
|
||||
int get_fd() const noexcept
|
||||
{
|
||||
if (!this->lp)
|
||||
return -1;
|
||||
|
||||
auto spa_callbacks = this->lp->control->iface.cb;
|
||||
auto spa_loop_methods = (const spa_loop_control_methods*)spa_callbacks.funcs;
|
||||
if (spa_loop_methods->get_fd)
|
||||
return spa_loop_methods->get_fd(spa_callbacks.data);
|
||||
else
|
||||
return -1;
|
||||
}
|
||||
|
||||
~pipewire_context()
|
||||
{
|
||||
if (this->registry)
|
||||
pw.proxy_destroy((pw_proxy*)this->registry);
|
||||
for (auto& [id, p, l] : this->port_listener)
|
||||
if (l)
|
||||
pw.proxy_destroy((pw_proxy*)p);
|
||||
if (this->core)
|
||||
pw.core_disconnect(this->core);
|
||||
if (this->context)
|
||||
pw.context_destroy(this->context);
|
||||
if (owns_main_loop && this->main_loop)
|
||||
pw.main_loop_destroy(this->main_loop);
|
||||
}
|
||||
|
||||
friend struct pipewire_filter;
|
||||
const libpipewire& pw = libpipewire::instance();
|
||||
std::shared_ptr<pipewire_instance> global_instance;
|
||||
|
||||
pw_main_loop* main_loop{};
|
||||
pw_loop* lp{};
|
||||
|
||||
pw_context* context{};
|
||||
pw_core* core{};
|
||||
|
||||
pw_registry* registry{};
|
||||
spa_hook registry_listener{};
|
||||
|
||||
std::function<void(const port_info&)> on_port_added;
|
||||
std::function<void(const port_info&)> on_port_removed;
|
||||
|
||||
std::vector<listened_port> port_listener{};
|
||||
|
||||
std::atomic<int> pending{};
|
||||
std::atomic<int> done{};
|
||||
bool owns_main_loop{true};
|
||||
int sync{};
|
||||
};
|
||||
|
||||
struct pipewire_filter
|
||||
{
|
||||
const libpipewire& pw = libpipewire::instance();
|
||||
std::shared_ptr<pipewire_context> loop{};
|
||||
pw_filter* filter{};
|
||||
std::vector<pw_proxy*> links{};
|
||||
|
||||
struct port
|
||||
{
|
||||
void* data;
|
||||
}* port{};
|
||||
|
||||
explicit pipewire_filter(std::shared_ptr<pipewire_context> loop)
|
||||
: loop{loop}
|
||||
{
|
||||
}
|
||||
|
||||
explicit pipewire_filter(std::shared_ptr<pipewire_context> loop, pw_filter* filter)
|
||||
: loop{loop}
|
||||
, filter{filter}
|
||||
{
|
||||
}
|
||||
|
||||
void create_filter(std::string_view filter_name, const pw_filter_events& events, void* context)
|
||||
{
|
||||
assert(!filter);
|
||||
|
||||
auto& pw = libpipewire::instance();
|
||||
// clang-format off
|
||||
this->filter = pw.filter_new_simple(
|
||||
loop->lp,
|
||||
filter_name.data(),
|
||||
pw.properties_new(
|
||||
PW_KEY_MEDIA_TYPE, "Midi",
|
||||
PW_KEY_MEDIA_CATEGORY, "Filter",
|
||||
PW_KEY_MEDIA_ROLE, "DSP",
|
||||
PW_KEY_MEDIA_NAME, "libremidi",
|
||||
#if defined(PW_KEY_NODE_LOCK_RATE)
|
||||
PW_KEY_NODE_LOCK_RATE, "true",
|
||||
#endif
|
||||
PW_KEY_NODE_ALWAYS_PROCESS, "true",
|
||||
PW_KEY_NODE_PAUSE_ON_IDLE, "false",
|
||||
#if defined(PW_KEY_NODE_SUSPEND_ON_IDLE)
|
||||
PW_KEY_NODE_SUSPEND_ON_IDLE, "false",
|
||||
#endif
|
||||
nullptr),
|
||||
&events,
|
||||
context);
|
||||
// clang-format on
|
||||
assert(filter);
|
||||
}
|
||||
|
||||
void destroy()
|
||||
{
|
||||
if (this->filter)
|
||||
pw.filter_destroy(this->filter);
|
||||
}
|
||||
|
||||
void create_local_port(std::string_view port_name, spa_direction direction)
|
||||
{
|
||||
// clang-format off
|
||||
this->port = (struct port*)pw.filter_add_port(
|
||||
this->filter,
|
||||
direction,
|
||||
PW_FILTER_PORT_FLAG_MAP_BUFFERS,
|
||||
sizeof(struct port),
|
||||
pw.properties_new(
|
||||
PW_KEY_FORMAT_DSP, "8 bit raw midi",
|
||||
PW_KEY_PORT_NAME, port_name.data(),
|
||||
nullptr),
|
||||
nullptr, 0);
|
||||
// clang-format on
|
||||
assert(port);
|
||||
}
|
||||
|
||||
void set_port_buffer(int bytes)
|
||||
{
|
||||
uint8_t buffer[1024];
|
||||
struct spa_pod_builder builder;
|
||||
spa_pod_builder_init(&builder, buffer, sizeof(buffer));
|
||||
|
||||
// clang-format off
|
||||
const struct spa_pod* params[1] = {
|
||||
(spa_pod*) spa_pod_builder_add_object(
|
||||
&builder,
|
||||
SPA_TYPE_OBJECT_ParamBuffers, SPA_PARAM_Buffers,
|
||||
SPA_PARAM_BUFFERS_buffers, SPA_POD_CHOICE_RANGE_Int(1, 1, 32),
|
||||
SPA_PARAM_BUFFERS_blocks, SPA_POD_Int(1),
|
||||
SPA_PARAM_BUFFERS_size, SPA_POD_CHOICE_RANGE_Int(bytes, 4096, INT32_MAX),
|
||||
SPA_PARAM_BUFFERS_stride, SPA_POD_Int(1)
|
||||
)
|
||||
};
|
||||
// clang-format on
|
||||
|
||||
pw.filter_update_params(this->filter, this->port, params, 1);
|
||||
}
|
||||
|
||||
void remove_port()
|
||||
{
|
||||
assert(this->port);
|
||||
pw.filter_remove_port(this->port);
|
||||
this->port = nullptr;
|
||||
}
|
||||
|
||||
void rename_port(std::string_view port_name)
|
||||
{
|
||||
assert(this->port);
|
||||
spa_dict_item items[1] = {
|
||||
SPA_DICT_ITEM_INIT(PW_KEY_PORT_NAME, port_name.data()),
|
||||
};
|
||||
|
||||
auto properties = SPA_DICT_INIT(items, 1);
|
||||
pw.filter_update_properties(this->filter, this->port, &properties);
|
||||
this->port = nullptr;
|
||||
}
|
||||
|
||||
void start_filter()
|
||||
{
|
||||
if (pw.filter_connect(this->filter, PW_FILTER_FLAG_RT_PROCESS, NULL, 0) < 0)
|
||||
{
|
||||
std::cerr << "can't connect\n";
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
uint32_t filter_node_id() { return this->loop->pw.filter_get_node_id(this->filter); }
|
||||
|
||||
void synchronize_node()
|
||||
{
|
||||
this->loop->synchronize();
|
||||
int k = 0;
|
||||
auto node_id = filter_node_id();
|
||||
while (node_id == 4294967295)
|
||||
{
|
||||
this->loop->synchronize();
|
||||
node_id = filter_node_id();
|
||||
|
||||
if (k++; k > 100)
|
||||
return;
|
||||
}
|
||||
}
|
||||
void synchronize_ports(const pipewire_context::node& this_node)
|
||||
{
|
||||
// Leave some time to resolve the ports
|
||||
int k = 0;
|
||||
const auto num_local_ins = 1;
|
||||
const auto num_local_outs = 0;
|
||||
while (this_node.inputs.size() < num_local_ins || this_node.outputs.size() < num_local_outs)
|
||||
{
|
||||
this->loop->synchronize();
|
||||
if (k++; k > 100)
|
||||
return;
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
#pragma GCC diagnostic pop
|
||||
@@ -0,0 +1,464 @@
|
||||
#pragma once
|
||||
|
||||
#include <libremidi/backends/linux/helpers.hpp>
|
||||
#include <libremidi/backends/linux/pipewire.hpp>
|
||||
#include <libremidi/backends/pipewire/context.hpp>
|
||||
#include <libremidi/detail/memory.hpp>
|
||||
#include <libremidi/detail/midi_in.hpp>
|
||||
#include <libremidi/detail/semaphore.hpp>
|
||||
|
||||
#include <atomic>
|
||||
#include <semaphore>
|
||||
#include <stop_token>
|
||||
#include <thread>
|
||||
|
||||
namespace libremidi
|
||||
{
|
||||
struct pipewire_helpers
|
||||
{
|
||||
struct port
|
||||
{
|
||||
void* data{};
|
||||
};
|
||||
|
||||
// All pipewire operations have to happen in the same thread
|
||||
// - and pipewire checks that internally.
|
||||
std::jthread main_loop_thread;
|
||||
const libpipewire& pw = libpipewire::instance();
|
||||
std::shared_ptr<pipewire_instance> global_instance;
|
||||
std::shared_ptr<pipewire_context> global_context;
|
||||
std::unique_ptr<pipewire_filter> filter;
|
||||
pw_proxy* link{};
|
||||
|
||||
int64_t this_instance{};
|
||||
|
||||
eventfd_notifier termination_event{};
|
||||
pollfd fds[2]{};
|
||||
|
||||
semaphore_pair_lock thread_lock;
|
||||
std::shared_ptr<void> canary = std::make_shared<int>();
|
||||
|
||||
enum poll_state
|
||||
{
|
||||
start_poll,
|
||||
in_poll,
|
||||
not_in_poll
|
||||
};
|
||||
std::atomic<poll_state> current_state{not_in_poll};
|
||||
|
||||
pipewire_helpers()
|
||||
{
|
||||
static std::atomic_int64_t instance{};
|
||||
this_instance = ++instance;
|
||||
|
||||
fds[1] = termination_event;
|
||||
}
|
||||
|
||||
template <typename Self>
|
||||
void create_filter(Self& self)
|
||||
{
|
||||
if (this->filter)
|
||||
return;
|
||||
|
||||
auto& configuration = self.configuration;
|
||||
if (configuration.context && configuration.filter && configuration.set_process_func)
|
||||
{
|
||||
this->filter = std::make_unique<pipewire_filter>(this->global_context, configuration.filter);
|
||||
|
||||
pipewire_callback cbs{
|
||||
.token = this_instance,
|
||||
.callback = [&self, p = std::weak_ptr{canary}](spa_io_position* nf) -> void {
|
||||
if (auto pt = p.lock())
|
||||
self.process(nf);
|
||||
|
||||
self.thread_lock.check_client_released();
|
||||
}};
|
||||
configuration.set_process_func(cbs);
|
||||
}
|
||||
else
|
||||
{
|
||||
this->filter = std::make_unique<pipewire_filter>(this->global_context);
|
||||
#pragma GCC diagnostic push
|
||||
#pragma GCC diagnostic ignored "-Wmissing-field-initializers"
|
||||
static constexpr struct pw_filter_events filter_events
|
||||
= {.version = PW_VERSION_FILTER_EVENTS,
|
||||
.process = +[](void* _data, struct spa_io_position* position) -> void {
|
||||
// FIXME likely we need the thread_lock check here too
|
||||
Self& self = *static_cast<Self*>(_data);
|
||||
self.process(position);
|
||||
}};
|
||||
#pragma GCC diagnostic pop
|
||||
|
||||
this->filter->create_filter(self.configuration.client_name, filter_events, &self);
|
||||
this->filter->start_filter();
|
||||
}
|
||||
}
|
||||
|
||||
template <typename Self>
|
||||
void destroy_filter(Self& self)
|
||||
{
|
||||
assert(global_context);
|
||||
if (!global_context->owns_main_loop)
|
||||
{
|
||||
if (self.configuration.clear_process_func)
|
||||
{
|
||||
self.configuration.clear_process_func(this_instance);
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
if (this->filter)
|
||||
{
|
||||
this->filter->destroy();
|
||||
}
|
||||
}
|
||||
|
||||
this->filter.reset();
|
||||
}
|
||||
|
||||
template <typename Self>
|
||||
int create_context(Self& self)
|
||||
{
|
||||
if (this->global_context)
|
||||
return 0;
|
||||
|
||||
// Initialize PipeWire client
|
||||
auto& configuration = self.configuration;
|
||||
if (configuration.context)
|
||||
{
|
||||
this->global_context = std::make_shared<pipewire_context>(configuration.context);
|
||||
}
|
||||
else
|
||||
{
|
||||
this->global_instance = std::make_shared<pipewire_instance>();
|
||||
this->global_context = std::make_shared<pipewire_context>(this->global_instance);
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
|
||||
void destroy_context()
|
||||
{
|
||||
assert(this->global_context);
|
||||
this->global_context.reset();
|
||||
this->global_instance.reset();
|
||||
}
|
||||
|
||||
void run_poll_loop()
|
||||
try
|
||||
{
|
||||
// Note: called from a std::jthread.
|
||||
assert(this->global_context);
|
||||
if (int fd = this->global_context->get_fd(); fd != -1)
|
||||
{
|
||||
fds[0] = {.fd = fd, .events = POLLIN, .revents = 0};
|
||||
current_state = poll_state::in_poll;
|
||||
|
||||
for (;;)
|
||||
{
|
||||
if (int err = poll(fds, 2, -1); err < 0)
|
||||
{
|
||||
if (err == -EAGAIN)
|
||||
continue;
|
||||
else
|
||||
break;
|
||||
}
|
||||
|
||||
// Check pipewire fd:
|
||||
if (fds[0].revents & POLLIN)
|
||||
{
|
||||
if (auto lp = this->global_context->lp)
|
||||
{
|
||||
int result = pw_loop_iterate(lp, 0);
|
||||
if (result < 0)
|
||||
std::cerr << "pw_loop_iterate: " << spa_strerror(result) << "\n";
|
||||
}
|
||||
fds[0].revents = 0;
|
||||
}
|
||||
|
||||
// Check exit fd:
|
||||
if (fds[1].revents & POLLIN)
|
||||
{
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
current_state = poll_state::not_in_poll;
|
||||
}
|
||||
catch (...)
|
||||
{
|
||||
current_state = poll_state::not_in_poll;
|
||||
}
|
||||
|
||||
template <typename Self>
|
||||
bool create_local_port(Self& self, std::string_view portName, spa_direction direction)
|
||||
{
|
||||
assert(this->global_context);
|
||||
assert(this->filter);
|
||||
|
||||
if (portName.empty())
|
||||
portName = direction == SPA_DIRECTION_INPUT ? "i" : "o";
|
||||
|
||||
if (!this->filter->port)
|
||||
{
|
||||
this->filter->create_local_port(portName.data(), direction);
|
||||
}
|
||||
|
||||
if (!this->filter->port)
|
||||
{
|
||||
self.template error<driver_error>(self.configuration, "PipeWire: error creating port");
|
||||
return false;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
void add_callbacks(const observer_configuration& conf)
|
||||
{
|
||||
assert(global_context);
|
||||
global_context->on_port_added = [&conf](const pipewire_context::port_info& port) {
|
||||
if (port.format.find("midi") == std::string::npos)
|
||||
return;
|
||||
|
||||
bool unfiltered = conf.track_any;
|
||||
unfiltered |= (port.physical && conf.track_hardware);
|
||||
unfiltered |= (!port.physical && conf.track_virtual);
|
||||
if (unfiltered)
|
||||
{
|
||||
if (port.direction == SPA_DIRECTION_INPUT)
|
||||
{
|
||||
if (conf.output_added)
|
||||
conf.output_added(to_port_info<SPA_DIRECTION_INPUT>(port));
|
||||
}
|
||||
else
|
||||
{
|
||||
if (conf.input_added)
|
||||
conf.input_added(to_port_info<SPA_DIRECTION_OUTPUT>(port));
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
global_context->on_port_removed = [&conf](const pipewire_context::port_info& port) {
|
||||
if (port.format.find("midi") == std::string::npos)
|
||||
return;
|
||||
|
||||
bool unfiltered = conf.track_any;
|
||||
unfiltered |= (port.physical && conf.track_hardware);
|
||||
unfiltered |= (!port.physical && conf.track_virtual);
|
||||
if (unfiltered)
|
||||
{
|
||||
if (port.direction == SPA_DIRECTION_INPUT)
|
||||
{
|
||||
if (conf.output_removed)
|
||||
conf.output_removed(to_port_info<SPA_DIRECTION_INPUT>(port));
|
||||
}
|
||||
else
|
||||
{
|
||||
if (conf.input_removed)
|
||||
conf.input_removed(to_port_info<SPA_DIRECTION_OUTPUT>(port));
|
||||
}
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
void start_thread()
|
||||
{
|
||||
if (!this->global_context->owns_main_loop)
|
||||
return;
|
||||
|
||||
current_state = poll_state::start_poll;
|
||||
main_loop_thread = std::jthread{[this]() { run_poll_loop(); }};
|
||||
}
|
||||
|
||||
void stop_thread()
|
||||
{
|
||||
assert(this->global_context);
|
||||
if (!this->global_context->owns_main_loop)
|
||||
return;
|
||||
|
||||
if (main_loop_thread.joinable() || current_state != poll_state::not_in_poll)
|
||||
{
|
||||
termination_event.notify();
|
||||
main_loop_thread.request_stop();
|
||||
|
||||
termination_event.notify();
|
||||
for (int i = 0; i < 100; i++)
|
||||
{
|
||||
if (current_state == poll_state::not_in_poll)
|
||||
break;
|
||||
std::this_thread::sleep_for(std::chrono::milliseconds(10));
|
||||
termination_event.notify();
|
||||
}
|
||||
|
||||
if (main_loop_thread.joinable())
|
||||
main_loop_thread.join();
|
||||
}
|
||||
}
|
||||
|
||||
void do_close_port()
|
||||
{
|
||||
if (!this->filter)
|
||||
return;
|
||||
if (!this->filter->port)
|
||||
return;
|
||||
|
||||
if (!this->global_context->owns_main_loop)
|
||||
{
|
||||
this->canary.reset();
|
||||
this->thread_lock.prepare_release_client();
|
||||
}
|
||||
|
||||
unlink_ports();
|
||||
this->filter->remove_port();
|
||||
}
|
||||
|
||||
void rename_port(std::string_view port_name)
|
||||
{
|
||||
if (this->filter)
|
||||
this->filter->rename_port(port_name);
|
||||
}
|
||||
|
||||
void unlink_ports()
|
||||
{
|
||||
if (link)
|
||||
{
|
||||
this->global_context->unlink_ports(link);
|
||||
link = nullptr;
|
||||
}
|
||||
}
|
||||
|
||||
bool link_ports(auto& self, const input_port& in_port)
|
||||
{
|
||||
// Wait for the pipewire server to send us back our node's info
|
||||
for (int i = 0; i < 1000; i++)
|
||||
this->filter->synchronize_node();
|
||||
|
||||
auto this_node = this->filter->filter_node_id();
|
||||
auto& midi = this->global_context->current_graph.software_midi;
|
||||
auto node_it = midi.find(this_node);
|
||||
if (node_it == midi.end())
|
||||
{
|
||||
std::cerr << "Node " << this_node << " not found! \n";
|
||||
return false;
|
||||
}
|
||||
|
||||
// Wait for the pipewire server to send us back our node's ports
|
||||
this->filter->synchronize_ports(node_it->second);
|
||||
|
||||
if (node_it->second.inputs.empty())
|
||||
{
|
||||
std::cerr << "Node " << this_node << " has no ports! \n";
|
||||
return false;
|
||||
}
|
||||
|
||||
// Link ports
|
||||
const auto& p = node_it->second.inputs.front();
|
||||
link = this->global_context->link_ports(in_port.port, p.id);
|
||||
pw_loop_iterate(this->global_context->lp, 1);
|
||||
if (!link)
|
||||
{
|
||||
self.template error<invalid_parameter_error>(
|
||||
self.configuration,
|
||||
"PipeWire: could not connect to port: " + in_port.port_name + " -> " + p.port_name);
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
bool link_ports(auto& self, const output_port& out_port)
|
||||
{
|
||||
// Wait for the pipewire server to send us back our node's info
|
||||
for (int i = 0; i < 1000; i++)
|
||||
this->filter->synchronize_node();
|
||||
|
||||
auto this_node = this->filter->filter_node_id();
|
||||
auto& midi = this->global_context->current_graph.software_midi;
|
||||
auto node_it = midi.find(this_node);
|
||||
if (node_it == midi.end())
|
||||
{
|
||||
std::cerr << "Node " << this_node << " not found! \n";
|
||||
return false;
|
||||
}
|
||||
|
||||
// Wait for the pipewire server to send us back our node's ports
|
||||
this->filter->synchronize_ports(node_it->second);
|
||||
|
||||
if (node_it->second.outputs.empty())
|
||||
{
|
||||
std::cerr << "Node " << this_node << " has no ports! \n";
|
||||
return false;
|
||||
}
|
||||
|
||||
// Link ports
|
||||
const auto& p = node_it->second.outputs.front();
|
||||
link = this->global_context->link_ports(p.id, out_port.port);
|
||||
pw_loop_iterate(this->global_context->lp, 1);
|
||||
if (!link)
|
||||
{
|
||||
self.template error<invalid_parameter_error>(
|
||||
self.configuration,
|
||||
"PipeWire: could not connect to port: " + p.port_name + " -> " + out_port.port_name);
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
template <spa_direction Direction>
|
||||
static auto to_port_info(const pipewire_context::port_info& port)
|
||||
-> std::conditional_t<Direction == SPA_DIRECTION_OUTPUT, input_port, output_port>
|
||||
{
|
||||
std::string device_name, port_name;
|
||||
auto name_colon = port.port_alias.find(':');
|
||||
if (name_colon != std::string::npos)
|
||||
{
|
||||
device_name = port.port_alias.substr(0, name_colon);
|
||||
port_name = port.port_alias.substr(name_colon + 1);
|
||||
}
|
||||
else
|
||||
{
|
||||
port_name = port.port_alias;
|
||||
}
|
||||
|
||||
return {{
|
||||
.client = 0,
|
||||
.port = port.id,
|
||||
.manufacturer = "",
|
||||
.device_name = device_name,
|
||||
.port_name = port.port_name,
|
||||
.display_name = port_name,
|
||||
}};
|
||||
}
|
||||
|
||||
// Note: keep in mind that an "input" port for us (e.g. a keyboard that goes to the computer)
|
||||
// is an "output" port from the point of view of pipewire as data will come out of it
|
||||
template <spa_direction Direction>
|
||||
static auto get_ports(const pipewire_context& ctx) noexcept -> std::vector<
|
||||
std::conditional_t<Direction == SPA_DIRECTION_OUTPUT, input_port, output_port>>
|
||||
{
|
||||
std::vector<std::conditional_t<Direction == SPA_DIRECTION_OUTPUT, input_port, output_port>>
|
||||
ret;
|
||||
|
||||
{
|
||||
std::lock_guard _{ctx.current_graph.mtx};
|
||||
for (auto& node : ctx.current_graph.physical_midi)
|
||||
{
|
||||
for (auto& p :
|
||||
(Direction == SPA_DIRECTION_INPUT ? node.second.inputs : node.second.outputs))
|
||||
{
|
||||
ret.push_back(to_port_info<Direction>(p));
|
||||
}
|
||||
}
|
||||
for (auto& node : ctx.current_graph.software_midi)
|
||||
{
|
||||
for (auto& p :
|
||||
(Direction == SPA_DIRECTION_INPUT ? node.second.inputs : node.second.outputs))
|
||||
{
|
||||
ret.push_back(to_port_info<Direction>(p));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return ret;
|
||||
}
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,125 @@
|
||||
#pragma once
|
||||
#include <libremidi/backends/pipewire/config.hpp>
|
||||
#include <libremidi/backends/pipewire/helpers.hpp>
|
||||
#include <libremidi/detail/midi_in.hpp>
|
||||
#include <libremidi/detail/midi_stream_decoder.hpp>
|
||||
|
||||
#include <chrono>
|
||||
|
||||
namespace libremidi
|
||||
{
|
||||
class midi_in_pipewire final
|
||||
: public midi1::in_api
|
||||
, public pipewire_helpers
|
||||
, public error_handler
|
||||
{
|
||||
public:
|
||||
struct
|
||||
: input_configuration
|
||||
, pipewire_input_configuration
|
||||
{
|
||||
} configuration;
|
||||
|
||||
explicit midi_in_pipewire(input_configuration&& conf, pipewire_input_configuration&& apiconf)
|
||||
: configuration{std::move(conf), std::move(apiconf)}
|
||||
{
|
||||
create_context(*this);
|
||||
create_filter(*this);
|
||||
}
|
||||
|
||||
~midi_in_pipewire() override
|
||||
{
|
||||
stop_thread();
|
||||
do_close_port();
|
||||
destroy_filter(*this);
|
||||
destroy_context();
|
||||
}
|
||||
|
||||
void set_client_name(std::string_view) override
|
||||
{
|
||||
warning(configuration, "midi_in_pipewire: set_client_name unsupported");
|
||||
}
|
||||
|
||||
libremidi::API get_current_api() const noexcept override { return libremidi::API::PIPEWIRE; }
|
||||
|
||||
bool open_port(const input_port& in_port, std::string_view name) override
|
||||
{
|
||||
if (!create_local_port(*this, name, SPA_DIRECTION_INPUT))
|
||||
return false;
|
||||
|
||||
if (!link_ports(*this, in_port))
|
||||
return false;
|
||||
|
||||
start_thread();
|
||||
return true;
|
||||
}
|
||||
|
||||
bool open_virtual_port(std::string_view name) override
|
||||
{
|
||||
if (!create_local_port(*this, name, SPA_DIRECTION_INPUT))
|
||||
return false;
|
||||
|
||||
start_thread();
|
||||
return true;
|
||||
}
|
||||
|
||||
void close_port() override
|
||||
{
|
||||
stop_thread();
|
||||
do_close_port();
|
||||
}
|
||||
|
||||
void set_port_name(std::string_view port_name) override { rename_port(port_name); }
|
||||
|
||||
timestamp absolute_timestamp() const noexcept override { return system_ns(); }
|
||||
|
||||
void process(struct spa_io_position* position)
|
||||
{
|
||||
static constexpr timestamp_backend_info timestamp_info{
|
||||
.has_absolute_timestamps = true,
|
||||
.absolute_is_monotonic = true,
|
||||
.has_samples = true,
|
||||
};
|
||||
|
||||
assert(this->filter);
|
||||
assert(this->filter->port);
|
||||
const auto b = pw.filter_dequeue_buffer(this->filter->port);
|
||||
if (!b)
|
||||
return;
|
||||
|
||||
const auto buf = b->buffer;
|
||||
const auto d = &buf->datas[0];
|
||||
|
||||
if (d->data == nullptr)
|
||||
return;
|
||||
|
||||
const auto pod
|
||||
= (spa_pod*)spa_pod_from_data(d->data, d->maxsize, d->chunk->offset, d->chunk->size);
|
||||
if (!pod)
|
||||
return;
|
||||
if (!spa_pod_is_sequence(pod))
|
||||
return;
|
||||
|
||||
struct spa_pod_control* c{};
|
||||
SPA_POD_SEQUENCE_FOREACH((struct spa_pod_sequence*)pod, c)
|
||||
{
|
||||
if (c->type != SPA_CONTROL_Midi)
|
||||
continue;
|
||||
|
||||
auto data = (uint8_t*)SPA_POD_BODY(&c->value);
|
||||
auto size = SPA_POD_BODY_SIZE(&c->value);
|
||||
|
||||
const auto to_ns = [=, clk = position->clock] {
|
||||
return 1e9 * ((clk.position + c->offset) / (double)clk.rate.denom);
|
||||
};
|
||||
|
||||
m_processing.on_bytes(
|
||||
{data, data + size}, m_processing.timestamp<timestamp_info>(to_ns, c->offset));
|
||||
}
|
||||
|
||||
pw.filter_queue_buffer(this->filter->port, b);
|
||||
}
|
||||
|
||||
midi1::input_state_machine m_processing{this->configuration};
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,171 @@
|
||||
#pragma once
|
||||
#include <libremidi/backends/pipewire/config.hpp>
|
||||
#include <libremidi/backends/pipewire/helpers.hpp>
|
||||
#include <libremidi/detail/midi_out.hpp>
|
||||
|
||||
#include <readerwriterqueue.h>
|
||||
|
||||
#include <semaphore>
|
||||
|
||||
namespace libremidi
|
||||
{
|
||||
class midi_out_pipewire
|
||||
: public midi1::out_api
|
||||
, public pipewire_helpers
|
||||
, public error_handler
|
||||
{
|
||||
public:
|
||||
struct
|
||||
: output_configuration
|
||||
, pipewire_output_configuration
|
||||
{
|
||||
} configuration;
|
||||
|
||||
midi_out_pipewire(output_configuration&& conf, pipewire_output_configuration&& apiconf)
|
||||
: configuration{std::move(conf), std::move(apiconf)}
|
||||
{
|
||||
create_context(*this);
|
||||
create_filter(*this);
|
||||
}
|
||||
|
||||
~midi_out_pipewire() override
|
||||
{
|
||||
stop_thread();
|
||||
do_close_port();
|
||||
destroy_filter(*this);
|
||||
destroy_context();
|
||||
}
|
||||
|
||||
void set_client_name(std::string_view) override
|
||||
{
|
||||
warning(configuration, "midi_out_pipewire: set_client_name unsupported");
|
||||
}
|
||||
|
||||
libremidi::API get_current_api() const noexcept override { return libremidi::API::PIPEWIRE; }
|
||||
|
||||
bool open_port(const output_port& out_port, std::string_view name) override
|
||||
{
|
||||
if (!create_local_port(*this, name, SPA_DIRECTION_OUTPUT))
|
||||
return false;
|
||||
|
||||
this->filter->set_port_buffer(configuration.output_buffer_size);
|
||||
|
||||
if (!link_ports(*this, out_port))
|
||||
return false;
|
||||
|
||||
start_thread();
|
||||
return true;
|
||||
}
|
||||
|
||||
bool open_virtual_port(std::string_view name) override
|
||||
{
|
||||
if (!create_local_port(*this, name, SPA_DIRECTION_OUTPUT))
|
||||
return false;
|
||||
|
||||
this->filter->set_port_buffer(configuration.output_buffer_size);
|
||||
|
||||
start_thread();
|
||||
return true;
|
||||
}
|
||||
|
||||
void close_port() override
|
||||
{
|
||||
stop_thread();
|
||||
do_close_port();
|
||||
}
|
||||
|
||||
void set_port_name(std::string_view port_name) override { rename_port(port_name); }
|
||||
|
||||
int process(spa_io_position* pos)
|
||||
{
|
||||
m_process_clock.store(pos->clock.nsec, std::memory_order_relaxed);
|
||||
const auto b = pw.filter_dequeue_buffer(this->filter->port);
|
||||
if (!b)
|
||||
return 1;
|
||||
|
||||
const auto buf = b->buffer;
|
||||
const auto d = &buf->datas[0];
|
||||
|
||||
if (d->data == nullptr)
|
||||
return 1;
|
||||
|
||||
spa_pod_builder build;
|
||||
spa_zero(build);
|
||||
spa_pod_builder_init(&build, d->data, d->maxsize);
|
||||
|
||||
spa_pod_frame f;
|
||||
spa_pod_builder_push_sequence(&build, &f, 0);
|
||||
|
||||
// for all events
|
||||
while (auto m_ptr = m_queue.peek())
|
||||
{
|
||||
auto& m = *m_ptr;
|
||||
if (m.empty())
|
||||
{
|
||||
m_queue.pop();
|
||||
continue;
|
||||
}
|
||||
|
||||
// TODO why
|
||||
if (m.bytes[0] == 0xff)
|
||||
{
|
||||
m_queue.pop();
|
||||
continue;
|
||||
}
|
||||
|
||||
spa_pod_builder_control(&build, m.timestamp, SPA_CONTROL_Midi);
|
||||
int res = spa_pod_builder_bytes(&build, m.bytes.data(), m.bytes.size());
|
||||
|
||||
// Try again next buffer
|
||||
if (res == -ENOSPC)
|
||||
break;
|
||||
|
||||
m_queue.pop();
|
||||
}
|
||||
spa_pod_builder_pop(&build, &f);
|
||||
|
||||
int n_fill_frames = build.state.offset;
|
||||
if (n_fill_frames > 0)
|
||||
{
|
||||
d->chunk->offset = 0;
|
||||
d->chunk->stride = 1;
|
||||
d->chunk->size = n_fill_frames;
|
||||
b->size = n_fill_frames;
|
||||
|
||||
pw.filter_queue_buffer(this->filter->port, b);
|
||||
return 0;
|
||||
}
|
||||
|
||||
pw.filter_flush(this->filter->filter, true);
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
void send_message(const unsigned char* message, size_t size) override
|
||||
{
|
||||
m_queue.enqueue(libremidi::message(midi_bytes{message, message + size}, 0));
|
||||
}
|
||||
|
||||
int convert_timestamp(int64_t user) const noexcept
|
||||
{
|
||||
switch (configuration.timestamps)
|
||||
{
|
||||
case timestamp_mode::AudioFrame:
|
||||
return static_cast<int>(user);
|
||||
|
||||
default:
|
||||
// TODO
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
||||
void schedule_message(int64_t ts, const unsigned char* message, size_t size) override
|
||||
{
|
||||
m_queue.enqueue(
|
||||
libremidi::message(midi_bytes{message, message + size}, convert_timestamp(ts)));
|
||||
}
|
||||
|
||||
moodycamel::ReaderWriterQueue<libremidi::message> m_queue;
|
||||
std::atomic_int64_t m_process_clock = 0;
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,74 @@
|
||||
#pragma once
|
||||
#include <libremidi/backends/pipewire/config.hpp>
|
||||
#include <libremidi/backends/pipewire/helpers.hpp>
|
||||
#include <libremidi/detail/observer.hpp>
|
||||
|
||||
#include <unordered_set>
|
||||
|
||||
namespace libremidi
|
||||
{
|
||||
class observer_pipewire final
|
||||
: public observer_api
|
||||
, private pipewire_helpers
|
||||
, private error_handler
|
||||
{
|
||||
public:
|
||||
struct
|
||||
: observer_configuration
|
||||
, pipewire_observer_configuration
|
||||
{
|
||||
} configuration;
|
||||
|
||||
explicit observer_pipewire(observer_configuration&& conf, pipewire_observer_configuration&& apiconf)
|
||||
: configuration{std::move(conf), std::move(apiconf)}
|
||||
{
|
||||
create_context(*this);
|
||||
|
||||
// FIXME notify_in_constructor
|
||||
// FIXME port rename callback
|
||||
#if 0
|
||||
// Initialize PipeWire client
|
||||
if (configuration.context)
|
||||
{
|
||||
this->client = configuration.context;
|
||||
set_callbacks();
|
||||
}
|
||||
else
|
||||
#endif
|
||||
{
|
||||
this->add_callbacks(configuration);
|
||||
this->start_thread();
|
||||
}
|
||||
}
|
||||
|
||||
libremidi::API get_current_api() const noexcept override { return libremidi::API::PIPEWIRE; }
|
||||
|
||||
std::vector<libremidi::input_port> get_input_ports() const noexcept override
|
||||
{
|
||||
return get_ports<SPA_DIRECTION_OUTPUT>(*this->global_context);
|
||||
}
|
||||
|
||||
std::vector<libremidi::output_port> get_output_ports() const noexcept override
|
||||
{
|
||||
return get_ports<SPA_DIRECTION_INPUT>(*this->global_context);
|
||||
}
|
||||
|
||||
~observer_pipewire()
|
||||
{
|
||||
stop_thread();
|
||||
destroy_context();
|
||||
#if 0
|
||||
if (client && !configuration.context)
|
||||
{
|
||||
// If we own the client, deactivate it
|
||||
pipewire_deactivate(this->client);
|
||||
pipewire_client_close(this->client);
|
||||
this->client = nullptr;
|
||||
}
|
||||
#endif
|
||||
}
|
||||
|
||||
std::unordered_set<std::string> seen_input_ports;
|
||||
std::unordered_set<std::string> seen_output_ports;
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,153 @@
|
||||
#pragma once
|
||||
#if __has_include(<boost/lockfree/spsc_queue.hpp>)
|
||||
#include <libremidi/backends/pipewire/config.hpp>
|
||||
#include <libremidi/backends/pipewire/helpers.hpp>
|
||||
#include <libremidi/shared_context.hpp>
|
||||
|
||||
#include <boost/lockfree/spsc_queue.hpp>
|
||||
|
||||
#include <variant>
|
||||
|
||||
namespace libremidi::pipewire
|
||||
{
|
||||
|
||||
// Create a PipeWire client which will be shared across objects
|
||||
struct shared_handler : public libremidi::shared_context
|
||||
{
|
||||
explicit shared_handler(std::string_view v)
|
||||
{
|
||||
midiin_callbacks.reserve(64);
|
||||
midiout_callbacks.reserve(64);
|
||||
|
||||
pipewire_status_t status{};
|
||||
client = pipewire_client_open(v.data(), PipewireNoStartServer, &status);
|
||||
assert(client);
|
||||
assert(status == 0);
|
||||
pipewire_set_process_callback(
|
||||
client,
|
||||
+[](pipewire_nframes_t cnt, void* ctx) -> int {
|
||||
((shared_handler*)ctx)->pipewire_callback(cnt);
|
||||
return 0;
|
||||
},
|
||||
this);
|
||||
}
|
||||
|
||||
virtual void start_processing() override { pipewire_activate(client); }
|
||||
virtual void stop_processing() override { pipewire_deactivate(client); }
|
||||
|
||||
static shared_configurations make(std::string_view client_name)
|
||||
{
|
||||
auto clt = std::make_shared<shared_handler>(client_name);
|
||||
auto add_in_cb = [client = std::weak_ptr{clt}](libremidi::pipewire_callback cb) {
|
||||
if (auto clt = client.lock())
|
||||
clt->events.push({shared_handler::event_type::in_callback_added, std::move(cb)});
|
||||
};
|
||||
auto clear_in_cb = [client = std::weak_ptr{clt}](int64_t index) {
|
||||
if (auto clt = client.lock())
|
||||
clt->events.push({shared_handler::event_type::in_callback_removed, index});
|
||||
};
|
||||
auto add_out_cb = [client = std::weak_ptr{clt}](libremidi::pipewire_callback cb) {
|
||||
if (auto clt = client.lock())
|
||||
clt->events.push({shared_handler::event_type::out_callback_added, std::move(cb)});
|
||||
};
|
||||
auto clear_out_cb = [client = std::weak_ptr{clt}](int64_t index) {
|
||||
if (auto clt = client.lock())
|
||||
clt->events.push({shared_handler::event_type::out_callback_removed, index});
|
||||
};
|
||||
return {
|
||||
.context = clt,
|
||||
.observer = pipewire_observer_configuration{.context = clt->client},
|
||||
.in
|
||||
= pipewire_input_configuration{.context = clt->client, .set_process_func = add_in_cb, .clear_process_func = clear_in_cb},
|
||||
.out
|
||||
= pipewire_output_configuration{.context = clt->client, .set_process_func = add_out_cb, .clear_process_func = clear_out_cb},
|
||||
};
|
||||
}
|
||||
|
||||
int pipewire_callback(pipewire_nframes_t cnt)
|
||||
{
|
||||
// 1. Process the events that will change the callback list
|
||||
event ev;
|
||||
while (events.pop(ev))
|
||||
{
|
||||
switch (ev.type)
|
||||
{
|
||||
case in_callback_added:
|
||||
midiin_callbacks.push_back(
|
||||
std::move(*std::get_if<libremidi::pipewire_callback>(&ev.payload)));
|
||||
break;
|
||||
case in_callback_removed: {
|
||||
auto idx = *std::get_if<int64_t>(&ev.payload);
|
||||
for (auto it = midiin_callbacks.begin(); it != midiin_callbacks.end();)
|
||||
{
|
||||
if (it->token == idx)
|
||||
{
|
||||
midiin_callbacks.erase(it);
|
||||
break;
|
||||
}
|
||||
else
|
||||
{
|
||||
++it;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
case out_callback_added:
|
||||
midiout_callbacks.push_back(
|
||||
std::move(*std::get_if<libremidi::pipewire_callback>(&ev.payload)));
|
||||
break;
|
||||
case out_callback_removed:
|
||||
auto idx = *std::get_if<int64_t>(&ev.payload);
|
||||
for (auto it = midiout_callbacks.begin(); it != midiout_callbacks.end();)
|
||||
{
|
||||
if (it->token == idx)
|
||||
{
|
||||
midiout_callbacks.erase(it);
|
||||
break;
|
||||
}
|
||||
else
|
||||
{
|
||||
++it;
|
||||
}
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
for (auto& cb : midiin_callbacks)
|
||||
cb.callback(cnt);
|
||||
|
||||
for (auto& cb : midiout_callbacks)
|
||||
cb.callback(cnt);
|
||||
|
||||
return 0;
|
||||
}
|
||||
|
||||
~shared_handler()
|
||||
{
|
||||
pipewire_deactivate(client);
|
||||
pipewire_client_close(client);
|
||||
}
|
||||
|
||||
pipewire_client_t* client{};
|
||||
|
||||
enum event_type
|
||||
{
|
||||
in_callback_added,
|
||||
in_callback_removed,
|
||||
out_callback_added,
|
||||
out_callback_removed,
|
||||
};
|
||||
struct event
|
||||
{
|
||||
event_type type;
|
||||
std::variant<libremidi::pipewire_callback, int64_t> payload;
|
||||
};
|
||||
|
||||
boost::lockfree::spsc_queue<event> events{16};
|
||||
|
||||
std::vector<libremidi::pipewire_callback> midiin_callbacks;
|
||||
std::vector<libremidi::pipewire_callback> midiout_callbacks;
|
||||
};
|
||||
}
|
||||
#endif
|
||||
Reference in New Issue
Block a user