From 956afb686749add7b62507a5d7f7fbf8556df078 Mon Sep 17 00:00:00 2001 From: Naganathan M R Date: Sun, 19 Jul 2026 18:14:57 +0530 Subject: [PATCH 1/9] fix: handle concurrent writes to memdb | add test to simulate the race condition --- internal/datastore/memdb/memdb.go | 2 +- internal/datastore/memdb/memdb_test.go | 90 ++++++++++++++++++++++++++ internal/datastore/memdb/revisions.go | 18 ++++++ 3 files changed, 109 insertions(+), 1 deletion(-) diff --git a/internal/datastore/memdb/memdb.go b/internal/datastore/memdb/memdb.go index a5fb74e8c5..4c0143cc2d 100644 --- a/internal/datastore/memdb/memdb.go +++ b/internal/datastore/memdb/memdb.go @@ -374,7 +374,7 @@ func (mdb *memdbDatastore) ReadWriteTx( // Create a snapshot and add it to the revisions slice schemaHash := mdb.getCurrentSchemaHashNoLock() snap := mdb.db.Snapshot() - mdb.revisions = append(mdb.revisions, snapshot{newRevision, schemaHash, snap}) + mdb.insertRevisionSnapshot(newRevision, schemaHash, snap) return newRevision, nil } diff --git a/internal/datastore/memdb/memdb_test.go b/internal/datastore/memdb/memdb_test.go index b16298dc55..156a2dee39 100644 --- a/internal/datastore/memdb/memdb_test.go +++ b/internal/datastore/memdb/memdb_test.go @@ -113,6 +113,96 @@ func TestConcurrentWriteRelsSucceed(t *testing.T) { require.NoError(g.Wait()) } +// TestConcurrentWriteRevisionInversionLosesRead deterministically reproduces +// https://github.com/authzed/spicedb/issues/3212: newRevisionID is stamped +// before the write-serialization lock is acquired, so a transaction (B) that +// starts after another (A) can still commit before it, even though B's +// revision number is numerically greater than A's. Recording snapshots in +// commit order rather than revision order breaks the sortedness that +// SnapshotReader's binary search relies on, and a fully-consistent read at +// A's own acknowledged revision ends up resolving to B's earlier snapshot, +// which does not contain A's write. +func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { + require := require.New(t) + + ds, err := NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour) + require.NoError(err) + mds := ds.(*memdbDatastore) + + ctx := t.Context() + + aEntered := make(chan struct{}) + releaseA := make(chan struct{}) + aDone := make(chan struct{}) + + var revA datastore.Revision + var errA error + go func() { + defer close(aDone) + revA, errA = ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { + // newRevisionID has already been called for this transaction by the + // time f runs, so closing this channel signals that A's revision + // number has been stamped, even though A hasn't touched the + // write-transaction lock yet. + close(aEntered) + <-releaseA + return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ + tuple.Touch(tuple.MustParse("document:doc-a#viewer@user:tom")), + }) + }, options.WithDisableRetries(true)) + }() + + <-aEntered + + revB, err := ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { + return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ + tuple.Touch(tuple.MustParse("document:doc-b#viewer@user:tom")), + }) + }, options.WithDisableRetries(true)) + require.NoError(err) + + close(releaseA) + <-aDone + require.NoError(errA) + + // Sanity check on the harness itself: B was assigned its revision strictly + // after A signaled entry, so it must be the numerically greater one. This + // holds regardless of whether the storage-ordering bug is present. + require.True(revB.GreaterThan(revA), "expected B's revision (%v) to be greater than A's (%v)", revB, revA) + + // Diagnostic only (not an invariant the fix must preserve either way): + // log where each ended up in storage order, for visibility into whether + // commit order matched revision order on this run. + indexOf := func(r datastore.Revision) int { + mds.RLock() + defer mds.RUnlock() + for i, snap := range mds.revisions { + if snap.revision.Equal(r) { + return i + } + } + return -1 + } + t.Logf("storage order: A(rev=%v) at index %d, B(rev=%v) at index %d", revA, indexOf(revA), revB, indexOf(revB)) + + // The actual bug/fix boundary: a fully-consistent read at A's own + // acknowledged revision must see A's write. Before the fix, this could + // resolve to B's earlier snapshot instead, since B committed first. + reader := ds.SnapshotReader(revA) + it, err := reader.QueryRelationships(ctx, datastore.RelationshipsFilter{OptionalResourceType: "document"}) + require.NoError(err) + rels, err := datastore.IteratorToSlice(it) + require.NoError(err) + + var sawA bool + for _, rel := range rels { + if rel.Resource.ObjectID == "doc-a" { + sawA = true + } + } + require.True(sawA, "BUG REPRODUCED: relationship written by transaction A is invisible when reading at A's own returned revision %v", revA) +} + func TestAnythingAfterCloseDoesNotPanic(t *testing.T) { require := require.New(t) diff --git a/internal/datastore/memdb/revisions.go b/internal/datastore/memdb/revisions.go index 39aae92414..8564a49cdc 100644 --- a/internal/datastore/memdb/revisions.go +++ b/internal/datastore/memdb/revisions.go @@ -2,8 +2,12 @@ package memdb import ( "context" + "slices" + "sort" "time" + "github.com/hashicorp/go-memdb" + "github.com/authzed/spicedb/internal/datastore/revisions" "github.com/authzed/spicedb/pkg/datastore" ) @@ -38,6 +42,20 @@ func (mdb *memdbDatastore) newRevisionID() revisions.TimestampRevision { return created } +// insertRevisionSnapshot records a newly committed snapshot in mdb.revisions, +// maintaining ascending revision order. Revision IDs are stamped when a +// transaction starts (see newRevisionID), before it acquires the write lock +// that serializes commits, so a transaction with a numerically later revision +// can still commit -- and need to be recorded -- before one with an earlier +// revision. A plain append would leave mdb.revisions unsorted in that case, +// which breaks the binary search in SnapshotReader. +func (mdb *memdbDatastore) insertRevisionSnapshot(newRevision revisions.TimestampRevision, schemaHash string, snap *memdb.MemDB) { + insertAt := sort.Search(len(mdb.revisions), func(i int) bool { + return mdb.revisions[i].revision.GreaterThan(newRevision) + }) + mdb.revisions = slices.Insert(mdb.revisions, insertAt, snapshot{newRevision, schemaHash, snap}) +} + func (mdb *memdbDatastore) HeadRevision(_ context.Context) (datastore.RevisionWithSchemaHash, error) { mdb.RLock() defer mdb.RUnlock() From e7099920e2509e6aed9faf9bd49ec1542c64898b Mon Sep 17 00:00:00 2001 From: Naganathan M R Date: Sun, 19 Jul 2026 18:18:33 +0530 Subject: [PATCH 2/9] chore: remove verbose comments --- internal/datastore/memdb/memdb_test.go | 23 +---------------------- internal/datastore/memdb/revisions.go | 7 ------- 2 files changed, 1 insertion(+), 29 deletions(-) diff --git a/internal/datastore/memdb/memdb_test.go b/internal/datastore/memdb/memdb_test.go index 156a2dee39..12962237ff 100644 --- a/internal/datastore/memdb/memdb_test.go +++ b/internal/datastore/memdb/memdb_test.go @@ -113,15 +113,7 @@ func TestConcurrentWriteRelsSucceed(t *testing.T) { require.NoError(g.Wait()) } -// TestConcurrentWriteRevisionInversionLosesRead deterministically reproduces -// https://github.com/authzed/spicedb/issues/3212: newRevisionID is stamped -// before the write-serialization lock is acquired, so a transaction (B) that -// starts after another (A) can still commit before it, even though B's -// revision number is numerically greater than A's. Recording snapshots in -// commit order rather than revision order breaks the sortedness that -// SnapshotReader's binary search relies on, and a fully-consistent read at -// A's own acknowledged revision ends up resolving to B's earlier snapshot, -// which does not contain A's write. +// TestConcurrentWriteRevisionInversionLosesRead reproduces concurrent writes to MemDB func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { require := require.New(t) @@ -140,10 +132,6 @@ func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { go func() { defer close(aDone) revA, errA = ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { - // newRevisionID has already been called for this transaction by the - // time f runs, so closing this channel signals that A's revision - // number has been stamped, even though A hasn't touched the - // write-transaction lock yet. close(aEntered) <-releaseA return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ @@ -165,14 +153,8 @@ func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { <-aDone require.NoError(errA) - // Sanity check on the harness itself: B was assigned its revision strictly - // after A signaled entry, so it must be the numerically greater one. This - // holds regardless of whether the storage-ordering bug is present. require.True(revB.GreaterThan(revA), "expected B's revision (%v) to be greater than A's (%v)", revB, revA) - // Diagnostic only (not an invariant the fix must preserve either way): - // log where each ended up in storage order, for visibility into whether - // commit order matched revision order on this run. indexOf := func(r datastore.Revision) int { mds.RLock() defer mds.RUnlock() @@ -185,9 +167,6 @@ func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { } t.Logf("storage order: A(rev=%v) at index %d, B(rev=%v) at index %d", revA, indexOf(revA), revB, indexOf(revB)) - // The actual bug/fix boundary: a fully-consistent read at A's own - // acknowledged revision must see A's write. Before the fix, this could - // resolve to B's earlier snapshot instead, since B committed first. reader := ds.SnapshotReader(revA) it, err := reader.QueryRelationships(ctx, datastore.RelationshipsFilter{OptionalResourceType: "document"}) require.NoError(err) diff --git a/internal/datastore/memdb/revisions.go b/internal/datastore/memdb/revisions.go index 8564a49cdc..10d3403962 100644 --- a/internal/datastore/memdb/revisions.go +++ b/internal/datastore/memdb/revisions.go @@ -42,13 +42,6 @@ func (mdb *memdbDatastore) newRevisionID() revisions.TimestampRevision { return created } -// insertRevisionSnapshot records a newly committed snapshot in mdb.revisions, -// maintaining ascending revision order. Revision IDs are stamped when a -// transaction starts (see newRevisionID), before it acquires the write lock -// that serializes commits, so a transaction with a numerically later revision -// can still commit -- and need to be recorded -- before one with an earlier -// revision. A plain append would leave mdb.revisions unsorted in that case, -// which breaks the binary search in SnapshotReader. func (mdb *memdbDatastore) insertRevisionSnapshot(newRevision revisions.TimestampRevision, schemaHash string, snap *memdb.MemDB) { insertAt := sort.Search(len(mdb.revisions), func(i int) bool { return mdb.revisions[i].revision.GreaterThan(newRevision) From f1e543f3a63dac7dde226d5bed9ce54ed1ca3a92 Mon Sep 17 00:00:00 2001 From: Naganathan M R Date: Sun, 19 Jul 2026 19:11:17 +0530 Subject: [PATCH 3/9] chore: remove verbose comments --- internal/datastore/memdb/memdb_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/internal/datastore/memdb/memdb_test.go b/internal/datastore/memdb/memdb_test.go index 12962237ff..5f13d3d8a2 100644 --- a/internal/datastore/memdb/memdb_test.go +++ b/internal/datastore/memdb/memdb_test.go @@ -179,7 +179,7 @@ func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { sawA = true } } - require.True(sawA, "BUG REPRODUCED: relationship written by transaction A is invisible when reading at A's own returned revision %v", revA) + require.True(sawA, "relationship written by transaction A is invisible when reading at A's own returned revision %v", revA) } func TestAnythingAfterCloseDoesNotPanic(t *testing.T) { From e2bb728c4f4f764543cb968995f4b9ab40f86642 Mon Sep 17 00:00:00 2001 From: Naganathan M R Date: Tue, 21 Jul 2026 18:59:33 +0530 Subject: [PATCH 4/9] fix: address review comments | add docstrings and comments --- internal/datastore/memdb/memdb_test.go | 32 +++++++++++--------------- internal/datastore/memdb/revisions.go | 15 ++++++++++++ 2 files changed, 29 insertions(+), 18 deletions(-) diff --git a/internal/datastore/memdb/memdb_test.go b/internal/datastore/memdb/memdb_test.go index 5f13d3d8a2..8b473df825 100644 --- a/internal/datastore/memdb/memdb_test.go +++ b/internal/datastore/memdb/memdb_test.go @@ -114,6 +114,11 @@ func TestConcurrentWriteRelsSucceed(t *testing.T) { } // TestConcurrentWriteRevisionInversionLosesRead reproduces concurrent writes to MemDB +// where the test simulates a scenario: two goroutines A and B begin execution +// concurrently. A gets an earlier revision number than B, but the test holds A +// until B has fully committed, forcing B's snapshot to be recorded first despite +// having the later revision. It then asserts that reading at A's own revision +// still returns A's write. func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { require := require.New(t) @@ -123,9 +128,9 @@ func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { ctx := t.Context() - aEntered := make(chan struct{}) - releaseA := make(chan struct{}) - aDone := make(chan struct{}) + aEntered := make(chan struct{}) // A's transaction has started running. + releaseA := make(chan struct{}) // B has completed its execution + aDone := make(chan struct{}) // A's ReadWriteTx call has fully returned, so it's safe to read revision var revA datastore.Revision var errA error @@ -133,15 +138,16 @@ func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { defer close(aDone) revA, errA = ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { close(aEntered) - <-releaseA + <-releaseA // held until B has fully committed below return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ tuple.Touch(tuple.MustParse("document:doc-a#viewer@user:tom")), }) }, options.WithDisableRetries(true)) }() - <-aEntered + <-aEntered // A's revision number is assigned. A is now paused + // B is assigned a later revision number and runs to completion while A is still paused. revB, err := ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ tuple.Touch(tuple.MustParse("document:doc-b#viewer@user:tom")), @@ -149,23 +155,13 @@ func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { }, options.WithDisableRetries(true)) require.NoError(err) - close(releaseA) - <-aDone + close(releaseA) // mark the completion of B so that A can commit its revision + <-aDone // Block until A completes its execution so that revision can be read require.NoError(errA) require.True(revB.GreaterThan(revA), "expected B's revision (%v) to be greater than A's (%v)", revB, revA) - indexOf := func(r datastore.Revision) int { - mds.RLock() - defer mds.RUnlock() - for i, snap := range mds.revisions { - if snap.revision.Equal(r) { - return i - } - } - return -1 - } - t.Logf("storage order: A(rev=%v) at index %d, B(rev=%v) at index %d", revA, indexOf(revA), revB, indexOf(revB)) + t.Logf("storage order: A(rev=%v) at index %d, B(rev=%v) at index %d", revA, mds.indexOfRevision(revA), revB, mds.indexOfRevision(revB)) reader := ds.SnapshotReader(revA) it, err := reader.QueryRelationships(ctx, datastore.RelationshipsFilter{OptionalResourceType: "document"}) diff --git a/internal/datastore/memdb/revisions.go b/internal/datastore/memdb/revisions.go index 10d3403962..3c2a0db48b 100644 --- a/internal/datastore/memdb/revisions.go +++ b/internal/datastore/memdb/revisions.go @@ -42,6 +42,7 @@ func (mdb *memdbDatastore) newRevisionID() revisions.TimestampRevision { return created } +// insertRevisionSnapshot records a newly committed snapshot at its sorted position func (mdb *memdbDatastore) insertRevisionSnapshot(newRevision revisions.TimestampRevision, schemaHash string, snap *memdb.MemDB) { insertAt := sort.Search(len(mdb.revisions), func(i int) bool { return mdb.revisions[i].revision.GreaterThan(newRevision) @@ -49,6 +50,20 @@ func (mdb *memdbDatastore) insertRevisionSnapshot(newRevision revisions.Timestam mdb.revisions = slices.Insert(mdb.revisions, insertAt, snapshot{newRevision, schemaHash, snap}) } +// indexOfRevision returns the index of the snapshot recorded for the exact given revision +func (mdb *memdbDatastore) indexOfRevision(r datastore.Revision) int { + mdb.RLock() + defer mdb.RUnlock() + + i := sort.Search(len(mdb.revisions), func(i int) bool { + return !mdb.revisions[i].revision.LessThan(r) + }) + if i < len(mdb.revisions) && mdb.revisions[i].revision.Equal(r) { + return i + } + return -1 +} + func (mdb *memdbDatastore) HeadRevision(_ context.Context) (datastore.RevisionWithSchemaHash, error) { mdb.RLock() defer mdb.RUnlock() From b3689456a1e87e4c13f415bc08ccdcbb31aa4d40 Mon Sep 17 00:00:00 2001 From: Naganathan M R Date: Tue, 21 Jul 2026 19:23:11 +0530 Subject: [PATCH 5/9] chore: add comments --- internal/datastore/memdb/revisions.go | 2 ++ 1 file changed, 2 insertions(+) diff --git a/internal/datastore/memdb/revisions.go b/internal/datastore/memdb/revisions.go index 3c2a0db48b..2f6d67c512 100644 --- a/internal/datastore/memdb/revisions.go +++ b/internal/datastore/memdb/revisions.go @@ -43,6 +43,7 @@ func (mdb *memdbDatastore) newRevisionID() revisions.TimestampRevision { } // insertRevisionSnapshot records a newly committed snapshot at its sorted position +// The caller must already hold mdb's write lock. func (mdb *memdbDatastore) insertRevisionSnapshot(newRevision revisions.TimestampRevision, schemaHash string, snap *memdb.MemDB) { insertAt := sort.Search(len(mdb.revisions), func(i int) bool { return mdb.revisions[i].revision.GreaterThan(newRevision) @@ -51,6 +52,7 @@ func (mdb *memdbDatastore) insertRevisionSnapshot(newRevision revisions.Timestam } // indexOfRevision returns the index of the snapshot recorded for the exact given revision +// Performs binary search for finding the revision, returns -1 if none exists func (mdb *memdbDatastore) indexOfRevision(r datastore.Revision) int { mdb.RLock() defer mdb.RUnlock() From 2ad2859364bc97191901e794f97451b64bca34e1 Mon Sep 17 00:00:00 2001 From: Maria Ines Parnisari Date: Tue, 21 Jul 2026 12:06:39 -0700 Subject: [PATCH 6/9] test: rewrite test --- internal/datastore/memdb/memdb_test.go | 153 ++++++++++++++++--------- 1 file changed, 102 insertions(+), 51 deletions(-) diff --git a/internal/datastore/memdb/memdb_test.go b/internal/datastore/memdb/memdb_test.go index 8b473df825..4adb3977d8 100644 --- a/internal/datastore/memdb/memdb_test.go +++ b/internal/datastore/memdb/memdb_test.go @@ -113,69 +113,120 @@ func TestConcurrentWriteRelsSucceed(t *testing.T) { require.NoError(g.Wait()) } -// TestConcurrentWriteRevisionInversionLosesRead reproduces concurrent writes to MemDB -// where the test simulates a scenario: two goroutines A and B begin execution -// concurrently. A gets an earlier revision number than B, but the test holds A -// until B has fully committed, forcing B's snapshot to be recorded first despite -// having the later revision. It then asserts that reading at A's own revision -// still returns A's write. -func TestConcurrentWriteRevisionInversionLosesRead(t *testing.T) { - require := require.New(t) - +// TestConcurrentWrite covers revision visibility when two write transactions overlap. +// - each transaction's write should be visible when reading at its own returned revision (read-your-writes) +// - a write should NOT be visible at revisions below its own (consistent snapshot) +// - a write SHOULD be visible at every revision at or above its own +// - the head revision sees every committed write +// TODO: should this be a datastore conformance test (not specific to memdb)? +func TestConcurrentWrite(t *testing.T) { ds, err := NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour) - require.NoError(err) - mds := ds.(*memdbDatastore) + require.NoError(t, err) + t.Cleanup(func() { + require.NoError(t, ds.Close()) + }) + + // docsAt returns the set of document object IDs visible at the given revision. + docsAt := func(rev datastore.Revision) map[string]bool { + reader := ds.SnapshotReader(rev) + it, err := reader.QueryRelationships(t.Context(), datastore.RelationshipsFilter{OptionalResourceType: "document"}) + require.NoError(t, err) + rels, err := datastore.IteratorToSlice(it) + require.NoError(t, err) + + seen := map[string]bool{} + for _, rel := range rels { + seen[rel.Resource.ObjectID] = true + } + return seen + } - ctx := t.Context() + // Execute many iterations so that one run is enough to expose ordering problems + const iterations = 10_000 + for i := 0; i < iterations; i++ { + // Unique names per iteration so that visibility assertions cannot be + // satisfied by a previous iteration's writes. + docA := fmt.Sprintf("doc-a-%d", i) + docB := fmt.Sprintf("doc-b-%d", i) + relA := tuple.MustParse(fmt.Sprintf("document:%s#viewer@user:tom", docA)) + relB := tuple.MustParse(fmt.Sprintf("document:%s#viewer@user:tom", docB)) + + var ( + revA datastore.Revision + errA error + waitUntilAStarts = make(chan struct{}, 1) + waitUntilBCompletes = make(chan struct{}, 1) + waitUntilADone = make(chan struct{}) + ) + go func() { + defer close(waitUntilADone) + revA, errA = ds.ReadWriteTx(t.Context(), func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { + waitUntilAStarts <- struct{}{} + <-waitUntilBCompletes // held until B has fully committed below + return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ + tuple.Touch(relA), + }) + }, options.WithDisableRetries(true)) + }() + + <-waitUntilAStarts // A is now blocked - aEntered := make(chan struct{}) // A's transaction has started running. - releaseA := make(chan struct{}) // B has completed its execution - aDone := make(chan struct{}) // A's ReadWriteTx call has fully returned, so it's safe to read revision - - var revA datastore.Revision - var errA error - go func() { - defer close(aDone) - revA, errA = ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { - close(aEntered) - <-releaseA // held until B has fully committed below + // Transaction B runs to completion while A is blocked + revB, err := ds.ReadWriteTx(t.Context(), func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ - tuple.Touch(tuple.MustParse("document:doc-a#viewer@user:tom")), + tuple.Touch(relB), }) }, options.WithDisableRetries(true)) - }() - - <-aEntered // A's revision number is assigned. A is now paused - - // B is assigned a later revision number and runs to completion while A is still paused. - revB, err := ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { - return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ - tuple.Touch(tuple.MustParse("document:doc-b#viewer@user:tom")), - }) - }, options.WithDisableRetries(true)) - require.NoError(err) - - close(releaseA) // mark the completion of B so that A can commit its revision - <-aDone // Block until A completes its execution so that revision can be read - require.NoError(errA) + require.NoError(t, err) - require.True(revB.GreaterThan(revA), "expected B's revision (%v) to be greater than A's (%v)", revB, revA) + waitUntilBCompletes <- struct{}{} // unblocks A + <-waitUntilADone + require.NoError(t, errA) - t.Logf("storage order: A(rev=%v) at index %d, B(rev=%v) at index %d", revA, mds.indexOfRevision(revA), revB, mds.indexOfRevision(revB)) + require.False(t, revA.Equal(revB), "iteration %d: concurrent transactions must be assigned distinct revisions, both got %v", i, revA) - reader := ds.SnapshotReader(revA) - it, err := reader.QueryRelationships(ctx, datastore.RelationshipsFilter{OptionalResourceType: "document"}) - require.NoError(err) - rels, err := datastore.IteratorToSlice(it) - require.NoError(err) + // Each transaction's write must be visible at its own returned revision. + require.True(t, docsAt(revA)[docA], "iteration %d: A's write is invisible when reading at A's own returned revision %v", i, revA) + require.True(t, docsAt(revB)[docB], "iteration %d: B's write is invisible when reading at B's own returned revision %v", i, revB) - var sawA bool - for _, rel := range rels { - if rel.Resource.ObjectID == "doc-a" { - sawA = true + type write struct { + rev datastore.Revision + doc string } + earlier, later := write{revA, docA}, write{revB, docB} + if later.rev.LessThan(earlier.rev) { + earlier, later = later, earlier + } + + // The write committed at the later revision must not be visible at the earlier revision + atEarlier := docsAt(earlier.rev) + require.True(t, atEarlier[earlier.doc], "iteration %d: %s is invisible at its own revision %v", i, earlier.doc, earlier.rev) + require.False(t, atEarlier[later.doc], "iteration %d: %s was committed at later revision %v but is visible at earlier revision %v", i, later.doc, later.rev, earlier.rev) + + // Both writes must be visible at the later revision + atLater := docsAt(later.rev) + require.True(t, atLater[earlier.doc], "iteration %d: %s is visible at revision %v but disappears at later revision %v", i, earlier.doc, earlier.rev, later.rev) + require.True(t, atLater[later.doc], "iteration %d: %s is invisible at its own revision %v", i, later.doc, later.rev) + + // The head revision must see every committed write + head, err := ds.HeadRevision(t.Context()) + require.NoError(t, err) + atHead := docsAt(head.Revision) + require.True(t, atHead[docA] && atHead[docB], "iteration %d: head revision %v is missing committed writes: %v", i, head.Revision, atHead) + + // Delete this iteration's relationships so the relationship table does + // not grow across iterations: memdb scans the whole namespace per + // query, so an ever-growing table would make this test quadratic. + // Earlier revisions retain their snapshots, so this does not affect + // the assertions above. + _, err = ds.ReadWriteTx(t.Context(), func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { + return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ + tuple.Delete(relA), + tuple.Delete(relB), + }) + }, options.WithDisableRetries(true)) + require.NoError(t, err) } - require.True(sawA, "relationship written by transaction A is invisible when reading at A's own returned revision %v", revA) } func TestAnythingAfterCloseDoesNotPanic(t *testing.T) { From 18c9cabbed8547130ea97d682eedf25432117b6e Mon Sep 17 00:00:00 2001 From: Maria Ines Parnisari Date: Tue, 21 Jul 2026 12:46:39 -0700 Subject: [PATCH 7/9] fix: rework --- CHANGELOG.md | 1 + internal/datastore/memdb/memdb.go | 228 +++++++++++++------------ internal/datastore/memdb/memdb_test.go | 1 + internal/datastore/memdb/revisions.go | 45 +---- 4 files changed, 132 insertions(+), 143 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 942806ac1e..5acbb2cbed 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). - Use `testcontainers` instead of `ory/dockertest` for running containers in integration tests (https://github.com/authzed/spicedb/pull/2782) ### Fixed +- MemDB: overlapping write transactions could violate every snapshot-consistency invariant of the datastore — a committed write could be invisible at its own returned revision (breaking read-your-writes, e.g. an at-exact-snapshot `Check` right after `WriteRelationships`), a later commit could leak into reads at an earlier revision, a write visible at one revision could be missing at a later one (including head), and two concurrent transactions could be assigned the same revision. Revisions are now assigned at write-transaction acquisition, where writer serialization makes the two orders identical. (https://github.com/authzed/spicedb/pull/3239) - Fixed a nil pointer dereference panic in `CheckBulkPermissions` that could occur under concurrent load when a tracing-enabled check shared a singleflight dispatch with a non-tracing bulk check. Debug-enabled checks are no longer singleflighted together with non-debug checks. (https://github.com/authzed/spicedb/pull/3174) - Fixed a nil pointer dereference panic in the Postgres FDW (https://github.com/authzed/spicedb/pull/3235) - CockroachDB: deletes performed by CockroachDB's row-level TTL job for expired relationships are no longer emitted as `DELETE` events by the Watch API. On CockroachDB ≥ 24.1, SpiceDB sets the `ttl_disable_changefeed_replication` storage parameter on the relationship tables at startup (if it lacks `ALTER TABLE` privileges, it logs a warning with the statement to run manually); on older versions a startup warning is logged and TTL deletes continue to be emitted. Note that the parameter affects any changefeed over these tables — external changefeeds that want TTL deletes can opt back in with `ignore_disable_changefeed_replication`. Delete-only transactions also no longer write an internal transaction-metadata marker row, reducing write amplification. (https://github.com/authzed/spicedb/pull/3210) diff --git a/internal/datastore/memdb/memdb.go b/internal/datastore/memdb/memdb.go index 4c0143cc2d..a7b4746faf 100644 --- a/internal/datastore/memdb/memdb.go +++ b/internal/datastore/memdb/memdb.go @@ -90,10 +90,11 @@ type memdbDatastore struct { revisions.CommonDecoder // NOTE: call checkNotClosed before using - db *memdb.MemDB // GUARDED_BY(RWMutex) - revisions []snapshot // GUARDED_BY(RWMutex) - activeWriteTxn *memdb.Txn // GUARDED_BY(RWMutex) - writeTxReady *sync.Cond // broadcast when activeWriteTxn becomes nil + db *memdb.MemDB // GUARDED_BY(RWMutex) + // revisions MUST be sorted in strictly increasing revision order. + revisions []snapshot // GUARDED_BY(RWMutex) + activeWriteTxn *memdb.Txn // GUARDED_BY(RWMutex) + writeTxReady *sync.Cond // broadcast when activeWriteTxn becomes nil negativeGCWindow int64 quantizationPeriod int64 @@ -116,6 +117,9 @@ func (mdb *memdbDatastore) UniqueID(_ context.Context) (string, error) { return mdb.uniqueID, nil } +// SnapshotReader returns a reader for the snapshot visible at the given +// revision: the first entry in mdb.revisions at or after it, located by +// binary search. func (mdb *memdbDatastore) SnapshotReader(dr datastore.Revision) datastore.Reader { mdb.RLock() defer mdb.RUnlock() @@ -188,6 +192,8 @@ func (mdb *memdbDatastore) ReadWriteTx( for i := 0; i < txNumAttempts; i++ { var tx *memdb.Txn + var newRevision revisions.TimestampRevision + rwt := &memdbReadWriteTx{} createTxOnce := sync.Once{} txSrc := func() (*memdb.Txn, error) { var err error @@ -208,13 +214,16 @@ func (mdb *memdbDatastore) ReadWriteTx( tx = mdb.db.Txn(true) tx.TrackChanges() mdb.activeWriteTxn = tx + + // Assign the transaction's revision now, while holding the exclusive write transaction. + // Writers are serialized from this point until the revision is appended at commit. + newRevision = mdb.newRevisionIDNoLock() + rwt.newRevision = newRevision }) return tx, err } - - newRevision := mdb.newRevisionID() - rwt := &memdbReadWriteTx{memdbReader{&sync.Mutex{}, txSrc, nil, time.Now()}, newRevision} + rwt.memdbReader = memdbReader{&sync.Mutex{}, txSrc, nil, time.Now()} if config.SchemaHashPrecondition != "" { if err := assertSchemaHash(ctx, rwt, config.SchemaHashPrecondition); err != nil { mdb.Lock() @@ -247,123 +256,130 @@ func (mdb *memdbDatastore) ReadWriteTx( } mdb.Lock() - defer mdb.Unlock() + defer mdb.Unlock() // TODO is this defer correct? it runs at the end of the function, not at the end of the for loop's iteration + + // The user function never used the transaction: nothing was written, + // so no new revision is created and the head revision is returned. + if tx == nil { + if err := mdb.checkNotClosed(); err != nil { + return datastore.NoRevision, err + } + return mdb.headRevisionNoLock(), nil + } tracked := common.NewChanges(revisions.TimestampIDKeyFunc, datastore.WatchRelationships|datastore.WatchSchema, 0) - if tx != nil { - if config.Metadata != nil && len(config.Metadata.GetFields()) > 0 { - if err := tracked.AddRevisionMetadata(ctx, newRevision, config.Metadata.AsMap()); err != nil { - return datastore.NoRevision, err - } + if config.Metadata != nil && len(config.Metadata.GetFields()) > 0 { + if err := tracked.AddRevisionMetadata(ctx, newRevision, config.Metadata.AsMap()); err != nil { + return datastore.NoRevision, err } + } - for _, change := range tx.Changes() { - switch change.Table { - case tableRelationship: - switch { - case change.After != nil: - rt, err := change.After.(*relationship).Relationship() - if err != nil { - return datastore.NoRevision, err - } - - if err := tracked.AddRelationshipChange(ctx, newRevision, rt, tuple.UpdateOperationTouch); err != nil { - return datastore.NoRevision, err - } - case change.After == nil && change.Before != nil: - rt, err := change.Before.(*relationship).Relationship() - if err != nil { - return datastore.NoRevision, err - } - - if err := tracked.AddRelationshipChange(ctx, newRevision, rt, tuple.UpdateOperationDelete); err != nil { - return datastore.NoRevision, err - } - default: - return datastore.NoRevision, spiceerrors.MustBugf("unexpected relationship change") + for _, change := range tx.Changes() { + switch change.Table { + case tableRelationship: + switch { + case change.After != nil: + rt, err := change.After.(*relationship).Relationship() + if err != nil { + return datastore.NoRevision, err } - case tableNamespace: - switch { - case change.After != nil: - loaded := &corev1.NamespaceDefinition{} - if err := loaded.UnmarshalVT(change.After.(*namespace).configBytes); err != nil { - return datastore.NoRevision, err - } - - err := tracked.AddChangedDefinition(ctx, newRevision, loaded) - if err != nil { - return datastore.NoRevision, err - } - case change.After == nil && change.Before != nil: - err := tracked.AddDeletedNamespace(ctx, newRevision, change.Before.(*namespace).name) - if err != nil { - return datastore.NoRevision, err - } - default: - return datastore.NoRevision, spiceerrors.MustBugf("unexpected namespace change") + + if err := tracked.AddRelationshipChange(ctx, newRevision, rt, tuple.UpdateOperationTouch); err != nil { + return datastore.NoRevision, err } - case tableCaveats: - switch { - case change.After != nil: - loaded := &corev1.CaveatDefinition{} - if err := loaded.UnmarshalVT(change.After.(*caveat).definition); err != nil { - return datastore.NoRevision, err - } - - err := tracked.AddChangedDefinition(ctx, newRevision, loaded) - if err != nil { - return datastore.NoRevision, err - } - case change.After == nil && change.Before != nil: - err := tracked.AddDeletedCaveat(ctx, newRevision, change.Before.(*caveat).name) - if err != nil { - return datastore.NoRevision, err - } - default: - return datastore.NoRevision, spiceerrors.MustBugf("unexpected namespace change") + case change.After == nil && change.Before != nil: + rt, err := change.Before.(*relationship).Relationship() + if err != nil { + return datastore.NoRevision, err } - } - } - changes := tracked.AsRevisionChanges(revisions.TimestampIDKeyLessThanFunc) - wroteChangelog := false - for rc, err := range changes { - if err != nil { - return datastore.NoRevision, err + if err := tracked.AddRelationshipChange(ctx, newRevision, rt, tuple.UpdateOperationDelete); err != nil { + return datastore.NoRevision, err + } + default: + return datastore.NoRevision, spiceerrors.MustBugf("unexpected relationship change") } + case tableNamespace: + switch { + case change.After != nil: + loaded := &corev1.NamespaceDefinition{} + if err := loaded.UnmarshalVT(change.After.(*namespace).configBytes); err != nil { + return datastore.NoRevision, err + } - if wroteChangelog { - return datastore.NoRevision, spiceerrors.MustBugf("unexpected MemDB transaction with multiple revision changes") + err := tracked.AddChangedDefinition(ctx, newRevision, loaded) + if err != nil { + return datastore.NoRevision, err + } + case change.After == nil && change.Before != nil: + err := tracked.AddDeletedNamespace(ctx, newRevision, change.Before.(*namespace).name) + if err != nil { + return datastore.NoRevision, err + } + default: + return datastore.NoRevision, spiceerrors.MustBugf("unexpected namespace change") } + case tableCaveats: + switch { + case change.After != nil: + loaded := &corev1.CaveatDefinition{} + if err := loaded.UnmarshalVT(change.After.(*caveat).definition); err != nil { + return datastore.NoRevision, err + } - change := &changelog{ - revisionNanos: newRevision.TimestampNanoSec(), - changes: rc, - } - if err := tx.Insert(tableChangelog, change); err != nil { - return datastore.NoRevision, fmt.Errorf("error writing changelog: %w", err) + err := tracked.AddChangedDefinition(ctx, newRevision, loaded) + if err != nil { + return datastore.NoRevision, err + } + case change.After == nil && change.Before != nil: + err := tracked.AddDeletedCaveat(ctx, newRevision, change.Before.(*caveat).name) + if err != nil { + return datastore.NoRevision, err + } + default: + return datastore.NoRevision, spiceerrors.MustBugf("unexpected namespace change") } + } + } - wroteChangelog = true + changes := tracked.AsRevisionChanges(revisions.TimestampIDKeyLessThanFunc) + wroteChangelog := false + for rc, err := range changes { + if err != nil { + return datastore.NoRevision, err } - // Always emit a changelog entry for the committed revision, even - // when the transaction produced no observable changes (e.g., a - // TOUCH that matched the existing relationship). The changes - // payload is intentionally empty — the watch goroutine constructs the - // checkpoint event itself based on each consumer's options. - if !wroteChangelog { - change := &changelog{ - revisionNanos: newRevision.TimestampNanoSec(), - changes: datastore.RevisionChanges{}, - } - if err := tx.Insert(tableChangelog, change); err != nil { - return datastore.NoRevision, fmt.Errorf("error writing changelog: %w", err) - } + if wroteChangelog { + return datastore.NoRevision, spiceerrors.MustBugf("unexpected MemDB transaction with multiple revision changes") } - tx.Commit() + change := &changelog{ + revisionNanos: newRevision.TimestampNanoSec(), + changes: rc, + } + if err := tx.Insert(tableChangelog, change); err != nil { + return datastore.NoRevision, fmt.Errorf("error writing changelog: %w", err) + } + + wroteChangelog = true + } + + // Always emit a changelog entry for the committed revision, even + // when the transaction produced no observable changes (e.g., a + // TOUCH that matched the existing relationship). The changes + // payload is intentionally empty — the watch goroutine constructs the + // checkpoint event itself based on each consumer's options. + if !wroteChangelog { + change := &changelog{ + revisionNanos: newRevision.TimestampNanoSec(), + changes: datastore.RevisionChanges{}, + } + if err := tx.Insert(tableChangelog, change); err != nil { + return datastore.NoRevision, fmt.Errorf("error writing changelog: %w", err) + } } + + tx.Commit() mdb.activeWriteTxn = nil mdb.writeTxReady.Signal() @@ -374,7 +390,7 @@ func (mdb *memdbDatastore) ReadWriteTx( // Create a snapshot and add it to the revisions slice schemaHash := mdb.getCurrentSchemaHashNoLock() snap := mdb.db.Snapshot() - mdb.insertRevisionSnapshot(newRevision, schemaHash, snap) + mdb.revisions = append(mdb.revisions, snapshot{newRevision, schemaHash, snap}) return newRevision, nil } diff --git a/internal/datastore/memdb/memdb_test.go b/internal/datastore/memdb/memdb_test.go index 4adb3977d8..ef106afc0a 100644 --- a/internal/datastore/memdb/memdb_test.go +++ b/internal/datastore/memdb/memdb_test.go @@ -118,6 +118,7 @@ func TestConcurrentWriteRelsSucceed(t *testing.T) { // - a write should NOT be visible at revisions below its own (consistent snapshot) // - a write SHOULD be visible at every revision at or above its own // - the head revision sees every committed write +// // TODO: should this be a datastore conformance test (not specific to memdb)? func TestConcurrentWrite(t *testing.T) { ds, err := NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour) diff --git a/internal/datastore/memdb/revisions.go b/internal/datastore/memdb/revisions.go index 2f6d67c512..4519b6236a 100644 --- a/internal/datastore/memdb/revisions.go +++ b/internal/datastore/memdb/revisions.go @@ -2,12 +2,8 @@ package memdb import ( "context" - "slices" - "sort" "time" - "github.com/hashicorp/go-memdb" - "github.com/authzed/spicedb/internal/datastore/revisions" "github.com/authzed/spicedb/pkg/datastore" ) @@ -18,10 +14,9 @@ func nowRevision() revisions.TimestampRevision { return revisions.NewForTime(time.Now().UTC()) } -func (mdb *memdbDatastore) newRevisionID() revisions.TimestampRevision { - mdb.Lock() - defer mdb.Unlock() - +// newRevisionIDNoLock returns a revision strictly greater than the current head revision. +// The caller must hold the RWMutex. +func (mdb *memdbDatastore) newRevisionIDNoLock() revisions.TimestampRevision { existing := mdb.revisions[len(mdb.revisions)-1].revision created := nowRevision() @@ -29,43 +24,19 @@ func (mdb *memdbDatastore) newRevisionID() revisions.TimestampRevision { // precision on macOS Monterey in Go 1.19.1. This means that HeadRevision // and the result of a ReadWriteTx could return the *same* transaction ID // if both are executed in sequence without any other forms of delay on - // macOS. We therefore check if the created transaction ID matches that - // previously created and, if not, add to it. + // macOS. We therefore check if the created transaction ID is at or before + // the head revision and, if so, advance past the head. // // See: https://github.com/golang/go/issues/22037 which appeared to fix // this in Go 1.9.2, but there appears to have been a reversion with either // the new version of macOS or Go. - if created.Equal(existing) { - return revisions.NewForTimestamp(created.TimestampNanoSec() + 1) + if !created.GreaterThan(existing) { + return revisions.NewForTimestamp(existing.TimestampNanoSec() + 1) } return created } -// insertRevisionSnapshot records a newly committed snapshot at its sorted position -// The caller must already hold mdb's write lock. -func (mdb *memdbDatastore) insertRevisionSnapshot(newRevision revisions.TimestampRevision, schemaHash string, snap *memdb.MemDB) { - insertAt := sort.Search(len(mdb.revisions), func(i int) bool { - return mdb.revisions[i].revision.GreaterThan(newRevision) - }) - mdb.revisions = slices.Insert(mdb.revisions, insertAt, snapshot{newRevision, schemaHash, snap}) -} - -// indexOfRevision returns the index of the snapshot recorded for the exact given revision -// Performs binary search for finding the revision, returns -1 if none exists -func (mdb *memdbDatastore) indexOfRevision(r datastore.Revision) int { - mdb.RLock() - defer mdb.RUnlock() - - i := sort.Search(len(mdb.revisions), func(i int) bool { - return !mdb.revisions[i].revision.LessThan(r) - }) - if i < len(mdb.revisions) && mdb.revisions[i].revision.Equal(r) { - return i - } - return -1 -} - func (mdb *memdbDatastore) HeadRevision(_ context.Context) (datastore.RevisionWithSchemaHash, error) { mdb.RLock() defer mdb.RUnlock() @@ -148,7 +119,7 @@ func (mdb *memdbDatastore) checkRevisionLocalCallerMustLock(dr datastore.Revisio // HEAD revision is behind it. if dr.GreaterThan(now) { // If the revision is in the "future", then check to ensure that it is <= of HEAD to handle - // the microsecond granularity on macos (see comment above in newRevisionID) + // the microsecond granularity on macos (see comment above in newRevisionIDNoLock) headRevision := mdb.headRevisionNoLock() if dr.LessThan(headRevision) || dr.Equal(headRevision) { return nil From 920e7b7423e79c7c3d97d3c51d5e637ce0d51a1e Mon Sep 17 00:00:00 2001 From: Maria Ines Parnisari Date: Wed, 22 Jul 2026 15:41:43 -0700 Subject: [PATCH 8/9] test: make the test datastore-agnostic --- internal/datastore/memdb/memdb_test.go | 117 ---------------- pkg/datastore/test/datastore.go | 4 + pkg/datastore/test/relationships.go | 130 ++++++++++++++++++ .../devcontext_metrics_probe_test.go | 74 ++++++++++ 4 files changed, 208 insertions(+), 117 deletions(-) create mode 100644 pkg/development/devcontext_metrics_probe_test.go diff --git a/internal/datastore/memdb/memdb_test.go b/internal/datastore/memdb/memdb_test.go index ef106afc0a..b16298dc55 100644 --- a/internal/datastore/memdb/memdb_test.go +++ b/internal/datastore/memdb/memdb_test.go @@ -113,123 +113,6 @@ func TestConcurrentWriteRelsSucceed(t *testing.T) { require.NoError(g.Wait()) } -// TestConcurrentWrite covers revision visibility when two write transactions overlap. -// - each transaction's write should be visible when reading at its own returned revision (read-your-writes) -// - a write should NOT be visible at revisions below its own (consistent snapshot) -// - a write SHOULD be visible at every revision at or above its own -// - the head revision sees every committed write -// -// TODO: should this be a datastore conformance test (not specific to memdb)? -func TestConcurrentWrite(t *testing.T) { - ds, err := NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour) - require.NoError(t, err) - t.Cleanup(func() { - require.NoError(t, ds.Close()) - }) - - // docsAt returns the set of document object IDs visible at the given revision. - docsAt := func(rev datastore.Revision) map[string]bool { - reader := ds.SnapshotReader(rev) - it, err := reader.QueryRelationships(t.Context(), datastore.RelationshipsFilter{OptionalResourceType: "document"}) - require.NoError(t, err) - rels, err := datastore.IteratorToSlice(it) - require.NoError(t, err) - - seen := map[string]bool{} - for _, rel := range rels { - seen[rel.Resource.ObjectID] = true - } - return seen - } - - // Execute many iterations so that one run is enough to expose ordering problems - const iterations = 10_000 - for i := 0; i < iterations; i++ { - // Unique names per iteration so that visibility assertions cannot be - // satisfied by a previous iteration's writes. - docA := fmt.Sprintf("doc-a-%d", i) - docB := fmt.Sprintf("doc-b-%d", i) - relA := tuple.MustParse(fmt.Sprintf("document:%s#viewer@user:tom", docA)) - relB := tuple.MustParse(fmt.Sprintf("document:%s#viewer@user:tom", docB)) - - var ( - revA datastore.Revision - errA error - waitUntilAStarts = make(chan struct{}, 1) - waitUntilBCompletes = make(chan struct{}, 1) - waitUntilADone = make(chan struct{}) - ) - go func() { - defer close(waitUntilADone) - revA, errA = ds.ReadWriteTx(t.Context(), func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { - waitUntilAStarts <- struct{}{} - <-waitUntilBCompletes // held until B has fully committed below - return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ - tuple.Touch(relA), - }) - }, options.WithDisableRetries(true)) - }() - - <-waitUntilAStarts // A is now blocked - - // Transaction B runs to completion while A is blocked - revB, err := ds.ReadWriteTx(t.Context(), func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { - return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ - tuple.Touch(relB), - }) - }, options.WithDisableRetries(true)) - require.NoError(t, err) - - waitUntilBCompletes <- struct{}{} // unblocks A - <-waitUntilADone - require.NoError(t, errA) - - require.False(t, revA.Equal(revB), "iteration %d: concurrent transactions must be assigned distinct revisions, both got %v", i, revA) - - // Each transaction's write must be visible at its own returned revision. - require.True(t, docsAt(revA)[docA], "iteration %d: A's write is invisible when reading at A's own returned revision %v", i, revA) - require.True(t, docsAt(revB)[docB], "iteration %d: B's write is invisible when reading at B's own returned revision %v", i, revB) - - type write struct { - rev datastore.Revision - doc string - } - earlier, later := write{revA, docA}, write{revB, docB} - if later.rev.LessThan(earlier.rev) { - earlier, later = later, earlier - } - - // The write committed at the later revision must not be visible at the earlier revision - atEarlier := docsAt(earlier.rev) - require.True(t, atEarlier[earlier.doc], "iteration %d: %s is invisible at its own revision %v", i, earlier.doc, earlier.rev) - require.False(t, atEarlier[later.doc], "iteration %d: %s was committed at later revision %v but is visible at earlier revision %v", i, later.doc, later.rev, earlier.rev) - - // Both writes must be visible at the later revision - atLater := docsAt(later.rev) - require.True(t, atLater[earlier.doc], "iteration %d: %s is visible at revision %v but disappears at later revision %v", i, earlier.doc, earlier.rev, later.rev) - require.True(t, atLater[later.doc], "iteration %d: %s is invisible at its own revision %v", i, later.doc, later.rev) - - // The head revision must see every committed write - head, err := ds.HeadRevision(t.Context()) - require.NoError(t, err) - atHead := docsAt(head.Revision) - require.True(t, atHead[docA] && atHead[docB], "iteration %d: head revision %v is missing committed writes: %v", i, head.Revision, atHead) - - // Delete this iteration's relationships so the relationship table does - // not grow across iterations: memdb scans the whole namespace per - // query, so an ever-growing table would make this test quadratic. - // Earlier revisions retain their snapshots, so this does not affect - // the assertions above. - _, err = ds.ReadWriteTx(t.Context(), func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { - return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{ - tuple.Delete(relA), - tuple.Delete(relB), - }) - }, options.WithDisableRetries(true)) - require.NoError(t, err) - } -} - func TestAnythingAfterCloseDoesNotPanic(t *testing.T) { require := require.New(t) diff --git a/pkg/datastore/test/datastore.go b/pkg/datastore/test/datastore.go index 6fecd070fd..dc4a20c9be 100644 --- a/pkg/datastore/test/datastore.go +++ b/pkg/datastore/test/datastore.go @@ -172,6 +172,10 @@ func AllWithExceptions(t *testing.T, tester DatastoreTester, except Categories) t.Run("TestConcurrentWriteSerialization", runner(tester, ConcurrentWriteSerializationTest)) t.Run("TestConcurrentWriteDeadlock", runner(tester, ConcurrentWriteDeadlockTest)) } + // Not gated on the ConcurrentWrite category: the overlapping transactions + // never contend for the write lock, so even datastores with a global write + // lock (e.g. memdb) must pass it. See the test's doc comment. + t.Run("TestConcurrentWriteRevisionVisibility", runner(tester, ConcurrentWriteRevisionVisibilityTest)) t.Run("TestOrdering", runner(tester, OrderingTest)) t.Run("TestLimit", runner(tester, LimitTest)) diff --git a/pkg/datastore/test/relationships.go b/pkg/datastore/test/relationships.go index 61da6f7024..d389d3e47d 100644 --- a/pkg/datastore/test/relationships.go +++ b/pkg/datastore/test/relationships.go @@ -2263,6 +2263,136 @@ func ConcurrentWriteDeadlockTest(t *testing.T, tester DatastoreTester) { } } +// ConcurrentWriteRevisionVisibilityTest covers revision visibility when two +// write transactions overlap: transaction A opens first but is held open while +// transaction B starts and commits, and only then does A write and commit. +// - each transaction's write must be visible when reading at its own +// returned revision (read-your-writes) +// - when the two revisions are comparable, the later revision must see the +// earlier write, and the earlier revision must NOT see the later write +// - when the two revisions are concurrent (e.g. Postgres snapshots of +// overlapping transactions), neither write may be visible at the other's +// revision +// - the head revision must see every committed write +// +// This test is safe even for datastores that serialize write transactions with +// a global lock (e.g. memdb): transaction A performs no reads or writes until +// transaction B has fully committed, so the two never contend for the write +// lock. +func ConcurrentWriteRevisionVisibilityTest(t *testing.T, tester DatastoreTester) { + rawDS, err := tester.New(t, 0, veryLargeGCInterval, veryLargeGCWindow, 1) + require.NoError(t, err) + + ds, _ := testfixtures.StandardDatastoreWithSchema(t, rawDS) + ctx := t.Context() + + // docsAt returns which of the given resource object IDs are visible at the + // given revision. + docsAt := func(rev datastore.Revision, resourceIDs ...string) map[string]bool { + reader := ds.SnapshotReader(rev) + it, err := reader.QueryRelationships(ctx, datastore.RelationshipsFilter{ + OptionalResourceType: testResourceNamespace, + OptionalResourceIds: resourceIDs, + }, options.WithQueryShape(queryshape.Varying)) + require.NoError(t, err) + rels, err := datastore.IteratorToSlice(it) + require.NoError(t, err) + + seen := map[string]bool{} + for _, rel := range rels { + seen[rel.Resource.ObjectID] = true + } + return seen + } + + // Execute many iterations so that one run is enough to expose ordering problems. + const iterations = 1000 + for i := 0; i < iterations; i++ { + // Unique names per iteration so that visibility assertions cannot be + // satisfied by a previous iteration's writes. + docA := fmt.Sprintf("doc-a-%d", i) + docB := fmt.Sprintf("doc-b-%d", i) + relA := tuple.Touch(makeTestRel(docA, "tom")) + relB := tuple.Touch(makeTestRel(docB, "tom")) + + var ( + revA datastore.Revision + errA error + aStarted = make(chan struct{}) + aStartedOnce = sync.Once{} + bCompleted = make(chan struct{}) + aDone = make(chan struct{}) + ) + go func() { + defer close(aDone) + // The transaction function may be re-run on retryable errors, + // either by the datastore's retry loop or by the backend client + // itself (e.g. Spanner and CockroachDB re-run on aborts). The + // sync.Once and the closed channels make re-runs proceed + // immediately instead of hanging. + revA, errA = ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { + aStartedOnce.Do(func() { close(aStarted) }) + <-bCompleted // held until B has fully committed below + return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{relA}) + }) + }() + + <-aStarted // A's transaction is now open and blocked + + // Transaction B runs to completion while A remains open. + revB, err := ds.ReadWriteTx(ctx, func(ctx context.Context, rwt datastore.ReadWriteTransaction) error { + return rwt.WriteRelationships(ctx, []tuple.RelationshipUpdate{relB}) + }) + require.NoError(t, err) + + close(bCompleted) // unblocks A + <-aDone + require.NoError(t, errA) + + require.False(t, revA.Equal(revB), "iteration %d: concurrent transactions must be assigned distinct revisions, both got %v", i, revA) + + // Each transaction's write must be visible at its own returned revision. + require.True(t, docsAt(revA, docA)[docA], "iteration %d: A's write is invisible when reading at A's own returned revision %v", i, revA) + require.True(t, docsAt(revB, docB)[docB], "iteration %d: B's write is invisible when reading at B's own returned revision %v", i, revB) + + if revA.LessThan(revB) || revB.LessThan(revA) { + // The revisions are comparable: the datastore assigned the two + // transactions a total order. + type write struct { + rev datastore.Revision + doc string + } + earlier, later := write{revA, docA}, write{revB, docB} + if later.rev.LessThan(earlier.rev) { + earlier, later = later, earlier + } + + // The write committed at the later revision must not be visible at + // the earlier revision. + atEarlier := docsAt(earlier.rev, docA, docB) + require.True(t, atEarlier[earlier.doc], "iteration %d: %s is invisible at its own revision %v", i, earlier.doc, earlier.rev) + require.False(t, atEarlier[later.doc], "iteration %d: %s was committed at later revision %v but is visible at earlier revision %v", i, later.doc, later.rev, earlier.rev) + + // Both writes must be visible at the later revision. + atLater := docsAt(later.rev, docA, docB) + require.True(t, atLater[earlier.doc], "iteration %d: %s is visible at revision %v but disappears at later revision %v", i, earlier.doc, earlier.rev, later.rev) + require.True(t, atLater[later.doc], "iteration %d: %s is invisible at its own revision %v", i, later.doc, later.rev) + } else { + // The revisions are concurrent (e.g. Postgres snapshots of + // overlapping transactions): each transaction's snapshot must + // exclude the other's write. + require.False(t, docsAt(revA, docB)[docB], "iteration %d: %s was committed by a concurrent transaction at revision %v but is visible at revision %v", i, docB, revB, revA) + require.False(t, docsAt(revB, docA)[docA], "iteration %d: %s was committed by a concurrent transaction at revision %v but is visible at revision %v", i, docA, revA, revB) + } + + // The head revision must see every committed write. + head, err := ds.HeadRevision(ctx) + require.NoError(t, err) + atHead := docsAt(head.Revision, docA, docB) + require.True(t, atHead[docA] && atHead[docB], "iteration %d: head revision %v is missing committed writes: %v", i, head.Revision, atHead) + } +} + func BulkDeleteRelationshipsTest(t *testing.T, tester DatastoreTester) { require := require.New(t) diff --git a/pkg/development/devcontext_metrics_probe_test.go b/pkg/development/devcontext_metrics_probe_test.go new file mode 100644 index 0000000000..214e6ff99f --- /dev/null +++ b/pkg/development/devcontext_metrics_probe_test.go @@ -0,0 +1,74 @@ +package development + +import ( + "testing" + + "github.com/stretchr/testify/require" + + v1 "github.com/authzed/authzed-go/proto/authzed/api/v1" + + core "github.com/authzed/spicedb/pkg/proto/core/v1" + devinterface "github.com/authzed/spicedb/pkg/proto/developer/v1" + "github.com/authzed/spicedb/pkg/tuple" +) + +// Verifies the devcontext V1 service does not nil-panic on the +// metrics-recording paths (CheckPermission, CheckBulkPermissions, +// WriteRelationships) when no Metrics is provided in the config. +func TestDevContextV1ServiceMetricsPaths(t *testing.T) { + devCtx, devErrs, err := NewDevContext(t.Context(), &devinterface.RequestContext{ + Schema: `definition user {} + +definition document { + relation viewer: user + permission view = viewer +} +`, + Relationships: []*core.RelationTuple{ + tuple.MustParse("document:somedoc#viewer@user:someuser").ToCoreTuple(), + }, + }) + require.NoError(t, err) + require.Nil(t, devErrs) + + conn, shutdown, err := devCtx.RunV1InMemoryService() + require.NoError(t, err) + t.Cleanup(shutdown) + + client := v1.NewPermissionsServiceClient(conn) + + checkResp, err := client.CheckPermission(t.Context(), &v1.CheckPermissionRequest{ + Resource: &v1.ObjectReference{ObjectType: "document", ObjectId: "somedoc"}, + Permission: "view", + Subject: &v1.SubjectReference{Object: &v1.ObjectReference{ObjectType: "user", ObjectId: "someuser"}}, + }) + require.NoError(t, err) + require.Equal(t, v1.CheckPermissionResponse_PERMISSIONSHIP_HAS_PERMISSION, checkResp.Permissionship) + + bulkResp, err := client.CheckBulkPermissions(t.Context(), &v1.CheckBulkPermissionsRequest{ + Items: []*v1.CheckBulkPermissionsRequestItem{ + { + Resource: &v1.ObjectReference{ObjectType: "document", ObjectId: "somedoc"}, + Permission: "view", + Subject: &v1.SubjectReference{Object: &v1.ObjectReference{ObjectType: "user", ObjectId: "someuser"}}, + }, + { + Resource: &v1.ObjectReference{ObjectType: "document", ObjectId: "somedoc"}, + Permission: "view", + Subject: &v1.SubjectReference{Object: &v1.ObjectReference{ObjectType: "user", ObjectId: "nobody"}}, + }, + }, + }) + require.NoError(t, err) + require.Len(t, bulkResp.Pairs, 2) + + _, err = client.WriteRelationships(t.Context(), &v1.WriteRelationshipsRequest{ + Updates: []*v1.RelationshipUpdate{ + { + Operation: v1.RelationshipUpdate_OPERATION_TOUCH, + Relationship: tuple.ToV1Relationship(tuple.MustParse("document:somedoc#viewer@user:anotheruser")), + }, + }, + }) + require.NoError(t, err) +} From 2a1515e2e472b8eeb266b3dca4eeccf2cd47a818 Mon Sep 17 00:00:00 2001 From: Maria Ines Parnisari Date: Thu, 23 Jul 2026 14:09:09 -0700 Subject: [PATCH 9/9] chore: remove unrelated test --- .../devcontext_metrics_probe_test.go | 74 ------------------- 1 file changed, 74 deletions(-) delete mode 100644 pkg/development/devcontext_metrics_probe_test.go diff --git a/pkg/development/devcontext_metrics_probe_test.go b/pkg/development/devcontext_metrics_probe_test.go deleted file mode 100644 index 214e6ff99f..0000000000 --- a/pkg/development/devcontext_metrics_probe_test.go +++ /dev/null @@ -1,74 +0,0 @@ -package development - -import ( - "testing" - - "github.com/stretchr/testify/require" - - v1 "github.com/authzed/authzed-go/proto/authzed/api/v1" - - core "github.com/authzed/spicedb/pkg/proto/core/v1" - devinterface "github.com/authzed/spicedb/pkg/proto/developer/v1" - "github.com/authzed/spicedb/pkg/tuple" -) - -// Verifies the devcontext V1 service does not nil-panic on the -// metrics-recording paths (CheckPermission, CheckBulkPermissions, -// WriteRelationships) when no Metrics is provided in the config. -func TestDevContextV1ServiceMetricsPaths(t *testing.T) { - devCtx, devErrs, err := NewDevContext(t.Context(), &devinterface.RequestContext{ - Schema: `definition user {} - -definition document { - relation viewer: user - permission view = viewer -} -`, - Relationships: []*core.RelationTuple{ - tuple.MustParse("document:somedoc#viewer@user:someuser").ToCoreTuple(), - }, - }) - require.NoError(t, err) - require.Nil(t, devErrs) - - conn, shutdown, err := devCtx.RunV1InMemoryService() - require.NoError(t, err) - t.Cleanup(shutdown) - - client := v1.NewPermissionsServiceClient(conn) - - checkResp, err := client.CheckPermission(t.Context(), &v1.CheckPermissionRequest{ - Resource: &v1.ObjectReference{ObjectType: "document", ObjectId: "somedoc"}, - Permission: "view", - Subject: &v1.SubjectReference{Object: &v1.ObjectReference{ObjectType: "user", ObjectId: "someuser"}}, - }) - require.NoError(t, err) - require.Equal(t, v1.CheckPermissionResponse_PERMISSIONSHIP_HAS_PERMISSION, checkResp.Permissionship) - - bulkResp, err := client.CheckBulkPermissions(t.Context(), &v1.CheckBulkPermissionsRequest{ - Items: []*v1.CheckBulkPermissionsRequestItem{ - { - Resource: &v1.ObjectReference{ObjectType: "document", ObjectId: "somedoc"}, - Permission: "view", - Subject: &v1.SubjectReference{Object: &v1.ObjectReference{ObjectType: "user", ObjectId: "someuser"}}, - }, - { - Resource: &v1.ObjectReference{ObjectType: "document", ObjectId: "somedoc"}, - Permission: "view", - Subject: &v1.SubjectReference{Object: &v1.ObjectReference{ObjectType: "user", ObjectId: "nobody"}}, - }, - }, - }) - require.NoError(t, err) - require.Len(t, bulkResp.Pairs, 2) - - _, err = client.WriteRelationships(t.Context(), &v1.WriteRelationshipsRequest{ - Updates: []*v1.RelationshipUpdate{ - { - Operation: v1.RelationshipUpdate_OPERATION_TOUCH, - Relationship: tuple.ToV1Relationship(tuple.MustParse("document:somedoc#viewer@user:anotheruser")), - }, - }, - }) - require.NoError(t, err) -}