Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 8 additions & 1 deletion server/cmd/api/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ import (
"github.com/kernel/kernel-images/server/lib/logger"
"github.com/kernel/kernel-images/server/lib/metrics"
"github.com/kernel/kernel-images/server/lib/nekoclient"
"github.com/kernel/kernel-images/server/lib/neterror"
oapi "github.com/kernel/kernel-images/server/lib/oapi"
"github.com/kernel/kernel-images/server/lib/recorder"
"github.com/kernel/kernel-images/server/lib/scaletozero"
Expand Down Expand Up @@ -297,6 +298,11 @@ func main() {
os.Exit(1)
}

// Tallies net::ERR_* failures seen on relayed CDP traffic. Shared by every
// proxy session and read by the metrics endpoint, so it lives for the
// process and its counts are cumulative per VM.
netErrorTracker := neterror.NewTracker(slogger)

rDevtools := chi.NewRouter()
rDevtools.Use(
chiMiddleware.Logger,
Expand All @@ -322,7 +328,7 @@ func main() {
rDevtools.Get("/json/list", jsonTargetHandler)
rDevtools.Get("/json/list/", jsonTargetHandler)
rDevtools.Get("/*", func(w http.ResponseWriter, r *http.Request) {
devtoolsproxy.WebSocketProxyHandler(upstreamMgr, slogger, config.LogCDPMessages, stz, telemetrySession.Publish, wsRegistry).ServeHTTP(w, r)
devtoolsproxy.WebSocketProxyHandler(upstreamMgr, slogger, config.LogCDPMessages, stz, telemetrySession.Publish, wsRegistry, netErrorTracker).ServeHTTP(w, r)
})

srvDevtools := &http.Server{
Expand Down Expand Up @@ -364,6 +370,7 @@ func main() {
rMetrics.Use(chiMiddleware.Recoverer)
metricsCollectors := []metrics.Collector{
metrics.NewChromeCollector(upstreamMgr),
metrics.NewNetErrorCollector(netErrorTracker),
metrics.NewGPUCollector(),
metrics.NewSystemCollector(),
}
Expand Down
17 changes: 15 additions & 2 deletions server/lib/devtoolsproxy/proxy.go
Original file line number Diff line number Diff line change
Expand Up @@ -305,19 +305,32 @@ func maybePauseAfterCurrentRead(ctx context.Context, logger *slog.Logger, r *htt
// so it can be wired directly; the proxy ignores the returns.
type EventPublisher func(ev events.Event) (events.Envelope, bool)

// NetErrorObserver is offered every CDP frame Chromium sends to the client so
// it can tally net::ERR_* failures. Satisfied by neterror.Tracker; nil
// disables the tap.
type NetErrorObserver interface {
Observe(msg []byte)
}

// WebSocketProxyHandler returns an http.Handler that upgrades incoming connections and
// proxies them to the current upstream websocket URL. It expects only websocket requests.
// If logCDPMessages is true, all CDP messages will be logged with their direction.
// publish is invoked on accept (cdp_connect) and on teardown (cdp_disconnect); pass
// nil to disable emission.
func WebSocketProxyHandler(mgr *UpstreamManager, logger *slog.Logger, logCDPMessages bool, ctrl scaletozero.Controller, publish EventPublisher, reg *wsdrain.Registry) http.Handler {
// nil to disable emission. netErrors, when non-nil, observes upstream frames; it
// must not mutate them.
func WebSocketProxyHandler(mgr *UpstreamManager, logger *slog.Logger, logCDPMessages bool, ctrl scaletozero.Controller, publish EventPublisher, reg *wsdrain.Registry, netErrors NetErrorObserver) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
// Counts every relayed message so cdp_disconnect can report message_count.
var msgCount atomic.Int64
var transform wsproxy.MessageTransform = func(direction string, mt websocket.MessageType, msg []byte) []byte {
if logCDPMessages {
logCDPMessage(logger, direction, mt, msg)
}
// Events only travel upstream-to-client, so client frames are
// skipped rather than scanned for a needle they cannot contain.
if netErrors != nil && direction == wsproxy.DirectionUpstreamToClient {
netErrors.Observe(msg)
}
msgCount.Add(1)
return msg
}
Expand Down
95 changes: 89 additions & 6 deletions server/lib/devtoolsproxy/proxy_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,13 +20,96 @@ import (
"time"

"github.com/coder/websocket"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"github.com/kernel/kernel-images/server/lib/events"
oapi "github.com/kernel/kernel-images/server/lib/oapi"
"github.com/kernel/kernel-images/server/lib/scaletozero"
"github.com/kernel/kernel-images/server/lib/wsdrain"
"github.com/kernel/kernel-images/server/lib/wsproxy"
)

// recordingObserver captures every frame handed to the net error tap.
type recordingObserver struct {
mu sync.Mutex
frames []string
}

func (o *recordingObserver) Observe(msg []byte) {
o.mu.Lock()
defer o.mu.Unlock()
o.frames = append(o.frames, string(msg))
}

func (o *recordingObserver) snapshot() []string {
o.mu.Lock()
defer o.mu.Unlock()
return append([]string(nil), o.frames...)
}

// The tap must see Chromium's frames and only Chromium's frames: a reversed
// direction check would silently scan client commands, which can never carry
// a CDP event, and count nothing.
func TestWebSocketProxyHandler_ObservesOnlyUpstreamFrames(t *testing.T) {
const clientFrame = "client-to-upstream"
const upstreamFrame = "upstream-to-client"

upstreamSrv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
c, err := websocket.Accept(w, r, &websocket.AcceptOptions{OriginPatterns: []string{"*"}})
if err != nil {
return
}
defer c.Close(websocket.StatusNormalClosure, "")
// Reply with a distinct payload so the two directions are never
// confusable in the recorded frames.
if _, _, err := c.Read(r.Context()); err != nil {
return
}
_ = c.Write(r.Context(), websocket.MessageText, []byte(upstreamFrame))
<-r.Context().Done()
}))
defer upstreamSrv.Close()

u, _ := url.Parse(upstreamSrv.URL)
logger := silentLogger()
mgr := NewUpstreamManager("/dev/null", logger)
mgr.setCurrent((&url.URL{Scheme: "ws", Host: u.Host, Path: "/"}).String())

obs := &recordingObserver{}
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), nil, nil, obs))
defer proxySrv.Close()

pu, _ := url.Parse(proxySrv.URL)
pu.Scheme = "ws"

ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
conn, _, err := websocket.Dial(ctx, pu.String(), nil)
require.NoError(t, err)
defer conn.Close(websocket.StatusNormalClosure, "")

require.NoError(t, conn.Write(ctx, websocket.MessageText, []byte(clientFrame)))

// The tap runs before the frame is written to the client, so a successful
// read means the observer has already been offered it.
_, got, err := conn.Read(ctx)
require.NoError(t, err)
require.Equal(t, upstreamFrame, string(got))

assert.Equal(t, []string{upstreamFrame}, obs.snapshot())
}

func TestWebSocketProxyHandler_NilObserverIsAllowed(t *testing.T) {
logger := silentLogger()
mgr := NewUpstreamManager("/dev/null", logger)
mgr.setCurrent("ws://127.0.0.1:1/")

assert.NotPanics(t, func() {
WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), nil, nil, nil)
})
}

func silentLogger() *slog.Logger {
return slog.New(slog.NewTextHandler(io.Discard, nil))
}
Expand Down Expand Up @@ -133,7 +216,7 @@ func TestWebSocketProxyHandler_ProxiesEcho(t *testing.T) {
// seed current upstream to echo server including path/query (bypass tailing)
mgr.setCurrent((&url.URL{Scheme: u.Scheme, Host: u.Host, Path: u.Path, RawQuery: u.RawQuery}).String())

proxy := WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), nil, nil)
proxy := WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), nil, nil, nil)
proxySrv := httptest.NewServer(proxy)
defer proxySrv.Close()

Expand Down Expand Up @@ -191,7 +274,7 @@ func TestWebSocketProxyHandler_RegistryClosesClientWithGoingAway(t *testing.T) {
mgr.setCurrent((&url.URL{Scheme: "ws", Host: u.Host, Path: "/echo"}).String())

reg := wsdrain.New()
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), nil, reg))
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), nil, reg, nil))
defer proxySrv.Close()

pu, _ := url.Parse(proxySrv.URL)
Expand Down Expand Up @@ -520,7 +603,7 @@ func TestWebSocketProxyHandler_EmitsConnectAndDisconnect(t *testing.T) {
mgr.setCurrent(u.String())

rp := &recordingPublisher{}
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil))
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil, nil))
defer proxySrv.Close()

pu, _ := url.Parse(proxySrv.URL)
Expand Down Expand Up @@ -695,7 +778,7 @@ func TestWebSocketProxyHandler_EmitsUpstreamChangedOnMidStreamRestart(t *testing
mgr.setCurrent(urlA.String())

rp := &recordingPublisher{}
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil))
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil, nil))
defer proxySrv.Close()

pu, _ := url.Parse(proxySrv.URL)
Expand Down Expand Up @@ -779,7 +862,7 @@ func TestWebSocketProxyHandler_KicksClientOffStaleUpstreamOnURLChange(t *testing
mgr.setCurrent(urlA.String())

rp := &recordingPublisher{}
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil))
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil, nil))
defer proxySrv.Close()

pu, _ := url.Parse(proxySrv.URL)
Expand Down Expand Up @@ -831,7 +914,7 @@ func TestWebSocketProxyHandler_EmitsUpstreamErrorOnDialFailure(t *testing.T) {
mgr.setCurrent(deadURL)

rp := &recordingPublisher{}
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil))
proxySrv := httptest.NewServer(WebSocketProxyHandler(mgr, logger, false, scaletozero.NewNoopController(), rp.publish, nil, nil))
defer proxySrv.Close()

pu, _ := url.Parse(proxySrv.URL)
Expand Down
52 changes: 52 additions & 0 deletions server/lib/metrics/neterror.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
package metrics

import (
"context"

"github.com/kernel/kernel-images/server/lib/neterror"
)

// NetErrorSource reports network failure counts observed on relayed CDP
// traffic. Satisfied by neterror.Tracker.
type NetErrorSource interface {
Snapshot() (map[neterror.Key]int64, int64)
}

// NetErrorCollector exposes Chromium net::ERR_* failures as counters. Unlike
// the Chrome collector it touches neither Chrome nor the network on scrape:
// the counts are accumulated by the CDP proxy as traffic flows and only read
// here.
type NetErrorCollector struct {
source NetErrorSource
}

func NewNetErrorCollector(source NetErrorSource) *NetErrorCollector {
return &NetErrorCollector{source: source}
}

func (c *NetErrorCollector) Name() string { return "neterror" }

func (c *NetErrorCollector) Collect(_ context.Context, w *Writer) error {
counts, overflow := c.source.Snapshot()

// A VM that has seen no failures emits no samples here, so this family is
// simply absent from queries until something fails. Use the unlabelled
// _dropped_total below to tell "no failures" apart from "not scraped".
w.Metric("kernel_chromium_net_errors_total",
"Chromium network request failures by error text and resource type, cumulative since VM start. Client-cancelled requests (net::ERR_ABORTED) are excluded; deliberate route.abort() blocking is not separable and lands in the net::ERR_FAILED series.",
"counter")
for _, k := range neterror.SortedKeys(counts) {
w.Sample("kernel_chromium_net_errors_total", []Label{
{"error", k.Error},
{"resource_type", k.ResourceType},
}, float64(counts[k]))
}

// Always sampled, so this doubles as the proof that the tap is running on
// a VM reporting no failures.
w.Metric("kernel_chromium_net_errors_dropped_total",
"Network failures not counted because the distinct error/resource-type cap was reached.",
"counter")
w.Sample("kernel_chromium_net_errors_dropped_total", nil, float64(overflow))
return nil
}
57 changes: 57 additions & 0 deletions server/lib/metrics/neterror_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
package metrics

import (
"context"
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"github.com/kernel/kernel-images/server/lib/neterror"
)

type stubNetErrorSource struct {
counts map[neterror.Key]int64
overflow int64
}

func (s stubNetErrorSource) Snapshot() (map[neterror.Key]int64, int64) {
return s.counts, s.overflow
}

func TestNetErrorCollector(t *testing.T) {
c := NewNetErrorCollector(stubNetErrorSource{
counts: map[neterror.Key]int64{
{Error: "net::ERR_TIMED_OUT", ResourceType: "XHR"}: 3,
{Error: "net::ERR_HTTP2_PROTOCOL_ERROR", ResourceType: "Document"}: 12,
},
overflow: 2,
})

w := &Writer{}
require.NoError(t, c.Collect(context.Background(), w))

out := string(w.Bytes())
assert.Contains(t, out, "# TYPE kernel_chromium_net_errors_total counter\n")
// Samples are ordered by error then resource type so an unchanged tracker
// scrapes byte-identically.
assert.Contains(t, out, `kernel_chromium_net_errors_total{error="net::ERR_HTTP2_PROTOCOL_ERROR",resource_type="Document"} 12
kernel_chromium_net_errors_total{error="net::ERR_TIMED_OUT",resource_type="XHR"} 3
`)
assert.Contains(t, out, "kernel_chromium_net_errors_dropped_total 2\n")
}

// With no failures the labelled family has no samples and so is absent from
// queries; the unlabelled dropped counter is what proves the VM reported.
func TestNetErrorCollectorWithNoFailures(t *testing.T) {
c := NewNetErrorCollector(stubNetErrorSource{})

w := &Writer{}
require.NoError(t, c.Collect(context.Background(), w))

out := string(w.Bytes())
assert.NotContains(t, out, "kernel_chromium_net_errors_total{")
assert.Contains(t, out, "kernel_chromium_net_errors_dropped_total 0\n")
}

var _ NetErrorSource = (*neterror.Tracker)(nil)
Loading
Loading