diff --git a/pkg/dotc1z/engine/pebble/adapter.go b/pkg/dotc1z/engine/pebble/adapter.go index fbd4aac28..0f1e17a9d 100644 --- a/pkg/dotc1z/engine/pebble/adapter.go +++ b/pkg/dotc1z/engine/pebble/adapter.go @@ -230,6 +230,18 @@ func (a *Adapter) EndSync(ctx context.Context) error { zap.Error(err), ) } + // Build the per-entitlement grant digests at seal time, after all + // grants are written and while still on the fresh-sync NoSync path + // (the EndFreshSync flush below hardens the nodes). Non-fatal, like + // the stats sidecar: a missing/stale digest only forces a diff + // consumer onto the on-demand index fold, and the on-Open migration + // backfills it next time the file opens writable. + if err := a.engine.BuildAllGrantDigests(ctx, existing.GetSyncId()); err != nil { + ctxzap.Extract(ctx).Warn("pebble: build grant digests failed; grant-diff callers will fall back to on-demand index folds until the next Open backfills them", + zap.String("sync_id", existing.GetSyncId()), + zap.Error(err), + ) + } // Single flush + WAL fsync at sync end. This is the durability // boundary — counterpart to MarkFreshSync at StartNewSync. After // this returns, all writes from the sync are on disk. diff --git a/pkg/dotc1z/engine/pebble/adapter_clone_sync.go b/pkg/dotc1z/engine/pebble/adapter_clone_sync.go index ef24f3864..93d54551b 100644 --- a/pkg/dotc1z/engine/pebble/adapter_clone_sync.go +++ b/pkg/dotc1z/engine/pebble/adapter_clone_sync.go @@ -110,6 +110,8 @@ func cloneSync( {GrantByPrincipalSyncLowerBound(syncIDBytes), GrantByPrincipalSyncUpperBound(syncIDBytes)}, {GrantByPrincipalResourceTypeSyncLowerBound(syncIDBytes), GrantByPrincipalResourceTypeSyncUpperBound(syncIDBytes)}, {GrantByNeedsExpansionSyncLowerBound(syncIDBytes), GrantByNeedsExpansionSyncUpperBound(syncIDBytes)}, + {GrantByEntPrincHashSyncLowerBound(syncIDBytes), GrantByEntPrincHashSyncUpperBound(syncIDBytes)}, + {DigestSyncLowerBound(syncIDBytes), DigestSyncUpperBound(syncIDBytes)}, {encodeAssetPrefix(syncIDBytes), upperBoundOf(encodeAssetPrefix(syncIDBytes))}, // Stats sidecar — single key per sync; copyRange's [lo, hi) // shape requires a half-open range, so we synthesize one diff --git a/pkg/dotc1z/engine/pebble/adapter_diff.go b/pkg/dotc1z/engine/pebble/adapter_diff.go index 5fe8b49e6..25040f63d 100644 --- a/pkg/dotc1z/engine/pebble/adapter_diff.go +++ b/pkg/dotc1z/engine/pebble/adapter_diff.go @@ -23,13 +23,17 @@ import ( // deletions are NOT captured — that matches the SQLite behavior // in pkg/dotc1z/diff.go (additions-only diff). // -// Strategy: for each record-type-bearing keyspace under -// appliedSyncID, iterate the source keys, recompute the same -// record's primary key under baseSyncID, Get from base; if base -// returns ErrNotFound, write the value under diffSyncID's -// keyspace in the same record type. Index entries are -// recomputed on write so the diff sync has its own (fresh) -// indexes that match the records that landed. +// Strategy: for the small record types (resource_types, resources, +// entitlements, assets), iterate the source keys under appliedSyncID, +// recompute the same record's primary key under baseSyncID, Get from +// base; if base returns ErrNotFound, write the value under +// diffSyncID's keyspace in the same record type. GRANTS — the type +// that dominates every real file — are diffed via the per-entitlement +// grant digests instead (see diffGrants): unchanged entitlements are +// skipped with a single root comparison, and only the principal-hash +// buckets that actually differ are scanned. Index entries are +// recomputed on write so the diff sync has its own (fresh) indexes +// that match the records that landed. // // Returns the diff sync's ID. func generateSyncDiff(ctx context.Context, a *Adapter, baseSyncID, appliedSyncID string) (string, error) { @@ -106,7 +110,7 @@ func generateSyncDiff(ctx context.Context, a *Adapter, baseSyncID, appliedSyncID if err := diffEntitlements(ctx, a, baseBytes, appliedBytes, diffSyncID); err != nil { return "", fmt.Errorf("generate-diff: entitlements: %w", err) } - if err := diffGrants(ctx, a, baseBytes, appliedBytes, diffSyncID); err != nil { + if err := diffGrants(ctx, a, baseBytes, appliedBytes, baseSyncID, appliedSyncID, diffSyncID); err != nil { return "", fmt.Errorf("generate-diff: grants: %w", err) } if err := diffAssets(ctx, a, baseBytes, appliedBytes, diffSyncID); err != nil { @@ -190,7 +194,84 @@ func diffEntitlements(ctx context.Context, a *Adapter, baseBytes, appliedBytes [ }) } -func diffGrants(ctx context.Context, a *Adapter, baseBytes, appliedBytes []byte, diffSyncID string) error { +// diffGrants computes the grants set difference using the +// per-entitlement grant digests instead of scanning every applied +// grant. +// +// Walk: enumerate the distinct entitlement_ids in the applied sync's +// hash index (one seek each), compare each entitlement's digest between +// base and applied — one root read per side when nothing changed, the +// overwhelmingly common case — and materialize only the principal-hash +// buckets the comparison flags as dirty. Grants in dirty buckets are +// probed against base's PRIMARY keyspace by external_id, which +// preserves the additions-only contract exactly: a grant whose +// external_id exists in base under a different entitlement/principal +// is still "present in base" (not emitted), which is why the probe +// targets the primary key rather than comparing index keys. +// +// Why pruning is sound: equal (count, digest) node pairs mean the two +// sides hold identical grant content-hash sets, and the content hash +// folds external_id — so a clean entitlement/bucket cannot contain an +// applied external_id that base lacks. Dirtiness over-approximates +// additions (it also fires for removals and source-set changes, which +// the base probe then filters out), never under-approximates them. +// +// NOTE: grants without an entitlement or principal ref have no +// hash-index entry and are invisible to the digest. They are silently +// skipped. A future O(1) coverage check (e.g. a stored grant count in +// the sync run record) will restore detection of this case. +func diffGrants(ctx context.Context, a *Adapter, baseBytes, appliedBytes []byte, baseSyncID, appliedSyncID, diffSyncID string) error { + eng := a.engine + + ents, err := eng.distinctDigestPartitions(ctx, grantDigestSpec, appliedBytes) + if err != nil { + return err + } + for _, ent := range ents { + if err := ctx.Err(); err != nil { + return err + } + // Entitlements only present in base (fully removed) are not in + // `ents` and are correctly skipped: they cannot contain + // additions. Digests missing on either side (e.g. a ghost + // entitlement with grants but no entitlement record) degrade + // to an on-demand index fold inside DirtyEntitlementBuckets. + dirty, err := eng.DirtyEntitlementBuckets(ctx, baseSyncID, eng, appliedSyncID, ent) + if err != nil { + return err + } + for _, bucket := range dirty { + var innerErr error + err := eng.IterateGrantsByEntitlementBucket(ctx, appliedSyncID, ent, bucket, func(rec *v3.GrantRecord) bool { + exists, probeErr := existsAt(eng.DB(), encodeGrantKey(baseBytes, rec.GetExternalId())) + if probeErr != nil { + innerErr = probeErr + return false + } + if exists { + return true + } + if putErr := eng.PutGrantRecord(ctx, rec); putErr != nil { + innerErr = putErr + return false + } + return true + }) + if err != nil { + return err + } + if innerErr != nil { + return innerErr + } + } + } + return nil +} + +// diffGrantsFullScan is the O(applied grants) path: iterate every +// applied grant, probe base by external_id, emit on miss. Used +// directly by benchmarks and available as a correctness reference. +func diffGrantsFullScan(ctx context.Context, a *Adapter, baseBytes, appliedBytes []byte, diffSyncID string) error { srcPrefix := encodeGrantPrefix(appliedBytes) return iterDiff(ctx, a.engine.DB(), srcPrefix, upperBoundOf(srcPrefix), func(_ []byte, val []byte) error { var rec v3.GrantRecord diff --git a/pkg/dotc1z/engine/pebble/adapter_diff_trie_test.go b/pkg/dotc1z/engine/pebble/adapter_diff_trie_test.go new file mode 100644 index 000000000..bee449d96 --- /dev/null +++ b/pkg/dotc1z/engine/pebble/adapter_diff_trie_test.go @@ -0,0 +1,201 @@ +package pebble + +import ( + "context" + "testing" + + "github.com/cockroachdb/pebble/v2" + + v2 "github.com/conductorone/baton-sdk/pb/c1/connector/v2" + v3 "github.com/conductorone/baton-sdk/pb/c1/storage/v3" + "github.com/conductorone/baton-sdk/pkg/connectorstore" +) + +func mkV2Ent(id string) *v2.Entitlement { + return v2.Entitlement_builder{ + Id: id, + Resource: v2.Resource_builder{ + Id: v2.ResourceId_builder{ + ResourceType: "app", + Resource: "github", + }.Build(), + }.Build(), + }.Build() +} + +// runSync starts a sync, writes the given entitlements + grants, and +// seals it (EndSync builds the per-entitlement tries). Returns the +// sync id. +func runSync(t *testing.T, a *Adapter, ents []*v2.Entitlement, grants []*v2.Grant) string { + t.Helper() + ctx := context.Background() + syncID, err := a.StartNewSync(ctx, connectorstore.SyncTypeFull, "") + if err != nil { + t.Fatalf("StartNewSync: %v", err) + } + if len(ents) > 0 { + if err := a.PutEntitlements(ctx, ents...); err != nil { + t.Fatalf("PutEntitlements: %v", err) + } + } + if len(grants) > 0 { + if err := a.PutGrants(ctx, grants...); err != nil { + t.Fatalf("PutGrants: %v", err) + } + } + if err := a.EndSync(ctx); err != nil { + t.Fatalf("EndSync: %v", err) + } + return syncID +} + +// diffGrantIDs runs generateSyncDiff and returns the set of grant +// external_ids that landed in the diff sync. +func diffGrantIDs(t *testing.T, a *Adapter, baseSyncID, appliedSyncID string) map[string]bool { + t.Helper() + ctx := context.Background() + diffID, err := generateSyncDiff(ctx, a, baseSyncID, appliedSyncID) + if err != nil { + t.Fatalf("generateSyncDiff: %v", err) + } + got := map[string]bool{} + if err := a.engine.IterateGrantsBySync(ctx, diffID, func(r *v3.GrantRecord) bool { + got[r.GetExternalId()] = true + return true + }); err != nil { + t.Fatalf("IterateGrantsBySync(diff): %v", err) + } + return got +} + +func requireDiffIDs(t *testing.T, got map[string]bool, want ...string) { + t.Helper() + wantSet := map[string]bool{} + for _, id := range want { + wantSet[id] = true + if !got[id] { + t.Errorf("diff missing expected grant %q", id) + } + } + for id := range got { + if !wantSet[id] { + t.Errorf("diff contains unexpected grant %q", id) + } + } +} + +// TestGenerateSyncDiffTrieMultiEntitlement exercises the trie-driven +// grants diff across the full semantic matrix in one pair of syncs: +// unchanged entitlements (pruned at the root), an addition inside an +// existing entitlement, an addition under a brand-new entitlement, a +// removal (must not be emitted), and a grant that MOVED to a different +// entitlement while keeping its external_id (dirty bucket on both +// entitlements, but per the additions-only-by-external_id contract it +// must not be emitted). +func TestGenerateSyncDiffTrieMultiEntitlement(t *testing.T) { + a := newAdapter(t) + + base := runSync(t, a, + []*v2.Entitlement{mkV2Ent("ent-A"), mkV2Ent("ent-B"), mkV2Ent("ent-C")}, + []*v2.Grant{ + mkV2Grant("g1", "ent-A", "user", "alice"), + mkV2Grant("g2", "ent-A", "user", "bob"), + mkV2Grant("g3", "ent-B", "user", "carol"), + mkV2Grant("g4", "ent-B", "user", "dave"), + mkV2Grant("g5", "ent-C", "user", "eve"), + mkV2Grant("g6", "ent-A", "user", "frank"), + }) + applied := runSync(t, a, + []*v2.Entitlement{mkV2Ent("ent-A"), mkV2Ent("ent-B"), mkV2Ent("ent-C"), mkV2Ent("ent-D")}, + []*v2.Grant{ + mkV2Grant("g1", "ent-A", "user", "alice"), // unchanged + mkV2Grant("g2", "ent-A", "user", "bob"), // unchanged + mkV2Grant("g3", "ent-B", "user", "carol"), // unchanged + mkV2Grant("g4", "ent-B", "user", "dave"), // unchanged + // g5 removed (ent-C empty in applied) — removals not emitted. + mkV2Grant("g6", "ent-B", "user", "frank"), // moved ent-A→ent-B, same external_id — NOT an addition + mkV2Grant("g7", "ent-B", "user", "grace"), // addition in an existing entitlement + mkV2Grant("g8", "ent-D", "user", "heidi"), // addition under a new entitlement + }) + + got := diffGrantIDs(t, a, base, applied) + requireDiffIDs(t, got, "g7", "g8") +} + +// TestGenerateSyncDiffOrphanGrantSkipped documents that a grant +// without entitlement/principal refs has no hash-index entry and is +// invisible to the trie-driven diff path. It is silently skipped. +// Normal grants in the same sync (g2) are still emitted correctly. +// TODO: restore detection via an O(1) coverage check (e.g. a stored +// grant count in the sync run record). +func TestGenerateSyncDiffOrphanGrantSkipped(t *testing.T) { + ctx := context.Background() + a := newAdapter(t) + + base := runSync(t, a, + []*v2.Entitlement{mkV2Ent("ent-A")}, + []*v2.Grant{mkV2Grant("g1", "ent-A", "user", "alice")}) + + applied, err := a.StartNewSync(ctx, connectorstore.SyncTypeFull, "") + if err != nil { + t.Fatalf("StartNewSync: %v", err) + } + if err := a.PutEntitlements(ctx, mkV2Ent("ent-A")); err != nil { + t.Fatalf("PutEntitlements: %v", err) + } + if err := a.PutGrants(ctx, + mkV2Grant("g1", "ent-A", "user", "alice"), + mkV2Grant("g2", "ent-A", "user", "bob"), + ); err != nil { + t.Fatalf("PutGrants: %v", err) + } + orphan := v3.GrantRecord_builder{ + ExternalId: "g-orphan", + }.Build() + if err := a.engine.PutGrantRecord(ctx, orphan); err != nil { + t.Fatalf("PutGrantRecord(orphan): %v", err) + } + if err := a.EndSync(ctx); err != nil { + t.Fatalf("EndSync: %v", err) + } + + got := diffGrantIDs(t, a, base, applied) + // g-orphan has no hash-index entry — silently skipped by trie diff. + requireDiffIDs(t, got, "g2") +} + +// TestGenerateSyncDiffTrieWithoutTrees simulates a file whose tries +// are missing (e.g. written before they existed and not yet +// backfilled): the hash index is intact, so the trie path runs, and +// DirtyEntitlementBuckets degrades to the on-demand index fold per +// entitlement. The diff must still be exact. +func TestGenerateSyncDiffTrieWithoutTrees(t *testing.T) { + a := newAdapter(t) + + base := runSync(t, a, + []*v2.Entitlement{mkV2Ent("ent-A"), mkV2Ent("ent-B")}, + []*v2.Grant{ + mkV2Grant("g1", "ent-A", "user", "alice"), + mkV2Grant("g2", "ent-B", "user", "bob"), + }) + applied := runSync(t, a, + []*v2.Entitlement{mkV2Ent("ent-A"), mkV2Ent("ent-B")}, + []*v2.Grant{ + mkV2Grant("g1", "ent-A", "user", "alice"), + mkV2Grant("g2", "ent-B", "user", "bob"), + mkV2Grant("g3", "ent-B", "user", "carol"), + }) + + for _, sid := range []string{base, applied} { + idBytes, err := a.engine.resolveSyncBytes(sid) + if err != nil { + t.Fatal(err) + } + if err := a.engine.db.DeleteRange(DigestSyncLowerBound(idBytes), DigestSyncUpperBound(idBytes), pebble.Sync); err != nil { + t.Fatalf("DeleteRange(digest %s): %v", sid, err) + } + } + + got := diffGrantIDs(t, a, base, applied) + requireDiffIDs(t, got, "g3") +} diff --git a/pkg/dotc1z/engine/pebble/cleanup.go b/pkg/dotc1z/engine/pebble/cleanup.go index 0547fafd3..eeb390ca0 100644 --- a/pkg/dotc1z/engine/pebble/cleanup.go +++ b/pkg/dotc1z/engine/pebble/cleanup.go @@ -35,6 +35,8 @@ func syncScopedRanges(syncIDBytes []byte) [][2][]byte { {GrantByPrincipalSyncLowerBound(syncIDBytes), GrantByPrincipalSyncUpperBound(syncIDBytes)}, {GrantByPrincipalResourceTypeSyncLowerBound(syncIDBytes), GrantByPrincipalResourceTypeSyncUpperBound(syncIDBytes)}, {GrantByNeedsExpansionSyncLowerBound(syncIDBytes), GrantByNeedsExpansionSyncUpperBound(syncIDBytes)}, + {GrantByEntPrincHashSyncLowerBound(syncIDBytes), GrantByEntPrincHashSyncUpperBound(syncIDBytes)}, + {DigestSyncLowerBound(syncIDBytes), DigestSyncUpperBound(syncIDBytes)}, {encodeAssetPrefix(syncIDBytes), upperBoundOf(encodeAssetPrefix(syncIDBytes))}, // Stats sidecar — single key per sync; the half-open range // shape contains exactly the one key for this sync. diff --git a/pkg/dotc1z/engine/pebble/diff_grants_bench_test.go b/pkg/dotc1z/engine/pebble/diff_grants_bench_test.go new file mode 100644 index 000000000..6820bf465 --- /dev/null +++ b/pkg/dotc1z/engine/pebble/diff_grants_bench_test.go @@ -0,0 +1,224 @@ +package pebble + +import ( + "context" + "fmt" + "os" + "testing" + + "github.com/segmentio/ksuid" + + v2 "github.com/conductorone/baton-sdk/pb/c1/connector/v2" + "github.com/conductorone/baton-sdk/pkg/connectorstore" + "github.com/conductorone/baton-sdk/pkg/dotc1z/engine/pebble/codec" +) + +// Diff-strategy benchmark — trie (diffGrants) vs full-scan (diffGrantsFullScan). +// +// Fixture shape: N entitlements × M grants in base, plus a small addition set in +// applied (diffBenchDirtyEnts existing ents each gain diffBenchAddPerDirty grants, +// plus diffBenchNewEnts brand-new entitlements). The fixture is scale-independent: +// only N×M changes between scale levels; the addition count stays fixed. +// +// Scales (grant counts in base): +// +// 10K — 50 ents × 200 grants, depth-0 trees. Always runs (also with -short). +// 1M — 1 000 ents × 1 000 grants, depth-1 trees. Default. +// 50M — 5 000 ents × 10 000 grants, depth-1 trees. Set BATONSDK_BENCH_DIFF_LONG=1. +// WARNING: seeding 2×50M grants takes ~20 min. Run overnight or on CI. +// +// Run examples: +// +// # default (1M) +// go test -run=^$ -bench=BenchmarkDiffGrants -benchtime=1x ./pkg/dotc1z/engine/pebble +// +// # quick smoke (10K) +// go test -run=^$ -bench=BenchmarkDiffGrants -benchtime=1x -short ./pkg/dotc1z/engine/pebble +// +// # full suite including 50M +// BATONSDK_BENCH_DIFF_LONG=1 go test -run=^$ -bench=BenchmarkDiffGrants -benchtime=1x \ +// ./pkg/dotc1z/engine/pebble +const ( + // additions injected into applied — same count at every scale. + diffBenchDirtyEnts = 5 // existing entitlements that gain new grants in applied + diffBenchAddPerDirty = 20 // new grants per dirty entitlement + diffBenchNewEnts = 5 // brand-new entitlements present only in applied + diffBenchGrantsNewEnt = 20 // grants per new entitlement + + diffBenchTotalAdded = diffBenchDirtyEnts*diffBenchAddPerDirty + + diffBenchNewEnts*diffBenchGrantsNewEnt // = 200 +) + +type diffScale struct { + tag string + ents int + grantsPerEnt int +} + +// diffBenchScales returns the scale levels to run. +// - -short → 10K only (fast smoke). +// - default → 1M. +// - BATONSDK_BENCH_DIFF_LONG=1 → 1M + 50M. +func diffBenchScales() []diffScale { + if testing.Short() { + return []diffScale{{tag: "10K", ents: 50, grantsPerEnt: 200}} + } + scales := []diffScale{{tag: "1M", ents: 1_000, grantsPerEnt: 1_000}} + if os.Getenv("BATONSDK_BENCH_DIFF_LONG") != "" { + scales = append(scales, diffScale{tag: "50M", ents: 5_000, grantsPerEnt: 10_000}) + } + return scales +} + +// seedGrantDiffBench builds the base and applied syncs. Seeding is the +// expensive part; call this once per benchmark sub-function before the loop. +// +// Grants are PUT entitlement-by-entitlement to keep peak memory per call at +// O(grantsPerEnt) rather than O(ents×grantsPerEnt). +// newAdapterNoSync opens an adapter with DurabilityNoSync so the 200 +// per-diff grant writes don't each pay a WAL fsync. +func newAdapterNoSync(t testing.TB) *Adapter { + t.Helper() + e, _ := newTestEngine(t, WithDurability(DurabilityNoSync)) + return NewAdapter(e) +} + +func seedGrantDiffBench(b *testing.B, a *Adapter, ents, grantsPerEnt int) (baseSyncID, appliedSyncID string) { + b.Helper() + ctx := context.Background() + + baseEnts := make([]*v2.Entitlement, ents) + for i := range baseEnts { + baseEnts[i] = mkV2Ent(fmt.Sprintf("ent-%04d", i)) + } + + // Base sync ---------------------------------------------------------------- + var err error + baseSyncID, err = a.StartNewSync(ctx, connectorstore.SyncTypeFull, "") + if err != nil { + b.Fatalf("StartNewSync(base): %v", err) + } + if err := a.PutEntitlements(ctx, baseEnts...); err != nil { + b.Fatalf("PutEntitlements(base): %v", err) + } + for e := range ents { + entID := fmt.Sprintf("ent-%04d", e) + batch := make([]*v2.Grant, grantsPerEnt) + for g := range batch { + batch[g] = mkV2Grant( + fmt.Sprintf("%s:g%06d", entID, g), + entID, "user", fmt.Sprintf("u%08d", g), + ) + } + if err := a.PutGrants(ctx, batch...); err != nil { + b.Fatalf("PutGrants(base ent %d): %v", e, err) + } + } + if err := a.EndSync(ctx); err != nil { + b.Fatalf("EndSync(base): %v", err) + } + + // Applied sync ------------------------------------------------------------- + appliedSyncID, err = a.StartNewSync(ctx, connectorstore.SyncTypeFull, "") + if err != nil { + b.Fatalf("StartNewSync(applied): %v", err) + } + appliedEnts := make([]*v2.Entitlement, 0, ents+diffBenchNewEnts) + appliedEnts = append(appliedEnts, baseEnts...) + for i := range diffBenchNewEnts { + appliedEnts = append(appliedEnts, mkV2Ent(fmt.Sprintf("new-ent-%04d", i))) + } + if err := a.PutEntitlements(ctx, appliedEnts...); err != nil { + b.Fatalf("PutEntitlements(applied): %v", err) + } + + // Existing entitlements: identical to base, plus additions in the first few. + for e := range ents { + entID := fmt.Sprintf("ent-%04d", e) + cap := grantsPerEnt + if e < diffBenchDirtyEnts { + cap += diffBenchAddPerDirty + } + batch := make([]*v2.Grant, 0, cap) + for g := range grantsPerEnt { + batch = append(batch, mkV2Grant( + fmt.Sprintf("%s:g%06d", entID, g), + entID, "user", fmt.Sprintf("u%08d", g), + )) + } + if e < diffBenchDirtyEnts { + for g := range diffBenchAddPerDirty { + batch = append(batch, mkV2Grant( + fmt.Sprintf("%s:added%04d", entID, g), + entID, "user", fmt.Sprintf("added-u%d-%08d", e, g), + )) + } + } + if err := a.PutGrants(ctx, batch...); err != nil { + b.Fatalf("PutGrants(applied ent %d): %v", e, err) + } + } + + // Brand-new entitlements present only in applied. + for i := range diffBenchNewEnts { + entID := fmt.Sprintf("new-ent-%04d", i) + batch := make([]*v2.Grant, diffBenchGrantsNewEnt) + for g := range batch { + batch[g] = mkV2Grant( + fmt.Sprintf("%s:g%06d", entID, g), + entID, "user", fmt.Sprintf("new-ent-u%d-%08d", i, g), + ) + } + if err := a.PutGrants(ctx, batch...); err != nil { + b.Fatalf("PutGrants(applied new-ent %d): %v", i, err) + } + } + + if err := a.EndSync(ctx); err != nil { + b.Fatalf("EndSync(applied): %v", err) + } + return baseSyncID, appliedSyncID +} + +// runDiffGrantsBench is the shared body for the two diff benchmarks. +func runDiffGrantsBench(b *testing.B, diffFn func(context.Context, *Adapter, []byte, []byte, string, string, string) error) { + b.Helper() + for _, sc := range diffBenchScales() { + b.Run(sc.tag, func(b *testing.B) { + a := newAdapterNoSync(b) + ctx := context.Background() + baseSyncID, appliedSyncID := seedGrantDiffBench(b, a, sc.ents, sc.grantsPerEnt) + + baseBytes, err := codec.EncodeSyncID(baseSyncID) + if err != nil { + b.Fatalf("encode base: %v", err) + } + appliedBytes, err := codec.EncodeSyncID(appliedSyncID) + if err != nil { + b.Fatalf("encode applied: %v", err) + } + + b.ReportMetric(float64(sc.ents*sc.grantsPerEnt), "base_grants") + b.ReportMetric(float64(diffBenchTotalAdded), "expected_additions") + + for b.Loop() { + diffID := ksuid.New().String() + if err := diffFn(ctx, a, baseBytes, appliedBytes, baseSyncID, appliedSyncID, diffID); err != nil { + b.Fatalf("diff: %v", err) + } + } + }) + } +} + +// fullScanAdapter wraps diffGrantsFullScan to match the diffGrants signature. +func fullScanAdapter(ctx context.Context, a *Adapter, baseBytes, appliedBytes []byte, _, _, diffID string) error { + return diffGrantsFullScan(ctx, a, baseBytes, appliedBytes, diffID) +} + +// BenchmarkDiffGrants_Trie / BenchmarkDiffGrants_FullScan use DurabilityNoSync +// so that the 200 per-diff grant writes don't each pay a WAL fsync. The writes +// still happen (churn is present), but sync latency is excluded from the +// measurement. This isolates the read/comparison cost of each strategy. +func BenchmarkDiffGrants_Trie(b *testing.B) { runDiffGrantsBench(b, diffGrants) } +func BenchmarkDiffGrants_FullScan(b *testing.B) { runDiffGrantsBench(b, fullScanAdapter) } diff --git a/pkg/dotc1z/engine/pebble/digest.go b/pkg/dotc1z/engine/pebble/digest.go new file mode 100644 index 000000000..5d216a900 --- /dev/null +++ b/pkg/dotc1z/engine/pebble/digest.go @@ -0,0 +1,915 @@ +package pebble + +import ( + "bytes" + "context" + "encoding/binary" + "errors" + "fmt" + + "github.com/cockroachdb/pebble/v2" + + "github.com/conductorone/baton-sdk/pkg/dotc1z/engine/pebble/codec" +) + +// Bucketed XOR set digests over bucket-hash indexes. +// +// Goal: answer "does this partition hold exactly the same records as +// some other sync/file?" with a single key read, and when the answer is +// no, identify which hash buckets differ so a caller can load only +// those records instead of re-reading the whole partition. +// +// The digest is generic over any secondary index matching the shape +// described on digestIndexSpec — a partition prefix followed by a raw +// fixed-width bucket hash, with a per-record content hash as the +// value. Two properties make such an index the right substrate: +// +// - hash-major order is identical across two files that hold the same +// records, so a streamed fold produces the same digest; and +// - a bucket is a contiguous bit-range of the hash, which is a +// contiguous range of the index key, so "all records in bucket P" +// is a single range scan. +// +// The grant instantiation (partition = entitlement, bucket hash = +// hash(principal identity)) lives in grant_digest.go. +// +// Combiner. A node's digest is the XOR of every content hash beneath +// it (the content hashes themselves are whatever the index writer +// chose — only the combiner is XOR). XOR is homomorphic (parent = XOR +// of children), order-independent, and invertible, which buys three +// things: +// +// - split-independence: a bucket's digest depends only on the records +// in its hash range, never on how the range is subdivided, so +// buckets from digests of different widths compare directly (a +// width-w bucket is the XOR of its two width-(w+1) halves); +// - O(1) incremental maintenance: post-seal insert/overwrite/delete +// XOR the record's content hash into/out of the root and the one +// leaf on its bucket path (see digestMutator); +// - the empty digest is all-zero (the XOR identity), so an absent +// leaf reads as {count: 0, digest: 0}. +// +// Every node also stores its record COUNT. Comparison always checks the +// (count, digest) pair, so a non-empty node whose hashes happened to XOR +// to zero can never be conflated with an empty/absent one. XOR +// set-hashing is not adversarially collision-resistant +// (Bellare–Micciancio); a digest is an optimization, not a trust +// boundary — see RFC 0003 §9. Index writers should fold every +// key-distinguishing field into the content hash so duplicate leaves +// (and thus in-partition self-cancellation) are impossible by +// construction. +// +// Shape. A flat hash table, not a multi-level radix tree: one root plus +// a single leaf level of 2^width buckets, where width — in BITS of the +// bucket hash, 0..digestMaxWidthBits — is chosen per partition so the +// average bucket holds at most digestTargetBucketSize records +// (chooseDigestWidth). Width 0 means root-only (small and empty +// partitions cost exactly one stored node). Growing capacity one bit at +// a time keeps realized bucket occupancy within a 2x band of the +// target; a byte-per-level radix could only pick capacities of 256^k, +// so a partition just past a boundary would store up to ~256x more +// nodes than needed. +// +// Interior levels are deliberately absent: hierarchical pruning only +// pays when the leaf level is too large to scan, and at <= 2^16 leaves +// a contiguous range scan beats a level-by-level descent. Comparison is +// a single merge scan of both sides' leaf levels (see +// dirtyPartitionBuckets); cross-width comparison folds the finer side's +// leaves down to the coarser width on the fly, which split-independence +// makes exact. For the same reason digestTargetBucketSize is a soft +// tuning knob, not an ABI constant: digests built with different +// targets still compare correctly. +// +// Leaves are stored sparsely — a leaf is materialized iff its bucket +// holds >=1 record. The root is always materialized (it is the "digest +// was built" marker — absence of the root means "never built", never +// "empty"). +// +// On-disk ABI. The bucket-hash width, each index's content-hash +// definition, the combiner, the leaf-prefix encoding, and the node +// value framing are all part of the stored format: changing +// digestBucketHashLen, an index's content-hash field set, the combiner, +// or the framing requires an index-migration version bump (see +// index_migrations.go). digestTargetBucketSize is NOT part of the ABI +// (see above). + +const ( + // digestBucketHashLen is the width, in bytes, of the raw bucket + // hash embedded in a digested index's key. Collisions in the bucket + // address are harmless (they only co-locate records — the index + // key's tail still distinguishes rows). + digestBucketHashLen = 8 + + // digestTargetBucketSize is the record count a single leaf bucket + // aims to hold. The width is grown until 2^width buckets bring the + // average bucket under this. Tunable without a migration: the width + // is read from the stored root and cross-width comparison folds to + // the coarser side. + digestTargetBucketSize = 512 + + // digestMaxWidthBits caps the leaf-level width. 2^16 buckets keeps + // the comparison's full leaf scan trivially cheap; past the cap + // (digestTargetBucketSize << 16 records) buckets simply grow beyond + // the target. + digestMaxWidthBits = 16 + + // digestLeafPrefixLen is the stored byte width of a leaf key's + // bucket prefix: the bucket index LEFT-ALIGNED in 16 bits, so leaf + // keys sort in bucket-hash order at every width and folding to a + // coarser width is "take the top bits". ABI. + digestLeafPrefixLen = 2 +) + +// Node-key levels: the root is level 0 (empty prefix); the single leaf +// level is 1 (digestLeafPrefixLen-byte prefix). See encodeDigestNodeKey. +const ( + digestLevelRoot byte = 0 + digestLevelLeaf byte = 1 +) + +// hashLen is the width of a content hash (index value) and of a node +// digest. +const hashLen = 8 + +// zeroDigest is the XOR identity — the digest of an empty/absent node. +var zeroDigest [hashLen]byte + +// xorInto XORs src into dst in place, over min(len(dst), len(src)). +func xorInto(dst, src []byte) { + for i := range min(len(dst), len(src)) { + dst[i] ^= src[i] + } +} + +// digestIndexSpec describes one bucket-hash index the digest core can +// fold. The index must have the shape +// +// index key = partitionPrefix(sync, partition) | | +// index value = content hash (hashLen bytes) +// +// and, for partition enumeration, +// +// index key = syncBounds(sync).lower | 0x00 | tuple(partition) | 0x00 | | … +// +// (i.e. the partition is the first tuple component after the sync-wide +// lower bound — the standard by-value index layout). +// +// The content hash defines record identity for the diff and must fold +// every key-distinguishing field (see the package comment); the bucket +// hash must be derived from fields that are stable across syncs. +type digestIndexSpec struct { + // indexID discriminates this index's nodes inside the typeDigest + // keyspace. Conventionally the digested index's own idx* byte. ABI. + indexID byte + + // partitionPrefix returns the index-key prefix covering one + // partition's entries, ending immediately before the raw bucket + // hash (trailing separator included). + partitionPrefix func(syncIDBytes []byte, partition string) []byte + + // syncBounds bounds the index's entire keyspace for one sync. + syncBounds func(syncIDBytes []byte) (lower, upper []byte) +} + +// DigestBucket addresses one hash bucket of a partition: the records +// whose bucket hash starts with the top Bits bits of Index. The zero +// value (Bits 0) addresses the whole partition. +type DigestBucket struct { + Index uint32 + Bits int +} + +// leafKeyPrefix returns the stored digestLeafPrefixLen-byte node-key +// prefix for a leaf bucket: the index left-aligned in 16 bits. +// Requires 1 <= Bits <= digestMaxWidthBits. +func (b DigestBucket) leafKeyPrefix() []byte { + out := make([]byte, digestLeafPrefixLen) + binary.BigEndian.PutUint16(out, uint16(b.Index)<<(16-b.Bits)) //nolint:gosec // Index < 2^Bits <= 2^16 by construction + return out +} + +// bucketOfHash returns the width-`bits` bucket holding bucket hash bh. +// Requires 1 <= bits <= digestMaxWidthBits. +func bucketOfHash(bh []byte, bits int) DigestBucket { + u := binary.BigEndian.Uint16(bh[:digestLeafPrefixLen]) + return DigestBucket{Index: uint32(u >> (16 - bits)), Bits: bits} +} + +// bucketBounds returns the index key range [lower, upper) covering a +// bucket's records in one partition. The raw bucket hash is a clean +// byte region of the index key, so a bit-granular bucket maps to plain +// uint64 arithmetic on that region. +func (s digestIndexSpec) bucketBounds(idBytes []byte, partition string, b DigestBucket) ([]byte, []byte) { + prefix := s.partitionPrefix(idBytes, partition) + if b.Bits == 0 { + return prefix, upperBoundOf(prefix) + } + boundAt := func(hash uint64) []byte { + out := append(append(make([]byte, 0, len(prefix)+digestBucketHashLen), prefix...), 0, 0, 0, 0, 0, 0, 0, 0) + binary.BigEndian.PutUint64(out[len(prefix):], hash) + return out + } + shift := uint(64 - b.Bits) //nolint:gosec // Bits in [1, digestMaxWidthBits] + lower := boundAt(uint64(b.Index) << shift) + if uint64(b.Index)+1 == uint64(1)< digestMaxWidthBits { + return 0, 0, nil, false + } + widthBits := int(val[0]) + count := int64(binary.BigEndian.Uint64(val[1:9])) //nolint:gosec // count is a non-negative row count + return widthBits, count, val[9:], true +} + +func packDigestLeaf(count int64, digest []byte) []byte { + buf := make([]byte, 0, 8+len(digest)) + var n [8]byte + binary.BigEndian.PutUint64(n[:], uint64(count)) //nolint:gosec // non-negative row count + buf = append(buf, n[:]...) + return append(buf, digest...) +} + +// unpackDigestLeaf returns (count, digest, ok) for a leaf node body. +func unpackDigestLeaf(val []byte) (int64, []byte, bool) { + if len(val) != 8+hashLen { + return 0, nil, false + } + count := int64(binary.BigEndian.Uint64(val[:8])) //nolint:gosec // non-negative count + return count, val[8:], true +} + +// buildPartitionDigest counts a partition's index entries (pass 1), +// picks the width from that count, and delegates the fold to +// buildPartitionDigestAtWidth (pass 2). +func (e *Engine) buildPartitionDigest(ctx context.Context, spec digestIndexSpec, idBytes []byte, partition string) error { + prefix := spec.partitionPrefix(idBytes, partition) + count, err := e.countKeysInRange(prefix, upperBoundOf(prefix)) + if err != nil { + return err + } + return e.buildPartitionDigestAtWidth(ctx, spec, idBytes, partition, chooseDigestWidth(count)) +} + +// buildPartitionDigestAtWidth folds the index for one partition into a +// root plus every non-empty leaf, in a single streaming pass — O(1) +// memory regardless of partition size. +// +// The pass starts by range-deleting the partition's whole node +// keyspace: the build only ever Sets nodes, so without the clear a +// rebuild that changes width or empties a bucket would leave stale +// leaves that the comparison merge scan (which enumerates leaves from +// the node keyspace) would read — and a stale digest that happens to +// match the peer prunes a real diff. Old and new framings are +// byte-length identical, so stale nodes are not detectable by +// inspection. +// +// Sorted index order means each bucket's entries are contiguous, so the +// single "open" leaf closes exactly when its prefix changes. Only +// non-empty leaves are ever opened, so sparsity is automatic, not a +// prune pass. +// +// The width is taken as a parameter rather than derived so the +// width-selection seam can be exercised directly: tests force a width +// that the natural count→width mapping would only produce at a very +// large record count, which is how the cross-width comparison path gets +// covered without seeding hundreds of thousands of records. +func (e *Engine) buildPartitionDigestAtWidth(ctx context.Context, spec digestIndexSpec, idBytes []byte, partition string, widthBits int) error { + prefix := spec.partitionPrefix(idBytes, partition) + upper := upperBoundOf(prefix) + nodeLower := encodeDigestPartitionPrefix(idBytes, spec.indexID, partition) + nodeUpper := upperBoundOf(nodeLower) + + return e.withWrite(func() error { + batch := e.db.NewBatch() + defer batch.Close() + + // Clear any prior build (see function comment). In-batch + // ordering makes this safe: the Sets below land after the + // tombstone and survive it. + if err := batch.DeleteRange(nodeLower, nodeUpper, nil); err != nil { + return err + } + + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: prefix, UpperBound: upper}) + if err != nil { + return err + } + defer iter.Close() + + // lowMask clears the sub-bucket bits of a left-aligned 16-bit + // prefix, leaving the leaf's stored key prefix value. + var lowMask uint16 + if widthBits > 0 { + lowMask = ^uint16(0) >> widthBits + } + + var ( + leafOpen bool + leafLV uint16 // left-aligned stored prefix of the open leaf + leafDigest [hashLen]byte + leafCount int64 + rootDigest [hashLen]byte + total int64 + ) + flushLeaf := func() error { + if !leafOpen { + return nil + } + var lp [digestLeafPrefixLen]byte + binary.BigEndian.PutUint16(lp[:], leafLV) + key := encodeDigestNodeKey(idBytes, spec.indexID, partition, digestLevelLeaf, lp[:]) + return batch.Set(key, packDigestLeaf(leafCount, leafDigest[:]), nil) + } + + for iter.First(); iter.Valid(); iter.Next() { + if err := ctx.Err(); err != nil { + return err + } + key := iter.Key() + if len(key) < len(prefix)+digestBucketHashLen { + continue // malformed; skip defensively + } + val := iter.Value() // per-record content hash + if len(val) != hashLen { + // Index writers always emit exactly hashLen bytes + // (grantContentHash); a wrong length is a corrupt or + // mis-encoded entry. xorInto would fold only a prefix and + // quietly corrupt the digest, so reject it. At seal the + // caller downgrades a build error to "no digest" (readers + // fall back to the on-demand fold), so this fails safe. + return fmt.Errorf("buildPartitionDigestAtWidth: content hash for %q is %d bytes, want %d", partition, len(val), hashLen) + } + if widthBits > 0 { + lv := binary.BigEndian.Uint16(key[len(prefix):]) &^ lowMask + if !leafOpen || lv != leafLV { + if err := flushLeaf(); err != nil { + return err + } + leafLV = lv + leafDigest = [hashLen]byte{} + leafCount = 0 + leafOpen = true + } + xorInto(leafDigest[:], val) + leafCount++ + } + xorInto(rootDigest[:], val) + total++ + } + if err := iter.Error(); err != nil { + return err + } + if err := flushLeaf(); err != nil { + return err + } + + // Root is written unconditionally — even at count 0 — as the + // "digest was built" marker. + rootKey := encodeDigestNodeKey(idBytes, spec.indexID, partition, digestLevelRoot, nil) + if err := batch.Set(rootKey, packDigestRoot(widthBits, total, rootDigest[:]), nil); err != nil { + return err + } + + opts := writeOpts(e.opts.durability) + if e.IsFreshSync() { + opts = pebble.NoSync + } + return batch.Commit(opts) + }) +} + +// countKeysInRange counts keys in [lower, upper) without materializing +// values. Used by the build's count→width pass and by the diff driver's +// primary-vs-index coverage guard. +func (e *Engine) countKeysInRange(lower, upper []byte) (int64, error) { + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: lower, UpperBound: upper}) + if err != nil { + return 0, err + } + defer iter.Close() + var n int64 + for iter.First(); iter.Valid(); iter.Next() { + n++ + } + return n, iter.Error() +} + +// distinctDigestPartitions returns the distinct partitions present in a +// sync's digested index, in index order. It seeks past each partition's +// whole range once its id is captured, so the cost is O(partitions) +// seeks, not O(entries). This is the diff driver's work list, and it +// deliberately comes from the INDEX rather than the partition-owning +// records: an entry whose partition has no owning record still has +// index entries (and a fold), so it still gets compared — only a digest +// was never built for it. +func (e *Engine) distinctDigestPartitions(ctx context.Context, spec digestIndexSpec, syncIDBytes []byte) ([]string, error) { + lower, upper := spec.syncBounds(syncIDBytes) + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: lower, UpperBound: upper}) + if err != nil { + return nil, err + } + defer iter.Close() + + // Keys are lower | 0x00 | tuple(partition) | 0x00 | bucketHash | … + // (the spec's shape contract). + partStart := len(lower) + 1 + var out []string + for iter.First(); iter.Valid(); { + if err := ctx.Err(); err != nil { + return nil, err + } + key := iter.Key() + if len(key) <= partStart { + iter.Next() + continue // malformed; skip defensively + } + partBytes, _, decErr := codec.DecodeTupleStringTo(nil, key[partStart:], 0) + if decErr != nil { + iter.Next() + continue + } + partition := string(partBytes) + out = append(out, partition) + // Seek past this partition's whole index range. + seekTo := upperBoundOf(spec.partitionPrefix(syncIDBytes, partition)) + if seekTo == nil { + break + } + iter.SeekGE(seekTo) + } + return out, iter.Error() +} + +// DigestRoot is a partition's stored root digest. +type DigestRoot struct { + Hash []byte + Bits int // leaf-level width in bits; 0 = root-only digest + Count int64 +} + +// getPartitionDigestRoot returns the stored root for a partition. ok is +// false when no digest has been built for it (the caller can fall back +// to computeBucketDigest, which derives the same digest from the index +// on demand). +func (e *Engine) getPartitionDigestRoot(spec digestIndexSpec, idBytes []byte, partition string) (DigestRoot, bool, error) { + val, closer, err := e.db.Get(encodeDigestNodeKey(idBytes, spec.indexID, partition, digestLevelRoot, nil)) + if err != nil { + if errors.Is(err, pebble.ErrNotFound) { + return DigestRoot{}, false, nil + } + return DigestRoot{}, false, err + } + defer closer.Close() + widthBits, count, h, valid := unpackDigestRoot(val) + if !valid { + return DigestRoot{}, false, fmt.Errorf("getPartitionDigestRoot: malformed root for %q", partition) + } + out := make([]byte, len(h)) + copy(out, h) + return DigestRoot{Hash: out, Bits: widthBits, Count: count}, true, nil +} + +// getDigestLeaf reads one stored leaf by its key prefix. An absent leaf +// returns (0, zero digest, present=false, nil) — the XOR identity. +func (e *Engine) getDigestLeaf(spec digestIndexSpec, idBytes []byte, partition string, leafPrefix []byte) (int64, []byte, bool, error) { + val, closer, err := e.db.Get(encodeDigestNodeKey(idBytes, spec.indexID, partition, digestLevelLeaf, leafPrefix)) + if err != nil { + if errors.Is(err, pebble.ErrNotFound) { + return 0, zeroDigest[:], false, nil + } + return 0, nil, false, err + } + defer closer.Close() + count, digest, ok := unpackDigestLeaf(val) + if !ok { + return 0, nil, false, fmt.Errorf("getDigestLeaf: malformed leaf for %q", partition) + } + out := make([]byte, hashLen) + copy(out, digest) + return count, out, true, nil +} + +// foldedBucket is one entry of a folded leaf scan: the (XOR, count) +// aggregate of the consecutive stored leaves sharing the top `bits` +// bits of their bucket index. +type foldedBucket struct { + idx uint32 + count int64 + digest [hashLen]byte +} + +// foldedLeafBuckets scans a partition's stored leaf level and folds it +// to width foldBits (which must be <= the width the digest was built +// at), returning the non-empty buckets in index order. One contiguous +// range scan; folding is exact because leaf prefixes are left-aligned +// (so keys sort in bucket-hash order at any width) and XOR digests are +// split-independent. +func (e *Engine) foldedLeafBuckets(ctx context.Context, spec digestIndexSpec, idBytes []byte, partition string, foldBits int) ([]foldedBucket, error) { + stem := encodeDigestNodeKey(idBytes, spec.indexID, partition, digestLevelLeaf, nil) + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: stem, UpperBound: upperBoundOf(stem)}) + if err != nil { + return nil, err + } + defer iter.Close() + var out []foldedBucket + for iter.First(); iter.Valid(); iter.Next() { + if err := ctx.Err(); err != nil { + return nil, err + } + key := iter.Key() + if len(key) != len(stem)+digestLeafPrefixLen { + continue // malformed; skip defensively + } + lv := binary.BigEndian.Uint16(key[len(stem):]) + idx := uint32(lv >> (16 - foldBits)) + count, digest, ok := unpackDigestLeaf(iter.Value()) + if !ok { + return nil, fmt.Errorf("foldedLeafBuckets: malformed leaf for %q", partition) + } + if n := len(out); n > 0 && out[n-1].idx == idx { + out[n-1].count += count + xorInto(out[n-1].digest[:], digest) + continue + } + fb := foldedBucket{idx: idx, count: count} + copy(fb.digest[:], digest) + out = append(out, fb) + } + return out, iter.Error() +} + +// computeBucketDigest folds the index over a single bucket (the zero +// bucket = whole partition = the root) and returns the content-defined +// XOR digest plus the record count. This is the authoritative +// definition of a node; stored nodes are a cache of it. +// Split-independent: the digest depends only on the records in the +// bucket's hash range, not on any digest's width. +func (e *Engine) computeBucketDigest(ctx context.Context, spec digestIndexSpec, idBytes []byte, partition string, bucket DigestBucket) ([]byte, int64, error) { + lower, upper := spec.bucketBounds(idBytes, partition, bucket) + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: lower, UpperBound: upper}) + if err != nil { + return nil, 0, err + } + defer iter.Close() + digest := make([]byte, hashLen) + var count int64 + for iter.First(); iter.Valid(); iter.Next() { + if err := ctx.Err(); err != nil { + return nil, 0, err + } + val := iter.Value() + if len(val) != hashLen { + // See buildPartitionDigestAtWidth: a mis-length index value is + // corruption; reject it rather than silently fold a prefix. + return nil, 0, fmt.Errorf("computeBucketDigest: content hash for %q is %d bytes, want %d", partition, len(val), hashLen) + } + xorInto(digest, val) + count++ + } + if err := iter.Error(); err != nil { + return nil, 0, err + } + return digest, count, nil +} + +// dirtyPartitionBuckets compares this engine's partition against +// other's and returns the buckets whose records differ. A single zero +// bucket (Bits 0) means "the whole partition differs" (used when the +// comparison granularity is the root, e.g. small partitions). A nil +// (empty) result means the two are identical. +// +// The fast path is a single root read per side. On mismatch both sides' +// leaf levels are folded to compareBits — the narrower digest's width, +// where both sides have directly comparable buckets (XOR digests are +// split-independent) — and merge-compared in one pass. Each side's fold +// is a single contiguous range scan of its stored leaves, so cost is +// O(leaves), bounded by count/digestTargetBucketSize per side. +// +// A missing root means the digest was never built on that side — NOT +// that the partition is empty — so both sides are compared via the +// authoritative on-demand fold instead. +func (e *Engine) dirtyPartitionBuckets(ctx context.Context, spec digestIndexSpec, idBytes []byte, other *Engine, otherIDBytes []byte, partition string) ([]DigestBucket, error) { + rootA, okA, err := e.getPartitionDigestRoot(spec, idBytes, partition) + if err != nil { + return nil, err + } + rootB, okB, err := other.getPartitionDigestRoot(spec, otherIDBytes, partition) + if err != nil { + return nil, err + } + + if !okA || !okB { + ha, ca, err := e.computeBucketDigest(ctx, spec, idBytes, partition, DigestBucket{}) + if err != nil { + return nil, err + } + hb, cb, err := other.computeBucketDigest(ctx, spec, otherIDBytes, partition, DigestBucket{}) + if err != nil { + return nil, err + } + if ca == cb && bytes.Equal(ha, hb) { + return nil, nil + } + return []DigestBucket{{}}, nil + } + + if rootA.Count == rootB.Count && bytes.Equal(rootA.Hash, rootB.Hash) { + return nil, nil + } + + // Roots differ. The comparison granularity is the narrower digest's + // width; at width 0 there is nothing below the root. + compareBits := min(rootA.Bits, rootB.Bits) + if compareBits == 0 { + return []DigestBucket{{}}, nil + } + + fa, err := e.foldedLeafBuckets(ctx, spec, idBytes, partition, compareBits) + if err != nil { + return nil, err + } + fb, err := other.foldedLeafBuckets(ctx, spec, otherIDBytes, partition, compareBits) + if err != nil { + return nil, err + } + + // Merge the two sorted folded-bucket streams. A bucket present on + // only one side is dirty by construction (stored leaves are never + // empty); a shared bucket is dirty iff its (count, digest) differs. + var dirty []DigestBucket + i, j := 0, 0 + for i < len(fa) || j < len(fb) { + switch { + case j == len(fb) || (i < len(fa) && fa[i].idx < fb[j].idx): + dirty = append(dirty, DigestBucket{Index: fa[i].idx, Bits: compareBits}) + i++ + case i == len(fa) || fb[j].idx < fa[i].idx: + dirty = append(dirty, DigestBucket{Index: fb[j].idx, Bits: compareBits}) + j++ + default: + if fa[i].count != fb[j].count || fa[i].digest != fb[j].digest { + dirty = append(dirty, DigestBucket{Index: fa[i].idx, Bits: compareBits}) + } + i++ + j++ + } + } + + // Roots differed but the merge found nothing: with consistent + // digests that's impossible (the root is the XOR of the leaves), so + // a stored node is stale or corrupt. Fail safe — whole partition + // dirty; the next rebuild heals the digest. + if len(dirty) == 0 { + return []DigestBucket{{}}, nil + } + return dirty, nil +} + +// --- Incremental maintenance (post-seal) --- + +// digestMutator accumulates per-node (XOR, count) deltas for a batch of +// post-seal record mutations against one digested index, and applies +// each touched node exactly once. +// +// Why an accumulator instead of read-modify-write per mutation: the +// updates target a plain pebble.Batch, which does NOT read through its +// own writes — and every mutation in a batch touches the root node, so +// naive per-mutation RMW would lose deltas. Accumulating also collapses +// N writes per node into one. All record writers run under withWrite's +// mutex, so reading current node values from the DB inside apply is +// race-free. +// +// Lifecycle: one mutator per write batch. Callers feed removeHash(old) +// / addHash(new) as they process records (an overwrite that moves a +// record to a different bucket is exactly remove+add), then call +// apply(batch) once before commit. +// +// Partitions whose digest was never built (no stored root) are skipped: +// the seal-time build or the on-Open backfill will construct them from +// the index. This also makes the mutator free during a fresh sync — but +// callers on the fresh-sync bulk path should skip constructing one +// anyway to avoid the per-partition root probe. +type digestMutator struct { + e *Engine + spec digestIndexSpec + parts map[string]*mutatorPartition // keyed by string(root node key) +} + +type mutatorPartition struct { + idBytes []byte + partition string + rootKey []byte + present bool // stored root exists; if false all deltas are dropped + bits int // leaf-level width read from the stored root + rootCount int64 // count read from the stored root + rootDigest [hashLen]byte // digest read from the stored root + xor [hashLen]byte // accumulated root delta + countDelta int64 + nodes map[string]*mutatorNode // leaves, keyed by string(node key) +} + +type mutatorNode struct { + key []byte + xor [hashLen]byte + countDelta int64 +} + +func newDigestMutator(e *Engine, spec digestIndexSpec) *digestMutator { + return &digestMutator{e: e, spec: spec, parts: make(map[string]*mutatorPartition)} +} + +// partFor returns the (cached) per-partition state, probing the stored +// root on first touch. The cached root snapshot stays valid for the +// mutator's lifetime because all writers serialize through withWrite. +func (m *digestMutator) partFor(idBytes []byte, partition string) (*mutatorPartition, error) { + rootKey := encodeDigestNodeKey(idBytes, m.spec.indexID, partition, digestLevelRoot, nil) + k := string(rootKey) + if mp, ok := m.parts[k]; ok { + return mp, nil + } + mp := &mutatorPartition{idBytes: idBytes, partition: partition, rootKey: rootKey, nodes: make(map[string]*mutatorNode)} + val, closer, err := m.e.db.Get(rootKey) + switch { + case err == nil: + widthBits, count, digest, ok := unpackDigestRoot(val) + closer.Close() + if ok { + mp.present = true + mp.bits = widthBits + mp.rootCount = count + copy(mp.rootDigest[:], digest) + } + // Malformed root: leave present=false so mutations are dropped; + // the stale root heals at the next rebuild. + case errors.Is(err, pebble.ErrNotFound): + // No digest — deltas for this partition are no-ops. + default: + return nil, err + } + m.parts[k] = mp + return mp, nil +} + +// addHash records the insertion of a record with the given bucket and +// content hashes into its partition's digest. +func (m *digestMutator) addHash(idBytes []byte, partition string, bucketHash, contentHash []byte) error { + return m.delta(idBytes, partition, bucketHash, contentHash, 1) +} + +// removeHash records the removal of a record from its partition's +// digest. +func (m *digestMutator) removeHash(idBytes []byte, partition string, bucketHash, contentHash []byte) error { + return m.delta(idBytes, partition, bucketHash, contentHash, -1) +} + +func (m *digestMutator) delta(idBytes []byte, partition string, bucketHash, contentHash []byte, sign int64) error { + mp, err := m.partFor(idBytes, partition) + if err != nil { + return err + } + if !mp.present { + return nil + } + xorInto(mp.xor[:], contentHash) + mp.countDelta += sign + if mp.bits == 0 { + return nil + } + key := encodeDigestNodeKey(idBytes, m.spec.indexID, partition, digestLevelLeaf, bucketOfHash(bucketHash, mp.bits).leafKeyPrefix()) + nk := string(key) + n, ok := mp.nodes[nk] + if !ok { + n = &mutatorNode{key: key} + mp.nodes[nk] = n + } + xorInto(n.xor[:], contentHash) + n.countDelta += sign + return nil +} + +// apply folds the accumulated deltas into the stored nodes via batch. +// Nodes whose delta cancelled to zero (e.g. an overwrite that changed +// only excluded fields) are skipped; a leaf whose count reaches zero is +// deleted (restoring sparsity); the root is rewritten in place. A count +// that would go negative means the stored digest disagrees with the +// mutation stream — the digest is dropped wholesale (DeleteRange), so +// readers fall back to the on-demand fold until the next rebuild. +func (m *digestMutator) apply(batch *pebble.Batch) error { + for _, mp := range m.parts { + if !mp.present { + continue + } + if err := m.applyPartition(batch, mp); err != nil { + return err + } + } + return nil +} + +func (m *digestMutator) applyPartition(batch *pebble.Batch, mp *mutatorPartition) error { + if mp.countDelta == 0 && mp.xor == zeroDigest && len(mp.nodes) == 0 { + return nil + } + dropDigest := func() error { + // In-batch ordering: this tombstone lands after any node Sets + // already staged for this partition and removes them too. + lo := encodeDigestPartitionPrefix(mp.idBytes, m.spec.indexID, mp.partition) + return batch.DeleteRange(lo, upperBoundOf(lo), nil) + } + if mp.rootCount+mp.countDelta < 0 { + return dropDigest() + } + for _, n := range mp.nodes { + if n.countDelta == 0 && n.xor == zeroDigest { + continue + } + var ( + curCount int64 + curDigest [hashLen]byte + ) + val, closer, err := m.e.db.Get(n.key) + switch { + case err == nil: + c, d, ok := unpackDigestLeaf(val) + closer.Close() + if !ok { + return dropDigest() + } + curCount = c + copy(curDigest[:], d) + case errors.Is(err, pebble.ErrNotFound): + // absent leaf = {0, zero} + default: + return err + } + newCount := curCount + n.countDelta + if newCount < 0 { + return dropDigest() + } + xorInto(curDigest[:], n.xor[:]) + if newCount == 0 { + // An emptied leaf's digest must cancel to exactly zero + // (count 0 ⇒ digest 0); anything else means the stored + // digest disagrees with the mutation stream. + if curDigest != zeroDigest { + return dropDigest() + } + if err := batch.Delete(n.key, nil); err != nil { + return err + } + continue + } + if err := batch.Set(n.key, packDigestLeaf(newCount, curDigest[:]), nil); err != nil { + return err + } + } + if mp.countDelta == 0 && mp.xor == zeroDigest { + return nil + } + newDigest := mp.rootDigest + xorInto(newDigest[:], mp.xor[:]) + return batch.Set(mp.rootKey, packDigestRoot(mp.bits, mp.rootCount+mp.countDelta, newDigest[:]), nil) +} diff --git a/pkg/dotc1z/engine/pebble/digest_test.go b/pkg/dotc1z/engine/pebble/digest_test.go new file mode 100644 index 000000000..819421ca2 --- /dev/null +++ b/pkg/dotc1z/engine/pebble/digest_test.go @@ -0,0 +1,927 @@ +package pebble + +import ( + "bytes" + "context" + "encoding/binary" + "fmt" + "testing" + + "github.com/cockroachdb/pebble/v2" + "github.com/segmentio/ksuid" + + v3 "github.com/conductorone/baton-sdk/pb/c1/storage/v3" +) + +// putEnt writes an entitlement record (under the engine's current +// sync) whose external_id is entID — the same string grants reference +// via EntitlementRef.EntitlementId, which is what BuildAllGrantDigests +// keys each digest on. +func putEnt(t testing.TB, e *Engine, ctx context.Context, entID string) { + t.Helper() + rec := v3.EntitlementRecord_builder{ + ExternalId: entID, + Resource: v3.ResourceRef_builder{ + ResourceTypeId: "app", + ResourceId: "github", + }.Build(), + }.Build() + if err := e.PutEntitlementRecord(ctx, rec); err != nil { + t.Fatalf("PutEntitlementRecord: %v", err) + } +} + +// makeGrantWithSources is makeGrant plus an optional source-entitlement +// set, which grantContentHash folds in — so two grants with the same +// (entitlement, principal, external_id) but different sources produce +// different content hashes while keeping the SAME index key. +func makeGrantWithSources(syncID, externalID, entID, principalID string, sources ...string) *v3.GrantRecord { + g := makeGrant(syncID, externalID, entID, principalID) + if len(sources) > 0 { + m := make(map[string]*v3.GrantSourceRecord, len(sources)) + for _, s := range sources { + m[s] = v3.GrantSourceRecord_builder{}.Build() + } + g.SetSources(m) + } + return g +} + +// digestNodeCount counts stored digest nodes for a sync (across all +// partitions and digested indexes). +func digestNodeCount(t testing.TB, e *Engine, syncID string) int { + t.Helper() + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatalf("resolveSyncBytes: %v", err) + } + iter, err := e.db.NewIter(&pebble.IterOptions{ + LowerBound: DigestSyncLowerBound(idBytes), + UpperBound: DigestSyncUpperBound(idBytes), + }) + if err != nil { + t.Fatalf("NewIter: %v", err) + } + defer iter.Close() + n := 0 + for iter.First(); iter.Valid(); iter.Next() { + n++ + } + if err := iter.Error(); err != nil { + t.Fatalf("iter: %v", err) + } + return n +} + +// rawLeafPrefixes returns the stored 2-byte leaf key prefixes for one +// entitlement's grant digest, in key order. +func rawLeafPrefixes(t testing.TB, e *Engine, idBytes []byte, entID string) [][]byte { + t.Helper() + stem := encodeDigestNodeKey(idBytes, grantDigestSpec.indexID, entID, digestLevelLeaf, nil) + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: stem, UpperBound: upperBoundOf(stem)}) + if err != nil { + t.Fatalf("NewIter: %v", err) + } + defer iter.Close() + var out [][]byte + for iter.First(); iter.Valid(); iter.Next() { + key := iter.Key() + if len(key) != len(stem)+digestLeafPrefixLen { + t.Fatalf("leaf key with prefix length %d, want %d", len(key)-len(stem), digestLeafPrefixLen) + } + out = append(out, append([]byte(nil), key[len(stem):]...)) + } + if err := iter.Error(); err != nil { + t.Fatalf("iter: %v", err) + } + return out +} + +// seedEntitlement writes the entitlement record + grants and builds the +// digest, returning the syncID. +func seedEntitlement(t testing.TB, e *Engine, entID string, grants []*v3.GrantRecord) string { + t.Helper() + ctx := context.Background() + syncID := ksuid.New().String() + if err := e.SetCurrentSync(syncID); err != nil { + t.Fatalf("SetCurrentSync: %v", err) + } + putEnt(t, e, ctx, entID) + if err := e.PutGrantRecords(ctx, grants...); err != nil { + t.Fatalf("PutGrantRecords: %v", err) + } + if err := e.BuildAllGrantDigests(ctx, syncID); err != nil { + t.Fatalf("BuildAllGrantDigests: %v", err) + } + return syncID +} + +// seedEntitlementAtWidth is seedEntitlement but forces a specific +// leaf-level width instead of deriving it from the grant count, so a +// test can build two digests of different widths over a small grant +// set. +func seedEntitlementAtWidth(t testing.TB, e *Engine, entID string, grants []*v3.GrantRecord, widthBits int) string { + t.Helper() + ctx := context.Background() + syncID := ksuid.New().String() + if err := e.SetCurrentSync(syncID); err != nil { + t.Fatalf("SetCurrentSync: %v", err) + } + putEnt(t, e, ctx, entID) + if err := e.PutGrantRecords(ctx, grants...); err != nil { + t.Fatalf("PutGrantRecords: %v", err) + } + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatalf("resolveSyncBytes: %v", err) + } + if err := e.buildPartitionDigestAtWidth(ctx, grantDigestSpec, idBytes, entID, widthBits); err != nil { + t.Fatalf("buildPartitionDigestAtWidth: %v", err) + } + return syncID +} + +// TestDigestDifferentWidthsComparison builds two digests of different +// widths (4 vs 8 bits) over the SAME entitlement and exercises +// DirtyEntitlementBuckets across them. It validates two things the +// equal-width tests cannot: +// +// - split-independence: identical grant content yields the same root +// hash regardless of digest width, and compares as zero dirty +// buckets; +// - the cross-width merge: after one principal's grant changes, the +// comparison (at compareBits = min(4,8) = 4) localizes the change to +// that principal's width-4 bucket — the finer (width-8) side's +// leaves fold down to width 4 during the scan — and leaves a known +// principal in a different bucket clean. +func TestDigestDifferentWidthsComparison(t *testing.T) { + ctx := context.Background() + + const nPrincipals = 40 + principals := make([]string, nPrincipals) + for i := range principals { + principals[i] = fmt.Sprintf("user-%03d", i) + } + mkGrants := func() []*v3.GrantRecord { + gs := make([]*v3.GrantRecord, 0, nPrincipals) + for i, p := range principals { + gs = append(gs, makeGrant("", fmt.Sprintf("g-%03d", i), "ent-A", p)) + } + return gs + } + + // bucket4 returns the width-4 bucket index for a principal, matching + // how grants are keyed (principal type "user" per makeGrant). + bucket4 := func(principalID string) uint32 { + return bucketOfHash(principalBucketHash("user", principalID), 4).Index + } + + ea, _ := newTestEngine(t) + eb, _ := newTestEngine(t) + syncA := seedEntitlementAtWidth(t, ea, "ent-A", mkGrants(), 4) + syncB := seedEntitlementAtWidth(t, eb, "ent-A", mkGrants(), 8) + + ra, okA, err := ea.GetEntitlementDigestRoot(ctx, syncA, "ent-A") + if err != nil || !okA { + t.Fatalf("root A: ok=%v err=%v", okA, err) + } + rb, okB, err := eb.GetEntitlementDigestRoot(ctx, syncB, "ent-A") + if err != nil || !okB { + t.Fatalf("root B: ok=%v err=%v", okB, err) + } + if ra.Bits != 4 || rb.Bits != 8 { + t.Fatalf("widths = %d, %d; want 4, 8", ra.Bits, rb.Bits) + } + // Split-independence: identical content -> identical root despite + // different digest widths. + if !bytes.Equal(ra.Hash, rb.Hash) { + t.Fatalf("different-width digests over identical content disagree on root:\n A(w4)=%x\n B(w8)=%x", ra.Hash, rb.Hash) + } + dirty, err := ea.DirtyEntitlementBuckets(ctx, syncA, eb, syncB, "ent-A") + if err != nil { + t.Fatalf("DirtyEntitlementBuckets (identical): %v", err) + } + if len(dirty) != 0 { + t.Fatalf("identical content across widths: dirty=%d, want 0", len(dirty)) + } + + // Pick the principal to change and a "clean" principal known to sit + // in a different width-4 bucket. + changed := principals[0] + changedBucket := bucket4(changed) + cleanP := "" + for _, p := range principals[1:] { + if bucket4(p) != changedBucket { + cleanP = p + break + } + } + if cleanP == "" { + t.Skip("no principal landed in a different width-4 bucket from the changed one; can't assert localization") + } + + // Mutate the changed principal's grant in B (same external_id -> + // same index key, new content hash via an added source) and rebuild + // B's digest at width 8. + g := makeGrantWithSources(syncB, "g-000", "ent-A", changed, "src-ent") + if err := eb.PutGrantRecord(ctx, g); err != nil { + t.Fatalf("PutGrantRecord (mutate): %v", err) + } + idBytesB, err := eb.resolveSyncBytes(syncB) + if err != nil { + t.Fatal(err) + } + if err := eb.buildPartitionDigestAtWidth(ctx, grantDigestSpec, idBytesB, "ent-A", 8); err != nil { + t.Fatalf("rebuild B: %v", err) + } + + rb2, _, _ := eb.GetEntitlementDigestRoot(ctx, syncB, "ent-A") + if bytes.Equal(ra.Hash, rb2.Hash) { + t.Fatal("mutation did not change B's root") + } + + dirty, err = ea.DirtyEntitlementBuckets(ctx, syncA, eb, syncB, "ent-A") + if err != nil { + t.Fatalf("DirtyEntitlementBuckets (changed): %v", err) + } + if len(dirty) == 0 { + t.Fatal("changed principal across widths produced no dirty buckets") + } + // Localization: every dirty entry is a width-4 bucket (not the + // whole-entitlement zero bucket). + for _, b := range dirty { + if b.Bits != 4 { + t.Fatalf("dirty bucket bits = %d, want 4 (compareBits); got whole-entitlement or wrong-width bucket", b.Bits) + } + } + // Loading the dirty buckets in B surfaces the changed principal and + // excludes the known-clean principal. + loaded := map[string]bool{} + for _, b := range dirty { + if err := eb.IterateGrantsByEntitlementBucket(ctx, syncB, "ent-A", b, func(g *v3.GrantRecord) bool { + loaded[g.GetPrincipal().GetResourceId()] = true + return true + }); err != nil { + t.Fatalf("IterateGrantsByEntitlementBucket: %v", err) + } + } + if !loaded[changed] { + t.Fatalf("dirty buckets did not include the changed principal %q; loaded=%v", changed, loaded) + } + if loaded[cleanP] { + t.Fatalf("dirty buckets wrongly included clean principal %q (different bucket); change was not localized", cleanP) + } +} + +func TestDigestEmptyEntitlementSingleRoot(t *testing.T) { + ctx := context.Background() + e, _ := newTestEngine(t) + syncID := seedEntitlement(t, e, "ent-empty", nil) + + if got := digestNodeCount(t, e, syncID); got != 1 { + t.Fatalf("empty entitlement: digest node count = %d, want 1 (root only)", got) + } + root, ok, err := e.GetEntitlementDigestRoot(ctx, syncID, "ent-empty") + if err != nil || !ok { + t.Fatalf("GetEntitlementDigestRoot: ok=%v err=%v", ok, err) + } + if root.Bits != 0 { + t.Fatalf("empty entitlement width = %d, want 0", root.Bits) + } + if root.Count != 0 { + t.Fatalf("empty entitlement count = %d, want 0", root.Count) + } +} + +func TestDigestIdenticalGrantsSameRoot(t *testing.T) { + ctx := context.Background() + mk := func() []*v3.GrantRecord { + return []*v3.GrantRecord{ + makeGrant("", "g1", "ent-A", "alice"), + makeGrant("", "g2", "ent-A", "bob"), + makeGrant("", "g3", "ent-A", "carol"), + } + } + ea, _ := newTestEngine(t) + eb, _ := newTestEngine(t) + syncA := seedEntitlement(t, ea, "ent-A", mk()) + syncB := seedEntitlement(t, eb, "ent-A", mk()) + + ra, okA, err := ea.GetEntitlementDigestRoot(ctx, syncA, "ent-A") + if err != nil || !okA { + t.Fatalf("root A: ok=%v err=%v", okA, err) + } + rb, okB, err := eb.GetEntitlementDigestRoot(ctx, syncB, "ent-A") + if err != nil || !okB { + t.Fatalf("root B: ok=%v err=%v", okB, err) + } + if !bytes.Equal(ra.Hash, rb.Hash) { + t.Fatalf("identical grants produced different roots:\n A=%x\n B=%x", ra.Hash, rb.Hash) + } + dirty, err := ea.DirtyEntitlementBuckets(ctx, syncA, eb, syncB, "ent-A") + if err != nil { + t.Fatalf("DirtyEntitlementBuckets: %v", err) + } + if len(dirty) != 0 { + t.Fatalf("identical grants: dirty buckets = %d, want 0", len(dirty)) + } +} + +func TestDigestContentChangeDirtyBucket(t *testing.T) { + ctx := context.Background() + // Base set: same in both engines except bob's grant gains a source + // in B. external_id is unchanged, so the index KEY is identical and + // only the content hash (and thus bob's bucket) differs. + baseA := []*v3.GrantRecord{ + makeGrant("", "g1", "ent-A", "alice"), + makeGrant("", "g2", "ent-A", "bob"), + makeGrant("", "g3", "ent-A", "carol"), + } + baseB := []*v3.GrantRecord{ + makeGrant("", "g1", "ent-A", "alice"), + makeGrantWithSources("", "g2", "ent-A", "bob", "src-ent"), + makeGrant("", "g3", "ent-A", "carol"), + } + ea, _ := newTestEngine(t) + eb, _ := newTestEngine(t) + syncA := seedEntitlement(t, ea, "ent-A", baseA) + syncB := seedEntitlement(t, eb, "ent-A", baseB) + + ra, _, _ := ea.GetEntitlementDigestRoot(ctx, syncA, "ent-A") + rb, _, _ := eb.GetEntitlementDigestRoot(ctx, syncB, "ent-A") + if bytes.Equal(ra.Hash, rb.Hash) { + t.Fatal("content change did not change the root hash") + } + + dirty, err := ea.DirtyEntitlementBuckets(ctx, syncA, eb, syncB, "ent-A") + if err != nil { + t.Fatalf("DirtyEntitlementBuckets: %v", err) + } + if len(dirty) == 0 { + t.Fatal("content change produced no dirty buckets") + } + + // Loading the dirty buckets in B must surface bob (the changed + // principal) and must NOT require touching alice/carol's buckets. + found := map[string]bool{} + for _, b := range dirty { + if err := eb.IterateGrantsByEntitlementBucket(ctx, syncB, "ent-A", b, func(g *v3.GrantRecord) bool { + found[g.GetPrincipal().GetResourceId()] = true + return true + }); err != nil { + t.Fatalf("IterateGrantsByEntitlementBucket: %v", err) + } + } + if !found["bob"] { + t.Fatalf("dirty buckets did not include the changed principal bob; found=%v", found) + } +} + +func TestDigestAddedGrantDirtyBucket(t *testing.T) { + ctx := context.Background() + baseA := []*v3.GrantRecord{ + makeGrant("", "g1", "ent-A", "alice"), + makeGrant("", "g2", "ent-A", "bob"), + } + baseB := []*v3.GrantRecord{ + makeGrant("", "g1", "ent-A", "alice"), + makeGrant("", "g2", "ent-A", "bob"), + makeGrant("", "g3", "ent-A", "dave"), // added in B + } + ea, _ := newTestEngine(t) + eb, _ := newTestEngine(t) + syncA := seedEntitlement(t, ea, "ent-A", baseA) + syncB := seedEntitlement(t, eb, "ent-A", baseB) + + dirty, err := ea.DirtyEntitlementBuckets(ctx, syncA, eb, syncB, "ent-A") + if err != nil { + t.Fatalf("DirtyEntitlementBuckets: %v", err) + } + if len(dirty) == 0 { + t.Fatal("added grant produced no dirty buckets") + } + found := map[string]bool{} + for _, b := range dirty { + if err := eb.IterateGrantsByEntitlementBucket(ctx, syncB, "ent-A", b, func(g *v3.GrantRecord) bool { + found[g.GetPrincipal().GetResourceId()] = true + return true + }); err != nil { + t.Fatalf("IterateGrantsByEntitlementBucket: %v", err) + } + } + if !found["dave"] { + t.Fatalf("dirty buckets did not include the added principal dave; found=%v", found) + } +} + +func TestDigestVariableWidth(t *testing.T) { + ctx := context.Background() + + // Small entitlement: under one target bucket -> width 0, single node. + small := make([]*v3.GrantRecord, 0, 10) + for i := 0; i < 10; i++ { + small = append(small, makeGrant("", ksuid.New().String(), "ent-small", ksuid.New().String())) + } + es, _ := newTestEngine(t) + syncS := seedEntitlement(t, es, "ent-small", small) + rootS, ok, err := es.GetEntitlementDigestRoot(ctx, syncS, "ent-small") + if err != nil || !ok { + t.Fatalf("small root: ok=%v err=%v", ok, err) + } + if rootS.Bits != 0 { + t.Fatalf("small entitlement width = %d, want 0", rootS.Bits) + } + if rootS.Count != 10 { + t.Fatalf("small entitlement count = %d, want 10", rootS.Count) + } + if got := digestNodeCount(t, es, syncS); got != 1 { + t.Fatalf("small entitlement node count = %d, want 1", got) + } + + // Large entitlement: well over the target bucket size -> the width + // grows one bit at a time, and the digest gains leaf nodes beyond + // the root. + const n = digestTargetBucketSize*3 + 7 + large := make([]*v3.GrantRecord, 0, n) + for i := 0; i < n; i++ { + large = append(large, makeGrant("", ksuid.New().String(), "ent-large", ksuid.New().String())) + } + el, _ := newTestEngine(t) + syncL := seedEntitlement(t, el, "ent-large", large) + rootL, ok, err := el.GetEntitlementDigestRoot(ctx, syncL, "ent-large") + if err != nil || !ok { + t.Fatalf("large root: ok=%v err=%v", ok, err) + } + if want := chooseDigestWidth(n); rootL.Bits != want { + t.Fatalf("large entitlement width = %d, want %d", rootL.Bits, want) + } + if rootL.Count != int64(n) { + t.Fatalf("large entitlement count = %d, want %d", rootL.Count, n) + } + // root + at least 2 leaves (width>=1 over n grants spreads across + // multiple buckets). + if got := digestNodeCount(t, el, syncL); got < 3 { + t.Fatalf("large entitlement node count = %d, want >= 3 (root + leaves)", got) + } + // Capacity invariant: 2^width buckets at the target size must cover + // the count, and width-1 must not (else the width is too large). + if int64(1)< 0 && int64(1)<<(rootL.Bits-1)*digestTargetBucketSize >= n { + t.Fatalf("width %d is one bit wider than the count %d needs", rootL.Bits, n) + } +} + +// TestHashIndexIsHashOrdered verifies the index iterates in +// hash(principal) order: the embedded bucket-hash region is +// non-decreasing across the entitlement's index range. +func TestHashIndexIsHashOrdered(t *testing.T) { + e, _ := newTestEngine(t) + grants := make([]*v3.GrantRecord, 0, 200) + for i := 0; i < 200; i++ { + grants = append(grants, makeGrant("", ksuid.New().String(), "ent-A", ksuid.New().String())) + } + syncID := seedEntitlement(t, e, "ent-A", grants) + + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatal(err) + } + entPrefix := encodeGrantByEntPrincHashEntPrefix(idBytes, "ent-A") + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: entPrefix, UpperBound: upperBoundOf(entPrefix)}) + if err != nil { + t.Fatal(err) + } + defer iter.Close() + var prev []byte + count := 0 + for iter.First(); iter.Valid(); iter.Next() { + bh, _, _, _, ok := decodeEntPrincHashTail(iter.Key(), entPrefix) + if !ok { + t.Fatal("failed to decode index tail") + } + if prev != nil && bytes.Compare(bh, prev) < 0 { + t.Fatalf("index not hash-ordered: %x < %x", bh, prev) + } + prev = append(prev[:0], bh...) + count++ + } + if count != 200 { + t.Fatalf("hash index entry count = %d, want 200", count) + } +} + +// dumpDigestNodes snapshots every digest node key/value for a sync. +// Used to byte-compare an incrementally-maintained digest against a +// from-scratch rebuild. +func dumpDigestNodes(t testing.TB, e *Engine, syncID string) map[string][]byte { + t.Helper() + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatalf("resolveSyncBytes: %v", err) + } + iter, err := e.db.NewIter(&pebble.IterOptions{ + LowerBound: DigestSyncLowerBound(idBytes), + UpperBound: DigestSyncUpperBound(idBytes), + }) + if err != nil { + t.Fatalf("NewIter: %v", err) + } + defer iter.Close() + out := map[string][]byte{} + for iter.First(); iter.Valid(); iter.Next() { + out[string(iter.Key())] = append([]byte(nil), iter.Value()...) + } + if err := iter.Error(); err != nil { + t.Fatalf("iter: %v", err) + } + return out +} + +// requireSameDigestNodes fails with a per-key diff when two node +// snapshots differ. +func requireSameDigestNodes(t *testing.T, got, want map[string][]byte) { + t.Helper() + for k, wv := range want { + gv, ok := got[k] + if !ok { + t.Errorf("missing node %x (want %x)", k, wv) + continue + } + if !bytes.Equal(gv, wv) { + t.Errorf("node %x differs:\n got %x\nwant %x", k, gv, wv) + } + } + for k, gv := range got { + if _, ok := want[k]; !ok { + t.Errorf("extra node %x = %x", k, gv) + } + } +} + +// TestDigestLeafFoldConsistent verifies the leaf-level build and the +// fold machinery the comparison rests on: every stored leaf is +// non-empty, the root is exactly the XOR (and count-sum) of the leaves, +// no nodes exist beyond root + leaves, folding the leaf level to a +// coarser width matches a manual regrouping, and a stored leaf +// byte-matches the authoritative on-demand fold of its bucket. +func TestDigestLeafFoldConsistent(t *testing.T) { + ctx := context.Background() + e, _ := newTestEngine(t) + const n = 60 + grants := make([]*v3.GrantRecord, 0, n) + for i := 0; i < n; i++ { + grants = append(grants, makeGrant("", fmt.Sprintf("g-%03d", i), "ent-A", fmt.Sprintf("user-%03d", i))) + } + syncID := seedEntitlementAtWidth(t, e, "ent-A", grants, 8) + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatal(err) + } + + root, ok, err := e.GetEntitlementDigestRoot(ctx, syncID, "ent-A") + if err != nil || !ok { + t.Fatalf("root: ok=%v err=%v", ok, err) + } + if root.Bits != 8 || root.Count != n { + t.Fatalf("root width=%d count=%d, want 8, %d", root.Bits, root.Count, n) + } + + // Folding at the build width returns the stored leaves one-to-one. + leaves, err := e.foldedLeafBuckets(ctx, grantDigestSpec, idBytes, "ent-A", 8) + if err != nil { + t.Fatal(err) + } + if len(leaves) == 0 { + t.Fatal("no leaf nodes stored") + } + var ( + rootXor [hashLen]byte + rootCount int64 + ) + for _, l := range leaves { + if l.count < 1 { + t.Fatalf("leaf %d stored with count %d; empty leaves must not be materialized", l.idx, l.count) + } + xorInto(rootXor[:], l.digest[:]) + rootCount += l.count + } + if rootCount != root.Count || !bytes.Equal(rootXor[:], root.Hash) { + t.Fatalf("root != fold of leaves: count %d vs %d", root.Count, rootCount) + } + + // Exactly root + leaves — nothing else in the keyspace. + if got, want := digestNodeCount(t, e, syncID), 1+len(leaves); got != want { + t.Fatalf("total node count = %d, want %d (root + leaves only)", got, want) + } + + // Folding to a coarser width matches a manual regroup of the + // build-width leaves. + leaves4, err := e.foldedLeafBuckets(ctx, grantDigestSpec, idBytes, "ent-A", 4) + if err != nil { + t.Fatal(err) + } + manual := map[uint32]*foldedBucket{} + var order []uint32 + for _, l := range leaves { + idx := l.idx >> 4 + fb, ok := manual[idx] + if !ok { + fb = &foldedBucket{idx: idx} + manual[idx] = fb + order = append(order, idx) + } + fb.count += l.count + xorInto(fb.digest[:], l.digest[:]) + } + if len(leaves4) != len(order) { + t.Fatalf("fold to width 4: %d buckets, want %d", len(leaves4), len(order)) + } + for i, idx := range order { + got, want := leaves4[i], manual[idx] + if got.idx != want.idx || got.count != want.count || got.digest != want.digest { + t.Fatalf("folded bucket %d mismatch: got {%d %d %x}, want {%d %d %x}", + i, got.idx, got.count, got.digest, want.idx, want.count, want.digest) + } + } + + // A stored leaf is a cache of the authoritative fold. + b := DigestBucket{Index: leaves[0].idx, Bits: 8} + h, c, err := e.ComputeEntitlementBucketDigest(ctx, syncID, "ent-A", b) + if err != nil { + t.Fatal(err) + } + lc, ld, present, err := e.getDigestLeaf(grantDigestSpec, idBytes, "ent-A", b.leafKeyPrefix()) + if err != nil || !present { + t.Fatalf("leaf %d: present=%v err=%v", b.Index, present, err) + } + if c != lc || !bytes.Equal(h, ld) { + t.Fatalf("stored leaf disagrees with ComputeEntitlementBucketDigest: count %d vs %d", lc, c) + } +} + +// TestDigestRebuildClearsStaleNodes verifies the leading DeleteRange in +// the build: a rebuild at a narrower width must remove the prior +// build's finer-grained leaves, or the comparison merge scan would read +// them. +func TestDigestRebuildClearsStaleNodes(t *testing.T) { + ctx := context.Background() + e, _ := newTestEngine(t) + grants := make([]*v3.GrantRecord, 0, 40) + for i := 0; i < 40; i++ { + grants = append(grants, makeGrant("", fmt.Sprintf("g-%03d", i), "ent-A", fmt.Sprintf("user-%03d", i))) + } + syncID := seedEntitlementAtWidth(t, e, "ent-A", grants, 8) + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatal(err) + } + + before := rawLeafPrefixes(t, e, idBytes, "ent-A") + if len(before) == 0 { + t.Fatal("width-8 build produced no leaves") + } + rootBefore, _, err := e.GetEntitlementDigestRoot(ctx, syncID, "ent-A") + if err != nil { + t.Fatal(err) + } + + if err := e.buildPartitionDigestAtWidth(ctx, grantDigestSpec, idBytes, "ent-A", 4); err != nil { + t.Fatalf("rebuild at width 4: %v", err) + } + + rootAfter, ok, err := e.GetEntitlementDigestRoot(ctx, syncID, "ent-A") + if err != nil || !ok { + t.Fatalf("root after rebuild: ok=%v err=%v", ok, err) + } + if rootAfter.Bits != 4 { + t.Fatalf("root width after rebuild = %d, want 4", rootAfter.Bits) + } + // Split-independence: same content, same root digest. + if !bytes.Equal(rootBefore.Hash, rootAfter.Hash) || rootBefore.Count != rootAfter.Count { + t.Fatal("rebuild at different width changed the root digest/count over identical content") + } + // Every surviving leaf prefix must be width-4 aligned (low 12 bits + // of the left-aligned prefix zero) — a width-8 leaf that escaped the + // range-clear would fail this. + after := rawLeafPrefixes(t, e, idBytes, "ent-A") + if len(after) == 0 || len(after) > 16 { + t.Fatalf("width-4 rebuild stored %d leaves, want 1..16", len(after)) + } + for _, p := range after { + if lv := binary.BigEndian.Uint16(p); lv&0x0FFF != 0 { + t.Fatalf("stale leaf prefix %x survived the width-4 rebuild", p) + } + } +} + +// TestDigestIncrementalEqualsRebuild is the §7 keystone invariant: after +// a sequence of post-seal inserts, content overwrites, a bucket-moving +// (principal-changing) overwrite, an excluded-field no-op overwrite, +// deletes, and a multi-record batch, the incrementally-maintained +// digest byte-equals a from-scratch rebuild. +func TestDigestIncrementalEqualsRebuild(t *testing.T) { + ctx := context.Background() + e, _ := newTestEngine(t) + grants := make([]*v3.GrantRecord, 0, 30) + for i := 0; i < 30; i++ { + grants = append(grants, makeGrant("", fmt.Sprintf("g-%03d", i), "ent-A", fmt.Sprintf("user-%03d", i))) + } + syncID := seedEntitlementAtWidth(t, e, "ent-A", grants, 8) + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatal(err) + } + + put := func(g *v3.GrantRecord) { + t.Helper() + if err := e.PutGrantRecord(ctx, g); err != nil { + t.Fatalf("PutGrantRecord: %v", err) + } + } + + // Post-seal inserts. + put(makeGrant(syncID, "g-100", "ent-A", "new-user-1")) + put(makeGrant(syncID, "g-101", "ent-A", "new-user-2")) + // Content overwrite: same index key, sources changed. + put(makeGrantWithSources(syncID, "g-005", "ent-A", "user-005", "src-ent")) + // Bucket-moving overwrite: same external_id, principal changed — + // must apply as remove(old path) + add(new path). + put(makeGrant(syncID, "g-006", "ent-A", "user-moved")) + // Excluded-field overwrite: needs_expansion is not part of the + // content hash, so this must leave the digest untouched. + noop := makeGrant(syncID, "g-007", "ent-A", "user-007") + noop.SetNeedsExpansion(true) + put(noop) + // Deletes. + for _, ext := range []string{"g-008", "g-009"} { + if err := e.DeleteGrantRecord(ctx, syncID, ext); err != nil { + t.Fatalf("DeleteGrantRecord(%s): %v", ext, err) + } + } + // Multi-record batch: inserts + an overwrite in one PutGrantRecords + // call, exercising the per-node delta accumulator (all of them + // share the root). + batch := []*v3.GrantRecord{ + makeGrant(syncID, "g-110", "ent-A", "batch-user-1"), + makeGrant(syncID, "g-111", "ent-A", "batch-user-2"), + makeGrant(syncID, "g-112", "ent-A", "batch-user-3"), + makeGrantWithSources(syncID, "g-010", "ent-A", "user-010", "src-2"), + } + if err := e.PutGrantRecords(ctx, batch...); err != nil { + t.Fatalf("PutGrantRecords: %v", err) + } + + // Sparsity restored on delete: unless another remaining principal + // shares user-008's width-8 bucket, its leaf must be gone. + remaining := []string{"user-moved", "new-user-1", "new-user-2", "batch-user-1", "batch-user-2", "batch-user-3"} + for i := 0; i < 30; i++ { + if i == 6 || i == 8 || i == 9 { + continue // moved or deleted + } + remaining = append(remaining, fmt.Sprintf("user-%03d", i)) + } + deletedLeaf := bucketOfHash(principalBucketHash("user", "user-008"), 8).leafKeyPrefix() + shared := false + for _, p := range remaining { + if bytes.Equal(bucketOfHash(principalBucketHash("user", p), 8).leafKeyPrefix(), deletedLeaf) { + shared = true + break + } + } + if !shared { + _, _, present, err := e.getDigestLeaf(grantDigestSpec, idBytes, "ent-A", deletedLeaf) + if err != nil { + t.Fatal(err) + } + if present { + t.Fatal("emptied leaf node survived an incremental delete; sparsity not restored") + } + } + + incremental := dumpDigestNodes(t, e, syncID) + if err := e.buildPartitionDigestAtWidth(ctx, grantDigestSpec, idBytes, "ent-A", 8); err != nil { + t.Fatalf("rebuild: %v", err) + } + rebuilt := dumpDigestNodes(t, e, syncID) + requireSameDigestNodes(t, incremental, rebuilt) +} + +// TestDigestSameBatchSameBucketRMW pins the accumulator behavior: one +// PutGrantRecords batch adds several grants that land in the SAME +// width-8 bucket. A naive read-modify-write against the batch would +// lose all but one delta (a plain pebble.Batch doesn't read through its +// own writes); the accumulator must fold all of them into one node +// write. +func TestDigestSameBatchSameBucketRMW(t *testing.T) { + ctx := context.Background() + e, _ := newTestEngine(t) + + // Find three principals whose bucket hashes share a first byte — + // the top 8 bits, i.e. the same width-8 bucket. + collide := []string{"seed-principal"} + target := principalBucketHash("user", collide[0])[0] + for i := 0; len(collide) < 3; i++ { + p := fmt.Sprintf("cand-%d", i) + if principalBucketHash("user", p)[0] == target { + collide = append(collide, p) + } + } + + base := make([]*v3.GrantRecord, 0, 5) + for i := 0; i < 5; i++ { + base = append(base, makeGrant("", fmt.Sprintf("b-%d", i), "ent-A", fmt.Sprintf("base-%d", i))) + } + syncID := seedEntitlementAtWidth(t, e, "ent-A", base, 8) + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + t.Fatal(err) + } + + gs := make([]*v3.GrantRecord, 0, len(collide)) + for i, p := range collide { + gs = append(gs, makeGrant(syncID, fmt.Sprintf("x-%d", i), "ent-A", p)) + } + if err := e.PutGrantRecords(ctx, gs...); err != nil { + t.Fatalf("PutGrantRecords: %v", err) + } + + want := int64(len(collide)) + for i := 0; i < 5; i++ { + if principalBucketHash("user", fmt.Sprintf("base-%d", i))[0] == target { + want++ + } + } + count, _, present, err := e.getDigestLeaf(grantDigestSpec, idBytes, "ent-A", []byte{target, 0}) + if err != nil || !present { + t.Fatalf("bucket leaf %x: present=%v err=%v", target, present, err) + } + if count != want { + t.Fatalf("same-batch deltas lost: bucket count = %d, want %d", count, want) + } + + incremental := dumpDigestNodes(t, e, syncID) + if err := e.buildPartitionDigestAtWidth(ctx, grantDigestSpec, idBytes, "ent-A", 8); err != nil { + t.Fatalf("rebuild: %v", err) + } + requireSameDigestNodes(t, incremental, dumpDigestNodes(t, e, syncID)) +} + +// TestDigestMissingRootFallback: a missing root means "digest never +// built", not "no grants". Comparison against a populated-but-unbuilt +// side must fall back to the authoritative fold — clean when content is +// identical, whole-entitlement dirty when it differs. +func TestDigestMissingRootFallback(t *testing.T) { + ctx := context.Background() + mk := func() []*v3.GrantRecord { + return []*v3.GrantRecord{ + makeGrant("", "g1", "ent-A", "alice"), + makeGrant("", "g2", "ent-A", "bob"), + makeGrant("", "g3", "ent-A", "carol"), + } + } + ea, _ := newTestEngine(t) + syncA := seedEntitlement(t, ea, "ent-A", mk()) + + // B holds the same grants but never builds a digest. + eb, _ := newTestEngine(t) + syncB := ksuid.New().String() + if err := eb.SetCurrentSync(syncB); err != nil { + t.Fatal(err) + } + putEnt(t, eb, ctx, "ent-A") + for _, g := range mk() { + if err := eb.PutGrantRecord(ctx, g); err != nil { + t.Fatal(err) + } + } + if _, ok, err := eb.GetEntitlementDigestRoot(ctx, syncB, "ent-A"); err != nil || ok { + t.Fatalf("B unexpectedly has a root: ok=%v err=%v", ok, err) + } + + for name, dirtyFn := range map[string]func() ([]DigestBucket, error){ + "A vs B": func() ([]DigestBucket, error) { return ea.DirtyEntitlementBuckets(ctx, syncA, eb, syncB, "ent-A") }, + "B vs A": func() ([]DigestBucket, error) { return eb.DirtyEntitlementBuckets(ctx, syncB, ea, syncA, "ent-A") }, + } { + dirty, err := dirtyFn() + if err != nil { + t.Fatalf("%s: %v", name, err) + } + if len(dirty) != 0 { + t.Fatalf("%s: identical content with one digest unbuilt: dirty=%d, want 0 (fold fallback)", name, len(dirty)) + } + } + + // Diverge B; the fallback must now flag the whole entitlement. + if err := eb.PutGrantRecord(ctx, makeGrant(syncB, "g9", "ent-A", "zed")); err != nil { + t.Fatal(err) + } + dirty, err := ea.DirtyEntitlementBuckets(ctx, syncA, eb, syncB, "ent-A") + if err != nil { + t.Fatal(err) + } + if len(dirty) != 1 || dirty[0].Bits != 0 { + t.Fatalf("diverged content with one digest unbuilt: dirty=%v, want one whole-entitlement bucket", dirty) + } +} diff --git a/pkg/dotc1z/engine/pebble/grant_digest.go b/pkg/dotc1z/engine/pebble/grant_digest.go new file mode 100644 index 000000000..5416b57eb --- /dev/null +++ b/pkg/dotc1z/engine/pebble/grant_digest.go @@ -0,0 +1,260 @@ +package pebble + +import ( + "context" + "encoding/binary" + "errors" + "fmt" + "sort" + + "github.com/cespare/xxhash/v2" + "github.com/cockroachdb/pebble/v2" + + v3 "github.com/conductorone/baton-sdk/pb/c1/storage/v3" +) + +// Grant instantiation of the digest core (digest.go): the per- +// entitlement grant digest, folded over the +// by_entitlement_principal_hash index. +// +// partition = entitlement_id +// bucket hash = principalBucketHash (identity of the principal) +// content hash = grantContentHash (the membership edge) +// +// This answers "does this entitlement have exactly the same grants as +// some other sync/file?" with a single root read, and localizes any +// difference to principal-hash buckets so the diff driver loads only +// those grants (see adapter_diff.go). + +// grantDigestSpec wires the grant hash index into the digest core. The +// index's key layout (encodeGrantByEntPrincHashIndexKey) satisfies the +// digestIndexSpec shape contract: partition prefix, then the raw 8-byte +// bucket hash, then the principal/external_id tail; the value is the +// grant content hash. +var grantDigestSpec = digestIndexSpec{ + indexID: idxGrantByEntitlementPrincipalHash, + partitionPrefix: encodeGrantByEntPrincHashEntPrefix, + syncBounds: func(syncIDBytes []byte) ([]byte, []byte) { + return GrantByEntPrincHashSyncLowerBound(syncIDBytes), GrantByEntPrincHashSyncUpperBound(syncIDBytes) + }, +} + +// principalBucketHash is the bucket address for a principal: the 8-byte +// xxHash64 of (rt + "\x00" + id). Identity only — never the principal's +// full object — so the address is stable across syncs even when the +// principal's attributes change. Returns a fresh slice. +func principalBucketHash(rt, id string) []byte { + h := xxhash.New() + _, _ = h.WriteString(rt) + _, _ = h.Write([]byte{0}) + _, _ = h.WriteString(id) + out := make([]byte, digestBucketHashLen) + binary.BigEndian.PutUint64(out, h.Sum64()) + return out +} + +// grantContentHash is the canonical content hash of a grant — the value +// stored in the hash index and the unit the grant digest folds. +// +// ABI: the field set below defines what "the same grant" means for the +// diff. It deliberately covers the membership EDGE (entitlement id, +// principal identity, external_id) plus the grant's source-entitlement +// set (expansion provenance), and deliberately EXCLUDES sync-relative +// and transient processing state — sync_id, discovered_at, +// needs_expansion, expansion, and annotations — none of which change +// "which principal holds which entitlement". This is a hand-rolled +// framing, NOT proto marshal: deterministic-proto output is not +// canonical across protobuf library versions, which would make two +// files written by different SDK builds hash identical grants +// differently. Changing this set requires an index-migration bump. +func grantContentHash(r *v3.GrantRecord) []byte { + h := xxhash.New() + ent := r.GetEntitlement() + princ := r.GetPrincipal() + _, _ = h.WriteString(ent.GetEntitlementId()) + _, _ = h.Write([]byte{0}) + _, _ = h.WriteString(princ.GetResourceTypeId()) + _, _ = h.Write([]byte{0}) + _, _ = h.WriteString(princ.GetResourceId()) + _, _ = h.Write([]byte{0}) + _, _ = h.WriteString(r.GetExternalId()) + _, _ = h.Write([]byte{0}) + + // Source-entitlement ids, sorted for order-independence. The map + // values (GrantSourceRecord) are not folded in v1 — only the set of + // source ids, which is the membership-composition signal. + sources := r.GetSources() + ids := make([]string, 0, len(sources)) + for k := range sources { + ids = append(ids, k) + } + sort.Strings(ids) + for _, id := range ids { + _, _ = h.WriteString(id) + _, _ = h.Write([]byte{0}) + } + out := make([]byte, hashLen) + binary.BigEndian.PutUint64(out, h.Sum64()) + return out +} + +// grantHashIndexKey returns the by_entitlement_principal_hash index key +// for r, or nil when the grant lacks the entitlement/principal needed to +// place it (mirrors the by_entitlement index's nil-guard). +func grantHashIndexKey(syncIDBytes []byte, r *v3.GrantRecord) []byte { + ent := r.GetEntitlement() + princ := r.GetPrincipal() + if ent == nil || princ == nil { + return nil + } + bh := principalBucketHash(princ.GetResourceTypeId(), princ.GetResourceId()) + return encodeGrantByEntPrincHashIndexKey( + syncIDBytes, ent.GetEntitlementId(), bh, + princ.GetResourceTypeId(), princ.GetResourceId(), r.GetExternalId(), + ) +} + +// BuildAllGrantDigests rebuilds the grant digest for every entitlement +// in syncID. Called at seal time (Adapter.EndSync) after all grants are +// written, and by the on-Open migration backfill. Every entitlement +// gets a digest — including those with zero grants, which store a +// single root node — so a reader can always distinguish "empty" from +// "never built". +func (e *Engine) BuildAllGrantDigests(ctx context.Context, syncID string) error { + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + return err + } + // Collect entitlement ids first: the build writes into the + // typeDigest keyspace while we'd otherwise be mid-iteration over + // typeEntitlement. Different keyspaces, but snapshotting the ids + // keeps the iterator and the writes cleanly separated. + var ents []string + if err := e.IterateEntitlementsBySync(ctx, syncID, func(r *v3.EntitlementRecord) bool { + ents = append(ents, r.GetExternalId()) + return true + }); err != nil { + return fmt.Errorf("BuildAllGrantDigests: list entitlements: %w", err) + } + for _, ent := range ents { + if err := ctx.Err(); err != nil { + return err + } + if err := e.buildPartitionDigest(ctx, grantDigestSpec, idBytes, ent); err != nil { + return fmt.Errorf("BuildAllGrantDigests: entitlement %q: %w", ent, err) + } + } + return nil +} + +// GetEntitlementDigestRoot returns the stored grant-digest root for an +// entitlement. ok is false when no digest has been built for it (the +// caller can fall back to ComputeEntitlementBucketDigest, which derives +// the same digest from the index on demand). +func (e *Engine) GetEntitlementDigestRoot(ctx context.Context, syncID, entitlementID string) (DigestRoot, bool, error) { + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + return DigestRoot{}, false, err + } + return e.getPartitionDigestRoot(grantDigestSpec, idBytes, entitlementID) +} + +// ComputeEntitlementBucketDigest folds the grant hash index over a +// single bucket of an entitlement (the zero bucket = the whole +// entitlement) — the authoritative on-demand counterpart of the stored +// digest nodes. +func (e *Engine) ComputeEntitlementBucketDigest(ctx context.Context, syncID, entitlementID string, bucket DigestBucket) ([]byte, int64, error) { + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + return nil, 0, err + } + return e.computeBucketDigest(ctx, grantDigestSpec, idBytes, entitlementID, bucket) +} + +// DirtyEntitlementBuckets compares this engine's entitlement against +// other's and returns the buckets whose grants differ — see +// dirtyPartitionBuckets for the comparison contract (zero bucket = +// whole entitlement; nil = identical). +func (e *Engine) DirtyEntitlementBuckets(ctx context.Context, syncID string, other *Engine, otherSyncID, entitlementID string) ([]DigestBucket, error) { + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + return nil, err + } + otherIDBytes, err := other.resolveSyncBytes(otherSyncID) + if err != nil { + return nil, err + } + return e.dirtyPartitionBuckets(ctx, grantDigestSpec, idBytes, other, otherIDBytes, entitlementID) +} + +// IterateGrantsByEntitlementBucket yields the grants in one +// principal-hash bucket of an entitlement (the zero bucket = the whole +// entitlement). This is the dirty-bucket loader: after a digest +// comparison flags a bucket, the caller materializes only those grants. +// Like the other index iterators it does a point Get per entry to fetch +// the primary; orphan index entries are skipped. +func (e *Engine) IterateGrantsByEntitlementBucket(ctx context.Context, syncID, entitlementID string, bucket DigestBucket, yield func(*v3.GrantRecord) bool) error { + idBytes, err := e.resolveSyncBytes(syncID) + if err != nil { + return err + } + entPrefix := encodeGrantByEntPrincHashEntPrefix(idBytes, entitlementID) + lower, upper := grantDigestSpec.bucketBounds(idBytes, entitlementID, bucket) + iter, err := e.db.NewIter(&pebble.IterOptions{LowerBound: lower, UpperBound: upper}) + if err != nil { + return err + } + defer iter.Close() + for iter.First(); iter.Valid(); iter.Next() { + if err := ctx.Err(); err != nil { + return err + } + _, _, _, externalID, ok := decodeEntPrincHashTail(iter.Key(), entPrefix) + if !ok { + continue + } + val, closer, getErr := e.db.Get(encodeGrantKey(idBytes, externalID)) + if getErr != nil { + if errors.Is(getErr, pebble.ErrNotFound) { + continue + } + return getErr + } + r := &v3.GrantRecord{} + uErr := unmarshalRecord(val, r) + closer.Close() + if uErr != nil { + return fmt.Errorf("IterateGrantsByEntitlementBucket: unmarshal: %w", uErr) + } + if !yield(r) { + return nil + } + } + return iter.Error() +} + +// newGrantDigestMutator returns a digestMutator bound to the grant +// digest; feed it with addGrant/removeGrant. +func newGrantDigestMutator(e *Engine) *digestMutator { + return newDigestMutator(e, grantDigestSpec) +} + +// addGrant records r's insertion into its entitlement's grant digest. +func (m *digestMutator) addGrant(idBytes []byte, r *v3.GrantRecord) error { + return m.grantDelta(idBytes, r, m.addHash) +} + +// removeGrant records r's removal from its entitlement's grant digest. +func (m *digestMutator) removeGrant(idBytes []byte, r *v3.GrantRecord) error { + return m.grantDelta(idBytes, r, m.removeHash) +} + +func (m *digestMutator) grantDelta(idBytes []byte, r *v3.GrantRecord, apply func(idBytes []byte, partition string, bucketHash, contentHash []byte) error) error { + ent := r.GetEntitlement() + princ := r.GetPrincipal() + if ent == nil || princ == nil { + return nil // not in the hash index → not in the digest + } + bh := principalBucketHash(princ.GetResourceTypeId(), princ.GetResourceId()) + return apply(idBytes, ent.GetEntitlementId(), bh, grantContentHash(r)) +} diff --git a/pkg/dotc1z/engine/pebble/grants.go b/pkg/dotc1z/engine/pebble/grants.go index 04998e8fc..0cae2cb62 100644 --- a/pkg/dotc1z/engine/pebble/grants.go +++ b/pkg/dotc1z/engine/pebble/grants.go @@ -87,6 +87,14 @@ func (e *Engine) PutGrantRecords(ctx context.Context, records ...*v3.GrantRecord return err } + // Incremental digest maintenance applies only on the non-fresh + // path: during a fresh sync the digests don't exist yet (built + // once at seal), so skip even the per-entitlement root probe. + var mm *digestMutator + if !fresh { + mm = newGrantDigestMutator(e) + } + // Dedup pre-pass: keep only the LAST occurrence of each // external_id. The map value is the records[] // index — when we re-iterate, we process record i only if @@ -127,6 +135,21 @@ func (e *Engine) PutGrantRecords(ctx context.Context, records ...*v3.GrantRecord closer.Close() return err } + if mm != nil { + // The digest removal needs the old record's content + // hash, which folds fields the raw index scan does + // not extract (sources) — unmarshal only on this + // non-fresh path. + old := &v3.GrantRecord{} + if err := unmarshalRecord(oldVal, old); err != nil { + closer.Close() + return fmt.Errorf("PutGrantRecords: unmarshal old %q: %w", r.GetExternalId(), err) + } + if err := mm.removeGrant(idBytes, old); err != nil { + closer.Close() + return err + } + } closer.Close() case errors.Is(getErr, pebble.ErrNotFound): // no prior record — write unconditionally @@ -140,6 +163,16 @@ func (e *Engine) PutGrantRecords(ctx context.Context, records ...*v3.GrantRecord if err := e.writeGrantIndexes(idxBatch, idBytes, r); err != nil { return err } + if mm != nil { + if err := mm.addGrant(idBytes, r); err != nil { + return err + } + } + } + if mm != nil { + if err := mm.apply(idxBatch); err != nil { + return err + } } opts := writeOpts(e.opts.durability) if fresh { @@ -183,6 +216,8 @@ func (e *Engine) UnsafePutUniqueGrantRecords(ctx context.Context, records ...*v3 priKey []byte priVal []byte idxKeys [][]byte + hashKey []byte + hashVal []byte } enc := make([]encoded, len(records)) @@ -231,6 +266,8 @@ func (e *Engine) UnsafePutUniqueGrantRecords(ctx context.Context, records ...*v3 priKey: encodeGrantKey(idBytes, r.GetExternalId()), priVal: val, idxKeys: grantIndexKeys(idBytes, r), + hashKey: grantHashIndexKey(idBytes, r), + hashVal: grantContentHash(r), } } }(lo, hi) @@ -256,6 +293,11 @@ func (e *Engine) UnsafePutUniqueGrantRecords(ctx context.Context, records ...*v3 return err } } + if enc[i].hashKey != nil { + if err := idxBatch.Set(enc[i].hashKey, enc[i].hashVal, nil); err != nil { + return err + } + } } opts := writeOpts(e.opts.durability) @@ -312,8 +354,23 @@ func (e *Engine) DeleteGrantRecord(ctx context.Context, syncID, externalID strin closer.Close() return err } + // The digest removal needs the old record's content hash, which + // folds fields the raw index scan does not extract (sources). + old := &v3.GrantRecord{} + if uErr := unmarshalRecord(oldVal, old); uErr != nil { + closer.Close() + return fmt.Errorf("DeleteGrantRecord: unmarshal old %q: %w", externalID, uErr) + } closer.Close() + mm := newGrantDigestMutator(e) + if err := mm.removeGrant(idBytes, old); err != nil { + return err + } + if err := mm.apply(batch); err != nil { + return err + } + if err := batch.Delete(key, nil); err != nil { return err } @@ -354,13 +411,21 @@ func grantIndexKeys(syncIDBytes []byte, r *v3.GrantRecord) [][]byte { return keys } -// writeGrantIndexes adds index entries for r to batch. +// writeGrantIndexes adds index entries for r to batch. The nil-valued +// indexes go in via grantIndexKeys; the by_entitlement_principal_hash +// index is the one grant index that carries a VALUE (the grant content +// hash), so it is written separately. func (e *Engine) writeGrantIndexes(batch *pebble.Batch, syncIDBytes []byte, r *v3.GrantRecord) error { for _, k := range grantIndexKeys(syncIDBytes, r) { if err := batch.Set(k, nil, nil); err != nil { return err } } + if hk := grantHashIndexKey(syncIDBytes, r); hk != nil { + if err := batch.Set(hk, grantContentHash(r), nil); err != nil { + return err + } + } return nil } diff --git a/pkg/dotc1z/engine/pebble/if_newer.go b/pkg/dotc1z/engine/pebble/if_newer.go index e7ff49b18..621920f50 100644 --- a/pkg/dotc1z/engine/pebble/if_newer.go +++ b/pkg/dotc1z/engine/pebble/if_newer.go @@ -37,6 +37,10 @@ func (e *Engine) PutGrantRecordsIfNewer(ctx context.Context, records ...*v3.Gran return e.withWrite(func() error { batch := e.db.NewBatch() defer batch.Close() + // IfNewer is by definition a post-seal mutation path, so the + // grant digests (cloned along with the sync's keyspace) are + // kept in step incrementally. + mm := newGrantDigestMutator(e) idBytes, err := e.resolveSyncBytes("") if err != nil { return err @@ -63,6 +67,18 @@ func (e *Engine) PutGrantRecordsIfNewer(ctx context.Context, records ...*v3.Gran closer.Close() return err } + // The digest removal needs the old record's content hash, + // which folds fields the raw index scan does not extract + // (sources). + old := &v3.GrantRecord{} + if err := unmarshalRecord(oldVal, old); err != nil { + closer.Close() + return fmt.Errorf("PutGrantRecordsIfNewer: unmarshal old: %w", err) + } + if err := mm.removeGrant(idBytes, old); err != nil { + closer.Close() + return err + } closer.Close() case errors.Is(getErr, pebble.ErrNotFound): // no existing record — write unconditionally @@ -79,11 +95,17 @@ func (e *Engine) PutGrantRecordsIfNewer(ctx context.Context, records ...*v3.Gran if err := e.writeGrantIndexes(batch, idBytes, r); err != nil { return err } + if err := mm.addGrant(idBytes, r); err != nil { + return err + } written++ } if written == 0 { return nil } + if err := mm.apply(batch); err != nil { + return err + } return batch.Commit(writeOpts(e.opts.durability)) }) } diff --git a/pkg/dotc1z/engine/pebble/index_migrations.go b/pkg/dotc1z/engine/pebble/index_migrations.go index cd7098b95..9f94c09e6 100644 --- a/pkg/dotc1z/engine/pebble/index_migrations.go +++ b/pkg/dotc1z/engine/pebble/index_migrations.go @@ -8,6 +8,7 @@ import ( "github.com/cockroachdb/pebble/v2" + v3 "github.com/conductorone/baton-sdk/pb/c1/storage/v3" "github.com/conductorone/baton-sdk/pkg/dotc1z/engine/pebble/codec" ) @@ -63,7 +64,85 @@ type indexMigration struct { // the corresponding index for any existing c1z that doesn't have // it yet. var indexMigrations = []indexMigration{ - // none yet, because we have no existing data. + { + // Backfill the by_entitlement_principal_hash index and the + // per-entitlement grant digests for files written before either + // existed. Idempotent: re-emitting an index entry is a Set over + // the same key/value, and each digest rebuild range-clears the + // partition's typeDigest keyspace before writing. + // New files persist this version at their initial (empty) Open, + // so the inline write path maintains both and the backfill never + // re-runs for them. + // + // v1: XOR combiner, 256-ary radix with all levels stored + // sparsely, count on every node (RFC 0003). Uses xxHash64 + // (8-byte digests) for both principalBucketHash and + // grantContentHash. + // v2: flat shape — root + a single leaf level of 2^width + // buckets (width in bits, chosen per entitlement), leaf keys + // carrying a 2-byte left-aligned bucket index; node keys carry + // the digested index's id after the sync_id (see + // encodeDigestNodeKey). Index entries are unchanged; the bump + // forces a digest rebuild, whose per-partition range-clear + // removes the v1 nodes. + Name: "grant_by_entitlement_principal_hash", + Version: 2, + Apply: func(ctx context.Context, e *Engine) error { + return e.backfillGrantHashIndexAndDigests(ctx) + }, + }, +} + +// backfillGrantHashIndexAndDigests reconstructs the +// by_entitlement_principal_hash index for every grant in every sync, +// then rebuilds the per-entitlement grant digests. The index must be +// committed before the digests are built because BuildAllGrantDigests +// folds over the committed index. +func (e *Engine) backfillGrantHashIndexAndDigests(ctx context.Context) error { + var syncIDs []string + if err := e.IterateAllSyncRuns(ctx, func(r *v3.SyncRunRecord) bool { + syncIDs = append(syncIDs, r.GetSyncId()) + return true + }); err != nil { + return fmt.Errorf("backfill hash index: list syncs: %w", err) + } + + for _, syncID := range syncIDs { + if err := ctx.Err(); err != nil { + return err + } + idBytes, err := codec.EncodeSyncID(syncID) + if err != nil { + return err + } + // Re-emit the hash index entry for each grant in this sync. + batch := e.db.NewBatch() + err = e.IterateGrantsBySync(ctx, syncID, func(r *v3.GrantRecord) bool { + hk := grantHashIndexKey(idBytes, r) + if hk == nil { + return true + } + if setErr := batch.Set(hk, grantContentHash(r), nil); setErr != nil { + err = setErr + return false + } + return true + }) + if err != nil { + batch.Close() + return fmt.Errorf("backfill hash index: sync %q: %w", syncID, err) + } + if err := batch.Commit(pebble.Sync); err != nil { + batch.Close() + return fmt.Errorf("backfill hash index: commit %q: %w", syncID, err) + } + batch.Close() + + if err := e.BuildAllGrantDigests(ctx, syncID); err != nil { + return fmt.Errorf("backfill grant digests: sync %q: %w", syncID, err) + } + } + return nil } // applyIndexMigrations runs on engine Open (writable opens only — diff --git a/pkg/dotc1z/engine/pebble/keys.go b/pkg/dotc1z/engine/pebble/keys.go index 4398a576c..341c53610 100644 --- a/pkg/dotc1z/engine/pebble/keys.go +++ b/pkg/dotc1z/engine/pebble/keys.go @@ -61,6 +61,7 @@ const ( typeIndex byte = 0x07 typeCounter byte = 0x08 typeSession byte = 0x09 + typeDigest byte = 0x0A typeEngineMeta byte = 0xFF ) @@ -76,6 +77,12 @@ const ( idxGrantByNeedsExpansion byte = 0x05 idxGrantByPrincipalResourceType byte = 0x06 idxGrantByEntitlementResource byte = 0x07 + // idxGrantByEntitlementPrincipalHash sorts grants by + // (entitlement_id, hash(principal)). Unlike every other grant + // index its entries carry a VALUE (the grant content hash). It is + // the substrate the per-entitlement grant digest (typeDigest) + // folds over; see digest.go and grant_digest.go. + idxGrantByEntitlementPrincipalHash byte = 0x08 ) // --- Grant --- @@ -275,6 +282,145 @@ func encodeGrantByEntitlementResourcePrefix(syncIDBytes []byte, entRT, entRID st return codec.AppendTupleSeparator(buf) } +// --- Grant by (entitlement, principal-hash) + digest nodes --- + +// encodeGrantByEntPrincHashIndexKey is the by_entitlement_principal_hash +// secondary index on GrantRecord. Unlike every other grant index it +// interposes a RAW, fixed-width principal-bucket hash between the +// entitlement_id and the principal tuple, so the keyspace sorts by +// (entitlement_id, hash(principal)). That hash-major order is what the +// per-entitlement grant digest folds over, and a bucket — a bit-range +// of the hash — is a contiguous key range, which is the property the +// digest bucket range scans rely on (see digestIndexSpec.bucketBounds). +// +// v3 | typeIndex | idxGrantByEntitlementPrincipalHash | sync_id | 0x00 | +// entitlement_id | 0x00 | | +// principal_rt | 0x00 | principal_id | 0x00 | external_id +// -> value: grant content hash (xxHash64, 8 bytes) +// +// Because the bucket hash is raw it can contain 0x00, so the generic +// tuple walkers (lastTupleComponent / decodeTwoTupleComponents) must NOT +// be pointed at a prefix that stops before the hash — their walk would +// derail on those bytes. Use decodeEntPrincHashTail, which accounts for +// the hash's fixed width positionally. +// +// Paired with encodeGrantByEntPrincHashEntPrefix (by-value prefix, with +// trailing sep) and digestIndexSpec.bucketBounds (digest.go). +func encodeGrantByEntPrincHashIndexKey(syncIDBytes []byte, entitlementID string, bucketHash []byte, principalRT, principalID, externalID string) []byte { + buf := make([]byte, 0, 8+len(syncIDBytes)+len(entitlementID)+len(bucketHash)+len(principalRT)+len(principalID)+len(externalID)) + buf = append(buf, versionV3, typeIndex, idxGrantByEntitlementPrincipalHash) + buf = append(buf, syncIDBytes...) + buf = codec.AppendTupleSeparator(buf) + buf = codec.AppendTupleStrings(buf, entitlementID) + buf = codec.AppendTupleSeparator(buf) + buf = append(buf, bucketHash...) + return codec.AppendTupleStrings(buf, principalRT, principalID, externalID) +} + +// encodeGrantByEntPrincHashEntPrefix is the by-value prefix for "all +// grants in this sync under this entitlement", in hash order. The +// trailing separator is load-bearing (see the keys.go convention doc); +// the raw bucket hash follows it. Its output length is also the offset +// the decoder uses to locate the raw hash region. +func encodeGrantByEntPrincHashEntPrefix(syncIDBytes []byte, entitlementID string) []byte { + buf := make([]byte, 0, 8+len(syncIDBytes)+len(entitlementID)) + buf = append(buf, versionV3, typeIndex, idxGrantByEntitlementPrincipalHash) + buf = append(buf, syncIDBytes...) + buf = codec.AppendTupleSeparator(buf) + buf = codec.AppendTupleStrings(buf, entitlementID) + return codec.AppendTupleSeparator(buf) +} + +// decodeEntPrincHashTail decodes one index key relative to entPrefix +// (which must be encodeGrantByEntPrincHashEntPrefix output — i.e. it +// ends right before the raw bucket hash). Returns the raw bucket hash +// (a sub-slice of key, valid only as long as key is), the principal +// rt/id and the external_id. ok is false if the key is shorter than +// prefix+hash or the tuple tail is malformed. +func decodeEntPrincHashTail(key, entPrefix []byte) ([]byte, string, string, string, bool) { + if len(key) < len(entPrefix)+digestBucketHashLen { + return nil, "", "", "", false + } + bucketHash := key[len(entPrefix) : len(entPrefix)+digestBucketHashLen] + tail := key[len(entPrefix)+digestBucketHashLen:] + rt, next, err := codec.DecodeTupleStringTo(nil, tail, 0) + if err != nil || next >= len(tail) { + return nil, "", "", "", false + } + id, next2, err := codec.DecodeTupleStringTo(nil, tail, next+1) + if err != nil || next2 >= len(tail) { + return nil, "", "", "", false + } + ext, _, err := codec.DecodeTupleStringTo(nil, tail, next2+1) + if err != nil { + return nil, "", "", "", false + } + return bucketHash, string(rt), string(id), string(ext), true +} + +// GrantByEntPrincHashSyncLowerBound / UpperBound bound the entire +// by_entitlement_principal_hash index for a sync. Exported for the +// cleanup/clone/compaction keyspace plans. +func GrantByEntPrincHashSyncLowerBound(syncIDBytes []byte) []byte { + buf := make([]byte, 0, 3+len(syncIDBytes)) + buf = append(buf, versionV3, typeIndex, idxGrantByEntitlementPrincipalHash) + return append(buf, syncIDBytes...) +} + +func GrantByEntPrincHashSyncUpperBound(syncIDBytes []byte) []byte { + return upperBoundOf(GrantByEntPrincHashSyncLowerBound(syncIDBytes)) +} + +// Digest node keys. +// +// v3 | typeDigest | sync_id | index_id(1 byte) | 0x00 | partition | 0x00 | level(1 byte) | bucket_prefix +// +// index_id discriminates WHICH digested index the node belongs to (the +// digested index's own idx* byte, see digestIndexSpec) — it sits right +// after the fixed-width sync_id so one sync's digests for all indexes +// are still a single contiguous range for the cleanup/clone plans. +// level 0 is the root (bucket_prefix empty); level 1 is the single leaf +// level, one node per non-empty bucket, whose bucket_prefix is the +// bucket index LEFT-ALIGNED in 2 raw bytes (digestLeafPrefixLen). The +// left alignment makes leaf keys sort in bucket-hash order at every +// digest width, so the comparison's fold-to-coarser-width merge is a +// single contiguous scan of this range. See digest.go for the node +// value framing. +func encodeDigestNodeKey(syncIDBytes []byte, indexID byte, partition string, level byte, bucketPrefix []byte) []byte { + buf := encodeDigestPartitionPrefix(syncIDBytes, indexID, partition) + buf = append(buf, level) + return append(buf, bucketPrefix...) +} + +// encodeDigestPartitionPrefix is the prefix of every digest node key +// for one (index, partition) — the range a rebuild clears before +// writing (the build only Sets nodes; without the leading DeleteRange a +// width change or an emptied bucket would leave stale nodes for the +// comparison merge scan to read) and the range digestMutator drops on +// detecting an inconsistent digest. +func encodeDigestPartitionPrefix(syncIDBytes []byte, indexID byte, partition string) []byte { + buf := make([]byte, 0, 7+len(syncIDBytes)+len(partition)) + buf = append(buf, versionV3, typeDigest) + buf = append(buf, syncIDBytes...) + buf = append(buf, indexID) + buf = codec.AppendTupleSeparator(buf) + buf = codec.AppendTupleStrings(buf, partition) + return codec.AppendTupleSeparator(buf) +} + +// DigestSyncLowerBound / UpperBound bound the entire digest keyspace +// (all digested indexes) for a sync. Exported for the +// cleanup/clone/compaction keyspace plans. +func DigestSyncLowerBound(syncIDBytes []byte) []byte { + buf := make([]byte, 0, 2+len(syncIDBytes)) + buf = append(buf, versionV3, typeDigest) + return append(buf, syncIDBytes...) +} + +func DigestSyncUpperBound(syncIDBytes []byte) []byte { + return upperBoundOf(DigestSyncLowerBound(syncIDBytes)) +} + // --- ResourceType --- // encodeResourceTypeKey returns the primary key for a resource_type: diff --git a/pkg/dotc1z/engine/pebble/raw_records.go b/pkg/dotc1z/engine/pebble/raw_records.go index a8d2eae32..1895d4000 100644 --- a/pkg/dotc1z/engine/pebble/raw_records.go +++ b/pkg/dotc1z/engine/pebble/raw_records.go @@ -154,6 +154,20 @@ func (e *Engine) deleteGrantIndexesRaw(batch *pebble.Batch, syncIDBytes []byte, return err } } + // by_entitlement_principal_hash: the bucket hash is derived from the + // principal identity, so the raw path reconstructs the same key + // writeGrantIndexes wrote, mirroring grantHashIndexKey's nil-guard. + // The entitlement's grant digest is kept in step separately: callers + // on the post-seal mutation paths feed the old record to + // digestMutator.removeGrant (and the new one to .addGrant), which + // folds the change into the stored nodes in the same batch. + if entID != "" && principalRT != "" && principalID != "" { + bh := principalBucketHash(principalRT, principalID) + hk := encodeGrantByEntPrincHashIndexKey(syncIDBytes, entID, bh, principalRT, principalID, externalID) + if err := batch.Delete(hk, nil); err != nil { + return err + } + } return batch.Delete(encodeGrantByNeedsExpansionIndexKey(syncIDBytes, externalID), nil) } diff --git a/pkg/synccompactor/pebble/bucket_plans.go b/pkg/synccompactor/pebble/bucket_plans.go index 0f5c0e239..204bbfe17 100644 --- a/pkg/synccompactor/pebble/bucket_plans.go +++ b/pkg/synccompactor/pebble/bucket_plans.go @@ -76,6 +76,21 @@ func buildBucketPlans(syncIDBytes []byte) []bucketPlan { lower: enginepkg.GrantByNeedsExpansionSyncLowerBound(syncIDBytes), upper: enginepkg.GrantByNeedsExpansionSyncUpperBound(syncIDBytes), }, + { + name: "grant_by_entitlement_principal_hash", + lower: enginepkg.GrantByEntPrincHashSyncLowerBound(syncIDBytes), + upper: enginepkg.GrantByEntPrincHashSyncUpperBound(syncIDBytes), + }, + { + // Digest nodes (all digested indexes, e.g. the per-entitlement + // grant digest). Copied byte-for-byte with the sync_id they + // belong to; because Compact preserves the sync_id and the + // records move verbatim, the digests stay valid in the + // destination without a rebuild. + name: "digest", + lower: enginepkg.DigestSyncLowerBound(syncIDBytes), + upper: enginepkg.DigestSyncUpperBound(syncIDBytes), + }, { name: "asset", lower: enginepkg.AssetSyncLowerBound(syncIDBytes),