#7 Ingest HTTPS endpoint: accept report POST, size limit

Add internal/ingest implementing POST /v1/reports, wired into the core
build-tag-free handler so the same route serves on the dev server and the
Cloudflare Worker.

Response contract (ADR #6 §2.4):
- 202 Accepted for valid JSON within the 256 KiB cap ({"status":"accepted"})
- 413 for oversized bodies (Content-Length fast path AND a MaxBytesReader
  hard cap, so a missing/lying Content-Length cannot bypass the limit)
- 415 when Content-Type is not application/json
- 400 for malformed JSON or failed schema validation (generic error body,
  never echoes request content)
- 405 with Allow: POST for any non-POST method
- 503 when the storage Sink fails

Storage is decoupled behind a small Sink interface (Store(ctx, raw)) with a
NopSink default and a MemorySink for tests, so PII scrubbing (#8) and
encrypted R2 storage (#9) can slot in without touching the HTTP contract.
Rate limiting (429) and volumetric shedding stay a Cloudflare-edge/Pulumi
concern per #2 and are intentionally not implemented in the Worker.

Tests:
- Go unit tests (net/http/httptest) for every response code, including 413
  via both Content-Length and an oversized streamed body, plus boundary,
  storage-failure, and no-content-echo cases.
- Bruno API tests in OpenCollection YAML format under api-tests/, asserting
  the full contract against the local dev server via @usebruno/cli.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
2026-07-02 13:52:04 -05:00
co-authored by Claude Opus 4.8
parent 8e1dc66c54
commit 4c22056e66
14 changed files with 3795 additions and 6 deletions
+5
View File
@@ -0,0 +1,5 @@
name: local
variables:
- name: baseUrl
value: http://localhost:8787
+27
View File
@@ -0,0 +1,27 @@
info:
name: 400 - malformed JSON is rejected
type: http
seq: 4
http:
method: POST
url: "{{baseUrl}}/v1/reports"
headers:
- name: Content-Type
value: application/json
# Sent as a raw text body (so the bytes are transmitted verbatim) but declared
# as application/json, so the endpoint tries to parse it and fails -> 400.
body:
type: text
data: '{"appVersion": "1.0.0", "platform": "android", "report":'
runtime:
scripts:
- type: tests
code: |-
test("malformed JSON returns 400", function () {
expect(res.getStatus()).to.equal(400);
});
test("400 body carries a generic error string", function () {
expect(res.getBody().error).to.be.a("string");
});
+11
View File
@@ -0,0 +1,11 @@
opencollection: "1.0.0"
info:
name: LibreMail Bug-Report Ingest API tests
# Contract tests for POST /v1/reports (GitHub issue #7).
# Run against the local dev server (cmd/devserver) on http://localhost:8787,
# which serves the exact same http.Handler as the deployed Cloudflare Worker.
#
# This collection is authored in the OpenCollection YAML format (the presence
# of this opencollection.yml file selects the "yml" format in @usebruno/cli),
# per the maintainer's requirement to avoid the .bru format.
+33
View File
@@ -0,0 +1,33 @@
info:
name: 413 - oversized report is rejected
type: http
seq: 2
http:
method: POST
url: "{{baseUrl}}/v1/reports"
headers:
- name: Content-Type
value: application/json
body:
type: json
data: "{{oversizedReport}}"
runtime:
scripts:
# Build a payload larger than the 256 KiB cap at run time so this file stays
# small. The before-request script runs before body interpolation, so the
# {{oversizedReport}} placeholder above is replaced with this ~300 KiB body.
- type: before-request
code: |-
const filler = "x".repeat(300 * 1024); // 300 KiB, over the 256 KiB cap
bru.setVar("oversizedReport", JSON.stringify({
appVersion: "1.0.0",
platform: "android",
report: filler
}));
- type: tests
code: |-
test("body larger than 256 KiB returns 413", function () {
expect(res.getStatus()).to.equal(413);
});
+33
View File
@@ -0,0 +1,33 @@
info:
name: 202 - valid report is accepted
type: http
seq: 1
http:
method: POST
url: "{{baseUrl}}/v1/reports"
headers:
- name: Content-Type
value: application/json
body:
type: json
data: |-
{
"appVersion": "1.4.2 (142)",
"platform": "android",
"osVersion": "Android 14",
"device": "Pixel 7",
"clientTimestamp": "2026-07-02T12:34:56Z",
"report": "NullPointerException in SyncService\n at line 42\n<attached logs>"
}
runtime:
scripts:
- type: tests
code: |-
test("valid JSON report within the size limit returns 202", function () {
expect(res.getStatus()).to.equal(202);
});
test("202 body reports accepted status", function () {
expect(res.getBody().status).to.equal("accepted");
});
+27
View File
@@ -0,0 +1,27 @@
info:
name: 415 - non-JSON content type is rejected
type: http
seq: 3
http:
method: POST
url: "{{baseUrl}}/v1/reports"
headers:
- name: Content-Type
value: text/plain
body:
type: text
data: |-
{
"appVersion": "1.0.0",
"platform": "android",
"report": "body is valid JSON but the Content-Type is not application/json"
}
runtime:
scripts:
- type: tests
code: |-
test("Content-Type other than application/json returns 415", function () {
expect(res.getStatus()).to.equal(415);
});
+19
View File
@@ -0,0 +1,19 @@
info:
name: 405 - non-POST method is rejected
type: http
seq: 5
http:
method: GET
url: "{{baseUrl}}/v1/reports"
runtime:
scripts:
- type: tests
code: |-
test("a non-POST method returns 405", function () {
expect(res.getStatus()).to.equal(405);
});
test("405 advertises the allowed method via the Allow header", function () {
expect(res.getHeader("allow")).to.equal("POST");
});
+13 -4
View File
@@ -11,20 +11,29 @@ package handler
import (
"encoding/json"
"net/http"
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/ingest"
)
// serviceName identifies this service in responses.
const serviceName = "libremail-bug-report-ingest"
// New returns an http.Handler serving the bootstrap health/hello endpoints:
// New returns an http.Handler serving the ingest Worker's endpoints:
//
// GET / -> 200, JSON service/status/message
// GET /healthz -> 200, JSON {"status":"ok"}
// GET / -> 200, JSON service/status/message
// GET /healthz -> 200, JSON {"status":"ok"}
// POST /v1/reports -> 202 on accept; 400/413/415/405/503 per the ingest contract
//
// Any other path returns 404 and any non-GET method returns 405.
// Any other path returns 404. On the health/hello endpoints any non-GET method
// returns 405; on /v1/reports any non-POST method returns 405 (Allow: POST).
//
// The ingest endpoint is wired with a NopSink for now: it enforces the full
// HTTP contract (size cap + schema validation) but discards accepted bodies
// until the real storage sink (scrub + encrypt + R2) lands in #8/#9.
func New() http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("/healthz", healthz)
mux.Handle("/v1/reports", ingest.NewHandler(ingest.NopSink{}))
mux.HandleFunc("/", root)
return mux
}
+136
View File
@@ -0,0 +1,136 @@
// Package ingest implements the LibreMail bug-report ingest endpoint,
// POST /v1/reports.
//
// It owns the Worker-side half of the abuse policy fixed in ADR #6
// (docs/decisions/labels-and-abuse.md): the 256 KiB payload cap and JSON schema
// validation, plus the 202/400/413/415/405 response contract. Rate limiting
// (429) and volumetric shedding are NOT implemented here: per ADR #6 those are
// enforced at the Cloudflare edge via Pulumi-provisioned Rate Limiting rules
// (ticket #2), before the Worker ever runs. The endpoint does answer 503 when
// its storage Sink fails, matching the "storage unavailable" row of the ADR.
//
// The package is plain, build-tag-free Go so it is host-testable with the
// standard toolchain (go test ./...) and reused verbatim by the Wasm Worker.
package ingest
import (
"encoding/json"
"errors"
"io"
"mime"
"net/http"
"strings"
)
// MaxBodyBytes is the hard cap on the request body: 256 KiB (262,144 bytes),
// per ADR #6 §2.1. Bodies larger than this are rejected with 413.
const MaxBodyBytes = 256 * 1024
// contentTypeJSON is the only accepted request media type (parameters such as
// "; charset=utf-8" are allowed and ignored).
const contentTypeJSON = "application/json"
// Handler serves POST /v1/reports. Construct it with NewHandler.
type Handler struct {
sink Sink
}
// NewHandler returns a Handler that stores accepted reports in sink. A nil sink
// is replaced with NopSink, so the endpoint is always safe to construct.
func NewHandler(sink Sink) *Handler {
if sink == nil {
sink = NopSink{}
}
return &Handler{sink: sink}
}
// ServeHTTP implements the ingest response contract from ADR #6 §2.4:
//
// POST + valid JSON within the size cap -> 202 {"status":"accepted"}
// non-POST method -> 405 (Allow: POST)
// Content-Type not application/json -> 415
// body over 256 KiB -> 413 (Content-Length fast path + hard read cap)
// malformed JSON / failed validation -> 400
// sink failure -> 503
//
// Error bodies are generic ({"error":"..."}) and never reflect request content.
func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
w.Header().Set("Allow", http.MethodPost)
writeError(w, http.StatusMethodNotAllowed, "method not allowed")
return
}
if !isJSONContentType(r.Header.Get("Content-Type")) {
writeError(w, http.StatusUnsupportedMediaType, "content-type must be application/json")
return
}
// Fast path: reject an honestly-declared oversized body before reading it.
if r.ContentLength > MaxBodyBytes {
writeError(w, http.StatusRequestEntityTooLarge, "payload too large")
return
}
// Hard cap: enforce the limit on the stream too, so a missing, chunked, or
// lying Content-Length cannot smuggle a larger body past the fast path.
// MaxBytesReader allows up to MaxBodyBytes and errors on the next byte.
r.Body = http.MaxBytesReader(w, r.Body, MaxBodyBytes)
raw, err := io.ReadAll(r.Body)
if err != nil {
var maxErr *http.MaxBytesError
if errors.As(err, &maxErr) {
writeError(w, http.StatusRequestEntityTooLarge, "payload too large")
return
}
writeError(w, http.StatusBadRequest, "could not read request body")
return
}
if _, err := parseReport(raw); err != nil {
writeError(w, http.StatusBadRequest, "invalid report")
return
}
// Store the raw, validated body. Scrubbing/encryption happen downstream
// (#8/#9) behind the Sink; this package stays contract-only.
if err := h.sink.Store(r.Context(), raw); err != nil {
writeError(w, http.StatusServiceUnavailable, "storage unavailable")
return
}
writeJSON(w, http.StatusAccepted, map[string]string{"status": "accepted"})
}
// isJSONContentType reports whether ct declares an application/json body,
// tolerating media-type parameters (e.g. "application/json; charset=utf-8") and
// case. A missing or unparseable Content-Type is rejected.
func isJSONContentType(ct string) bool {
if strings.TrimSpace(ct) == "" {
return false
}
mediaType, _, err := mime.ParseMediaType(ct)
if err != nil {
return false
}
return mediaType == contentTypeJSON
}
// writeJSON writes body as a JSON response with the given status code.
func writeJSON(w http.ResponseWriter, status int, body any) {
w.Header().Set("Content-Type", "application/json; charset=utf-8")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(body)
}
// writeError writes a generic {"error": reason} body. reason is a fixed,
// caller-supplied string and must never include request-derived content.
func writeError(w http.ResponseWriter, status int, reason string) {
writeJSON(w, status, map[string]string{"error": reason})
}
// compile-time assurance that the stub sinks satisfy Sink.
var (
_ Sink = NopSink{}
_ Sink = (*MemorySink)(nil)
)
+233
View File
@@ -0,0 +1,233 @@
package ingest
import (
"bytes"
"context"
"encoding/json"
"errors"
"io"
"net/http"
"net/http/httptest"
"strings"
"testing"
)
// validBody is a minimal well-formed v1 report.
const validBody = `{"appVersion":"1.4.2 (142)","platform":"android","osVersion":"Android 14","device":"Pixel 7","clientTimestamp":"2026-07-02T12:34:56Z","report":"NPE in SyncService\n<logs>"}`
// newHandler returns a Handler backed by a MemorySink for assertions.
func newHandler(t *testing.T) (*Handler, *MemorySink) {
t.Helper()
sink := &MemorySink{}
return NewHandler(sink), sink
}
// do runs one request through h and returns the recorder.
func do(t *testing.T, h http.Handler, method, contentType, body string) *httptest.ResponseRecorder {
t.Helper()
var r io.Reader
if body != "" {
r = strings.NewReader(body)
}
req := httptest.NewRequest(method, "/v1/reports", r)
if contentType != "" {
req.Header.Set("Content-Type", contentType)
}
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
return rec
}
// decodeMap parses the JSON response body into a map.
func decodeMap(t *testing.T, rec *httptest.ResponseRecorder) map[string]string {
t.Helper()
var got map[string]string
if err := json.Unmarshal(rec.Body.Bytes(), &got); err != nil {
t.Fatalf("response body is not valid JSON: %v (body=%q)", err, rec.Body.String())
}
return got
}
func TestAccepted202(t *testing.T) {
h, sink := newHandler(t)
rec := do(t, h, http.MethodPost, "application/json", validBody)
if rec.Code != http.StatusAccepted {
t.Fatalf("status = %d, want %d (body=%q)", rec.Code, http.StatusAccepted, rec.Body.String())
}
if ct := rec.Header().Get("Content-Type"); ct != "application/json; charset=utf-8" {
t.Errorf("Content-Type = %q, want application/json; charset=utf-8", ct)
}
if got := decodeMap(t, rec); got["status"] != "accepted" {
t.Errorf("status field = %q, want accepted", got["status"])
}
if sink.Len() != 1 {
t.Fatalf("sink stored %d reports, want 1", sink.Len())
}
if stored := sink.Reports()[0]; !bytes.Equal(stored, []byte(validBody)) {
t.Errorf("sink stored %q, want the raw body %q", stored, validBody)
}
}
func TestAcceptedContentTypeWithCharset(t *testing.T) {
h, _ := newHandler(t)
for _, ct := range []string{"application/json", "application/json; charset=utf-8", "application/JSON", " application/json "} {
rec := do(t, h, http.MethodPost, ct, validBody)
if rec.Code != http.StatusAccepted {
t.Errorf("Content-Type %q: status = %d, want %d", ct, rec.Code, http.StatusAccepted)
}
}
}
func TestAcceptedAtExactLimit(t *testing.T) {
h, _ := newHandler(t)
prefix := `{"appVersion":"1.0.0","platform":"android","report":"`
suffix := `"}`
pad := MaxBodyBytes - len(prefix) - len(suffix)
body := prefix + strings.Repeat("x", pad) + suffix
if len(body) != MaxBodyBytes {
t.Fatalf("test setup: body len = %d, want exactly %d", len(body), MaxBodyBytes)
}
rec := do(t, h, http.MethodPost, "application/json", body)
if rec.Code != http.StatusAccepted {
t.Fatalf("body of exactly %d bytes: status = %d, want %d", MaxBodyBytes, rec.Code, http.StatusAccepted)
}
}
func TestOversizedViaContentLength413(t *testing.T) {
h, sink := newHandler(t)
big := strings.Repeat("x", MaxBodyBytes+1024)
body := `{"appVersion":"1.0.0","platform":"android","report":"` + big + `"}`
req := httptest.NewRequest(http.MethodPost, "/v1/reports", strings.NewReader(body))
req.Header.Set("Content-Type", "application/json")
if req.ContentLength <= MaxBodyBytes {
t.Fatalf("test setup: ContentLength = %d, want > %d (fast path must trigger)", req.ContentLength, MaxBodyBytes)
}
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusRequestEntityTooLarge {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusRequestEntityTooLarge)
}
if sink.Len() != 0 {
t.Errorf("sink stored %d reports, want 0 (oversized must not be stored)", sink.Len())
}
}
func TestOversizedViaStreamedBody413(t *testing.T) {
h, sink := newHandler(t)
big := strings.Repeat("x", MaxBodyBytes+1024)
body := `{"appVersion":"1.0.0","platform":"android","report":"` + big + `"}`
req := httptest.NewRequest(http.MethodPost, "/v1/reports", strings.NewReader(body))
req.Header.Set("Content-Type", "application/json")
// Simulate a missing/lying Content-Length so the fast path is bypassed and
// only the streamed MaxBytesReader cap can catch the oversized body.
req.ContentLength = -1
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusRequestEntityTooLarge {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusRequestEntityTooLarge)
}
if sink.Len() != 0 {
t.Errorf("sink stored %d reports, want 0 (oversized must not be stored)", sink.Len())
}
}
func TestWrongContentType415(t *testing.T) {
h, _ := newHandler(t)
for _, ct := range []string{"text/plain", "application/xml", "multipart/form-data", ""} {
rec := do(t, h, http.MethodPost, ct, validBody)
if rec.Code != http.StatusUnsupportedMediaType {
t.Errorf("Content-Type %q: status = %d, want %d", ct, rec.Code, http.StatusUnsupportedMediaType)
}
}
}
func TestMalformedJSON400(t *testing.T) {
h, _ := newHandler(t)
cases := map[string]string{
"truncated object": `{"appVersion":"1.0.0","platform":`,
"not json": `this is not json`,
"empty body": ``,
"trailing garbage": `{"appVersion":"1.0.0","platform":"android","report":"x"} garbage`,
"json array": `["not","an","object"]`,
}
for name, body := range cases {
t.Run(name, func(t *testing.T) {
rec := do(t, h, http.MethodPost, "application/json", body)
if rec.Code != http.StatusBadRequest {
t.Errorf("status = %d, want %d (body=%q)", rec.Code, http.StatusBadRequest, rec.Body.String())
}
})
}
}
func TestSchemaValidation400(t *testing.T) {
h, _ := newHandler(t)
cases := map[string]string{
"missing appVersion": `{"platform":"android","report":"boom"}`,
"empty appVersion": `{"appVersion":" ","platform":"android","report":"boom"}`,
"missing platform": `{"appVersion":"1.0.0","report":"boom"}`,
"missing report": `{"appVersion":"1.0.0","platform":"android"}`,
"empty report": `{"appVersion":"1.0.0","platform":"android","report":" "}`,
"bad clientTimestamp": `{"appVersion":"1.0.0","platform":"android","report":"boom","clientTimestamp":"not-a-time"}`,
}
for name, body := range cases {
t.Run(name, func(t *testing.T) {
rec := do(t, h, http.MethodPost, "application/json", body)
if rec.Code != http.StatusBadRequest {
t.Errorf("status = %d, want %d (body=%q)", rec.Code, http.StatusBadRequest, rec.Body.String())
}
})
}
}
func TestWrongMethod405(t *testing.T) {
h, _ := newHandler(t)
for _, method := range []string{http.MethodGet, http.MethodPut, http.MethodDelete, http.MethodPatch, http.MethodHead} {
rec := do(t, h, method, "application/json", validBody)
if rec.Code != http.StatusMethodNotAllowed {
t.Errorf("%s: status = %d, want %d", method, rec.Code, http.StatusMethodNotAllowed)
}
if allow := rec.Header().Get("Allow"); allow != http.MethodPost {
t.Errorf("%s: Allow header = %q, want POST", method, allow)
}
}
}
// failSink always fails, to exercise the 503 storage-unavailable path.
type failSink struct{}
func (failSink) Store(context.Context, []byte) error { return errors.New("boom") }
func TestStorageFailure503(t *testing.T) {
h := NewHandler(failSink{})
rec := do(t, h, http.MethodPost, "application/json", validBody)
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusServiceUnavailable)
}
}
func TestErrorBodiesDoNotEchoRequest(t *testing.T) {
h, _ := newHandler(t)
secret := "SUPER-SECRET-TOKEN-DO-NOT-LEAK"
body := `{"appVersion":"1.0.0","platform":"android","report":"boom","clientTimestamp":"` + secret + `"}`
rec := do(t, h, http.MethodPost, "application/json", body)
if rec.Code != http.StatusBadRequest {
t.Fatalf("status = %d, want %d", rec.Code, http.StatusBadRequest)
}
if strings.Contains(rec.Body.String(), secret) {
t.Errorf("error body echoed request content: %q", rec.Body.String())
}
}
func TestNilSinkDefaultsToNop(t *testing.T) {
h := NewHandler(nil)
rec := do(t, h, http.MethodPost, "application/json", validBody)
if rec.Code != http.StatusAccepted {
t.Fatalf("status = %d, want %d (nil sink should default to NopSink)", rec.Code, http.StatusAccepted)
}
}
+108
View File
@@ -0,0 +1,108 @@
package ingest
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
)
// Report is the v1 bug-report payload contract that the LibreMail app
// (LibreMail#33) sends to POST /v1/reports as a JSON object.
//
// Required fields: appVersion, platform, report. Everything else is optional.
// Validation is deliberately loose (see validate): the endpoint rejects only
// clearly-invalid payloads (missing required fields, absurd field lengths, a
// non-RFC3339 timestamp) so that newer app versions can add fields without
// breaking ingest. Unknown fields are ignored, not rejected.
//
// Example:
//
// {
// "appVersion": "1.4.2 (142)",
// "platform": "android",
// "osVersion": "Android 14",
// "device": "Pixel 7",
// "clientTimestamp": "2026-07-02T12:34:56Z",
// "report": "NullPointerException in SyncService...\n<logs>"
// }
type Report struct {
// AppVersion is the LibreMail app version, e.g. "1.4.2" or "1.4.2 (142)". Required.
AppVersion string `json:"appVersion"`
// Platform is the client platform, e.g. "android". Required.
Platform string `json:"platform"`
// Report is the free-text bug report: user description, logs, stack traces. Required.
Report string `json:"report"`
// OSVersion is the client OS version, e.g. "Android 14". Optional.
OSVersion string `json:"osVersion,omitempty"`
// Device is the device model/descriptor, e.g. "Pixel 7". Optional.
Device string `json:"device,omitempty"`
// ClientTimestamp is when the report was captured on the client, RFC 3339. Optional.
ClientTimestamp string `json:"clientTimestamp,omitempty"`
}
// Per-field length caps. These are defensive only: the 256 KiB body cap
// (MaxBodyBytes) is the real ceiling, and these just reject an absurd single
// metadata value with a clear 400 rather than storing it.
const (
maxAppVersionLen = 256
maxPlatformLen = 64
maxOSVersionLen = 128
maxDeviceLen = 256
)
// parseReport decodes raw into a Report and validates it. A non-nil error means
// the caller should respond 400; the error text is for internal logging only
// and must never be echoed to the client (it could reflect request contents).
//
// raw is expected to already be within MaxBodyBytes (the handler enforces the
// size cap before calling this), so no read limiting happens here.
func parseReport(raw []byte) (*Report, error) {
dec := json.NewDecoder(bytes.NewReader(raw))
// Note: DisallowUnknownFields is intentionally NOT set, so forward-compatible
// clients may add fields. We do reject trailing garbage after the object.
var rep Report
if err := dec.Decode(&rep); err != nil {
return nil, fmt.Errorf("json decode: %w", err)
}
if dec.More() {
return nil, errors.New("unexpected trailing data after JSON value")
}
if err := rep.validate(); err != nil {
return nil, err
}
return &rep, nil
}
// validate enforces the loose v1 schema rules and returns the first problem
// found, or nil if the report is acceptable.
func (r *Report) validate() error {
appVersion := strings.TrimSpace(r.AppVersion)
platform := strings.TrimSpace(r.Platform)
switch {
case appVersion == "":
return errors.New("appVersion is required")
case len(appVersion) > maxAppVersionLen:
return errors.New("appVersion is too long")
case platform == "":
return errors.New("platform is required")
case len(platform) > maxPlatformLen:
return errors.New("platform is too long")
case strings.TrimSpace(r.Report) == "":
return errors.New("report is required")
case len(r.OSVersion) > maxOSVersionLen:
return errors.New("osVersion is too long")
case len(r.Device) > maxDeviceLen:
return errors.New("device is too long")
}
if r.ClientTimestamp != "" {
if _, err := time.Parse(time.RFC3339, r.ClientTimestamp); err != nil {
return fmt.Errorf("clientTimestamp must be RFC3339: %w", err)
}
}
return nil
}
+64
View File
@@ -0,0 +1,64 @@
package ingest
import (
"context"
"sync"
)
// Sink is the storage seam for accepted reports. The ingest endpoint calls
// Store exactly once, with the raw (already size-capped and schema-validated)
// request body, after it has decided to accept a report.
//
// It is intentionally tiny and knows nothing about scrubbing, encryption, or
// R2: those land downstream (PII redaction is #8, encrypted R2 storage is #9).
// Keeping storage behind this interface lets #9 supply a real implementation
// without touching the HTTP-contract code in this package.
type Sink interface {
// Store persists a raw accepted report body. Returning a non-nil error
// causes the endpoint to answer 503 (storage unavailable) per the ADR #6
// response contract; the client is expected to retry with backoff.
Store(ctx context.Context, raw []byte) error
}
// NopSink is a Sink that accepts and discards every report. It is the default
// wired into the handler until the real storage sink (#9) exists, so the HTTP
// contract (202/400/413/415/405) is fully exercisable today without any
// storage backend.
type NopSink struct{}
// Store discards raw and always succeeds.
func (NopSink) Store(context.Context, []byte) error { return nil }
// MemorySink is an in-memory Sink that retains a copy of every stored report.
// It is intended for tests and local experimentation only: it grows without
// bound and is not safe to use as a production backend.
type MemorySink struct {
mu sync.Mutex
reports [][]byte
}
// Store appends a copy of raw to the in-memory slice.
func (m *MemorySink) Store(_ context.Context, raw []byte) error {
cp := make([]byte, len(raw))
copy(cp, raw)
m.mu.Lock()
defer m.mu.Unlock()
m.reports = append(m.reports, cp)
return nil
}
// Reports returns a snapshot copy of the report bodies stored so far.
func (m *MemorySink) Reports() [][]byte {
m.mu.Lock()
defer m.mu.Unlock()
out := make([][]byte, len(m.reports))
copy(out, m.reports)
return out
}
// Len reports how many reports have been stored.
func (m *MemorySink) Len() int {
m.mu.Lock()
defer m.mu.Unlock()
return len(m.reports)
}
+3 -1
View File
@@ -7,9 +7,11 @@
"build": "go run github.com/syumai/workers/cmd/workers-assets-gen && tinygo build -o ./build/app.wasm -target wasm -no-debug ./worker",
"dev": "wrangler dev",
"start": "wrangler dev",
"deploy": "wrangler deploy"
"deploy": "wrangler deploy",
"test:api": "cd api-tests && bru run --env local"
},
"devDependencies": {
"@usebruno/cli": "^3.5.0",
"wrangler": "^4.106.0"
}
}
+3083 -1
View File
File diff suppressed because it is too large Load Diff