diff --git a/storage/gcp/gcp.go b/storage/gcp/gcp.go index 8f95ae1f7..aa81978e5 100644 --- a/storage/gcp/gcp.go +++ b/storage/gcp/gcp.go @@ -35,6 +35,7 @@ import ( "errors" "fmt" "io" + "iter" "math" "net/http" "os" @@ -1455,6 +1456,75 @@ func (s *gcsStorage) deleteObjectsWithPrefix(ctx context.Context, objPrefix stri }) } +// MirrorWriter creates a new GCP storage for the MirrorTarget lifecycle mode. +func (s *Storage) MirrorWriter(ctx context.Context, opts *tessera.MirrorOptions) (tessera.MirrorWriter, tessera.LogReader, error) { + var err error + if s.cfg.GCSClient == nil { + s.cfg.GCSClient, err = gcs.NewClient(ctx, gcs.WithJSONReads()) + if err != nil { + return nil, nil, fmt.Errorf("failed to create GCS client: %v", err) + } + } + if s.cfg.SpannerClient == nil { + s.cfg.SpannerClient, err = spanner.NewClient(ctx, s.cfg.Spanner) + if err != nil { + return nil, nil, fmt.Errorf("failed to connect to Spanner: %v", err) + } + } + table := func(t string) string { + return s.cfg.SpannerTablePrefix + t + } + if err := initDB(ctx, s.cfg.Spanner, table); err != nil { + return nil, nil, fmt.Errorf("failed to verify/init Spanner schema: %v", err) + } + seq, err := newSpannerCoordinator(ctx, s.cfg.SpannerClient, table, 0) + if err != nil { + return nil, nil, fmt.Errorf("failed to create Spanner coordinator: %v", err) + } + gs := &gcsStorage{ + gcsClient: s.cfg.GCSClient, + bucket: s.cfg.Bucket, + bucketPrefix: s.cfg.BucketPrefix, + } + m := &mirrorWriter{ + seq: seq, + logStore: &logResourceStore{ + objStore: gs, + entriesPath: opts.EntriesPath(), + }, + } + lr := &LogReader{ + lrs: *m.logStore, + integratedSize: m.IntegratedSize, + nextIndex: seq.nextIndex, + } + return m, lr, nil +} + +// mirrorWriter implements the tessera.MirrorWriter lifecycle contract. +type mirrorWriter struct { + seq *spannerCoordinator + logStore *logResourceStore +} + +var _ tessera.MirrorWriter = &mirrorWriter{} + +func (w *mirrorWriter) IntegrateBundles(ctx context.Context, from uint64, bundles iter.Seq2[*api.EntryBundle, error]) (uint64, []byte, error) { + return 0, nil, errors.New("not implemented") +} + +func (w *mirrorWriter) IntegratedSize(ctx context.Context) (uint64, error) { + size, _, err := w.seq.currentTree(ctx) + if err != nil { + return 0, fmt.Errorf("failed to read integrated size: %w", err) + } + return size, nil +} + +func (w *mirrorWriter) UpdateCheckpoint(ctx context.Context, g func(oldCP []byte) (newCP []byte, err error)) error { + return errors.New("not implemented") +} + // MigrationWriter creates a new GCP storage for the MigrationTarget lifecycle mode. func (s *Storage) MigrationWriter(ctx context.Context, opts *tessera.MigrationOptions) (migrate.MigrationWriter, tessera.LogReader, error) { var err error