< Summary

Line coverage
84%
Covered lines: 84
Uncovered lines: 16
Coverable lines: 100
Total lines: 162
Line coverage: 84%
Branch coverage
N/A
Covered branches: 0
Total branches: 0
Branch coverage: N/A
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

File(s)

C:\Users\Kremenchuk\actions-runner\_work\Holoflow\Holoflow\src\holoflow_event\src\router.cc

#LineLine coverage
 1// Copyright 2025 Digital Holography Foundation
 2//
 3// Licensed under the Apache License, Version 2.0 (the "License");
 4// you may not use this file except in compliance with the License.
 5// You may obtain a copy of the License at
 6//
 7//     http://www.apache.org/licenses/LICENSE-2.0
 8//
 9// Unless required by applicable law or agreed to in writing, software
 10// distributed under the License is distributed on an "AS IS" BASIS,
 11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 12// See the License for the specific language governing permissions and
 13// limitations under the License.
 14
 15#include "holoflow_event/router.hh"
 16
 17#include <iostream>
 18
 19namespace holoflow_event {
 20
 121MailboxCounters::Snapshot MailboxCounters::snapshot() const noexcept {
 122  return Snapshot{received_events.load(std::memory_order_relaxed),
 23                  dropped_events.load(std::memory_order_relaxed),
 24                  sent_events.load(std::memory_order_relaxed)};
 125}
 26
 27EventReader::EventReader(BoundedMPMC<Event> *queue, MailboxCounters *counters)
 128    : queue_(queue), counters_(counters) {}
 29
 130std::optional<Event> EventReader::try_pop() noexcept {
 131  auto event = queue_->try_pop();
 132  if (event) {
 133    counters_->received_events.fetch_add(1, std::memory_order_relaxed);
 34  }
 135  return event;
 136}
 37
 138size_t EventReader::drop_all() noexcept {
 139  size_t dropped = 0;
 140  while (true) {
 141    auto event = queue_->try_pop();
 142    if (!event) {
 143      break;
 44    }
 145    dropped++;
 146  }
 147  counters_->dropped_events.fetch_add(dropped, std::memory_order_relaxed);
 148  return dropped;
 149}
 50
 151MailboxCounters::Snapshot EventReader::counters() const noexcept { return counters_->snapshot(); }
 52
 53EventWriter::EventWriter(BoundedMPMC<Event> *queue, MailboxCounters *counters)
 154    : queue_(queue), counters_(counters) {}
 55
 156bool EventWriter::try_push(Event &&event) noexcept {
 157  bool ok = queue_->try_push(std::move(event));
 158  counters_->sent_events.fetch_add(ok, std::memory_order_relaxed);
 159  counters_->dropped_events.fetch_add(!ok, std::memory_order_relaxed);
 160  return ok;
 161}
 62
 063MailboxCounters::Snapshot EventWriter::counters() const noexcept { return counters_->snapshot(); }
 64
 65Router::Router(const Config &config)
 166    : config_(config), ui_to_router_(config.ui_to_router_capacity),
 167      router_to_ui_(config.router_to_ui_capacity),
 168      nodes_to_router_(config.nodes_to_router_capacity) {}
 69
 170Router::NodeHandles Router::bind_node(const NodeId &node_id) {
 171  auto queue          = std::make_unique<BoundedMPMC<Event>>(config_.router_to_node_capacity);
 172  auto [it, inserted] = router_to_nodes_.try_emplace(node_id, std::move(queue));
 173  auto &counters      = router_to_node_counters_[node_id];
 74
 175  if (!inserted) {
 076    std::cerr << "[holoflow_event::Router] Error: Node '" << node_id
 77              << "' is already bound to the router." << std::endl;
 078    std::abort();
 79  }
 80
 181  EventWriter w{&nodes_to_router_, &nodes_to_router_counters_};
 182  EventReader r{it->second.get(), &counters};
 183  return NodeHandles{r, w};
 184}
 85
 186bool Router::ui_try_send(const NodeId &node_id, nlohmann::json &&data, time_point ts) noexcept {
 187  Event event{EventDirection::ToNode, node_id, std::move(data), ts};
 188  bool  ok = ui_to_router_.try_push(std::move(event));
 189  ui_to_router_counters_.received_events.fetch_add(ok, std::memory_order_relaxed);
 190  ui_to_router_counters_.dropped_events.fetch_add(!ok, std::memory_order_relaxed);
 191  return ok;
 192}
 93
 194std::optional<Event> Router::ui_try_receive() noexcept {
 195  auto event = router_to_ui_.try_pop();
 196  if (event) {
 197    router_to_ui_counters_.sent_events.fetch_add(1, std::memory_order_relaxed);
 98  }
 199  return event;
 1100}
 101
 1102void Router::tick(size_t budget) {
 103  // Nodes->Router to Router->UI
 1104  for (size_t i = 0; i < budget; i++) {
 1105    auto event = nodes_to_router_.try_pop();
 1106    if (!event) {
 1107      break;
 108    }
 1109    nodes_to_router_counters_.received_events.fetch_add(1, std::memory_order_relaxed);
 110
 1111    bool ok = router_to_ui_.try_push(std::move(*event));
 1112    router_to_ui_counters_.sent_events.fetch_add(ok, std::memory_order_relaxed);
 1113    router_to_ui_counters_.dropped_events.fetch_add(!ok, std::memory_order_relaxed);
 1114  }
 115
 116  // UI->Router to Router->Nodes
 1117  for (size_t i = 0; i < budget; i++) {
 1118    auto event = ui_to_router_.try_pop();
 1119    if (!event) {
 1120      break;
 121    }
 1122    ui_to_router_counters_.received_events.fetch_add(1, std::memory_order_relaxed);
 123
 1124    auto it = router_to_nodes_.find(event->node_id);
 1125    if (it == router_to_nodes_.end()) {
 1126      std::cerr << "[holoflow_event::Router] Warning: Dropping event for unknown node '"
 127                << event->node_id << "'." << std::endl;
 1128      ui_to_router_counters_.dropped_events.fetch_add(1, std::memory_order_relaxed);
 1129      continue;
 130    }
 131
 1132    auto &queue    = it->second;
 1133    auto &counters = router_to_node_counters_[event->node_id];
 1134    bool  ok       = queue->try_push(std::move(*event));
 1135    counters.sent_events.fetch_add(ok, std::memory_order_relaxed);
 1136    counters.dropped_events.fetch_add(!ok, std::memory_order_relaxed);
 1137  }
 1138}
 139
 1140MailboxCounters::Snapshot Router::ui_to_router_counters() const noexcept {
 1141  return ui_to_router_counters_.snapshot();
 1142}
 143
 0144MailboxCounters::Snapshot Router::router_to_ui_counters() const noexcept {
 0145  return router_to_ui_counters_.snapshot();
 0146}
 147
 0148MailboxCounters::Snapshot Router::nodes_to_router_counters() const noexcept {
 0149  return nodes_to_router_counters_.snapshot();
 0150}
 151
 0152MailboxCounters::Snapshot Router::router_to_node_counters(const NodeId &node_id) const noexcept {
 0153  auto it = router_to_node_counters_.find(node_id);
 0154  if (it == router_to_node_counters_.end()) {
 0155    std::cerr << "[holoflow_event::Router] Warning: No counters for unknown node '" << node_id
 156              << "'." << std::endl;
 0157    return MailboxCounters::Snapshot{0, 0, 0};
 158  }
 0159  return it->second.snapshot();
 0160}
 161
 162} // namespace holoflow_event

Methods/Properties