From 141468be7b0435ef5d019e982e5a53888603b94d Mon Sep 17 00:00:00 2001 From: Erin Perrine Date: Fri, 21 Aug 2026 16:40:42 +0000 Subject: [PATCH] storage/gcp: optionally set object retention on immutable log resources Add Config.ObjectRetentionPeriod and Config.ObjectRetentionLocked. When a period is set, the driver writes the log's immutable resources - full tiles and full entry bundles - with a GCS object retention configuration expiring that period after the write, in Locked or Unlocked mode, so that on a bucket with object retention enabled they cannot be overwritten or deleted until then. The checkpoint and partial tiles/bundles are always written without retention so that they can still be updated and garbage collected. Both the Appender and MigrationWriter lifecycles honour the option; left unset, behaviour is unchanged. The decision is made in logResourceStore, which knows from the layout which resources are partial, and passed to objStore.setObject as a new retention argument which gcsStorage sets on the Writer. Tests cover the layout decision (checkpoint and partials never retained, full tiles and bundles retained with the configured mode and expiry, and nothing retained when no period is set), that gcsStorage sends a non-nil retention as the object's retention metadata on upload and none for nil (against a minimal fake of the GCS JSON API), and validation of a negative period in New. --- storage/gcp/gcp.go | 48 +++++++++-- storage/gcp/gcp_test.go | 172 +++++++++++++++++++++++++++++++++++++++- 2 files changed, 208 insertions(+), 12 deletions(-) diff --git a/storage/gcp/gcp.go b/storage/gcp/gcp.go index 7725470cf..9b523d823 100644 --- a/storage/gcp/gcp.go +++ b/storage/gcp/gcp.go @@ -165,6 +165,14 @@ type Config struct { // TODO: consider providing a mechanism to set both BucketPrefix and SpannerTablePrefix // from a single string to avoid misconfiguration foot-guns. SpannerTablePrefix string + + // ObjectRetentionPeriod, if non-zero, writes full tiles and entry bundles with a GCS object + // retention (https://cloud.google.com/storage/docs/object-lock) expiring this long after the + // write. Checkpoints and partial tiles/bundles are never retained. Requires object retention + // enabled on the Bucket, and the JSON API (the gRPC API cannot set retention). + ObjectRetentionPeriod time.Duration + // ObjectRetentionLocked selects "Locked" rather than "Unlocked" retention mode. + ObjectRetentionLocked bool } // tablePrefixRE matches valid values for a Spanner table prefix: empty, or a leading @@ -176,6 +184,9 @@ func New(ctx context.Context, cfg Config) (tessera.Driver, error) { if !tablePrefixRE.MatchString(cfg.SpannerTablePrefix) { return nil, fmt.Errorf("invalid SpannerTablePrefix %q: must start with a letter, contain only letters, digits, or underscores, and be at most 64 characters long", cfg.SpannerTablePrefix) } + if cfg.ObjectRetentionPeriod < 0 { + return nil, fmt.Errorf("invalid ObjectRetentionPeriod %v: must not be negative", cfg.ObjectRetentionPeriod) + } if cfg.HTTPClient == nil { cfg.HTTPClient = http.DefaultClient } @@ -281,8 +292,10 @@ func (s *Storage) newAppender(ctx context.Context, o objStore, seq *spannerCoord a := &Appender{ logStore: &logResourceStore{ - objStore: o, - entriesPath: opts.EntriesPath(), + objStore: o, + entriesPath: opts.EntriesPath(), + retentionPeriod: s.cfg.ObjectRetentionPeriod, + retentionLocked: s.cfg.ObjectRetentionLocked, }, sequencer: seq, cpUpdated: make(chan struct{}), @@ -518,7 +531,7 @@ func (a *Appender) updateCheckpoint(ctx context.Context, size uint64, root []byt // objStore describes a type which can store and retrieve objects. type objStore interface { getObject(ctx context.Context, obj string) ([]byte, *gcs.ReaderObjectAttrs, error) - setObject(ctx context.Context, obj string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string) error + setObject(ctx context.Context, obj string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string, retention *gcs.ObjectRetention) error deleteObjectsWithPrefix(ctx context.Context, prefix string) error } @@ -526,10 +539,25 @@ type objStore interface { type logResourceStore struct { objStore objStore entriesPath func(uint64, uint8) string + // See Config.ObjectRetentionPeriod. + retentionPeriod time.Duration + retentionLocked bool +} + +// retention returns the retention to set on a full (partial == 0) tile or entry bundle, or nil. +func (lrs *logResourceStore) retention(partial uint8) *gcs.ObjectRetention { + if lrs.retentionPeriod == 0 || partial > 0 { + return nil + } + mode := "Unlocked" + if lrs.retentionLocked { + mode = "Locked" + } + return &gcs.ObjectRetention{Mode: mode, RetainUntil: time.Now().Add(lrs.retentionPeriod)} } func (lrs *logResourceStore) setCheckpoint(ctx context.Context, cpRaw []byte) error { - return lrs.objStore.setObject(ctx, layout.CheckpointPath, cpRaw, nil, ckptContType, ckptCacheControl) + return lrs.objStore.setObject(ctx, layout.CheckpointPath, cpRaw, nil, ckptContType, ckptCacheControl, nil) } func (lrs *logResourceStore) getCheckpoint(ctx context.Context) ([]byte, error) { @@ -549,7 +577,7 @@ func (s *logResourceStore) setTile(ctx context.Context, level, index uint64, par start := time.Now() tPath := layout.TilePath(level, index, partial) - err := s.objStore.setObject(ctx, tPath, data, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl) + err := s.objStore.setObject(ctx, tPath, data, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl, s.retention(partial)) opsHistogram.Record(ctx, time.Since(start).Milliseconds(), metric.WithAttributes(opNameKey.String("writeTile"))) return err } @@ -627,7 +655,7 @@ func (s *logResourceStore) setEntryBundle(ctx context.Context, bundleIndex uint6 // Note that setObject does an idempotent interpretation of DoesNotExist - it only // returns an error if the named object exists _and_ contains different data to what's // passed in here. - if err := s.objStore.setObject(ctx, objName, bundleRaw, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl); err != nil { + if err := s.objStore.setObject(ctx, objName, bundleRaw, &gcs.Conditions{DoesNotExist: true}, logContType, logCacheControl, s.retention(p)); err != nil { return fmt.Errorf("setObject(%q): %v", objName, err) } @@ -1349,11 +1377,12 @@ func (s *gcsStorage) getObject(ctx context.Context, obj string) ([]byte, *gcs.Re // // cond can be used to specify preconditions for the write (e.g. write iff not exists, write iff // current generation is X, etc.), or nil can be passed if no preconditions are desired. +// retention, if non-nil, is applied to the written object. // // Note that when preconditions are specified and are not met, an error will be returned *unless* // the currently stored data is bit-for-bit identical to the data to-be-written. // This is intended to provide idempotentency for writes. -func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string) error { +func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte, cond *gcs.Conditions, contType string, cacheCtl string, retention *gcs.ObjectRetention) error { return otel.TraceErr(ctx, "tessera.storage.gcp.setObject", tracer, func(ctx context.Context, span trace.Span) error { if s.bucketPrefix != "" { objName = filepath.Join(s.bucketPrefix, objName) @@ -1373,6 +1402,7 @@ func (s *gcsStorage) setObject(ctx context.Context, objName string, data []byte, } w.ContentType = contType w.CacheControl = cacheCtl + w.Retention = retention // Limit the amount of memory used for buffers, see https://pkg.go.dev/cloud.google.com/go/storage#Writer w.ChunkSize = len(data) + 1024 if _, err := w.Write(data); err != nil { @@ -1490,7 +1520,9 @@ func (s *Storage) MigrationWriter(ctx context.Context, opts *tessera.MigrationOp bucket: s.cfg.Bucket, bucketPrefix: s.cfg.BucketPrefix, }, - entriesPath: opts.EntriesPath(), + entriesPath: opts.EntriesPath(), + retentionPeriod: s.cfg.ObjectRetentionPeriod, + retentionLocked: s.cfg.ObjectRetentionLocked, }, } diff --git a/storage/gcp/gcp_test.go b/storage/gcp/gcp_test.go index 33540a268..7aa869944 100644 --- a/storage/gcp/gcp_test.go +++ b/storage/gcp/gcp_test.go @@ -18,10 +18,15 @@ import ( "bytes" "context" "crypto/sha256" + "encoding/json" "errors" "fmt" "log" "log/slog" + "mime" + "mime/multipart" + "net/http" + "net/http/httptest" "os" "reflect" "strings" @@ -175,6 +180,26 @@ func TestSpannerTablePrefixValidation(t *testing.T) { } } +func TestObjectRetentionPeriodValidation(t *testing.T) { + for _, test := range []struct { + period time.Duration + wantErr bool + }{ + {period: 0}, + {period: 24 * time.Hour}, + {period: -time.Second, wantErr: true}, + } { + _, err := New(t.Context(), Config{ + Bucket: "bucket", + Spanner: "projects/p/instances/i/databases/d", + ObjectRetentionPeriod: test.period, + }) + if gotErr := err != nil; gotErr != test.wantErr { + t.Errorf("New with ObjectRetentionPeriod %v: got err %v, want err %t", test.period, err, test.wantErr) + } + } +} + func TestSpannerSequencerPushback(t *testing.T) { ctx := t.Context() @@ -607,7 +632,7 @@ func TestPublishTree(t *testing.T) { t.Fatalf("publishTree: %v", err) } cpOld := []byte("bananas") - if err := m.setObject(ctx, layout.CheckpointPath, cpOld, nil, "", ""); err != nil { + if err := m.setObject(ctx, layout.CheckpointPath, cpOld, nil, "", "", nil); err != nil { t.Fatalf("setObject(bananas): %v", err) } updatesSeen := 0 @@ -880,12 +905,14 @@ func expectedPartialPrefixes(size uint64, entriesPath func(uint64, uint8) string type memObjStore struct { sync.RWMutex - mem map[string][]byte + mem map[string][]byte + retention map[string]*gcs.ObjectRetention } func newMemObjStore() *memObjStore { return &memObjStore{ - mem: make(map[string][]byte), + mem: make(map[string][]byte), + retention: make(map[string]*gcs.ObjectRetention), } } @@ -901,7 +928,7 @@ func (m *memObjStore) getObject(_ context.Context, obj string) ([]byte, *gcs.Rea } // TODO(phboneff): add content type tests -func (m *memObjStore) setObject(_ context.Context, obj string, data []byte, cond *gcs.Conditions, _, _ string) error { +func (m *memObjStore) setObject(_ context.Context, obj string, data []byte, cond *gcs.Conditions, _, _ string, retention *gcs.ObjectRetention) error { m.Lock() defer m.Unlock() @@ -915,6 +942,7 @@ func (m *memObjStore) setObject(_ context.Context, obj string, data []byte, cond } } m.mem[obj] = data + m.retention[obj] = retention return nil } @@ -941,6 +969,142 @@ func (m *memObjStore) deleteObjectsWithPrefix(_ context.Context, prefix string) return nil } +// TestImmutableResourceRetention checks that only full tiles and entry bundles are retained. +func TestImmutableResourceRetention(t *testing.T) { + const period = 24 * time.Hour + + for _, test := range []struct { + name string + period time.Duration + locked bool + wantMode string + }{ + {name: "no retention configured"}, + {name: "unlocked", period: period, wantMode: "Unlocked"}, + {name: "locked", period: period, locked: true, wantMode: "Locked"}, + } { + t.Run(test.name, func(t *testing.T) { + ctx := t.Context() + m := newMemObjStore() + s := &logResourceStore{ + objStore: m, + entriesPath: layout.EntriesPath, + retentionPeriod: test.period, + retentionLocked: test.locked, + } + + start := time.Now() + if err := s.setCheckpoint(ctx, []byte("checkpoint")); err != nil { + t.Fatalf("setCheckpoint: %v", err) + } + for _, p := range []uint8{0, 10} { + if err := s.setTile(ctx, 0, 0, p, []byte("tile")); err != nil { + t.Fatalf("setTile(partial=%d): %v", p, err) + } + if err := s.setEntryBundle(ctx, 0, p, []byte("bundle")); err != nil { + t.Fatalf("setEntryBundle(partial=%d): %v", p, err) + } + } + + for _, mutable := range []string{layout.CheckpointPath, layout.TilePath(0, 0, 10), layout.EntriesPath(0, 10)} { + if got := m.retention[mutable]; got != nil { + t.Errorf("%s: got retention %+v, want none", mutable, got) + } + } + for _, immutable := range []string{layout.TilePath(0, 0, 0), layout.EntriesPath(0, 0)} { + got := m.retention[immutable] + if test.period == 0 { + if got != nil { + t.Errorf("%s: got retention %+v, want none", immutable, got) + } + continue + } + if got == nil { + t.Fatalf("%s: got no retention, want mode %s", immutable, test.wantMode) + } + if got.Mode != test.wantMode { + t.Errorf("%s: got mode %q, want %q", immutable, got.Mode, test.wantMode) + } + if lo, hi := start.Add(test.period), time.Now().Add(test.period); got.RetainUntil.Before(lo) || got.RetainUntil.After(hi) { + t.Errorf("%s: got RetainUntil %v, want within [%v, %v]", immutable, got.RetainUntil, lo, hi) + } + } + }) + } +} + +// TestSetObjectRetention checks the retention gcsStorage sends on upload, using a fake GCS JSON API. +func TestSetObjectRetention(t *testing.T) { + const bucket = "test-bucket" + until := time.Date(2030, 1, 2, 3, 4, 5, 0, time.UTC) + + var mu sync.Mutex + gotMeta := map[string]map[string]any{} + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost || !strings.HasPrefix(r.URL.Path, "/upload/") { + t.Errorf("unexpected request: %s %s", r.Method, r.URL) + http.Error(w, "unexpected", http.StatusNotImplemented) + return + } + // Multipart upload: JSON metadata part first. + _, params, err := mime.ParseMediaType(r.Header.Get("Content-Type")) + if err != nil { + t.Errorf("upload Content-Type: %v", err) + return + } + part, err := multipart.NewReader(r.Body, params["boundary"]).NextPart() + if err != nil { + t.Errorf("multipart metadata part: %v", err) + return + } + var meta map[string]any + if err := json.NewDecoder(part).Decode(&meta); err != nil { + t.Errorf("decode metadata: %v", err) + return + } + name, _ := meta["name"].(string) + mu.Lock() + gotMeta[name] = meta + mu.Unlock() + w.Header().Set("Content-Type", "application/json") + _, _ = fmt.Fprintf(w, `{"bucket":%q,"name":%q,"generation":"1"}`, bucket, name) + })) + defer srv.Close() + + t.Setenv("STORAGE_EMULATOR_HOST", strings.TrimPrefix(srv.URL, "http://")) + ctx := t.Context() + c, err := gcs.NewClient(ctx, gcs.WithJSONReads()) + if err != nil { + t.Fatalf("NewClient: %v", err) + } + defer func() { + if err := c.Close(); err != nil { + t.Logf("Close: %v", err) + } + }() + + s := &gcsStorage{gcsClient: c, bucket: bucket} + if err := s.setObject(ctx, "retained", []byte("data"), nil, "application/octet-stream", "", &gcs.ObjectRetention{Mode: "Locked", RetainUntil: until}); err != nil { + t.Fatalf("setObject(retained): %v", err) + } + if err := s.setObject(ctx, "unretained", []byte("data"), nil, "application/octet-stream", "", nil); err != nil { + t.Fatalf("setObject(unretained): %v", err) + } + + mu.Lock() + defer mu.Unlock() + retained, unretained := gotMeta["retained"], gotMeta["unretained"] + if retained == nil || unretained == nil { + t.Fatalf("expected uploads not seen, got %v", gotMeta) + } + if ret, ok := retained["retention"].(map[string]any); !ok || ret["mode"] != "Locked" || ret["retainUntilTime"] != until.Format(time.RFC3339) { + t.Errorf("retained object: got retention %v, want mode Locked retainUntilTime %s", retained["retention"], until.Format(time.RFC3339)) + } + if ret, ok := unretained["retention"]; ok { + t.Errorf("unretained object: got retention %v, want none", ret) + } +} + func mustGenerateKeys(t *testing.T) (note.Signer, note.Verifier) { sk, vk, err := note.GenerateKey(nil, "testlog") if err != nil {