Skip to content

Commit c29a279

Browse files
committed
feat(metrics): add receive channel blocking latency and logs
1 parent 1abdcde commit c29a279

4 files changed

Lines changed: 44 additions & 4 deletions

File tree

pkg/agent/client.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -577,9 +577,16 @@ func (a *Client) remoteToProxy(connID int64, eConn *endpointConn) {
577577
Data: buf[:n],
578578
ConnectID: connID,
579579
}}
580+
581+
klog.V(4).InfoS("sending data to kube-apiserver", "bytes", n, "connectionID", connID)
582+
sendStart := time.Now()
580583
if err := a.Send(resp); err != nil {
581584
klog.ErrorS(err, "could not send DATA", "connectionID", connID)
582585
}
586+
sendLatency := time.Since(sendStart)
587+
if sendLatency > 10*time.Millisecond {
588+
klog.V(3).InfoS("slow send to kube-apiserver", "latency", sendLatency, "bytes", n, "connectionID", connID)
589+
}
583590
klog.V(4).InfoS("send data to server successfully", "bytes", n, "connectionID", connID)
584591
}
585592
}

pkg/server/metrics/metrics.go

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@ type ServerMetrics struct {
5454
pendingDials *prometheus.GaugeVec
5555
establishedConns *prometheus.GaugeVec
5656
fullRecvChannels *prometheus.GaugeVec
57+
fullRecvChannelWaits *prometheus.HistogramVec
5758
dialFailures *prometheus.CounterVec
5859
streamPackets *prometheus.CounterVec
5960
streamErrors *prometheus.CounterVec
@@ -143,6 +144,18 @@ func newServerMetrics() *ServerMetrics {
143144
"service_method",
144145
},
145146
)
147+
fullRecvChannelWaits := prometheus.NewHistogramVec(
148+
prometheus.HistogramOpts{
149+
Namespace: Namespace,
150+
Subsystem: Subsystem,
151+
Name: "full_receive_channel_wait_seconds",
152+
Help: "Time spent blocked on a full receive channel before the send completed, partitioned by service method.",
153+
Buckets: latencyBuckets,
154+
},
155+
[]string{
156+
"service_method",
157+
},
158+
)
146159
dialFailures := prometheus.NewCounterVec(
147160
prometheus.CounterOpts{
148161
Namespace: Namespace,
@@ -206,6 +219,7 @@ func newServerMetrics() *ServerMetrics {
206219
prometheus.MustRegister(pendingDials)
207220
prometheus.MustRegister(establishedConns)
208221
prometheus.MustRegister(fullRecvChannels)
222+
prometheus.MustRegister(fullRecvChannelWaits)
209223
prometheus.MustRegister(dialFailures)
210224
prometheus.MustRegister(streamPackets)
211225
prometheus.MustRegister(streamErrors)
@@ -223,6 +237,7 @@ func newServerMetrics() *ServerMetrics {
223237
pendingDials: pendingDials,
224238
establishedConns: establishedConns,
225239
fullRecvChannels: fullRecvChannels,
240+
fullRecvChannelWaits: fullRecvChannelWaits,
226241
dialFailures: dialFailures,
227242
streamPackets: streamPackets,
228243
streamErrors: streamErrors,
@@ -243,6 +258,7 @@ func (s *ServerMetrics) Reset() {
243258
s.pendingDials.Reset()
244259
s.establishedConns.Reset()
245260
s.fullRecvChannels.Reset()
261+
s.fullRecvChannelWaits.Reset()
246262
s.dialFailures.Reset()
247263
s.streamPackets.Reset()
248264
s.streamErrors.Reset()
@@ -299,6 +315,11 @@ func (s *ServerMetrics) FullRecvChannel(serviceMethod string) prometheus.Gauge {
299315
return s.fullRecvChannels.With(prometheus.Labels{"service_method": serviceMethod})
300316
}
301317

318+
// ObserveFullRecvChannelWait records how long a send blocked on a full receive channel.
319+
func (s *ServerMetrics) ObserveFullRecvChannelWait(serviceMethod string, elapsed time.Duration) {
320+
s.fullRecvChannelWaits.With(prometheus.Labels{"service_method": serviceMethod}).Observe(elapsed.Seconds())
321+
}
322+
302323
type DialFailureReason string
303324

304325
const (

pkg/server/server.go

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -497,7 +497,9 @@ func (s *ProxyServer) readFrontendToChannel(frontend *GrpcFrontend, userAgent []
497497
klog.V(2).InfoS("Receive channel from frontend is full", "userAgent", userAgent)
498498
fullRecvChannelMetric := metrics.Metrics.FullRecvChannel(metrics.Proxy)
499499
fullRecvChannelMetric.Inc()
500+
start := time.Now()
500501
recvCh <- in
502+
metrics.Metrics.ObserveFullRecvChannelWait(metrics.Proxy, time.Since(start))
501503
fullRecvChannelMetric.Dec()
502504
}
503505
}
@@ -833,7 +835,11 @@ func (s *ProxyServer) readBackendToChannel(backend *Backend, recvCh chan *client
833835
klog.V(2).InfoS("Receive channel from agent is full", "agentID", agentID)
834836
fullRecvChannelMetric := metrics.Metrics.FullRecvChannel(metrics.Connect)
835837
fullRecvChannelMetric.Inc()
838+
start := time.Now()
836839
recvCh <- in
840+
latency := time.Since(start)
841+
metrics.Metrics.ObserveFullRecvChannelWait(metrics.Connect, latency)
842+
klog.V(2).InfoS("Latency: Receive channel from agent is full", "agentID", agentID, "latency", latency.Milliseconds())
837843
fullRecvChannelMetric.Dec()
838844
}
839845
}
@@ -982,7 +988,8 @@ func (s *ProxyServer) serveRecvBackend(backend *Backend, agentID string, recvCh
982988
break
983989
}
984990
if err := frontend.send(pkt); err != nil {
985-
klog.ErrorS(err, "send to client stream failure", "agentID", agentID, "connectionID", resp.ConnectID)
991+
// Likely cause: Kube-apiserver already terminated the connection with k-server.
992+
klog.V(4).InfoS("send to client stream failure (receiving DATA)", "err", err.Error(), "agentID", agentID, "connectionID", resp.ConnectID)
986993
} else {
987994
klog.V(5).InfoS("DATA sent to frontend")
988995
}
@@ -993,7 +1000,8 @@ func (s *ProxyServer) serveRecvBackend(backend *Backend, agentID string, recvCh
9931000
klog.V(4).InfoS("Received data ACK from agent", "agentID", agentID, "connectionID", resp.ConnectID)
9941001
frontend, err := s.getFrontend(agentID, resp.ConnectID)
9951002
if err != nil {
996-
klog.ErrorS(err, "could not get frontent client")
1003+
// Likely cause: Kube-apiserver already terminated the connection with k-server.
1004+
klog.V(4).InfoS("could not get frontend client (receiving DATA_ACK)", "err", err.Error(), "agentID", agentID, "connectionID", resp.ConnectID)
9971005
break
9981006
}
9991007

pkg/server/tunnel.go

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -159,6 +159,8 @@ func (t *Tunnel) ServeHTTP(w http.ResponseWriter, r *http.Request) {
159159
}
160160
if err != nil {
161161
klog.ErrorS(err, "Received failure on connection")
162+
// Likely cause: Kube-apiserver already terminated the connection with k-server.
163+
klog.V(4).InfoS("Received failure on connection: frontent likely closed connection", "err", err.Error(), "host", r.Host, "agentID", agentID, "connectionID", connID)
162164
break
163165
}
164166

@@ -179,13 +181,15 @@ func (t *Tunnel) ServeHTTP(w http.ResponseWriter, r *http.Request) {
179181
if !acquired {
180182
start := time.Now()
181183

182-
klog.InfoS("Semaphore full, waiting for client receive window > 0", "start", start.String(), "host", r.Host, "agentID", agentID, "connectionID", connID)
184+
klog.V(4).InfoS("Semaphore full, waiting for client receive window > 0", "start", start.String(), "host", r.Host, "agentID", agentID, "connectionID", connID)
183185
// Blocking: if semaphore is full (waits till server.go serveRecvBackend() - which receives packets via the grpc stream from an agent - receives
184186
// an ACK packet which releases 1 from the semaphore.
185187
connection.flow.Acquire(context.Background(), 1)
186188
latency := time.Now().Sub(start)
187189

188-
klog.V(3).InfoS("Latency when waiting for client receive window > 0", "latency", latency.Milliseconds(), "start", start.String(), "host", r.Host, "agentID", agentID, "connectionID", connID)
190+
if latency > time.Millisecond {
191+
klog.InfoS("Latency when waiting for client receive window > 0", "latency", latency.Milliseconds(), "start", start.String(), "host", r.Host, "agentID", agentID, "connectionID", connID)
192+
}
189193
}
190194

191195
} else {

0 commit comments

Comments
 (0)