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
longer logged (only its length), and `client_id`, `timestamp` and decode error
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
plain-text stderr
added; panic recovery now routes the stack through slog instead of chi's
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
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
+67 -5
View File
@@ -2,6 +2,8 @@ package middleware_test
import (
"bytes"
"encoding/json"
"errors"
"log/slog"
"net/http"
"net/http/httptest"
@@ -180,6 +182,15 @@ func TestMaxBodyBytesRejectsOversizeOnNonReadingRoute(t *testing.T) {
if called {
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) {
@@ -229,12 +240,63 @@ func TestRecovererReturns500AndLogsThroughSlog(t *testing.T) {
rec.Code, http.StatusInternalServerError)
}
out := logbuf.String()
if !strings.Contains(out, "panic recovered") {
t.Fatalf("panic was not logged through slog: %q", out)
var record map[string]any
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"`) {
t.Fatalf("panic log was not structured JSON at error level: %q", out)
if record["msg"] != "panic recovered" || record["level"] != "ERROR" {
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
// invocation must not close an already-closed channel
// (which would panic) or flush again.
var err error
b.stopOnce.Do(func() {
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()
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
}
@@ -128,7 +137,10 @@ func (b *Buffer) flushLoop() {
for {
select {
case <-ticker.C:
b.flushLocked()
err := b.flushLocked()
if err != nil {
b.log.Error("flush reports failed", "error", err)
}
case <-b.done:
return
}
@@ -137,19 +149,19 @@ func (b *Buffer) flushLoop() {
// flushLocked acquires the lock, drains the buffer, and
// writes the data to a compressed file.
func (b *Buffer) flushLocked() {
func (b *Buffer) flushLocked() error {
b.mu.Lock()
if b.buf.Len() == 0 {
b.mu.Unlock()
return
return nil
}
data := b.drainBuf()
b.mu.Unlock()
b.writeFile(data)
return b.writeFile(data)
}
// 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
// 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")
name := fmt.Sprintf("reports-%s.jsonl.zst", ts)
path := filepath.Join(b.dataDir, name)
@@ -177,31 +189,35 @@ func (b *Buffer) writeFile(data []byte) {
filePerms,
)
if err != nil {
b.log.Error("create report file", "error", err)
return
return fmt.Errorf("create report file: %w", err)
}
// 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() }()
enc, err := zstd.NewWriter(f)
if err != nil {
b.log.Error("create zstd encoder", "error", err)
return
return fmt.Errorf("create zstd encoder: %w", err)
}
_, writeErr := enc.Write(data)
if writeErr != nil {
b.log.Error("write compressed data", "error", writeErr)
_, err = enc.Write(data)
if err != nil {
_ = enc.Close()
return
return fmt.Errorf("write compressed data: %w", err)
}
closeErr := enc.Close()
if closeErr != nil {
b.log.Error("close zstd encoder", "error", closeErr)
err = enc.Close()
if err != nil {
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
import (
"errors"
"io/fs"
"os"
"strings"
"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 {
t.Helper()
+9
View File
@@ -64,4 +64,13 @@ func TestHealthCheckRejectsOversizeBody(t *testing.T) {
t.Fatalf("status = %d, want %d",
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)
}
}