1 Commits
Author SHA1 Message Date
clawbot 8cdc68e5ba fix(backend): report ingest correctness — propagate storage failure, 413 on oversize, global body cap (closes #23)
check / check (push) Successful in 45s
A buffer failure on POST /api/v1/reports now returns 500 instead of a
false `ok`: the failure is server-side and a client can retry. An
over-limit body returns 413 (errors.As on `*http.MaxBytesError`);
malformed JSON stays 400. A MaxBodyBytes middleware (1 MiB) caps every
route; a route group can only lower that limit. The raw geo blob is no
longer logged, only its length; client_id, timestamp and decode error
text are length-bounded before logging. A decodeJSON handler helper is
added. Panic recovery routes the stack through slog as structured
JSON. Writing a report file now returns its error, so a failed final
flush fails the stop and the process exits non-zero.

Model: opus-5-5
2026-09-28 18:05:57 +00:00
5 changed files with 160 additions and 29 deletions
+4 -2
View File
@@ -29,8 +29,10 @@ latest run passes.
route, not just the report route; the raw attacker-controlled `geo` blob is no route, not just the report route; the raw attacker-controlled `geo` blob is no
longer logged (only its length), and `client_id`, `timestamp` and decode error longer logged (only its length), and `client_id`, `timestamp` and decode error
text are length-bounded before logging; a `decodeJSON` handler helper was text are length-bounded before logging; a `decodeJSON` handler helper was
added; and panic recovery now routes the stack through slog instead of chi's added; panic recovery now routes the stack through slog instead of chi's
plain-text stderr plain-text stderr; and writing a report file now returns its error, so a
failed final flush on shutdown makes the process exit non-zero instead of
losing the buffered reports silently
- 2026-09-21: shutdown lifecycle correctness. The process now shuts down through - 2026-09-21: shutdown lifecycle correctness. The process now shuts down through
fx instead of `os.Exit`, so every component's `OnStop` runs and buffered fx instead of `os.Exit`, so every component's `OnStop` runs and buffered
reports are flushed to disk on `SIGTERM` — previously a full flush window of reports are flushed to disk on `SIGTERM` — previously a full flush window of
+67 -5
View File
@@ -2,6 +2,8 @@ package middleware_test
import ( import (
"bytes" "bytes"
"encoding/json"
"errors"
"log/slog" "log/slog"
"net/http" "net/http"
"net/http/httptest" "net/http/httptest"
@@ -180,6 +182,15 @@ func TestMaxBodyBytesRejectsOversizeOnNonReadingRoute(t *testing.T) {
if called { if called {
t.Fatal("handler ran despite oversize body") t.Fatal("handler ran despite oversize body")
} }
if got := rec.Body.String(); got != "{\"status\":\"error\"}\n" {
t.Errorf("body = %q, want %q", got, "{\"status\":\"error\"}\n")
}
got := rec.Header().Get("Content-Type")
if got != "application/json; charset=utf-8" {
t.Errorf("Content-Type = %q, want a JSON content type", got)
}
} }
func TestMaxBodyBytesAllowsWithinLimit(t *testing.T) { func TestMaxBodyBytesAllowsWithinLimit(t *testing.T) {
@@ -229,12 +240,63 @@ func TestRecovererReturns500AndLogsThroughSlog(t *testing.T) {
rec.Code, http.StatusInternalServerError) rec.Code, http.StatusInternalServerError)
} }
out := logbuf.String() var record map[string]any
if !strings.Contains(out, "panic recovered") {
t.Fatalf("panic was not logged through slog: %q", out) err := json.Unmarshal(logbuf.Bytes(), &record)
if err != nil {
t.Fatalf("panic log is not one JSON record: %v (%q)", err, logbuf.String())
} }
if !strings.Contains(out, `"level":"ERROR"`) { if record["msg"] != "panic recovered" || record["level"] != "ERROR" {
t.Fatalf("panic log was not structured JSON at error level: %q", out) t.Errorf("log record = %v, want msg %q at level ERROR",
record, "panic recovered")
}
if record["panic"] != "boom" {
t.Errorf("panic field = %v, want %q", record["panic"], "boom")
}
stack, _ := record["stack"].(string)
if !strings.HasPrefix(stack, "goroutine ") {
t.Errorf("stack field = %q, want a stack trace", stack)
}
}
// TestRecovererRepanicsOnAbortHandler checks that a handler aborting
// with http.ErrAbortHandler is not treated as a crash: Recoverer
// panics again so the server aborts the response, and logs nothing.
func TestRecovererRepanicsOnAbortHandler(t *testing.T) {
t.Parallel()
var logbuf bytes.Buffer
mw := middleware.NewWithLogger(
slog.New(slog.NewJSONHandler(&logbuf, nil)),
)
handler := mw.Recoverer()(
http.HandlerFunc(func(_ http.ResponseWriter, _ *http.Request) {
panic(http.ErrAbortHandler)
}),
)
rec := httptest.NewRecorder()
req := httptest.NewRequestWithContext(t.Context(), http.MethodGet, "/", http.NoBody)
var recovered any
func() {
defer func() { recovered = recover() }()
handler.ServeHTTP(rec, req)
}()
err, _ := recovered.(error)
if !errors.Is(err, http.ErrAbortHandler) {
t.Errorf("Recoverer panicked with %v, want http.ErrAbortHandler", recovered)
}
if logbuf.Len() != 0 {
t.Errorf("abort was logged: %q", logbuf.String())
} }
} }
+38 -22
View File
@@ -80,12 +80,16 @@ func New(
// stopOnce makes OnStop idempotent: a second // stopOnce makes OnStop idempotent: a second
// invocation must not close an already-closed channel // invocation must not close an already-closed channel
// (which would panic) or flush again. // (which would panic) or flush again.
var err error
b.stopOnce.Do(func() { b.stopOnce.Do(func() {
close(b.done) close(b.done)
b.flushLocked() err = b.flushLocked()
}) })
return nil // A failed final flush fails the stop, so the process
// exits non-zero.
return err
}, },
}) })
@@ -109,7 +113,12 @@ func (b *Buffer) Append(v any) error {
data := b.drainBuf() data := b.drainBuf()
b.mu.Unlock() b.mu.Unlock()
go b.writeFile(data) go func() {
writeErr := b.writeFile(data)
if writeErr != nil {
b.log.Error("flush reports failed", "error", writeErr)
}
}()
return nil return nil
} }
@@ -128,7 +137,10 @@ func (b *Buffer) flushLoop() {
for { for {
select { select {
case <-ticker.C: case <-ticker.C:
b.flushLocked() err := b.flushLocked()
if err != nil {
b.log.Error("flush reports failed", "error", err)
}
case <-b.done: case <-b.done:
return return
} }
@@ -137,19 +149,19 @@ func (b *Buffer) flushLoop() {
// flushLocked acquires the lock, drains the buffer, and // flushLocked acquires the lock, drains the buffer, and
// writes the data to a compressed file. // writes the data to a compressed file.
func (b *Buffer) flushLocked() { func (b *Buffer) flushLocked() error {
b.mu.Lock() b.mu.Lock()
if b.buf.Len() == 0 { if b.buf.Len() == 0 {
b.mu.Unlock() b.mu.Unlock()
return return nil
} }
data := b.drainBuf() data := b.drainBuf()
b.mu.Unlock() b.mu.Unlock()
b.writeFile(data) return b.writeFile(data)
} }
// drainBuf copies the buffer contents and resets it. // drainBuf copies the buffer contents and resets it.
@@ -164,7 +176,7 @@ func (b *Buffer) drainBuf() []byte {
// writeFile creates a timestamped zstd-compressed JSONL file // writeFile creates a timestamped zstd-compressed JSONL file
// in the data directory. // in the data directory.
func (b *Buffer) writeFile(data []byte) { func (b *Buffer) writeFile(data []byte) error {
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z") ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
name := fmt.Sprintf("reports-%s.jsonl.zst", ts) name := fmt.Sprintf("reports-%s.jsonl.zst", ts)
path := filepath.Join(b.dataDir, name) path := filepath.Join(b.dataDir, name)
@@ -177,31 +189,35 @@ func (b *Buffer) writeFile(data []byte) {
filePerms, filePerms,
) )
if err != nil { if err != nil {
b.log.Error("create report file", "error", err) return fmt.Errorf("create report file: %w", err)
return
} }
// Closes the file on the early returns below. The success
// path closes it explicitly to check the error; closing it
// a second time here is harmless.
defer func() { _ = f.Close() }() defer func() { _ = f.Close() }()
enc, err := zstd.NewWriter(f) enc, err := zstd.NewWriter(f)
if err != nil { if err != nil {
b.log.Error("create zstd encoder", "error", err) return fmt.Errorf("create zstd encoder: %w", err)
return
} }
_, writeErr := enc.Write(data) _, err = enc.Write(data)
if writeErr != nil { if err != nil {
b.log.Error("write compressed data", "error", writeErr)
_ = enc.Close() _ = enc.Close()
return return fmt.Errorf("write compressed data: %w", err)
} }
closeErr := enc.Close() err = enc.Close()
if closeErr != nil { if err != nil {
b.log.Error("close zstd encoder", "error", closeErr) return fmt.Errorf("close zstd encoder: %w", err)
} }
err = f.Close()
if err != nil {
return fmt.Errorf("close report file: %w", err)
}
return nil
} }
@@ -1,6 +1,8 @@
package reportbuf_test package reportbuf_test
import ( import (
"errors"
"io/fs"
"os" "os"
"strings" "strings"
"testing" "testing"
@@ -52,6 +54,46 @@ func TestFlushOnShutdown(t *testing.T) {
} }
} }
// TestFailedFinalFlushFailsStop proves a final flush that cannot
// write its file makes the stop fail, which makes the process
// exit non-zero instead of dropping the buffered reports silently.
func TestFailedFinalFlushFailsStop(t *testing.T) {
dir := t.TempDir()
t.Setenv("DATA_DIR", dir)
var buf *reportbuf.Buffer
app := fxtest.New(t,
fx.Provide(
globals.New,
logger.New,
config.New,
reportbuf.New,
),
fx.Populate(&buf),
)
app.RequireStart()
err := buf.Append(map[string]string{"probe": "shutdown"})
if err != nil {
t.Fatalf("append report: %v", err)
}
// Removing the data directory leaves the final flush nowhere to
// write. A read-only directory would not do: tests run as root
// in the backend image, and root ignores the read-only bit.
err = os.RemoveAll(dir)
if err != nil {
t.Fatalf("remove data dir: %v", err)
}
err = app.Stop(t.Context())
if !errors.Is(err, fs.ErrNotExist) {
t.Fatalf("stop error = %v, want the final flush's error", err)
}
}
func hasReportFile(t *testing.T, dir string) bool { func hasReportFile(t *testing.T, dir string) bool {
t.Helper() t.Helper()
+9
View File
@@ -64,4 +64,13 @@ func TestHealthCheckRejectsOversizeBody(t *testing.T) {
t.Fatalf("status = %d, want %d", t.Fatalf("status = %d, want %d",
rec.Code, http.StatusRequestEntityTooLarge) rec.Code, http.StatusRequestEntityTooLarge)
} }
if got := rec.Body.String(); got != "{\"status\":\"error\"}\n" {
t.Errorf("body = %q, want %q", got, "{\"status\":\"error\"}\n")
}
got := rec.Header().Get("Content-Type")
if got != "application/json; charset=utf-8" {
t.Errorf("Content-Type = %q, want a JSON content type", got)
}
} }