libpromeki 1.0.0-alpha
PROfessional MEdia toolKIt
 
Loading...
Searching...
No Matches
pcapflowrouter.h
Go to the documentation of this file.
1
8#pragma once
9
10
11#include <promeki/config.h>
12#if PROMEKI_ENABLE_NETWORK
13#include <cstdint>
14#include <utility>
15#include <promeki/namespace.h>
16#include <promeki/ancpacket.h>
17#include <promeki/buffer.h>
18#include <promeki/datetime.h>
19#include <promeki/duration.h>
20#include <promeki/error.h>
21#include <promeki/function.h>
22#include <promeki/list.h>
23#include <promeki/packetdemux.h>
24#include <promeki/pcapsdpmap.h>
25#include <promeki/rtppacket.h>
29#include <promeki/string.h>
30#include <promeki/uniqueptr.h>
31
32PROMEKI_NAMESPACE_BEGIN
33
34class SdpSession;
35
62class PcapFlowRouter {
63 public:
65 struct RoutedAncFrame {
67 SocketAddress src;
69 SocketAddress dst;
71 uint32_t ssrc = 0;
73 DateTime captureTime;
76 RxAncFrame anc;
77 };
78
80 using AncFrameCallback = Function<void(const RoutedAncFrame &)>;
81
83 struct FlowStat {
84 SocketAddress dst;
85 uint32_t ssrc = 0;
86 uint8_t payloadType = 0;
87 PcapFlowKind kind = PcapFlowKind::Unknown;
88 uint64_t packets = 0;
89 uint64_t bytes = 0;
90 // -- RFC 3550 health (via RtpSeqTracker), cumulative over this SSRC --
91 uint64_t lostPackets = 0;
92 uint64_t duplicatePackets = 0;
93 uint64_t reorderedPackets = 0;
94 uint32_t timestampRegressions = 0;
95 Duration maxJitter = Duration::zero();
96 };
97
108 struct RtpAnomaly {
110 enum class Kind {
111 SsrcChange,
112 PayloadTypeChange,
113 PacketLoss,
114 Reorder,
115 Duplicate,
116 TimestampRegression,
117 JitterExceeded,
118 };
119 Kind kind = Kind::SsrcChange;
120 PcapFlowKind flowKind = PcapFlowKind::Unknown;
121 SocketAddress dst;
122 uint32_t ssrc = 0;
127 uint32_t previous = 0;
128 uint32_t count = 0;
129 Duration jitter = Duration::zero();
130 DateTime captureTime;
131 };
132
134 using RtpAnomalyCallback = Function<void(const RtpAnomaly &)>;
135
136 PcapFlowRouter() = default;
137
143 Error setSdp(const SdpSession &sdp);
144
156 void addAncFlow(const SocketAddress &dst, int payloadType = -1) {
157 _map.addAncFlow(dst.address(), dst.port(), payloadType);
158 }
159
161 void onAncFrame(AncFrameCallback cb) { _ancCb = std::move(cb); }
162
164 void onRtpAnomaly(RtpAnomalyCallback cb) { _anomalyCb = std::move(cb); }
165
177 void setJitterWarnThreshold(const Duration &threshold) { _jitterWarn = threshold; }
178
185 Error processFile(const String &path);
186
191 Error processBuffer(const Buffer &buf);
192
194 const List<FlowStat> &flowStats() const { return _stats; }
195
197 const PcapSdpMap &sdpMap() const { return _map; }
198
200 void reset();
201
202 private:
204 struct AncReasm {
205 SocketAddress dst;
206 SocketAddress src;
207 uint32_t ssrc = 0;
208 bool haveSsrc = false;
209 bool haveTs = false;
210 uint32_t timestamp = 0;
211 uint8_t payloadType = 0;
212 AncDesc desc;
213 RtpPacket::List packets;
214 DateTime captureTime;
215 };
216
221 struct FlowHealth {
222 SocketAddress dst;
223 uint32_t ssrc = 0;
224 bool haveSsrc = false;
225 uint8_t payloadType = 0;
226 bool havePt = false;
227 uint32_t lastTimestamp = 0;
228 bool haveTs = false;
229 uint32_t lastExtendedSeq = 0;
230 bool haveSeq = false;
231 // §A.8 jitter state (computed directly in ns; see trackRtpHealth).
232 int64_t prevArrivalNs = 0;
233 uint32_t prevRtpTs = 0;
234 int64_t jitterNs = 0;
235 bool haveJitterPrev = false;
236 Duration maxJitter = Duration::zero();
237 bool jitterOver = false;
238 uint32_t timestampRegressions = 0;
239 RtpSeqTracker tracker;
240 };
241
242 void handleDatagram(const UdpDatagram &dg, const DateTime &captureTime);
243 void routeAnc(const UdpDatagram &dg, const PcapFlow &flow, const RtpPacket &pkt, const DateTime &captureTime);
244 void flushAnc(AncReasm &r);
245 FlowStat &statFor(const SocketAddress &dst, uint32_t ssrc, uint8_t pt, PcapFlowKind kind);
246 AncReasm &ancReasmFor(const UdpDatagram &dg, const PcapFlow &flow);
247 FlowHealth &healthFor(const SocketAddress &dst);
248 void trackRtpHealth(const SocketAddress &dst, PcapFlowKind kind, const RtpPacket &pkt, FlowStat &st,
249 const DateTime &captureTime);
250 void emitAnomaly(RtpAnomaly::Kind kind, PcapFlowKind flowKind, const SocketAddress &dst, uint32_t ssrc,
251 uint32_t previous, uint32_t count, const DateTime &captureTime,
252 const Duration &jitter = Duration::zero());
253
254 Error runReader(class PcapReader &reader);
255
256 PcapSdpMap _map;
257 PacketDemux _demux;
258 AncFrameCallback _ancCb;
259 RtpAnomalyCallback _anomalyCb;
260 Duration _jitterWarn = Duration::zero();
261 List<FlowStat> _stats;
262 List<AncReasm> _ancFlows;
263 List<UniquePtr<FlowHealth>> _health;
264};
265
266PROMEKI_NAMESPACE_END
267
268#endif // PROMEKI_ENABLE_NETWORK