srdusr
aboutsummaryrefslogtreecommitdiffstats
path: root/include/wireframe/capture_queue.hpp
diff options
context:
space:
mode:
Diffstat (limited to 'include/wireframe/capture_queue.hpp')
-rw-r--r--include/wireframe/capture_queue.hpp98
1 files changed, 98 insertions, 0 deletions
diff --git a/include/wireframe/capture_queue.hpp b/include/wireframe/capture_queue.hpp
new file mode 100644
index 0000000..14794ba
--- /dev/null
+++ b/include/wireframe/capture_queue.hpp
@@ -0,0 +1,98 @@
+#pragma once
+
+#include <condition_variable>
+#include <cstdint>
+#include <mutex>
+#include <optional>
+#include <queue>
+#include <vector>
+
+// Bounded queue between the capture thread and the render/analysis
+// thread (PLAN.md's architecture sketch). Owns a copy of each packet's
+// bytes since the buffer libpcap hands the callback is only valid for
+// the duration of that call.
+namespace wireframe {
+
+struct CapturedPacket {
+ std::uint32_t ts_sec;
+ std::uint32_t ts_usec;
+ std::uint32_t original_len;
+ std::vector<unsigned char> data; // caplen bytes
+};
+
+// Single-producer / single-consumer. Two producer-side push variants
+// for two different producers with different constraints: a live
+// capture thread can't be allowed to stall (PLAN.md is explicit that a
+// traffic spike should drop packets, not block), but a replay-from-file
+// producer has no such real-time pressure, and dropping from a fixed
+// historical record would defeat the point of "faithfully replaying
+// what was captured" - so it blocks for room instead.
+class CaptureQueue {
+public:
+ explicit CaptureQueue(std::size_t capacity) : capacity_(capacity) {}
+
+ // Never blocks: drops the packet and counts it if the queue is full.
+ bool try_push(CapturedPacket&& packet) {
+ {
+ std::lock_guard<std::mutex> lock(mutex_);
+ if (queue_.size() >= capacity_) {
+ ++dropped_;
+ return false;
+ }
+ queue_.push(std::move(packet));
+ }
+ cv_.notify_all();
+ return true;
+ }
+
+ // Blocks until there's room, then pushes. Returns false without
+ // pushing if stop() is called while waiting - the consumer side is
+ // going away, so nothing will ever pop it.
+ bool push(CapturedPacket&& packet) {
+ {
+ std::unique_lock<std::mutex> lock(mutex_);
+ cv_.wait(lock, [this] { return queue_.size() < capacity_ || stopped_; });
+ if (stopped_) return false;
+ queue_.push(std::move(packet));
+ }
+ cv_.notify_all();
+ return true;
+ }
+
+ // Blocks until a packet is available. Returns nullopt only once
+ // stop() has been called and the queue has fully drained - so a
+ // consumer loop on pop() processes everything queued before the
+ // capture side stopped, rather than discarding it.
+ std::optional<CapturedPacket> pop() {
+ std::unique_lock<std::mutex> lock(mutex_);
+ cv_.wait(lock, [this] { return !queue_.empty() || stopped_; });
+ if (queue_.empty()) return std::nullopt;
+ CapturedPacket packet = std::move(queue_.front());
+ queue_.pop();
+ cv_.notify_all(); // wake a push() blocked on room, if any
+ return packet;
+ }
+
+ void stop() {
+ {
+ std::lock_guard<std::mutex> lock(mutex_);
+ stopped_ = true;
+ }
+ cv_.notify_all();
+ }
+
+ std::uint64_t dropped() const {
+ std::lock_guard<std::mutex> lock(mutex_);
+ return dropped_;
+ }
+
+private:
+ mutable std::mutex mutex_;
+ std::condition_variable cv_;
+ std::queue<CapturedPacket> queue_;
+ std::size_t capacity_;
+ bool stopped_ = false;
+ std::uint64_t dropped_ = 0;
+};
+
+} // namespace wireframe