Files
LibreMail-Bug-Report-Ingest/internal/storage/sink.go
T
JMR-devandClaude Opus 4.8 cb7d9c67a0 #10 Report lifecycle/status metadata (pending/removed/published)
Add a status layer over the ObjectStore so each stored report has a
lifecycle state, encoded in its object key as reports/<status>/<id>:
pending (new reports), removed (#11), published (#15). Encoding status in
the key prefix means "list pending" is a single prefix listing with no
secondary index to drift, so it returns exactly the pending reports by
construction.

Storage:
- Extend ObjectStore with List(ctx, prefix) and Delete(ctx, key); implement
  in MemoryStore (host) and the js/wasm R2Store. R2Store.List drives the R2
  binding's list() directly to page a prefix (the syumai helper takes no
  options), so a status with >1000 objects is still enumerated exactly.
- The ingest Sink now writes new reports under reports/pending/<id>, so
  accepted reports enter the lifecycle as pending. The <id> is stable across
  transitions.

lifecycle package:
- Manager over an ObjectStore: ListPending, GetPending(id), MarkRemoved(id),
  MarkPublished(id). A transition copies the opaque ciphertext frame to the
  destination status key and deletes the source key — bytes are never
  decrypted or re-encrypted; no key is needed to change status.
- Copy-then-delete is idempotent and retry-safe: Put(dest) before Delete(src)
  never loses a report, a retry converges (re-Put identical bytes, Delete the
  leftover source), and a transition of an id not in the source status returns
  ErrUnknownReport (unless it is already at the destination -> idempotent nil).

Tests (host, MemoryStore, no TinyGo):
- List-pending exactness across a mix of pending/removed/published.
- pending->removed and pending->published leave the pending set, appear under
  the target, and move byte-identical ciphertext that still decrypts.
- Idempotent retry and convergence from an interrupted (both-keys) state.
- Unknown/terminal-state ids error sensibly; new Sink reports list as pending.

Closes #10

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-02 15:25:11 -05:00

85 lines
2.9 KiB
Go

package storage
import (
"context"
"crypto/rand"
"encoding/hex"
"fmt"
"time"
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/crypto"
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/ingest"
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/scrub"
)
// Sink is the real ingest.Sink. For each accepted report it:
//
// 1. scrubs the raw body (best-effort PII redaction, #8);
// 2. seals the scrubbed bytes with the keyring's active AES-256 key
// (AES-256-GCM, ADR #5) into a self-describing frame;
// 3. writes the frame to the ObjectStore under a unique, unguessable key.
//
// The store only ever receives ciphertext. Sink is provider-agnostic: paired
// with a MemoryStore it runs on the host (tests, cmd/devserver); paired with an
// R2Store it runs in the Worker. The Worker's keyring loading (from Secrets
// Store) is handled by WorkerSink, which delegates the pipeline to a Sink.
type Sink struct {
store ObjectStore
keyring *crypto.Keyring
keyFn func() string
}
// Option customises a Sink.
type Option func(*Sink)
// WithKeyFunc overrides the object-key generator. Intended for tests that need a
// deterministic key; production uses the default random key.
func WithKeyFunc(fn func() string) Option {
return func(s *Sink) { s.keyFn = fn }
}
// NewSink returns a Sink that stores into store, encrypting under kr's active
// key. kr must be non-nil.
func NewSink(store ObjectStore, kr *crypto.Keyring, opts ...Option) *Sink {
s := &Sink{store: store, keyring: kr, keyFn: defaultObjectKey}
for _, o := range opts {
o(s)
}
return s
}
// Store scrubs, encrypts, and persists one accepted report. A non-nil return
// makes the ingest endpoint answer 503 (per ADR #6), so callers should retry.
func (s *Sink) Store(ctx context.Context, raw []byte) error {
scrubbed := scrub.Scrub(raw)
sealed, err := crypto.Seal(s.keyring, scrubbed)
if err != nil {
return fmt.Errorf("storage: seal: %w", err)
}
if err := s.store.Put(ctx, s.keyFn(), sealed); err != nil {
return fmt.Errorf("storage: put: %w", err)
}
return nil
}
// defaultObjectKey builds the key for a freshly accepted report. Reports enter
// the lifecycle (#10) as pending, so the key lives under the pending status
// prefix: reports/pending/<id>. The id is a UTC timestamp (for rough
// lexicographic ordering, convenient for the weekly publish job) plus 80 bits of
// CSPRNG randomness (so ids are unguessable and collision-free within a second),
// and is stable as the report later transitions to published/removed.
func defaultObjectKey() string {
return ReportKey(StatusPending, newReportID())
}
// newReportID mints a stable, unguessable report id: <utc-ts>-<80-bit-rand>.
func newReportID() string {
var b [10]byte
_, _ = rand.Read(b[:])
ts := time.Now().UTC().Format("20060102T150405")
return ts + "-" + hex.EncodeToString(b[:])
}
// Sink satisfies the ingest storage seam.
var _ ingest.Sink = (*Sink)(nil)