diff --git a/CMakeLists.txt b/CMakeLists.txt index 281a41f..f227e85 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -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) diff --git a/gui/WarpBezierConnectionPainter.cpp b/gui/BezierConnectionPainter.cpp similarity index 83% rename from gui/WarpBezierConnectionPainter.cpp rename to gui/BezierConnectionPainter.cpp index 75d23c0..54155c6 100644 --- a/gui/WarpBezierConnectionPainter.cpp +++ b/gui/BezierConnectionPainter.cpp @@ -1,4 +1,4 @@ -#include "WarpBezierConnectionPainter.h" +#include "BezierConnectionPainter.h" #include "WarpGraphModel.h" #include @@ -14,10 +14,11 @@ #include #include -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(&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); } diff --git a/gui/WarpBezierConnectionPainter.h b/gui/BezierConnectionPainter.h similarity index 82% rename from gui/WarpBezierConnectionPainter.h rename to gui/BezierConnectionPainter.h index a9377a3..1c51473 100644 --- a/gui/WarpBezierConnectionPainter.h +++ b/gui/BezierConnectionPainter.h @@ -2,7 +2,7 @@ #include -class WarpBezierConnectionPainter : public QtNodes::AbstractConnectionPainter { +class BezierConnectionPainter : public QtNodes::AbstractConnectionPainter { public: void paint(QPainter *painter, QtNodes::ConnectionGraphicsObject const &cgo) const override; diff --git a/gui/GraphEditorWidget.cpp b/gui/GraphEditorWidget.cpp index 57047ce..d8a4b18 100644 --- a/gui/GraphEditorWidget.cpp +++ b/gui/GraphEditorWidget.cpp @@ -1,4 +1,5 @@ #include "AudioLevelMeter.h" +#include "BezierConnectionPainter.h" #include "GraphEditorWidget.h" #include "PresetManager.h" #include "SquareConnectionPainter.h" @@ -9,7 +10,7 @@ #include #include #include -#include "WarpBezierConnectionPainter.h" + #include #include #include @@ -203,6 +204,8 @@ GraphEditorWidget::GraphEditorWidget(warppipe::Client *client, "UseDataDefinedColors": false }})"); + m_scene->setConnectionPainter(std::make_unique()); + 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()); } else { m_scene->setConnectionPainter( - std::make_unique()); + std::make_unique()); } for (auto *item : m_scene->items()) { diff --git a/gui/GraphEditorWidget.h b/gui/GraphEditorWidget.h index b7e0dbe..6f57dd3 100644 --- a/gui/GraphEditorWidget.h +++ b/gui/GraphEditorWidget.h @@ -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 m_mixerStrips; + bool m_sidebarRebuildPending = false; QTimer *m_meterTimer = nullptr; AudioLevelMeter *m_masterMeterL = nullptr; diff --git a/gui/SquareConnectionPainter.cpp b/gui/SquareConnectionPainter.cpp index 763aabb..226e897 100644 --- a/gui/SquareConnectionPainter.cpp +++ b/gui/SquareConnectionPainter.cpp @@ -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(&sceneForChannel->graphModel()) + : nullptr; + bool connectionAlive = mdl ? mdl->connectionExists(cId) : true; double spread = static_cast(cId.outPortIndex) * kSpacing; - auto *sceneForChannel = cgo.nodeScene(); - if (sceneForChannel) { - auto *mdl = dynamic_cast(&sceneForChannel->graphModel()); - if (mdl) { - auto ch = mdl->connectionChannel(cId); - spread = (static_cast(ch.index) - (ch.count - 1) / 2.0) - * kSpacing; - } + if (mdl && connectionAlive) { + auto ch = mdl->connectionChannel(cId); + spread = (static_cast(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(&sceneForChannel->graphModel()); - if (mdl2) { - auto ch = mdl2->connectionChannel(cId); - railOffset = static_cast(ch.index) * kSpacing; - } + if (mdl && connectionAlive) { + auto ch = mdl->connectionChannel(cId); + railOffset = static_cast(ch.index) * kSpacing; } double rightX = out.x() + kMinStub + railOffset; @@ -205,7 +203,19 @@ void SquareConnectionPainter::paint( auto *model = dynamic_cast(&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)); + } + } + } } } diff --git a/gui/WarpGraphModel.cpp b/gui/WarpGraphModel.cpp index 87cf2d2..101e97e 100644 --- a/gui/WarpGraphModel.cpp +++ b/gui/WarpGraphModel.cpp @@ -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(); diff --git a/gui/WarpGraphModel.h b/gui/WarpGraphModel.h index dcc2c3d..1142a53 100644 --- a/gui/WarpGraphModel.h +++ b/gui/WarpGraphModel.h @@ -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; diff --git a/src/warppipe.cpp b/src/warppipe.cpp index 6ee2aab..b8d6806 100644 --- a/src/warppipe.cpp +++ b/src/warppipe.cpp @@ -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(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(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(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(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(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(data); - if (!np || !info || np->params_subscribed) return; + if (!np || !info) return; + + auto* impl = static_cast(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 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 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 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 Client::Impl::CreateVirtualStreamLocked(std::string_view name, { std::lock_guard 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(); + 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 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_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(); - 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 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 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 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 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 lock(cache_mutex); saved_link_proxies.push_back(std::move(link_data)); @@ -1730,7 +1802,7 @@ void Client::Impl::AutoSave() { std::lock_guard lock(cache_mutex); std::vector 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 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 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 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 Client::CreateLink(PortId output, PortId input, const LinkOptions& link.input_port = input; { std::lock_guard 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 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); } diff --git a/tests/gui/warppipe_gui_tests.cpp b/tests/gui/warppipe_gui_tests.cpp index b248f0e..fd1d227 100644 --- a/tests/gui/warppipe_gui_tests.cpp +++ b/tests/gui/warppipe_gui_tests.cpp @@ -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; }