diff --git a/CMakeLists.txt b/CMakeLists.txt index f227e85..281a41f 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/GraphEditorWidget.cpp b/gui/GraphEditorWidget.cpp index d8a4b18..57047ce 100644 --- a/gui/GraphEditorWidget.cpp +++ b/gui/GraphEditorWidget.cpp @@ -1,5 +1,4 @@ #include "AudioLevelMeter.h" -#include "BezierConnectionPainter.h" #include "GraphEditorWidget.h" #include "PresetManager.h" #include "SquareConnectionPainter.h" @@ -10,7 +9,7 @@ #include #include #include - +#include "WarpBezierConnectionPainter.h" #include #include #include @@ -204,8 +203,6 @@ 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); @@ -569,17 +566,21 @@ GraphEditorWidget::GraphEditorWidget(warppipe::Client *client, connect(m_model, &QtNodes::AbstractGraphModel::nodeCreated, this, [this](QtNodes::NodeId nodeId) { wireVolumeWidget(nodeId); - this->scheduleSidebarRebuild(); + rebuildMixerStrips(); + rebuildNodeMeters(); + rebuildRulesList(); }); 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, @@ -1543,7 +1544,7 @@ void GraphEditorWidget::rebuildMixerStrips() { while (layout->count() > 0) { auto *item = layout->takeAt(0); if (item->widget()) - delete item->widget(); + item->widget()->deleteLater(); delete item; } m_mixerStrips.clear(); @@ -1721,7 +1722,7 @@ void GraphEditorWidget::rebuildNodeMeters() { while (layout->count() > 0) { auto *item = layout->takeAt(0); if (item->widget()) - delete item->widget(); + item->widget()->deleteLater(); delete item; } m_nodeMeters.clear(); @@ -1791,20 +1792,6 @@ 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; @@ -1816,7 +1803,7 @@ void GraphEditorWidget::rebuildRulesList() { while (layout->count() > 0) { auto *item = layout->takeAt(0); if (item->widget()) - delete item->widget(); + item->widget()->deleteLater(); delete item; } @@ -1940,7 +1927,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 6f57dd3..b7e0dbe 100644 --- a/gui/GraphEditorWidget.h +++ b/gui/GraphEditorWidget.h @@ -96,7 +96,6 @@ private: void rebuildMixerStrips(); void updateMeters(); void rebuildNodeMeters(); - void scheduleSidebarRebuild(); void rebuildRulesList(); void showAddRuleDialog(const std::string &prefillApp = {}, const std::string &prefillBin = {}, @@ -136,7 +135,6 @@ 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 226e897..763aabb 100644 --- a/gui/SquareConnectionPainter.cpp +++ b/gui/SquareConnectionPainter.cpp @@ -66,17 +66,16 @@ 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; - if (mdl && connectionAlive) { - auto ch = mdl->connectionChannel(cId); - spread = (static_cast(ch.index) - (ch.count - 1) / 2.0) - * 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; + } } double const dy = in.y() - out.y(); @@ -126,9 +125,12 @@ QPainterPath SquareConnectionPainter::orthogonalPath( // double railOffset = 0.0; - if (mdl && connectionAlive) { - auto ch = mdl->connectionChannel(cId); - railOffset = static_cast(ch.index) * kSpacing; + if (sceneForChannel) { + auto *mdl2 = dynamic_cast(&sceneForChannel->graphModel()); + if (mdl2) { + auto ch = mdl2->connectionChannel(cId); + railOffset = static_cast(ch.index) * kSpacing; + } } double rightX = out.x() + kMinStub + railOffset; @@ -203,19 +205,7 @@ void SquareConnectionPainter::paint( auto *model = dynamic_cast(&scene->graphModel()); if (model) { auto cId = cgo.connectionId(); - 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)); - } - } - } + peakLevel = model->connectionPeakLevel(cId); } } diff --git a/gui/BezierConnectionPainter.cpp b/gui/WarpBezierConnectionPainter.cpp similarity index 83% rename from gui/BezierConnectionPainter.cpp rename to gui/WarpBezierConnectionPainter.cpp index 54155c6..75d23c0 100644 --- a/gui/BezierConnectionPainter.cpp +++ b/gui/WarpBezierConnectionPainter.cpp @@ -1,4 +1,4 @@ -#include "BezierConnectionPainter.h" +#include "WarpBezierConnectionPainter.h" #include "WarpGraphModel.h" #include @@ -14,11 +14,10 @@ #include #include -QPainterPath BezierConnectionPainter::cubicPath( +QPainterPath WarpBezierConnectionPainter::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); @@ -26,7 +25,7 @@ QPainterPath BezierConnectionPainter::cubicPath( return cubic; } -void BezierConnectionPainter::paint( +void WarpBezierConnectionPainter::paint( QPainter *painter, QtNodes::ConnectionGraphicsObject const &cgo) const { auto const &style = QtNodes::StyleCollection::connectionStyle(); @@ -43,19 +42,7 @@ void BezierConnectionPainter::paint( auto *model = dynamic_cast(&scene->graphModel()); if (model) { auto cId = cgo.connectionId(); - 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)); - } - } - } + peakLevel = model->connectionPeakLevel(cId); } } @@ -111,7 +98,7 @@ void BezierConnectionPainter::paint( painter->drawEllipse(cgo.in(), pointRadius, pointRadius); } -QPainterPath BezierConnectionPainter::getPainterStroke( +QPainterPath WarpBezierConnectionPainter::getPainterStroke( QtNodes::ConnectionGraphicsObject const &cgo) const { auto cubic = cubicPath(cgo); @@ -126,6 +113,5 @@ QPainterPath BezierConnectionPainter::getPainterStroke( QPainterPathStroker stroker; stroker.setWidth(10.0); - return stroker.createStroke(result); } diff --git a/gui/BezierConnectionPainter.h b/gui/WarpBezierConnectionPainter.h similarity index 82% rename from gui/BezierConnectionPainter.h rename to gui/WarpBezierConnectionPainter.h index 1c51473..a9377a3 100644 --- a/gui/BezierConnectionPainter.h +++ b/gui/WarpBezierConnectionPainter.h @@ -2,7 +2,7 @@ #include -class BezierConnectionPainter : public QtNodes::AbstractConnectionPainter { +class WarpBezierConnectionPainter : public QtNodes::AbstractConnectionPainter { public: void paint(QPainter *painter, QtNodes::ConnectionGraphicsObject const &cgo) const override; diff --git a/gui/WarpGraphModel.cpp b/gui/WarpGraphModel.cpp index 101e97e..87cf2d2 100644 --- a/gui/WarpGraphModel.cpp +++ b/gui/WarpGraphModel.cpp @@ -343,8 +343,6 @@ bool WarpGraphModel::deleteConnection( } m_connections.erase(it); - m_connectionChannels.erase(connectionId); - recomputeConnectionChannels(); Q_EMIT connectionDeleted(connectionId); return true; } @@ -376,7 +374,6 @@ 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); @@ -484,14 +481,8 @@ 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) { @@ -1482,6 +1473,23 @@ 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 1142a53..dcc2c3d 100644 --- a/gui/WarpGraphModel.h +++ b/gui/WarpGraphModel.h @@ -109,6 +109,7 @@ 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 b8d6806..6ee2aab 100644 --- a/src/warppipe.cpp +++ b/src/warppipe.cpp @@ -107,12 +107,10 @@ 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; @@ -126,6 +124,7 @@ 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; @@ -135,6 +134,20 @@ 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) { @@ -152,6 +165,8 @@ 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); @@ -163,6 +178,8 @@ 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) { @@ -177,42 +194,6 @@ 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) { @@ -329,12 +310,10 @@ void NodeMeterProcess(void* data) { had_data = true; pw_stream_queue_buffer(meter->stream, buf); } - if (!had_data) { - left = 0.0f; - right = 0.0f; + if (had_data) { + meter->peak_left.store(left, std::memory_order_relaxed); + meter->peak_right.store(right, std::memory_order_relaxed); } - meter->peak_left.store(left, std::memory_order_relaxed); - meter->peak_right.store(right, std::memory_order_relaxed); } static const pw_stream_events kNodeMeterEvents = { @@ -475,29 +454,7 @@ struct Client::Impl { void Client::Impl::NodeInfoChanged(void* data, const struct pw_node_info* info) { auto* np = static_cast(data); - 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; + if (!np || !info || np->params_subscribed) return; for (uint32_t i = 0; i < info->n_params; ++i) { if (info->params[i].id == SPA_PARAM_Props && @@ -570,18 +527,15 @@ void Client::Impl::NodeParamChanged(void* data, int, uint32_t id, } } - bool uses_monitor_volume = false; + bool is_virtual = false; { std::lock_guard lock(impl->cache_mutex); - auto node_it = impl->nodes.find(np->node_id); - if (node_it != impl->nodes.end()) { - uses_monitor_volume = node_it->second.is_virtual; - } + is_virtual = impl->virtual_streams.count(np->node_id) > 0; } - 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; + 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; bool changed = false; { @@ -762,7 +716,6 @@ 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); @@ -920,25 +873,14 @@ Result Client::Impl::CreateVirtualStreamLocked(std::string_view name, { std::lock_guard lock(cache_mutex); - for (const auto& entry : virtual_streams) { - if (entry.second && entry.second->name == stream_name) { + for (const auto& entry : nodes) { + if (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}; + for (const auto& entry : virtual_streams) { + if (entry.second && entry.second->name == stream_name) { + return {Status::Error(StatusCode::kInvalidArgument, "duplicate node name"), 0}; } } if (options.behavior == VirtualBehavior::kLoopback && options.target_node) { @@ -1052,53 +994,29 @@ Result Client::Impl::CreateVirtualStreamLocked(std::string_view name, return {Status::Ok(), node_id}; } - // 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 + 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); if (!props) { - return {Status::Error(StatusCode::kInternal, "failed to allocate node properties"), 0}; + return {Status::Error(StatusCode::kInternal, "failed to allocate stream properties"), 0}; } if (node_group) { pw_properties_set(props, PW_KEY_NODE_GROUP, node_group); } - 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}; + 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}; } auto stream_data = std::make_unique(); - stream_data->proxy = proxy; -#ifndef WARPPIPE_TESTING - stream_data->linger = true; -#endif + stream_data->stream = stream; stream_data->loop = thread_loop; stream_data->is_source = is_source; stream_data->loopback = false; @@ -1113,11 +1031,28 @@ Result Client::Impl::CreateVirtualStreamLocked(std::string_view name, stream_data->channels = options.format.channels; } - pw_proxy_add_listener(proxy, &stream_data->listener, &kNodeProxyEvents, stream_data.get()); + pw_stream_add_listener(stream, &stream_data->listener, &kStreamEvents, stream_data.get()); - uint32_t node_id = SPA_ID_INVALID; + 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); int wait_attempts = 0; - while (node_id == SPA_ID_INVALID && !stream_data->failed && !stream_data->ready && wait_attempts < 3) { + while (node_id == SPA_ID_INVALID && !stream_data->failed && wait_attempts < 3) { int wait_res = pw_thread_loop_timed_wait(thread_loop, kSyncWaitSeconds); if (wait_res == -ETIMEDOUT) { break; @@ -1127,17 +1062,17 @@ Result Client::Impl::CreateVirtualStreamLocked(std::string_view name, } if (stream_data->failed) { - std::string error = stream_data->error.empty() ? "node creation failed" : stream_data->error; - pw_proxy_destroy(proxy); + std::string error = stream_data->error.empty() ? "stream entered error state" : stream_data->error; + pw_stream_destroy(stream); return {Status::Error(StatusCode::kUnavailable, std::move(error)), 0}; } - node_id = stream_data->node_id; if (node_id == SPA_ID_INVALID) { - pw_proxy_destroy(proxy); - return {Status::Error(StatusCode::kTimeout, "timed out waiting for virtual node id"), 0}; + pw_stream_destroy(stream); + return {Status::Error(StatusCode::kTimeout, "timed out waiting for stream node id"), 0}; } + stream_data->node_id = node_id; stream_data->ready = true; { std::lock_guard lock(cache_mutex); @@ -1214,24 +1149,14 @@ void Client::Impl::DisconnectLocked() { streams.swap(virtual_streams); } for (auto& entry : links) { - LinkProxy* link = entry.second.get(); - if (link) { - spa_hook_remove(&link->listener); - link->proxy = nullptr; - } + DetachLinkProxy(entry.second.get(), false); } for (auto& entry : streams) { StreamData* stream_data = entry.second.get(); if (!stream_data) continue; - if (stream_data->linger && stream_data->proxy) { - spa_hook_remove(&stream_data->listener); - stream_data->proxy = nullptr; - } else if (stream_data->module) { + 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); @@ -1239,17 +1164,11 @@ void Client::Impl::DisconnectLocked() { } } for (auto& entry : auto_link_proxies) { - if (entry) { - spa_hook_remove(&entry->listener); - entry->proxy = nullptr; - } + DetachLinkProxy(entry.get(), false); } auto_link_proxies.clear(); for (auto& entry : saved_link_proxies) { - if (entry) { - spa_hook_remove(&entry->listener); - entry->proxy = nullptr; - } + DetachLinkProxy(entry.get(), false); } saved_link_proxies.clear(); for (auto& entry : node_proxies) { @@ -1372,13 +1291,17 @@ void Client::Impl::EnforceRulesForLink(uint32_t link_id, uint32_t out_port, if (!should_destroy) return; - if (link_proxies.count(link_id)) 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; + } for (const auto& proxy : auto_link_proxies) { - if (proxy && proxy->output_port == out_port && + if (proxy && proxy->proxy && proxy->output_port == out_port && proxy->input_port == in_port) return; } for (const auto& proxy : saved_link_proxies) { - if (proxy && proxy->output_port == out_port && + if (proxy && proxy->proxy && proxy->output_port == out_port && proxy->input_port == in_port) return; } for (const auto& pair : auto_link_claimed_pairs) { @@ -1519,12 +1442,14 @@ void Client::Impl::ProcessPendingAutoLinks() { if (target_in == in_port) { is_ours = true; break; } } if (!is_ours) { - if (link_proxies.count(link_id)) is_ours = true; + auto owned = link_proxies.find(link_id); + if (owned != link_proxies.end() && owned->second && owned->second->proxy) + 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->output_port == out_port && + if (proxy && proxy->proxy && proxy->output_port == out_port && proxy->input_port == in_port) { is_ours = true; break; } } } @@ -1539,7 +1464,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->output_port == out_port && + if (proxy && proxy->proxy && proxy->output_port == out_port && proxy->input_port == in_port) { is_ours = true; break; } } } @@ -1587,6 +1512,7 @@ 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)); @@ -1677,14 +1603,15 @@ void Client::Impl::ProcessSavedLinks() { if (saved_in == in_port) { is_ours = true; break; } } if (!is_ours) { - if (link_proxies.count(link_id)) { + auto owned = link_proxies.find(link_id); + if (owned != link_proxies.end() && owned->second && owned->second->proxy) { 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->output_port == out_port && + if (proxy && proxy->proxy && proxy->output_port == out_port && proxy->input_port == in_port) { is_ours = true; break; } } } @@ -1699,7 +1626,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->output_port == out_port && + if (proxy && proxy->proxy && proxy->output_port == out_port && proxy->input_port == in_port) { is_ours = true; break; } } } @@ -1743,6 +1670,7 @@ 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)); @@ -1802,7 +1730,7 @@ void Client::Impl::AutoSave() { std::lock_guard lock(cache_mutex); std::vector live; for (const auto& entry : link_proxies) { - if (!entry.second) { + if (!entry.second || !entry.second->proxy) { continue; } auto link_it = links.find(entry.first); @@ -1831,7 +1759,7 @@ void Client::Impl::AutoSave() { links_array.push_back(std::move(link_obj)); } for (const auto& lp : saved_link_proxies) { - if (!lp || lp->id == SPA_ID_INVALID) continue; + if (!lp || !lp->proxy || 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; @@ -2202,9 +2130,6 @@ 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); @@ -2229,15 +2154,12 @@ 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); - auto node_it = impl_->nodes.find(node.value); - if (node_it == impl_->nodes.end()) { + if (impl_->nodes.find(node.value) == 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; @@ -2249,8 +2171,7 @@ 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; - } - if (!own_stream) { + } else { auto proxyIt = impl_->node_proxies.find(node.value); if (proxyIt == impl_->node_proxies.end() || !proxyIt->second->proxy) { pw_thread_loop_unlock(impl_->thread_loop); @@ -2266,8 +2187,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 = is_virtual ? SPA_PROP_monitorVolumes : SPA_PROP_channelVolumes; - uint32_t mute_prop = is_virtual ? SPA_PROP_monitorMute : SPA_PROP_mute; + uint32_t vol_prop = own_stream ? SPA_PROP_monitorVolumes : SPA_PROP_channelVolumes; + uint32_t mute_prop = own_stream ? SPA_PROP_monitorMute : SPA_PROP_mute; spa_pod_frame obj_frame; spa_pod_builder_push_object(&builder, &obj_frame, @@ -2545,6 +2466,7 @@ 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) { @@ -2565,13 +2487,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(); - pw_proxy_destroy(proxy); + DetachLinkProxy(link_proxy.get(), true); pw_thread_loop_unlock(impl_->thread_loop); return {Status::Error(StatusCode::kUnavailable, std::move(error)), {}}; } if (link_proxy->id == SPA_ID_INVALID) { remove_pending(); - pw_proxy_destroy(proxy); + DetachLinkProxy(link_proxy.get(), true); pw_thread_loop_unlock(impl_->thread_loop); return {Status::Error(StatusCode::kTimeout, "timed out waiting for link id"), {}}; } @@ -2582,7 +2504,7 @@ Result Client::CreateLink(PortId output, PortId input, const LinkOptions& link.input_port = input; { std::lock_guard lock(impl_->cache_mutex); - impl_->link_proxies.emplace(link_proxy->id, std::move(link_proxy)); + impl_->link_proxies[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; @@ -2668,14 +2590,9 @@ 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; @@ -2705,26 +2622,6 @@ 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; }); @@ -2888,18 +2785,7 @@ Status Client::RemoveRouteRule(RuleId id) { } { std::lock_guard lock(impl_->cache_mutex); - 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; - }); - } + (void)pairs_to_remove; } pw_thread_loop_unlock(impl_->thread_loop); } diff --git a/tests/gui/warppipe_gui_tests.cpp b/tests/gui/warppipe_gui_tests.cpp index fd1d227..b248f0e 100644 --- a/tests/gui/warppipe_gui_tests.cpp +++ b/tests/gui/warppipe_gui_tests.cpp @@ -1671,6 +1671,42 @@ 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; }