Compare commits

..

4 commits

Author SHA1 Message Date
0a9cd54f1a Fix stuff 2026-02-12 17:47:51 -07:00
2ac234ecca Levels work 2026-02-12 17:18:36 -07:00
2208483123 Almost working, no green to sink 2026-02-12 17:07:33 -07:00
3c1c86f952 Fix nodes 2026-02-12 16:52:00 -07:00
10 changed files with 301 additions and 193 deletions

View file

@ -95,8 +95,8 @@ if(WARPPIPE_BUILD_GUI)
gui/PresetManager.cpp
gui/VolumeWidgets.cpp
gui/AudioLevelMeter.cpp
gui/BezierConnectionPainter.cpp
gui/SquareConnectionPainter.cpp
gui/WarpBezierConnectionPainter.cpp
)
target_link_libraries(warppipe-gui PRIVATE
@ -114,8 +114,8 @@ if(WARPPIPE_BUILD_GUI)
gui/PresetManager.cpp
gui/VolumeWidgets.cpp
gui/AudioLevelMeter.cpp
gui/BezierConnectionPainter.cpp
gui/SquareConnectionPainter.cpp
gui/WarpBezierConnectionPainter.cpp
)
target_compile_definitions(warppipe-gui-tests PRIVATE WARPPIPE_TESTING)

View file

@ -1,4 +1,4 @@
#include "WarpBezierConnectionPainter.h"
#include "BezierConnectionPainter.h"
#include "WarpGraphModel.h"
#include <QtNodes/internal/BasicGraphicsScene.hpp>
@ -14,10 +14,11 @@
#include <algorithm>
#include <cmath>
QPainterPath WarpBezierConnectionPainter::cubicPath(
QPainterPath BezierConnectionPainter::cubicPath(
QtNodes::ConnectionGraphicsObject const &cgo) const {
QPointF const &in = cgo.endPoint(QtNodes::PortType::In);
QPointF const &out = cgo.endPoint(QtNodes::PortType::Out);
auto const c1c2 = cgo.pointsC1C2();
QPainterPath cubic(out);
@ -25,7 +26,7 @@ QPainterPath WarpBezierConnectionPainter::cubicPath(
return cubic;
}
void WarpBezierConnectionPainter::paint(
void BezierConnectionPainter::paint(
QPainter *painter,
QtNodes::ConnectionGraphicsObject const &cgo) const {
auto const &style = QtNodes::StyleCollection::connectionStyle();
@ -42,7 +43,19 @@ void WarpBezierConnectionPainter::paint(
auto *model = dynamic_cast<WarpGraphModel *>(&scene->graphModel());
if (model) {
auto cId = cgo.connectionId();
peakLevel = model->connectionPeakLevel(cId);
if (model->connectionExists(cId)) {
peakLevel = model->nodePeakLevel(cId.outNodeId);
if (peakLevel < 0.005f) {
auto const *outData = model->warpNodeData(cId.outNodeId);
if (outData &&
WarpGraphModel::classifyNode(outData->info) ==
WarpNodeType::kApplication) {
peakLevel = std::max(peakLevel,
model->nodePeakLevel(cId.inNodeId));
}
}
}
}
}
@ -98,7 +111,7 @@ void WarpBezierConnectionPainter::paint(
painter->drawEllipse(cgo.in(), pointRadius, pointRadius);
}
QPainterPath WarpBezierConnectionPainter::getPainterStroke(
QPainterPath BezierConnectionPainter::getPainterStroke(
QtNodes::ConnectionGraphicsObject const &cgo) const {
auto cubic = cubicPath(cgo);
@ -113,5 +126,6 @@ QPainterPath WarpBezierConnectionPainter::getPainterStroke(
QPainterPathStroker stroker;
stroker.setWidth(10.0);
return stroker.createStroke(result);
}

View file

@ -2,7 +2,7 @@
#include <QtNodes/internal/AbstractConnectionPainter.hpp>
class WarpBezierConnectionPainter : public QtNodes::AbstractConnectionPainter {
class BezierConnectionPainter : public QtNodes::AbstractConnectionPainter {
public:
void paint(QPainter *painter,
QtNodes::ConnectionGraphicsObject const &cgo) const override;

View file

@ -1,4 +1,5 @@
#include "AudioLevelMeter.h"
#include "BezierConnectionPainter.h"
#include "GraphEditorWidget.h"
#include "PresetManager.h"
#include "SquareConnectionPainter.h"
@ -9,7 +10,7 @@
#include <QtNodes/BasicGraphicsScene>
#include <QtNodes/ConnectionStyle>
#include <QtNodes/GraphicsView>
#include "WarpBezierConnectionPainter.h"
#include <QtNodes/internal/NodeGraphicsObject.hpp>
#include <QtNodes/internal/ConnectionGraphicsObject.hpp>
#include <QtNodes/internal/UndoCommands.hpp>
@ -203,6 +204,8 @@ GraphEditorWidget::GraphEditorWidget(warppipe::Client *client,
"UseDataDefinedColors": false
}})");
m_scene->setConnectionPainter(std::make_unique<BezierConnectionPainter>());
m_view = new ZoomGraphicsView(m_scene);
m_view->setFocusPolicy(Qt::StrongFocus);
m_view->viewport()->setFocusPolicy(Qt::StrongFocus);
@ -566,21 +569,17 @@ GraphEditorWidget::GraphEditorWidget(warppipe::Client *client,
connect(m_model, &QtNodes::AbstractGraphModel::nodeCreated, this,
[this](QtNodes::NodeId nodeId) {
wireVolumeWidget(nodeId);
rebuildMixerStrips();
rebuildNodeMeters();
rebuildRulesList();
this->scheduleSidebarRebuild();
});
connect(m_model, &QtNodes::AbstractGraphModel::nodeDeleted, this,
[this](QtNodes::NodeId nodeId) {
m_mixerStrips.erase(nodeId);
m_nodeMeters.erase(nodeId);
rebuildMixerStrips();
rebuildNodeMeters();
rebuildRulesList();
if (nodeId == m_selectedNodeId) {
m_selectedNodeId = 0;
clearNodeDetailsPanel();
}
this->scheduleSidebarRebuild();
});
connect(m_scene, &QGraphicsScene::selectionChanged, this,
@ -1544,7 +1543,7 @@ void GraphEditorWidget::rebuildMixerStrips() {
while (layout->count() > 0) {
auto *item = layout->takeAt(0);
if (item->widget())
item->widget()->deleteLater();
delete item->widget();
delete item;
}
m_mixerStrips.clear();
@ -1722,7 +1721,7 @@ void GraphEditorWidget::rebuildNodeMeters() {
while (layout->count() > 0) {
auto *item = layout->takeAt(0);
if (item->widget())
item->widget()->deleteLater();
delete item->widget();
delete item;
}
m_nodeMeters.clear();
@ -1792,6 +1791,20 @@ void GraphEditorWidget::rebuildNodeMeters() {
}
}
void GraphEditorWidget::scheduleSidebarRebuild() {
if (m_sidebarRebuildPending) {
return;
}
m_sidebarRebuildPending = true;
QTimer::singleShot(0, this, [this]() {
m_sidebarRebuildPending = false;
rebuildMixerStrips();
rebuildNodeMeters();
rebuildRulesList();
});
}
void GraphEditorWidget::rebuildRulesList() {
if (!m_rulesContainer || !m_client)
return;
@ -1803,7 +1816,7 @@ void GraphEditorWidget::rebuildRulesList() {
while (layout->count() > 0) {
auto *item = layout->takeAt(0);
if (item->widget())
item->widget()->deleteLater();
delete item->widget();
delete item;
}
@ -1927,7 +1940,7 @@ void GraphEditorWidget::setConnectionStyle(ConnectionStyleType style) {
m_scene->setConnectionPainter(std::make_unique<SquareConnectionPainter>());
} else {
m_scene->setConnectionPainter(
std::make_unique<WarpBezierConnectionPainter>());
std::make_unique<BezierConnectionPainter>());
}
for (auto *item : m_scene->items()) {

View file

@ -96,6 +96,7 @@ private:
void rebuildMixerStrips();
void updateMeters();
void rebuildNodeMeters();
void scheduleSidebarRebuild();
void rebuildRulesList();
void showAddRuleDialog(const std::string &prefillApp = {},
const std::string &prefillBin = {},
@ -135,6 +136,7 @@ private:
QWidget *m_mixerContainer = nullptr;
QScrollArea *m_mixerScroll = nullptr;
std::unordered_map<QtNodes::NodeId, QWidget *> m_mixerStrips;
bool m_sidebarRebuildPending = false;
QTimer *m_meterTimer = nullptr;
AudioLevelMeter *m_masterMeterL = nullptr;

View file

@ -66,16 +66,17 @@ QPainterPath SquareConnectionPainter::orthogonalPath(
constexpr double kNodePad = 15.0;
auto const cId = cgo.connectionId();
auto *sceneForChannel = cgo.nodeScene();
auto *mdl = sceneForChannel
? dynamic_cast<WarpGraphModel *>(&sceneForChannel->graphModel())
: nullptr;
bool connectionAlive = mdl ? mdl->connectionExists(cId) : true;
double spread = static_cast<double>(cId.outPortIndex) * kSpacing;
auto *sceneForChannel = cgo.nodeScene();
if (sceneForChannel) {
auto *mdl = dynamic_cast<WarpGraphModel *>(&sceneForChannel->graphModel());
if (mdl) {
auto ch = mdl->connectionChannel(cId);
spread = (static_cast<double>(ch.index) - (ch.count - 1) / 2.0)
* kSpacing;
}
if (mdl && connectionAlive) {
auto ch = mdl->connectionChannel(cId);
spread = (static_cast<double>(ch.index) - (ch.count - 1) / 2.0)
* kSpacing;
}
double const dy = in.y() - out.y();
@ -125,12 +126,9 @@ QPainterPath SquareConnectionPainter::orthogonalPath(
//
double railOffset = 0.0;
if (sceneForChannel) {
auto *mdl2 = dynamic_cast<WarpGraphModel *>(&sceneForChannel->graphModel());
if (mdl2) {
auto ch = mdl2->connectionChannel(cId);
railOffset = static_cast<double>(ch.index) * kSpacing;
}
if (mdl && connectionAlive) {
auto ch = mdl->connectionChannel(cId);
railOffset = static_cast<double>(ch.index) * kSpacing;
}
double rightX = out.x() + kMinStub + railOffset;
@ -205,7 +203,19 @@ void SquareConnectionPainter::paint(
auto *model = dynamic_cast<WarpGraphModel *>(&scene->graphModel());
if (model) {
auto cId = cgo.connectionId();
peakLevel = model->connectionPeakLevel(cId);
if (model->connectionExists(cId)) {
peakLevel = model->nodePeakLevel(cId.outNodeId);
if (peakLevel < 0.005f) {
auto const *outData = model->warpNodeData(cId.outNodeId);
if (outData &&
WarpGraphModel::classifyNode(outData->info) ==
WarpNodeType::kApplication) {
peakLevel = std::max(peakLevel,
model->nodePeakLevel(cId.inNodeId));
}
}
}
}
}

View file

@ -343,6 +343,8 @@ bool WarpGraphModel::deleteConnection(
}
m_connections.erase(it);
m_connectionChannels.erase(connectionId);
recomputeConnectionChannels();
Q_EMIT connectionDeleted(connectionId);
return true;
}
@ -374,6 +376,7 @@ bool WarpGraphModel::deleteNode(QtNodes::NodeId const nodeId) {
m_positions.erase(nodeId);
m_sizes.erase(nodeId);
m_volumeStates.erase(nodeId);
m_peakLevels.erase(nodeId);
m_styleCache.erase(nodeId);
m_volumeWidgets.erase(nodeId);
Q_EMIT nodeDeleted(nodeId);
@ -481,8 +484,14 @@ void WarpGraphModel::refreshFromClient() {
if (existing != m_pwToQt.end()) {
QtNodes::NodeId qtId = existing->second;
auto &data = m_nodes[qtId];
bool typeChanged = (data.info.is_virtual != nodeInfo.is_virtual);
data.info = nodeInfo;
if (typeChanged) {
m_styleCache.erase(qtId);
Q_EMIT nodeUpdated(qtId);
}
bool portsMissing =
data.inputPorts.empty() && data.outputPorts.empty();
if (portsMissing) {
@ -1473,23 +1482,6 @@ float WarpGraphModel::nodePeakLevel(QtNodes::NodeId nodeId) const {
return it != m_peakLevels.end() ? it->second : 0.0f;
}
float WarpGraphModel::connectionPeakLevel(QtNodes::ConnectionId cId) const {
constexpr float kSourceReliableThreshold = 0.005f;
float outPeak = nodePeakLevel(cId.outNodeId);
if (outPeak >= kSourceReliableThreshold)
return outPeak;
auto outNodeIt = m_nodes.find(cId.outNodeId);
if (outNodeIt == m_nodes.end())
return outPeak;
if (classifyNode(outNodeIt->second.info) != WarpNodeType::kApplication)
return outPeak;
return std::max(outPeak, nodePeakLevel(cId.inNodeId));
}
void WarpGraphModel::recomputeConnectionChannels() {
m_connectionChannels.clear();

View file

@ -109,7 +109,6 @@ public:
void setNodePeakLevel(QtNodes::NodeId nodeId, float level);
float nodePeakLevel(QtNodes::NodeId nodeId) const;
float connectionPeakLevel(QtNodes::ConnectionId cId) const;
struct ConnectionChannel {
int index = 0;

View file

@ -107,10 +107,12 @@ bool MatchesRule(const NodeInfo& node, const RuleMatch& match) {
struct StreamData {
pw_stream* stream = nullptr;
pw_impl_module* module = nullptr;
pw_proxy* proxy = nullptr;
spa_hook listener{};
pw_thread_loop* loop = nullptr;
bool is_source = false;
bool loopback = false;
bool linger = false;
std::string target_node;
std::string name;
bool ready = false;
@ -124,7 +126,6 @@ struct StreamData {
struct LinkProxy {
pw_proxy* proxy = nullptr;
spa_hook listener{};
bool listener_attached = false;
pw_thread_loop* loop = nullptr;
bool done = false;
bool failed = false;
@ -134,20 +135,6 @@ struct LinkProxy {
uint32_t input_port = 0;
};
void DetachLinkProxy(LinkProxy* link, bool destroy_proxy) {
if (!link) {
return;
}
if (link->listener_attached) {
spa_hook_remove(&link->listener);
link->listener_attached = false;
}
if (destroy_proxy && link->proxy) {
pw_proxy_destroy(link->proxy);
}
link->proxy = nullptr;
}
void LinkProxyBound(void* data, uint32_t global_id) {
auto* link = static_cast<LinkProxy*>(data);
if (!link) {
@ -165,8 +152,6 @@ void LinkProxyRemoved(void* data) {
if (!link) {
return;
}
link->listener_attached = false;
link->proxy = nullptr;
link->done = true;
if (link->loop) {
pw_thread_loop_signal(link->loop, false);
@ -178,8 +163,6 @@ void LinkProxyError(void* data, int, int res, const char* message) {
if (!link) {
return;
}
link->listener_attached = false;
link->proxy = nullptr;
link->failed = true;
link->error = message ? message : spa_strerror(res);
if (link->loop) {
@ -194,6 +177,42 @@ static const pw_proxy_events kLinkProxyEvents = {
.error = LinkProxyError,
};
void NodeProxyBound(void* data, uint32_t global_id) {
auto* sd = static_cast<StreamData*>(data);
if (!sd) return;
sd->node_id = global_id;
sd->ready = true;
if (sd->loop) {
pw_thread_loop_signal(sd->loop, false);
}
}
void NodeProxyRemoved(void* data) {
auto* sd = static_cast<StreamData*>(data);
if (!sd) return;
sd->ready = true;
if (sd->loop) {
pw_thread_loop_signal(sd->loop, false);
}
}
void NodeProxyError(void* data, int, int res, const char* message) {
auto* sd = static_cast<StreamData*>(data);
if (!sd) return;
sd->failed = true;
sd->error = message ? message : spa_strerror(res);
if (sd->loop) {
pw_thread_loop_signal(sd->loop, false);
}
}
static const pw_proxy_events kNodeProxyEvents = {
.version = PW_VERSION_PROXY_EVENTS,
.bound = NodeProxyBound,
.removed = NodeProxyRemoved,
.error = NodeProxyError,
};
void StreamProcess(void* data) {
auto* stream_data = static_cast<StreamData*>(data);
if (!stream_data || !stream_data->stream) {
@ -310,10 +329,12 @@ void NodeMeterProcess(void* data) {
had_data = true;
pw_stream_queue_buffer(meter->stream, buf);
}
if (had_data) {
meter->peak_left.store(left, std::memory_order_relaxed);
meter->peak_right.store(right, std::memory_order_relaxed);
if (!had_data) {
left = 0.0f;
right = 0.0f;
}
meter->peak_left.store(left, std::memory_order_relaxed);
meter->peak_right.store(right, std::memory_order_relaxed);
}
static const pw_stream_events kNodeMeterEvents = {
@ -454,7 +475,29 @@ struct Client::Impl {
void Client::Impl::NodeInfoChanged(void* data, const struct pw_node_info* info) {
auto* np = static_cast<NodeProxyData*>(data);
if (!np || !info || np->params_subscribed) return;
if (!np || !info) return;
auto* impl = static_cast<Client::Impl*>(np->impl_ptr);
bool notify = false;
if (impl && info->props) {
const char* virt_str = spa_dict_lookup(info->props, PW_KEY_NODE_VIRTUAL);
if (virt_str) {
std::lock_guard<std::mutex> lock(impl->cache_mutex);
auto it = impl->nodes.find(np->node_id);
if (it != impl->nodes.end()) {
bool new_virtual = spa_streq(virt_str, "true");
if (it->second.is_virtual != new_virtual) {
it->second.is_virtual = new_virtual;
notify = true;
}
}
}
}
if (notify) {
impl->NotifyChange();
}
if (np->params_subscribed) return;
for (uint32_t i = 0; i < info->n_params; ++i) {
if (info->params[i].id == SPA_PARAM_Props &&
@ -527,15 +570,18 @@ void Client::Impl::NodeParamChanged(void* data, int, uint32_t id,
}
}
bool is_virtual = false;
bool uses_monitor_volume = false;
{
std::lock_guard<std::mutex> lock(impl->cache_mutex);
is_virtual = impl->virtual_streams.count(np->node_id) > 0;
auto node_it = impl->nodes.find(np->node_id);
if (node_it != impl->nodes.end()) {
uses_monitor_volume = node_it->second.is_virtual;
}
}
float effective_vol = is_virtual && mon_volume >= 0.0f ? mon_volume : volume;
bool effective_mute = is_virtual && found_mon_mute ? mon_mute : mute;
bool effective_found_mute = is_virtual ? found_mon_mute : found_mute;
float effective_vol = uses_monitor_volume && mon_volume >= 0.0f ? mon_volume : volume;
bool effective_mute = uses_monitor_volume && found_mon_mute ? mon_mute : mute;
bool effective_found_mute = uses_monitor_volume ? found_mon_mute : found_mute;
bool changed = false;
{
@ -716,6 +762,7 @@ void Client::Impl::RegistryGlobalRemove(void* data, uint32_t id) {
{
std::lock_guard<std::mutex> lock(impl->cache_mutex);
impl->virtual_streams.erase(id);
impl->link_proxies.erase(id);
auto node_it = impl->nodes.find(id);
if (node_it != impl->nodes.end()) {
impl->nodes.erase(node_it);
@ -873,16 +920,27 @@ Result<uint32_t> Client::Impl::CreateVirtualStreamLocked(std::string_view name,
{
std::lock_guard<std::mutex> lock(cache_mutex);
for (const auto& entry : nodes) {
if (entry.second.name == stream_name) {
return {Status::Error(StatusCode::kInvalidArgument, "duplicate node name"), 0};
}
}
for (const auto& entry : virtual_streams) {
if (entry.second && entry.second->name == stream_name) {
return {Status::Error(StatusCode::kInvalidArgument, "duplicate node name"), 0};
}
}
for (const auto& entry : nodes) {
if (entry.second.name == stream_name) {
uint32_t existing_id = entry.first;
auto adopted = std::make_unique<StreamData>();
adopted->linger = true;
adopted->loop = thread_loop;
adopted->is_source = is_source;
adopted->name = stream_name;
adopted->node_id = existing_id;
adopted->ready = true;
if (options.format.rate != 0) adopted->rate = options.format.rate;
if (options.format.channels != 0) adopted->channels = options.format.channels;
virtual_streams.emplace(existing_id, std::move(adopted));
return {Status::Ok(), existing_id};
}
}
if (options.behavior == VirtualBehavior::kLoopback && options.target_node) {
bool found_target = false;
for (const auto& entry : nodes) {
@ -994,29 +1052,53 @@ Result<uint32_t> Client::Impl::CreateVirtualStreamLocked(std::string_view name,
return {Status::Ok(), node_id};
}
pw_properties* props = pw_properties_new(PW_KEY_MEDIA_TYPE, "Audio",
PW_KEY_MEDIA_CATEGORY, media_category,
PW_KEY_MEDIA_ROLE, "Music",
PW_KEY_MEDIA_CLASS, media_class_value.c_str(),
PW_KEY_NODE_NAME, stream_name.c_str(),
PW_KEY_MEDIA_NAME, display_name.c_str(),
PW_KEY_NODE_DESCRIPTION, display_name.c_str(),
PW_KEY_NODE_VIRTUAL, "true",
nullptr);
// Build audio position string for channels (e.g. "FL,FR" for stereo).
std::string audio_position;
{
static const char* const kChannelNames[] = {
"FL", "FR", "FC", "LFE", "RL", "RR", "SL", "SR"
};
uint32_t ch = options.format.channels;
for (uint32_t i = 0; i < ch && i < 8; ++i) {
if (i > 0) audio_position += ',';
audio_position += kChannelNames[i];
}
if (audio_position.empty()) audio_position = "FL,FR";
}
pw_properties* props = pw_properties_new(
"factory.name", "support.null-audio-sink",
PW_KEY_NODE_NAME, stream_name.c_str(),
PW_KEY_NODE_DESCRIPTION, display_name.c_str(),
PW_KEY_MEDIA_CLASS, media_class_value.c_str(),
PW_KEY_NODE_VIRTUAL, "true",
"audio.position", audio_position.c_str(),
nullptr);
#ifndef WARPPIPE_TESTING
pw_properties_set(props, PW_KEY_OBJECT_LINGER, "true");
#endif
if (!props) {
return {Status::Error(StatusCode::kInternal, "failed to allocate stream properties"), 0};
return {Status::Error(StatusCode::kInternal, "failed to allocate node properties"), 0};
}
if (node_group) {
pw_properties_set(props, PW_KEY_NODE_GROUP, node_group);
}
pw_stream* stream = pw_stream_new(core, stream_name.c_str(), props);
if (!stream) {
return {Status::Error(StatusCode::kUnavailable, "failed to create pipewire stream"), 0};
pw_proxy* proxy = reinterpret_cast<pw_proxy*>(
pw_core_create_object(core, "adapter",
PW_TYPE_INTERFACE_Node,
PW_VERSION_NODE,
&props->dict, 0));
pw_properties_free(props);
if (!proxy) {
return {Status::Error(StatusCode::kUnavailable, "failed to create virtual node"), 0};
}
auto stream_data = std::make_unique<StreamData>();
stream_data->stream = stream;
stream_data->proxy = proxy;
#ifndef WARPPIPE_TESTING
stream_data->linger = true;
#endif
stream_data->loop = thread_loop;
stream_data->is_source = is_source;
stream_data->loopback = false;
@ -1031,28 +1113,11 @@ Result<uint32_t> Client::Impl::CreateVirtualStreamLocked(std::string_view name,
stream_data->channels = options.format.channels;
}
pw_stream_add_listener(stream, &stream_data->listener, &kStreamEvents, stream_data.get());
pw_proxy_add_listener(proxy, &stream_data->listener, &kNodeProxyEvents, stream_data.get());
const struct spa_pod* params[1];
uint8_t buffer[1024];
spa_pod_builder builder = SPA_POD_BUILDER_INIT(buffer, sizeof(buffer));
spa_audio_info_raw audio_info{};
audio_info.format = SPA_AUDIO_FORMAT_F32;
audio_info.rate = stream_data->rate;
audio_info.channels = stream_data->channels;
params[0] = spa_format_audio_raw_build(&builder, SPA_PARAM_EnumFormat, &audio_info);
enum pw_direction direction = is_source ? PW_DIRECTION_OUTPUT : PW_DIRECTION_INPUT;
enum pw_stream_flags flags = PW_STREAM_FLAG_MAP_BUFFERS;
int res = pw_stream_connect(stream, direction, PW_ID_ANY, flags, params, 1);
if (res < 0) {
pw_stream_destroy(stream);
return {Status::Error(StatusCode::kUnavailable, "failed to connect pipewire stream"), 0};
}
uint32_t node_id = pw_stream_get_node_id(stream);
uint32_t node_id = SPA_ID_INVALID;
int wait_attempts = 0;
while (node_id == SPA_ID_INVALID && !stream_data->failed && wait_attempts < 3) {
while (node_id == SPA_ID_INVALID && !stream_data->failed && !stream_data->ready && wait_attempts < 3) {
int wait_res = pw_thread_loop_timed_wait(thread_loop, kSyncWaitSeconds);
if (wait_res == -ETIMEDOUT) {
break;
@ -1062,17 +1127,17 @@ Result<uint32_t> Client::Impl::CreateVirtualStreamLocked(std::string_view name,
}
if (stream_data->failed) {
std::string error = stream_data->error.empty() ? "stream entered error state" : stream_data->error;
pw_stream_destroy(stream);
std::string error = stream_data->error.empty() ? "node creation failed" : stream_data->error;
pw_proxy_destroy(proxy);
return {Status::Error(StatusCode::kUnavailable, std::move(error)), 0};
}
node_id = stream_data->node_id;
if (node_id == SPA_ID_INVALID) {
pw_stream_destroy(stream);
return {Status::Error(StatusCode::kTimeout, "timed out waiting for stream node id"), 0};
pw_proxy_destroy(proxy);
return {Status::Error(StatusCode::kTimeout, "timed out waiting for virtual node id"), 0};
}
stream_data->node_id = node_id;
stream_data->ready = true;
{
std::lock_guard<std::mutex> lock(cache_mutex);
@ -1149,14 +1214,24 @@ void Client::Impl::DisconnectLocked() {
streams.swap(virtual_streams);
}
for (auto& entry : links) {
DetachLinkProxy(entry.second.get(), false);
LinkProxy* link = entry.second.get();
if (link) {
spa_hook_remove(&link->listener);
link->proxy = nullptr;
}
}
for (auto& entry : streams) {
StreamData* stream_data = entry.second.get();
if (!stream_data) continue;
if (stream_data->module) {
if (stream_data->linger && stream_data->proxy) {
spa_hook_remove(&stream_data->listener);
stream_data->proxy = nullptr;
} else if (stream_data->module) {
pw_impl_module_destroy(stream_data->module);
stream_data->module = nullptr;
} else if (stream_data->proxy) {
pw_proxy_destroy(stream_data->proxy);
stream_data->proxy = nullptr;
} else if (stream_data->stream) {
pw_stream_disconnect(stream_data->stream);
pw_stream_destroy(stream_data->stream);
@ -1164,11 +1239,17 @@ void Client::Impl::DisconnectLocked() {
}
}
for (auto& entry : auto_link_proxies) {
DetachLinkProxy(entry.get(), false);
if (entry) {
spa_hook_remove(&entry->listener);
entry->proxy = nullptr;
}
}
auto_link_proxies.clear();
for (auto& entry : saved_link_proxies) {
DetachLinkProxy(entry.get(), false);
if (entry) {
spa_hook_remove(&entry->listener);
entry->proxy = nullptr;
}
}
saved_link_proxies.clear();
for (auto& entry : node_proxies) {
@ -1291,17 +1372,13 @@ void Client::Impl::EnforceRulesForLink(uint32_t link_id, uint32_t out_port,
if (!should_destroy) return;
auto active_link_proxy = link_proxies.find(link_id);
if (active_link_proxy != link_proxies.end() &&
active_link_proxy->second && active_link_proxy->second->proxy) {
return;
}
if (link_proxies.count(link_id)) return;
for (const auto& proxy : auto_link_proxies) {
if (proxy && proxy->proxy && proxy->output_port == out_port &&
if (proxy && proxy->output_port == out_port &&
proxy->input_port == in_port) return;
}
for (const auto& proxy : saved_link_proxies) {
if (proxy && proxy->proxy && proxy->output_port == out_port &&
if (proxy && proxy->output_port == out_port &&
proxy->input_port == in_port) return;
}
for (const auto& pair : auto_link_claimed_pairs) {
@ -1442,14 +1519,12 @@ void Client::Impl::ProcessPendingAutoLinks() {
if (target_in == in_port) { is_ours = true; break; }
}
if (!is_ours) {
auto owned = link_proxies.find(link_id);
if (owned != link_proxies.end() && owned->second && owned->second->proxy)
is_ours = true;
if (link_proxies.count(link_id)) is_ours = true;
}
if (!is_ours) {
uint32_t out_port = link_entry.second.output_port.value;
for (const auto& proxy : auto_link_proxies) {
if (proxy && proxy->proxy && proxy->output_port == out_port &&
if (proxy && proxy->output_port == out_port &&
proxy->input_port == in_port) { is_ours = true; break; }
}
}
@ -1464,7 +1539,7 @@ void Client::Impl::ProcessPendingAutoLinks() {
if (!is_ours) {
uint32_t out_port = link_entry.second.output_port.value;
for (const auto& proxy : saved_link_proxies) {
if (proxy && proxy->proxy && proxy->output_port == out_port &&
if (proxy && proxy->output_port == out_port &&
proxy->input_port == in_port) { is_ours = true; break; }
}
}
@ -1512,7 +1587,6 @@ void Client::Impl::CreateAutoLinkAsync(uint32_t output_port, uint32_t input_port
link_data->output_port = output_port;
link_data->input_port = input_port;
pw_proxy_add_listener(proxy, &link_data->listener, &kLinkProxyEvents, link_data.get());
link_data->listener_attached = true;
std::lock_guard<std::mutex> lock(cache_mutex);
auto_link_proxies.push_back(std::move(link_data));
@ -1603,15 +1677,14 @@ void Client::Impl::ProcessSavedLinks() {
if (saved_in == in_port) { is_ours = true; break; }
}
if (!is_ours) {
auto owned = link_proxies.find(link_id);
if (owned != link_proxies.end() && owned->second && owned->second->proxy) {
if (link_proxies.count(link_id)) {
is_ours = true;
}
}
if (!is_ours) {
uint32_t out_port = link_entry.second.output_port.value;
for (const auto& proxy : auto_link_proxies) {
if (proxy && proxy->proxy && proxy->output_port == out_port &&
if (proxy && proxy->output_port == out_port &&
proxy->input_port == in_port) { is_ours = true; break; }
}
}
@ -1626,7 +1699,7 @@ void Client::Impl::ProcessSavedLinks() {
if (!is_ours) {
uint32_t out_port = link_entry.second.output_port.value;
for (const auto& proxy : saved_link_proxies) {
if (proxy && proxy->proxy && proxy->output_port == out_port &&
if (proxy && proxy->output_port == out_port &&
proxy->input_port == in_port) { is_ours = true; break; }
}
}
@ -1670,7 +1743,6 @@ void Client::Impl::CreateSavedLinkAsync(uint32_t output_port,
link_data->input_port = input_port;
pw_proxy_add_listener(proxy, &link_data->listener, &kLinkProxyEvents,
link_data.get());
link_data->listener_attached = true;
std::lock_guard<std::mutex> lock(cache_mutex);
saved_link_proxies.push_back(std::move(link_data));
@ -1730,7 +1802,7 @@ void Client::Impl::AutoSave() {
std::lock_guard<std::mutex> lock(cache_mutex);
std::vector<SavedLink> live;
for (const auto& entry : link_proxies) {
if (!entry.second || !entry.second->proxy) {
if (!entry.second) {
continue;
}
auto link_it = links.find(entry.first);
@ -1759,7 +1831,7 @@ void Client::Impl::AutoSave() {
links_array.push_back(std::move(link_obj));
}
for (const auto& lp : saved_link_proxies) {
if (!lp || !lp->proxy || lp->id == SPA_ID_INVALID) continue;
if (!lp || lp->id == SPA_ID_INVALID) continue;
auto link_it = links.find(lp->id);
if (link_it == links.end()) continue;
const Link& link = link_it->second;
@ -2130,6 +2202,9 @@ Status Client::RemoveNode(NodeId node) {
}
pw_impl_module_destroy(owned_stream->module);
owned_stream->module = nullptr;
} else if (owned_stream->proxy) {
pw_proxy_destroy(owned_stream->proxy);
owned_stream->proxy = nullptr;
} else if (owned_stream->stream) {
pw_stream_disconnect(owned_stream->stream);
pw_stream_destroy(owned_stream->stream);
@ -2154,12 +2229,15 @@ Status Client::SetNodeVolume(NodeId node, float volume, bool mute) {
pw_thread_loop_lock(impl_->thread_loop);
bool is_virtual = false;
{
std::lock_guard<std::mutex> lock(impl_->cache_mutex);
if (impl_->nodes.find(node.value) == impl_->nodes.end()) {
auto node_it = impl_->nodes.find(node.value);
if (node_it == impl_->nodes.end()) {
pw_thread_loop_unlock(impl_->thread_loop);
return Status::Error(StatusCode::kNotFound, "node not found");
}
is_virtual = node_it->second.is_virtual;
}
uint32_t n_channels;
@ -2171,7 +2249,8 @@ Status Client::SetNodeVolume(NodeId node, float volume, bool mute) {
if (streamIt != impl_->virtual_streams.end() && streamIt->second->stream) {
own_stream = streamIt->second->stream;
n_channels = streamIt->second->channels;
} else {
}
if (!own_stream) {
auto proxyIt = impl_->node_proxies.find(node.value);
if (proxyIt == impl_->node_proxies.end() || !proxyIt->second->proxy) {
pw_thread_loop_unlock(impl_->thread_loop);
@ -2187,8 +2266,8 @@ Status Client::SetNodeVolume(NodeId node, float volume, bool mute) {
uint8_t buffer[512];
spa_pod_builder builder = SPA_POD_BUILDER_INIT(buffer, sizeof(buffer));
uint32_t vol_prop = own_stream ? SPA_PROP_monitorVolumes : SPA_PROP_channelVolumes;
uint32_t mute_prop = own_stream ? SPA_PROP_monitorMute : SPA_PROP_mute;
uint32_t vol_prop = is_virtual ? SPA_PROP_monitorVolumes : SPA_PROP_channelVolumes;
uint32_t mute_prop = is_virtual ? SPA_PROP_monitorMute : SPA_PROP_mute;
spa_pod_frame obj_frame;
spa_pod_builder_push_object(&builder, &obj_frame,
@ -2466,7 +2545,6 @@ Result<Link> Client::CreateLink(PortId output, PortId input, const LinkOptions&
link_proxy->output_port = output.value;
link_proxy->input_port = input.value;
pw_proxy_add_listener(proxy, &link_proxy->listener, &kLinkProxyEvents, link_proxy.get());
link_proxy->listener_attached = true;
int wait_attempts = 0;
while (link_proxy->id == SPA_ID_INVALID && !link_proxy->failed && wait_attempts < 3) {
@ -2487,13 +2565,13 @@ Result<Link> Client::CreateLink(PortId output, PortId input, const LinkOptions&
if (link_proxy->failed) {
std::string error = link_proxy->error.empty() ? "link creation failed" : link_proxy->error;
remove_pending();
DetachLinkProxy(link_proxy.get(), true);
pw_proxy_destroy(proxy);
pw_thread_loop_unlock(impl_->thread_loop);
return {Status::Error(StatusCode::kUnavailable, std::move(error)), {}};
}
if (link_proxy->id == SPA_ID_INVALID) {
remove_pending();
DetachLinkProxy(link_proxy.get(), true);
pw_proxy_destroy(proxy);
pw_thread_loop_unlock(impl_->thread_loop);
return {Status::Error(StatusCode::kTimeout, "timed out waiting for link id"), {}};
}
@ -2504,7 +2582,7 @@ Result<Link> Client::CreateLink(PortId output, PortId input, const LinkOptions&
link.input_port = input;
{
std::lock_guard<std::mutex> lock(impl_->cache_mutex);
impl_->link_proxies[link_proxy->id] = std::move(link_proxy);
impl_->link_proxies.emplace(link_proxy->id, std::move(link_proxy));
impl_->links[link.id.value] = link;
std::erase_if(impl_->pending_link_pairs, [&](const auto& p) {
return p.first == output.value && p.second == input.value;
@ -2590,9 +2668,14 @@ Status Client::RemoveLink(LinkId link) {
}
auto it = impl_->link_proxies.find(link.value);
if (it != impl_->link_proxies.end()) {
if (it->second && it->second->proxy) {
spa_hook_remove(&it->second->listener);
pw_proxy_destroy(it->second->proxy);
}
if (impl_->registry) {
pw_registry_destroy(impl_->registry, link.value);
}
impl_->link_proxies.erase(it);
auto link_it2 = impl_->links.find(link.value);
if (link_it2 != impl_->links.end()) {
uint32_t op = link_it2->second.output_port.value;
@ -2622,6 +2705,26 @@ Status Client::RemoveLink(LinkId link) {
impl_->links.erase(link_it);
}
if (out_port && in_port) {
for (auto& p : impl_->saved_link_proxies) {
if (p && p->output_port == out_port && p->input_port == in_port) {
spa_hook_remove(&p->listener);
if (p->proxy) pw_proxy_destroy(p->proxy);
p->proxy = nullptr;
}
}
std::erase_if(impl_->saved_link_proxies, [&](const auto& p) {
return p && p->output_port == out_port && p->input_port == in_port;
});
for (auto& p : impl_->auto_link_proxies) {
if (p && p->output_port == out_port && p->input_port == in_port) {
spa_hook_remove(&p->listener);
if (p->proxy) pw_proxy_destroy(p->proxy);
p->proxy = nullptr;
}
}
std::erase_if(impl_->auto_link_proxies, [&](const auto& p) {
return p && p->output_port == out_port && p->input_port == in_port;
});
std::erase_if(impl_->auto_link_claimed_pairs, [&](const auto& pair) {
return pair.first == out_port && pair.second == in_port;
});
@ -2785,7 +2888,18 @@ Status Client::RemoveRouteRule(RuleId id) {
}
{
std::lock_guard<std::mutex> lock(impl_->cache_mutex);
(void)pairs_to_remove;
for (const auto& pair : pairs_to_remove) {
for (auto& p : impl_->auto_link_proxies) {
if (p && p->output_port == pair.first &&
p->input_port == pair.second) {
spa_hook_remove(&p->listener);
}
}
std::erase_if(impl_->auto_link_proxies, [&](const auto& p) {
return p && p->output_port == pair.first &&
p->input_port == pair.second;
});
}
}
pw_thread_loop_unlock(impl_->thread_loop);
}

View file

@ -1671,42 +1671,6 @@ TEST_CASE("findPwNodeIdByName returns 0 for ghost nodes without pw mapping") {
REQUIRE(model.findPwNodeIdByName("ghost-lookup") == 100220);
}
TEST_CASE("connectionPeakLevel uses source node activity, with application fallback") {
auto tc = TestClient::Create();
if (!tc.available()) { SUCCEED("PipeWire unavailable"); return; }
ensureApp();
REQUIRE(tc.client->Test_InsertNode(
MakeNode(100240, "app-out", "Stream/Output/Audio", "Firefox")).ok());
REQUIRE(tc.client->Test_InsertNode(
MakeNode(100241, "sink-out", "Audio/Sink")).ok());
REQUIRE(tc.client->Test_InsertNode(
MakeNode(100242, "hw-in", "Audio/Sink")).ok());
WarpGraphModel model(tc.client.get());
model.refreshFromClient();
auto appQt = model.qtNodeIdForPw(100240);
auto sinkQt = model.qtNodeIdForPw(100241);
auto hwQt = model.qtNodeIdForPw(100242);
REQUIRE(appQt != 0);
REQUIRE(sinkQt != 0);
REQUIRE(hwQt != 0);
model.setNodePeakLevel(appQt, 0.0f);
model.setNodePeakLevel(sinkQt, 0.0f);
model.setNodePeakLevel(hwQt, 0.8f);
REQUIRE(model.connectionPeakLevel(QtNodes::ConnectionId{sinkQt, 0u, hwQt, 0u}) ==
Catch::Approx(0.0f));
REQUIRE(model.connectionPeakLevel(QtNodes::ConnectionId{appQt, 0u, hwQt, 0u}) ==
Catch::Approx(0.8f));
model.setNodePeakLevel(appQt, 0.4f);
REQUIRE(model.connectionPeakLevel(QtNodes::ConnectionId{appQt, 0u, hwQt, 0u}) ==
Catch::Approx(0.4f));
}
TEST_CASE("saveLayout stores and loadLayout restores view state") {
auto tc = TestClient::Create();
if (!tc.available()) { SUCCEED("PipeWire unavailable"); return; }