blob: 14794badfe4c8cd225aa42f2f698f7e7016cee17 (
plain) (
blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
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
|