Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
161 changes: 159 additions & 2 deletions storage/gcp/gcp.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,12 +31,14 @@ package gcp
import (
"bytes"
"context"
"encoding/binary"
"encoding/gob"
"errors"
"fmt"
"io"
"iter"
"math"
"math/bits"
"net/http"
"os"
"path/filepath"
Expand Down Expand Up @@ -1525,7 +1527,15 @@ func (w *mirrorWriter) IntegratedSize(ctx context.Context) (uint64, error) {
}

func (w *mirrorWriter) UpdateCheckpoint(ctx context.Context, fn func(oldCP []byte) (newCP []byte, err error)) error {
_, err := w.seq.dbPool.ReadWriteTransactionWithOptions(ctx, func(ctx context.Context, txn *spanner.ReadWriteTransaction) error {
// Read the integrated size before starting the PubCoord transaction: the Spanner client rejects nested
// transactions, and reading IntCoord within that transaction would make publishing contend with integration.
// A stale value is fine since the integrated size never shrinks and the resources it implies are immutable.
treeSize, err := w.IntegratedSize(ctx)
if err != nil {
return fmt.Errorf("failed to read integrated size: %v", err)
}

_, err = w.seq.dbPool.ReadWriteTransactionWithOptions(ctx, func(ctx context.Context, txn *spanner.ReadWriteTransaction) error {
if _, err := txn.ReadRowWithOptions(ctx, w.seq.table("PubCoord"), spanner.Key{0}, []string{"publishedAt"}, &spanner.ReadOptions{LockHint: spannerpb.ReadRequest_LOCK_HINT_EXCLUSIVE}); err != nil {
return fmt.Errorf("failed to lock PubCoord: %w", err)
}
Expand All @@ -1547,7 +1557,10 @@ func (w *mirrorWriter) UpdateCheckpoint(ctx context.Context, fn func(oldCP []byt
return fmt.Errorf("checkpoint size %d exceeds maximum supported size %d", newSize, int64(math.MaxInt64))
}

// TODO: Ensure partial tiles and entry bundles for the new checkpoint exist.
// Ensure that all the partial resources implied by the new size are present.
if err := w.ensureGeometry(ctx, newSize, treeSize); err != nil {
return fmt.Errorf("failed to ensure geometry: %w", err)
}

if err := w.logStore.setCheckpoint(ctx, newCP); err != nil {
return fmt.Errorf("failed to set checkpoint: %w", err)
Expand All @@ -1559,6 +1572,150 @@ func (w *mirrorWriter) UpdateCheckpoint(ctx context.Context, fn func(oldCP []byt
return err
}

// ensureGeometry checks that the partial tiles and entry bundles implied by the new size are
// present in the stored resources.
//
// If an implied partial resource is not already present, this function will attempt to create
// it from a strictly larger resource whose presence is implied by treeSize.
func (w *mirrorWriter) ensureGeometry(ctx context.Context, cpSize, treeSize uint64) error {
// If cpSize is zero then no tree exists.
// If cpSize == treeSize the resources are guaranteed present by integration.
// In both cases there's nothing to do.
if cpSize == 0 || cpSize == treeSize {
return nil
}
if cpSize > treeSize {
return fmt.Errorf("new size %d is greater than integrated size %d", cpSize, treeSize)
}
ml := maxLevel(cpSize)
idx := (cpSize - 1) >> layout.TileHeight
for l := uint64(0); l <= uint64(ml); l, idx = l+1, idx>>layout.TileHeight {
treeP := layout.PartialTileSize(l, idx, treeSize)
cpP := layout.PartialTileSize(l, idx, cpSize)

if cpP == treeP {
// Nothing to be done at this level.
continue
}

if l == 0 {
if err := w.ensurePartialBundle(ctx, idx, cpP, treeP); err != nil {
return err
}
}
if err := w.ensurePartialTile(ctx, l, idx, cpP, treeP); err != nil {
return err
}
}
return nil
}

// ensurePartialBundle ensures that the partial entry bundle implied by the new size is
// present in the stored resources.
//
// If the implied partial entry bundle is not already present, this function will attempt to create
// it from the entry bundle implied by treeSize.
func (w *mirrorWriter) ensurePartialBundle(ctx context.Context, idx uint64, cpP, treeP uint8) error {
if cpP == 0 {
return nil
}

// Check if the entry bundle already exists.
switch _, err := w.logStore.getEntryBundle(ctx, idx, cpP); {
case err == nil:
// Already exists, we're done!
return nil
case !errors.Is(err, os.ErrNotExist):
return fmt.Errorf("failed to read entry bundle %d.%d: %v", idx, cpP, err)
}

// Read bundle implied by the tree size...
d, err := w.logStore.getEntryBundle(ctx, idx, treeP)
if err != nil {
return fmt.Errorf("failed to read entry bundle @%d.%d: %v", idx, treeP, err)
}
eb := &api.EntryBundle{}
if err := eb.UnmarshalText(d); err != nil {
return fmt.Errorf("failed to unmarshal entry bundle @%d.%d: %v", idx, treeP, err)
}

eb.Entries = eb.Entries[:cpP]
d, err = marshalTlogEntryBundle(eb)
if err != nil {
return fmt.Errorf("failed to marshal entry bundle @%d.%d: %v", idx, cpP, err)
}
if err := w.logStore.setEntryBundle(ctx, idx, cpP, d); err != nil {
return fmt.Errorf("failed to write entry bundle @%d.%d: %v", idx, cpP, err)
}
return nil
}

// ensurePartialTile ensures that the partial tile implied by the new size is
// present in the stored resources.
//
// If the implied partial tile is not already present, this function will attempt to create
// it from the tile implied by treeSize.
func (w *mirrorWriter) ensurePartialTile(ctx context.Context, l uint64, idx uint64, cpP, treeP uint8) error {
if cpP == 0 {
return nil
}

// Check if the tile already exists.
switch _, err := w.logStore.getTile(ctx, l, idx, cpP); {
case err == nil:
// Already exists, we're done!
return nil
case !errors.Is(err, os.ErrNotExist):
return fmt.Errorf("failed to read tile @%d/%d.%d: %v", l, idx, cpP, err)
}

// Read tile implied by the tree size...
d, err := w.logStore.getTile(ctx, l, idx, treeP)
if err != nil {
return fmt.Errorf("failed to read tile @%d/%d.%d: %v", l, idx, treeP, err)
}
t := &api.HashTile{}
if err := t.UnmarshalText(d); err != nil {
return fmt.Errorf("failed to unmarshal tile @%d/%d.%d: %v", l, idx, treeP, err)
}

t.Nodes = t.Nodes[:cpP]
d, err = t.MarshalText()
if err != nil {
return fmt.Errorf("failed to marshal tile @%d/%d.%d: %v", l, idx, cpP, err)
}
if err := w.logStore.setTile(ctx, l, idx, cpP, d); err != nil {
return fmt.Errorf("failed to write tile @%d/%d.%d: %v", l, idx, cpP, err)
}
return nil
}

// maxLevel returns the maximum tile level for a tree of the given size.
func maxLevel(sz uint64) int {
if sz == 0 {
return 0
}
return (bits.Len64(sz) - 1) / layout.TileHeight
}

// marshalTlogEntryBundle returns a tlog-tiles compatible serialization of the provided entry bundle.
func marshalTlogEntryBundle(b *api.EntryBundle) ([]byte, error) {
total := 0
for i, e := range b.Entries {
l := len(e)
if l >= 1<<16 {
return nil, fmt.Errorf("entry #%d has length %d >= 1<<16", i, l)
}
total += 2 + l
}
data := make([]byte, 0, total)
for _, e := range b.Entries {
data = binary.BigEndian.AppendUint16(data, uint16(len(e)))
data = append(data, e...)
}
return data, nil
}

// 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
Expand Down
Loading