Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1f9e192001 |
@@ -29,10 +29,8 @@ 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; 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
|
||||
added; and panic recovery now routes the stack through slog instead of chi's
|
||||
plain-text stderr
|
||||
- 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
|
||||
|
||||
@@ -2,8 +2,6 @@ package middleware_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
@@ -182,15 +180,6 @@ 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) {
|
||||
@@ -240,63 +229,12 @@ func TestRecovererReturns500AndLogsThroughSlog(t *testing.T) {
|
||||
rec.Code, http.StatusInternalServerError)
|
||||
}
|
||||
|
||||
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())
|
||||
out := logbuf.String()
|
||||
if !strings.Contains(out, "panic recovered") {
|
||||
t.Fatalf("panic was not logged through slog: %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())
|
||||
if !strings.Contains(out, `"level":"ERROR"`) {
|
||||
t.Fatalf("panic log was not structured JSON at error level: %q", out)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -80,16 +80,12 @@ 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)
|
||||
err = b.flushLocked()
|
||||
b.flushLocked()
|
||||
})
|
||||
|
||||
// A failed final flush fails the stop, so the process
|
||||
// exits non-zero.
|
||||
return err
|
||||
return nil
|
||||
},
|
||||
})
|
||||
|
||||
@@ -113,12 +109,7 @@ func (b *Buffer) Append(v any) error {
|
||||
data := b.drainBuf()
|
||||
b.mu.Unlock()
|
||||
|
||||
go func() {
|
||||
writeErr := b.writeFile(data)
|
||||
if writeErr != nil {
|
||||
b.log.Error("flush reports failed", "error", writeErr)
|
||||
}
|
||||
}()
|
||||
go b.writeFile(data)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -137,10 +128,7 @@ func (b *Buffer) flushLoop() {
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
err := b.flushLocked()
|
||||
if err != nil {
|
||||
b.log.Error("flush reports failed", "error", err)
|
||||
}
|
||||
b.flushLocked()
|
||||
case <-b.done:
|
||||
return
|
||||
}
|
||||
@@ -149,19 +137,19 @@ func (b *Buffer) flushLoop() {
|
||||
|
||||
// flushLocked acquires the lock, drains the buffer, and
|
||||
// writes the data to a compressed file.
|
||||
func (b *Buffer) flushLocked() error {
|
||||
func (b *Buffer) flushLocked() {
|
||||
b.mu.Lock()
|
||||
|
||||
if b.buf.Len() == 0 {
|
||||
b.mu.Unlock()
|
||||
|
||||
return nil
|
||||
return
|
||||
}
|
||||
|
||||
data := b.drainBuf()
|
||||
b.mu.Unlock()
|
||||
|
||||
return b.writeFile(data)
|
||||
b.writeFile(data)
|
||||
}
|
||||
|
||||
// drainBuf copies the buffer contents and resets it.
|
||||
@@ -176,7 +164,7 @@ func (b *Buffer) drainBuf() []byte {
|
||||
|
||||
// writeFile creates a timestamped zstd-compressed JSONL file
|
||||
// in the data directory.
|
||||
func (b *Buffer) writeFile(data []byte) error {
|
||||
func (b *Buffer) writeFile(data []byte) {
|
||||
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)
|
||||
@@ -189,35 +177,31 @@ func (b *Buffer) writeFile(data []byte) error {
|
||||
filePerms,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create report file: %w", err)
|
||||
b.log.Error("create report file", "error", 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() }()
|
||||
|
||||
enc, err := zstd.NewWriter(f)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create zstd encoder: %w", err)
|
||||
b.log.Error("create zstd encoder", "error", err)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
_, err = enc.Write(data)
|
||||
if err != nil {
|
||||
_, writeErr := enc.Write(data)
|
||||
if writeErr != nil {
|
||||
b.log.Error("write compressed data", "error", writeErr)
|
||||
|
||||
_ = enc.Close()
|
||||
|
||||
return fmt.Errorf("write compressed data: %w", err)
|
||||
return
|
||||
}
|
||||
|
||||
err = enc.Close()
|
||||
if err != nil {
|
||||
return fmt.Errorf("close zstd encoder: %w", err)
|
||||
closeErr := enc.Close()
|
||||
if closeErr != nil {
|
||||
b.log.Error("close zstd encoder", "error", closeErr)
|
||||
}
|
||||
|
||||
err = f.Close()
|
||||
if err != nil {
|
||||
return fmt.Errorf("close report file: %w", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,8 +1,6 @@
|
||||
package reportbuf_test
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"io/fs"
|
||||
"os"
|
||||
"strings"
|
||||
"testing"
|
||||
@@ -54,46 +52,6 @@ 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()
|
||||
|
||||
|
||||
@@ -64,13 +64,4 @@ 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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user