Skip to content

Commit ccf1a2e

Browse files
committed
fix(packetparser): Fix under reporting of TCP flags and packet metrics, improve scalability
Signed-off-by: Matthew McKeen <matthew.mckeen@fastly.com>
1 parent 2eca8c0 commit ccf1a2e

20 files changed

Lines changed: 1567 additions & 1016 deletions

pkg/metrics/metrics.go

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -188,6 +188,12 @@ func InitializeMetrics() {
188188
ConntrackTotalConnectionsDescription,
189189
)
190190

191+
ParsedPacketsCounter = exporter.CreatePrometheusCounterVecForControlPlaneMetric(
192+
exporter.DefaultRegistry,
193+
parsedPacketsCounterName,
194+
parsedPacketsCounterDescription,
195+
)
196+
191197
isInitialized = true
192198
metricsLogger.Info("Metrics initialized")
193199
}

pkg/metrics/types.go

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ const (
1212
// Control plane metrics
1313
pluginManagerFailedToReconcileCounterName = "plugin_manager_failed_to_reconcile"
1414
lostEventsCounterName = "lost_events_counter"
15+
parsedPacketsCounterName = "parsed_packets_counter"
1516

1617
// Windows
1718
hnsStats = "windows_hns_stats"
@@ -43,6 +44,7 @@ const (
4344
// Control plane metrics
4445
pluginManagerFailedToReconcileCounterDescription = "Number of times the plugin manager failed to reconcile the plugins"
4546
lostEventsCounterDescription = "Number of events lost in control plane"
47+
parsedPacketsCounterDescription = "Number of packets parsed by the packetparser plugin"
4648

4749
// Conntrack metrics
4850
ConntrackPacketTxDescription = "Number of tx packets"
@@ -90,6 +92,7 @@ var (
9092
// Control Plane Metrics
9193
PluginManagerFailedToReconcileCounter CounterVec
9294
LostEventsCounter CounterVec
95+
ParsedPacketsCounter CounterVec
9396

9497
// DNS Metrics.
9598
DNSRequestCounter CounterVec

pkg/module/metrics/forward.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -170,8 +170,8 @@ func (f *ForwardMetrics) processLocalCtxFlow(flow *v1.Flow) {
170170
func (f *ForwardMetrics) update(fl *v1.Flow, labels []string) {
171171
switch f.metricName {
172172
case utils.ForwardPacketsGaugeName:
173-
f.forwardMetric.WithLabelValues(labels...).Inc()
173+
f.forwardMetric.WithLabelValues(labels...).Add(float64(utils.PreviouslyObservedPackets(fl) + 1))
174174
case utils.ForwardBytesGaugeName:
175-
f.forwardMetric.WithLabelValues(labels...).Add(float64(utils.PacketSize(fl)))
175+
f.forwardMetric.WithLabelValues(labels...).Add(float64(utils.PacketSize(fl) + utils.PreviouslyObservedBytes(fl)))
176176
}
177177
}

pkg/module/metrics/latency_test.go

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -124,7 +124,7 @@ func TestProcessFlow(t *testing.T) {
124124
f1 := utils.ToFlow(l, t1, apiSeverIp, nodeIp, 80, 443, 6, 3, 0)
125125
metaf1 := &utils.RetinaMetadata{}
126126
utils.AddTCPID(metaf1, 1234)
127-
utils.AddTCPFlags(f1, 1, 0, 0, 0, 0, 0)
127+
utils.AddTCPFlags(f1, 1, 0, 0, 0, 0, 0, 0, 0, 0)
128128
utils.AddRetinaMetadata(f1, metaf1)
129129
f1.Destination = &flow.Endpoint{
130130
PodName: "kubernetes-apiserver",
@@ -134,7 +134,7 @@ func TestProcessFlow(t *testing.T) {
134134
f2 := utils.ToFlow(l, t2, nodeIp, apiSeverIp, 443, 80, 6, 2, 0)
135135
metaf2 := &utils.RetinaMetadata{}
136136
utils.AddTCPID(metaf2, 1234)
137-
utils.AddTCPFlags(f2, 1, 1, 0, 0, 0, 0)
137+
utils.AddTCPFlags(f2, 1, 1, 0, 0, 0, 0, 0, 0, 0)
138138
utils.AddRetinaMetadata(f2, metaf2)
139139
f2.Source = &flow.Endpoint{
140140
PodName: "kubernetes-apiserver",
@@ -147,9 +147,9 @@ func TestProcessFlow(t *testing.T) {
147147
* Test case 2: Existing TCP connection.
148148
*/
149149
// Node -> Api server.
150-
utils.AddTCPFlags(f1, 1, 0, 0, 0, 0, 0)
150+
utils.AddTCPFlags(f1, 1, 0, 0, 0, 0, 0, 0, 0, 0)
151151
// Api server -> Node.
152-
utils.AddTCPFlags(f2, 0, 1, 0, 0, 0, 0)
152+
utils.AddTCPFlags(f2, 0, 1, 0, 0, 0, 0, 0, 0, 0)
153153
// Process flow.
154154
lm.ProcessFlow(f1)
155155
lm.ProcessFlow(f2)

pkg/module/metrics/tcpflags.go

Lines changed: 37 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,27 @@ func (t *TCPMetrics) getLabels() []string {
6565
return labels
6666
}
6767

68+
func combineFlagsWithPrevious(flags []string, flow *v1.Flow) map[string]uint32 {
69+
var combinedFlags map[string]uint32
70+
71+
previous := utils.PreviouslyObservedTCPFlags(flow)
72+
if previous != nil {
73+
combinedFlags = previous
74+
} else {
75+
combinedFlags = map[string]uint32{}
76+
}
77+
78+
for _, flag := range flags {
79+
if _, ok := combinedFlags[flag]; !ok {
80+
combinedFlags[flag] = 1
81+
} else {
82+
combinedFlags[flag]++
83+
}
84+
}
85+
86+
return combinedFlags
87+
}
88+
6889
func (t *TCPMetrics) ProcessFlow(flow *v1.Flow) {
6990
if flow == nil {
7091
return
@@ -100,11 +121,11 @@ func (t *TCPMetrics) ProcessFlow(flow *v1.Flow) {
100121
dstLabels = t.dstCtx.getValues(flow)
101122
}
102123

103-
for _, flag := range flags {
124+
for flag, count := range combineFlagsWithPrevious(flags, flow) {
104125
labels := append([]string{flag}, srcLabels...)
105126
labels = append(labels, dstLabels...)
106-
t.tcpFlagsMetrics.WithLabelValues(labels...).Inc()
107-
t.l.Debug("TCP flag metric", zap.String("flag", flag), zap.Strings("labels", labels))
127+
t.tcpFlagsMetrics.WithLabelValues(labels...).Add(float64(count))
128+
t.l.Debug("TCP flag metric", zap.String("flag", flag), zap.Strings("labels", labels), zap.Uint32("count", count))
108129
}
109130
}
110131

@@ -113,20 +134,23 @@ func (t *TCPMetrics) processLocalCtxFlow(flow *v1.Flow, flags []string) {
113134
if labelValuesMap == nil {
114135
return
115136
}
137+
138+
combinedFlags := combineFlagsWithPrevious(flags, flow)
139+
116140
// Ingress values
117141
if l := len(labelValuesMap[ingress]); l > 0 {
118-
for _, flag := range flags {
142+
for flag, count := range combinedFlags {
119143
labels := append([]string{flag}, labelValuesMap[ingress]...)
120-
t.tcpFlagsMetrics.WithLabelValues(labels...).Inc()
121-
t.l.Debug("TCP flag metric", zap.String("flag", flag), zap.Strings("labels", labels))
144+
t.tcpFlagsMetrics.WithLabelValues(labels...).Add(float64(count))
145+
t.l.Debug("TCP flag metric", zap.String("flag", flag), zap.Strings("labels", labels), zap.Uint32("count", count))
122146
}
123147
}
124148

125149
if l := len(labelValuesMap[egress]); l > 0 {
126-
for _, flag := range flags {
150+
for flag, count := range combinedFlags {
127151
labels := append([]string{flag}, labelValuesMap[egress]...)
128-
t.tcpFlagsMetrics.WithLabelValues(labels...).Inc()
129-
t.l.Debug("TCP flag metric", zap.String("flag", flag), zap.Strings("labels", labels))
152+
t.tcpFlagsMetrics.WithLabelValues(labels...).Add(float64(count))
153+
t.l.Debug("TCP flag metric", zap.String("flag", flag), zap.Strings("labels", labels), zap.Uint32("count", count))
130154
}
131155
}
132156
}
@@ -171,6 +195,10 @@ func (t *TCPMetrics) getFlagValues(flags *v1.TCPFlags) []string {
171195
f = append(f, utils.CWR)
172196
}
173197

198+
if flags.GetNS() {
199+
f = append(f, utils.NS)
200+
}
201+
174202
return f
175203
}
176204

0 commit comments

Comments
 (0)