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
This commit is contained in:
2026-09-28 18:05:57 +00:00
parent 7a1ee6e5a8
commit 8cdc68e5ba
13 changed files with 737 additions and 55 deletions
+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
}