Skip to content

Commit a79d6b4

Browse files
committed
Add browser OTEL bootstrap and trace proxy
1 parent 97d2a7d commit a79d6b4

26 files changed

Lines changed: 1287 additions & 49 deletions

cmd/wl/cmd_runtime_callbacks_test.go

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -17,8 +17,10 @@ import (
1717

1818
type fakeSelfHostedServer struct {
1919
client *sdk.Client
20+
env string
2021
}
2122

23+
func (s *fakeSelfHostedServer) SetEnvironment(environment string) { s.env = environment }
2224
func (s *fakeSelfHostedServer) SetScoreboard(*api.CachedEndpoint) {}
2325
func (s *fakeSelfHostedServer) SetScoreboardDetail(*api.CachedEndpoint) {}
2426
func (s *fakeSelfHostedServer) SetScoreboardDump(*api.CachedEndpoint) {}
@@ -424,10 +426,10 @@ func TestRunServe_RemoteClientCallbacks(t *testing.T) {
424426
}
425427

426428
func TestRunServe_UsesInjectedSelfHostedServer(t *testing.T) {
427-
var capturedClient *sdk.Client
429+
var capturedServer *fakeSelfHostedServer
428430
withSelfHostedAPIServerOverride(t, func(client *sdk.Client) selfHostedAPIServer {
429-
capturedClient = client
430-
return &fakeSelfHostedServer{client: client}
431+
capturedServer = &fakeSelfHostedServer{client: client}
432+
return capturedServer
431433
})
432434
withServeListenOverride(t, func(*http.Server) error { return nil })
433435
withLocalWorkflowDBOverride(t, func(string, string) localWorkflowDB {
@@ -458,7 +460,10 @@ func TestRunServe_UsesInjectedSelfHostedServer(t *testing.T) {
458460
if err := runServe(cmd, io.Discard, io.Discard); err != nil {
459461
t.Fatalf("runServe() error = %v", err)
460462
}
461-
if capturedClient == nil {
463+
if capturedServer == nil || capturedServer.client == nil {
462464
t.Fatal("self-hosted API server did not receive client")
463465
}
466+
if capturedServer.env != "self-sovereign" {
467+
t.Fatalf("environment = %q, want %q", capturedServer.env, "self-sovereign")
468+
}
464469
}

cmd/wl/cmd_serve.go

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,7 @@ type remoteWorkflowDB interface {
8282

8383
type selfHostedAPIServer interface {
8484
http.Handler
85+
SetEnvironment(string)
8586
SetScoreboard(*api.CachedEndpoint)
8687
SetScoreboardDetail(*api.CachedEndpoint)
8788
SetScoreboardDump(*api.CachedEndpoint)
@@ -298,6 +299,7 @@ func runServe(cmd *cobra.Command, stdout, stderr io.Writer) error {
298299
})
299300

300301
server := newSelfHostedAPIServer(client)
302+
server.SetEnvironment("self-sovereign")
301303

302304
scoreboardCache := api.NewScoreboardCache(db, 5*time.Minute)
303305
server.SetScoreboard(scoreboardCache)
@@ -317,7 +319,9 @@ func runServe(cmd *cobra.Command, stdout, stderr io.Writer) error {
317319
rateLimiter := api.NewRateLimiter(120, 120, time.Minute)
318320
defer rateLimiter.Stop()
319321
generalRL := api.RateLimit(rateLimiter)
320-
bodyLimit := api.MaxBytesBody(64 << 10) // 64 KB
322+
bodyLimit := api.MaxBytesBodyByPath(64<<10, map[string]int64{
323+
observability.BrowserTraceIngressPath: 1 << 20,
324+
})
321325
sentryMiddleware := sentryhttp.New(sentryhttp.Options{Repanic: true})
322326
handler := sentryMiddleware.Handle(observability.NewHTTPHandler(api.RequestLog(logger)(api.SecurityHeaders(generalRL(bodyLimit(api.SPAHandler(server, web.Assets)))))))
323327
if devMode {
@@ -392,6 +396,7 @@ func runServeHosted(cmd *cobra.Command, stdout, _ io.Writer) error {
392396

393397
// Build the API server with hosted workspace resolution.
394398
apiServer := api.NewHostedWorkspace(hosted.NewClientFunc(), hosted.NewWorkspaceFunc())
399+
apiServer.SetEnvironment(environment)
395400

396401
// Public read-only RemoteDB against the canonical hosted upstream (no token needed).
397402
publicDB := newHostedPublicDB()
@@ -438,7 +443,9 @@ func runServeHosted(cmd *cobra.Command, stdout, _ io.Writer) error {
438443
hostedRateLimiter := api.NewRateLimiter(120, 120, time.Minute)
439444
defer hostedRateLimiter.Stop()
440445
generalRL := api.RateLimit(hostedRateLimiter)
441-
bodyLimit := api.MaxBytesBody(64 << 10) // 64 KB
446+
bodyLimit := api.MaxBytesBodyByPath(64<<10, map[string]int64{
447+
observability.BrowserTraceIngressPath: 1 << 20,
448+
})
442449
sentryMiddleware := sentryhttp.New(sentryhttp.Options{Repanic: true})
443450
handler := sentryMiddleware.Handle(observability.NewHTTPHandler(api.RequestLog(logger)(api.SecurityHeaders(generalRL(bodyLimit(hostedServer.Handler(apiServer, web.Assets)))))))
444451
if devMode {

internal/api/bodylimit.go

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,9 +5,19 @@ import "net/http"
55
// MaxBytesBody returns middleware that limits request body size to n bytes.
66
// Requests exceeding the limit receive a 413 Request Entity Too Large response.
77
func MaxBytesBody(n int64) func(http.Handler) http.Handler {
8+
return MaxBytesBodyByPath(n, nil)
9+
}
10+
11+
// MaxBytesBodyByPath returns middleware that limits request body size to a
12+
// default number of bytes, with optional per-path overrides.
13+
func MaxBytesBodyByPath(defaultLimit int64, overrides map[string]int64) func(http.Handler) http.Handler {
814
return func(next http.Handler) http.Handler {
915
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
10-
r.Body = http.MaxBytesReader(w, r.Body, n)
16+
limit := defaultLimit
17+
if override, ok := overrides[r.URL.Path]; ok {
18+
limit = override
19+
}
20+
r.Body = http.MaxBytesReader(w, r.Body, limit)
1121
next.ServeHTTP(w, r)
1222
})
1323
}

internal/api/bodylimit_test.go

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -48,3 +48,30 @@ func TestMaxBytesBody_RejectsLargeBody(t *testing.T) {
4848
t.Fatalf("large body: got %d, want 413", rec.Code)
4949
}
5050
}
51+
52+
func TestMaxBytesBodyByPath_UsesPerPathOverride(t *testing.T) {
53+
called := false
54+
inner := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
55+
called = true
56+
_, err := io.ReadAll(r.Body)
57+
if err != nil {
58+
http.Error(w, "read failed", http.StatusRequestEntityTooLarge)
59+
return
60+
}
61+
w.WriteHeader(http.StatusOK)
62+
})
63+
64+
handler := MaxBytesBodyByPath(16, map[string]int64{
65+
"/api/telemetry/v1/traces": 128,
66+
})(inner)
67+
req := httptest.NewRequest(http.MethodPost, "/api/telemetry/v1/traces", strings.NewReader(strings.Repeat("x", 64)))
68+
rec := httptest.NewRecorder()
69+
handler.ServeHTTP(rec, req)
70+
71+
if rec.Code != http.StatusOK {
72+
t.Fatalf("override body: got %d, want 200", rec.Code)
73+
}
74+
if !called {
75+
t.Fatal("expected inner handler to be called")
76+
}
77+
}

internal/api/browser_traces.go

Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,104 @@
1+
package api
2+
3+
import (
4+
"bytes"
5+
"errors"
6+
"io"
7+
"mime"
8+
"net/http"
9+
"strings"
10+
"time"
11+
12+
"github.com/gastownhall/wasteland/internal/observability"
13+
)
14+
15+
const maxBrowserTracePayloadBytes = 1 << 20
16+
17+
var browserTraceProxyClient = &http.Client{
18+
Timeout: 10 * time.Second,
19+
Transport: observability.NewTransport(http.DefaultTransport),
20+
}
21+
22+
func (s *Server) handleBrowserTraces(w http.ResponseWriter, r *http.Request) {
23+
if r.Method == http.MethodOptions {
24+
w.WriteHeader(http.StatusNoContent)
25+
return
26+
}
27+
if r.Method != http.MethodPost {
28+
writeError(w, http.StatusMethodNotAllowed, "method not allowed")
29+
return
30+
}
31+
32+
target := observability.BrowserTraceProxyTarget()
33+
if target == "" || !observability.BrowserTracingEnabled() {
34+
writeError(w, http.StatusNotFound, "browser tracing not enabled")
35+
return
36+
}
37+
if !allowedBrowserTraceContentType(r.Header.Get("Content-Type")) {
38+
writeError(w, http.StatusUnsupportedMediaType, "unsupported browser trace content type")
39+
return
40+
}
41+
42+
reader := http.MaxBytesReader(w, r.Body, maxBrowserTracePayloadBytes)
43+
defer reader.Close() //nolint:errcheck // request cleanup
44+
45+
payload, err := io.ReadAll(reader)
46+
if err != nil {
47+
status := http.StatusBadRequest
48+
if strings.Contains(err.Error(), "http: request body too large") {
49+
status = http.StatusRequestEntityTooLarge
50+
}
51+
writeError(w, status, "invalid browser trace payload")
52+
return
53+
}
54+
55+
req, err := http.NewRequestWithContext(r.Context(), http.MethodPost, target, bytes.NewReader(payload))
56+
if err != nil {
57+
writeError(w, http.StatusInternalServerError, "failed to prepare browser trace proxy request")
58+
return
59+
}
60+
req.Header.Set("Content-Type", contentTypeOrDefault(r.Header.Get("Content-Type"), "application/x-protobuf"))
61+
if enc := r.Header.Get("Content-Encoding"); enc != "" {
62+
req.Header.Set("Content-Encoding", enc)
63+
}
64+
65+
resp, err := browserTraceProxyClient.Do(req)
66+
if err != nil {
67+
writeError(w, http.StatusBadGateway, "browser trace collector unavailable")
68+
return
69+
}
70+
defer resp.Body.Close() //nolint:errcheck // upstream cleanup
71+
72+
for _, header := range []string{"Content-Type"} {
73+
if value := resp.Header.Get(header); value != "" {
74+
w.Header().Set(header, value)
75+
}
76+
}
77+
w.WriteHeader(resp.StatusCode)
78+
if _, copyErr := io.Copy(w, resp.Body); copyErr != nil && !errors.Is(copyErr, io.EOF) {
79+
return
80+
}
81+
}
82+
83+
func contentTypeOrDefault(value, fallback string) string {
84+
if strings.TrimSpace(value) == "" {
85+
return fallback
86+
}
87+
return value
88+
}
89+
90+
func allowedBrowserTraceContentType(value string) bool {
91+
if strings.TrimSpace(value) == "" {
92+
return true
93+
}
94+
mediaType, _, err := mime.ParseMediaType(value)
95+
if err != nil {
96+
return false
97+
}
98+
switch mediaType {
99+
case "application/x-protobuf", "application/json":
100+
return true
101+
default:
102+
return false
103+
}
104+
}

0 commit comments

Comments
 (0)