Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions internal/datastore/crdb/crdb.go
Original file line number Diff line number Diff line change
Expand Up @@ -320,6 +320,10 @@ func (cds *crdbDatastore) MetricsID() (string, error) {
return common.MetricsIDFromURL(cds.dburl)
}

func (cds *crdbDatastore) EngineName() string {
return Engine
}

func (cds *crdbDatastore) ReadWriteTx(
ctx context.Context,
f datastore.TxUserFunc,
Expand Down
4 changes: 2 additions & 2 deletions internal/datastore/crdb/crdb_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -206,7 +206,7 @@ func TestCRDBDatastoreWithIntegrity(t *testing.T) { //nolint:tparallel
t.Parallel()
b := testdatastore.RunCRDBForTesting(t, crdbTestVersion())

test.All(t, crdbFactory.NewTester(test.DatastoreTesterFunc(func(_ testing.TB, revisionQuantization, gcInterval, gcWindow time.Duration, watchBufferLength uint16) (datastore.Datastore, error) {
test.AllWithExceptions(t, crdbFactory.NewTester(test.DatastoreTesterFunc(func(_ testing.TB, revisionQuantization, gcInterval, gcWindow time.Duration, watchBufferLength uint16) (datastore.Datastore, error) {
ctx := t.Context()
ds := b.NewDatastore(t, func(engine, uri string) datastore.Datastore {
ds, err := NewCRDBDatastore(
Expand All @@ -231,7 +231,7 @@ func TestCRDBDatastoreWithIntegrity(t *testing.T) { //nolint:tparallel
})

return ds, nil
})))
})), test.WithCategories(test.MigrationCategory))

unwrappedTester := test.DatastoreTesterFunc(func(_ testing.TB, revisionQuantization, gcInterval, gcWindow time.Duration, watchBufferLength uint16) (datastore.Datastore, error) {
ctx := t.Context()
Expand Down
79 changes: 79 additions & 0 deletions internal/datastore/crdb/enginebuilder.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
package crdb

import (
"context"
"errors"

"github.com/ccoveille/go-safecast/v2"

"github.com/authzed/spicedb/internal/datastore/common"
"github.com/authzed/spicedb/internal/datastore/crdb/migrations"
datastorecfg "github.com/authzed/spicedb/pkg/cmd/datastore/dsconfig"
"github.com/authzed/spicedb/pkg/datastore"
"github.com/authzed/spicedb/pkg/datastore/migration"
)

func init() {
datastorecfg.RegisterEngine(Engine, newDatastoreFromConfig)
migration.RegisterMigratableEngine(Engine, migrations.CRDBMigrations, newMigrationDriverFromConfig, "add-schema-tables")
}

func newMigrationDriverFromConfig(_ context.Context, cfg *migration.Config) (*migrations.CRDBDriver, error) {
return migrations.NewCRDBDriver(cfg.DatastoreURI)
}

func newDatastoreFromConfig(ctx context.Context, opts datastorecfg.Config) (datastore.Datastore, error) {
if len(opts.ReadReplicaURIs) > 0 {
return nil, errors.New("read replicas are not supported for the CockroachDB datastore engine")
}

Check warning on line 28 in internal/datastore/crdb/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/crdb/enginebuilder.go#L27-L28

Added lines #L27 - L28 were not covered by tests

maxRetries, err := safecast.Convert[uint8](opts.MaxRetries)
if err != nil {
return nil, errors.New("max-retries could not be cast to uint8")
}

Check warning on line 33 in internal/datastore/crdb/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/crdb/enginebuilder.go#L32-L33

Added lines #L32 - L33 were not covered by tests

watchChangeBufferMaximumSize, err := common.WatchBufferSize(opts.WatchChangeBufferMaximumSize)
if err != nil {
return nil, err
}

Check warning on line 38 in internal/datastore/crdb/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/crdb/enginebuilder.go#L37-L38

Added lines #L37 - L38 were not covered by tests

return NewCRDBDatastore(
ctx,
opts.URI,
GCWindow(opts.GCWindow),
RevisionQuantization(opts.RevisionQuantization),
MaxRevisionStalenessPercent(opts.MaxRevisionStalenessPercent),
ReadConnsMaxOpen(opts.ReadConnPool.MaxOpenConns),
ReadConnsMinOpen(opts.ReadConnPool.MinOpenConns),
ReadConnMaxIdleTime(opts.ReadConnPool.MaxIdleTime),
ReadConnMaxLifetime(opts.ReadConnPool.MaxLifetime),
ReadConnMaxLifetimeJitter(opts.ReadConnPool.MaxLifetimeJitter),
ReadConnHealthCheckInterval(opts.ReadConnPool.HealthCheckInterval),
ReadConnPingTimeout(opts.ReadConnPool.PingTimeout),
WithAcquireTimeout(opts.WriteAcquisitionTimeout),
WriteConnsMaxOpen(opts.WriteConnPool.MaxOpenConns),
WriteConnsMinOpen(opts.WriteConnPool.MinOpenConns),
WriteConnMaxIdleTime(opts.WriteConnPool.MaxIdleTime),
WriteConnMaxLifetime(opts.WriteConnPool.MaxLifetime),
WriteConnMaxLifetimeJitter(opts.WriteConnPool.MaxLifetimeJitter),
WriteConnHealthCheckInterval(opts.WriteConnPool.HealthCheckInterval),
WriteConnPingTimeout(opts.WriteConnPool.PingTimeout),
FollowerReadDelay(opts.FollowerReadDelay),
MaxRetries(maxRetries),
OverlapKey(opts.OverlapKey),
OverlapStrategy(opts.OverlapStrategy),
WatchBufferLength(opts.WatchBufferLength),
WatchChangeBufferMaximumSize(watchChangeBufferMaximumSize),
WatchBufferWriteTimeout(opts.WatchBufferWriteTimeout),
WatchConnectTimeout(opts.WatchConnectTimeout),
WithEnablePrometheusStats(opts.EnableDatastoreMetrics),
WithEnableConnectionBalancing(opts.EnableConnectionBalancing),
ConnectRate(opts.ConnectRate),
FilterMaximumIDCount(opts.FilterMaximumIDCount),
WithIntegrity(opts.RelationshipIntegrityEnabled),
AllowedMigrations(opts.AllowedMigrations),
WithColumnOptimization(opts.ExperimentalColumnOptimization),
IncludeQueryParametersInTraces(opts.IncludeQueryParametersInTraces),
WithWatchDisabled(opts.DisableWatchSupport),
)
}
13 changes: 13 additions & 0 deletions internal/datastore/engines/engines.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
// Package engines registers, via side effect of import, every datastore
// engine defined in this repository with pkg/cmd/datastore's engine registry.
// Blank-import it from any binary or test that constructs datastores by engine
// name through pkg/cmd/datastore.NewDatastore.
package engines

import (
_ "github.com/authzed/spicedb/internal/datastore/crdb"
_ "github.com/authzed/spicedb/internal/datastore/memdb"
_ "github.com/authzed/spicedb/internal/datastore/mysql"
_ "github.com/authzed/spicedb/internal/datastore/postgres"
_ "github.com/authzed/spicedb/internal/datastore/spanner"
)
23 changes: 23 additions & 0 deletions internal/datastore/memdb/enginebuilder.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
package memdb

import (
"context"
"errors"

log "github.com/authzed/spicedb/internal/logging"
datastorecfg "github.com/authzed/spicedb/pkg/cmd/datastore/dsconfig"
"github.com/authzed/spicedb/pkg/datastore"
)

func init() {
datastorecfg.RegisterEngine(Engine, newDatastoreFromConfig)
}

func newDatastoreFromConfig(_ context.Context, opts datastorecfg.Config) (datastore.Datastore, error) {
if len(opts.ReadReplicaURIs) > 0 {
return nil, errors.New("read replicas are not supported for the in-memory datastore engine")
}

Check warning on line 19 in internal/datastore/memdb/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/memdb/enginebuilder.go#L18-L19

Added lines #L18 - L19 were not covered by tests

log.Warn().Msg("in-memory datastore is not persistent and not feasible to run in a high availability fashion")
return NewMemdbDatastore(opts.WatchBufferLength, opts.RevisionQuantization, opts.GCWindow)
}
4 changes: 4 additions & 0 deletions internal/datastore/memdb/memdb.go
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,10 @@
return "memdb", nil
}

func (mdb *memdbDatastore) EngineName() string {
return Engine

Check warning on line 116 in internal/datastore/memdb/memdb.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/memdb/memdb.go#L115-L116

Added lines #L115 - L116 were not covered by tests
}

func (mdb *memdbDatastore) UniqueID(_ context.Context) (string, error) {
return mdb.uniqueID, nil
}
Expand Down
26 changes: 14 additions & 12 deletions internal/datastore/memdb/memdb_test.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package memdb
package memdb_test

import (
"context"
Expand All @@ -11,6 +11,7 @@ import (
"github.com/stretchr/testify/require"
"golang.org/x/sync/errgroup"

"github.com/authzed/spicedb/internal/datastore/memdb"
"github.com/authzed/spicedb/pkg/datastore"
"github.com/authzed/spicedb/pkg/datastore/options"
test "github.com/authzed/spicedb/pkg/datastore/test"
Expand All @@ -19,24 +20,25 @@ import (
"github.com/authzed/spicedb/pkg/tuple"
)

var memdbFactory = test.NewTesterFactory(ErrSerialization)
var memdbFactory = test.NewTesterFactory(memdb.ErrSerialization)

type memDBTest struct{}

func (memDBTest) New(_ testing.TB, revisionQuantization, _, gcWindow time.Duration, watchBufferLength uint16) (datastore.Datastore, error) {
return NewMemdbDatastore(watchBufferLength, revisionQuantization, gcWindow)
return memdb.NewMemdbDatastore(watchBufferLength, revisionQuantization, gcWindow)
}

func TestMemdbDatastore(t *testing.T) {
// ConcurrentWrite tests require row-level locking; memdb uses a global write lock
// and would deadlock if two write transactions were opened concurrently.
test.AllWithExceptions(t, memdbFactory.NewTester(memDBTest{}), test.WithCategories(test.ConcurrentWriteCategory))
// Migration tests are excluded because memdb has no schema migrations.
test.AllWithExceptions(t, memdbFactory.NewTester(memDBTest{}), test.WithCategories(test.ConcurrentWriteCategory, test.MigrationCategory))
}

func TestConcurrentWritePanic(t *testing.T) {
require := require.New(t)

ds, err := NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour)
ds, err := memdb.NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour)
require.NoError(err)

ctx := t.Context()
Expand Down Expand Up @@ -87,7 +89,7 @@ func TestConcurrentWritePanic(t *testing.T) {
func TestConcurrentWriteRelsSucceed(t *testing.T) {
require := require.New(t)

ds, err := NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour)
ds, err := memdb.NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour)
require.NoError(err)

ctx := t.Context()
Expand Down Expand Up @@ -116,7 +118,7 @@ func TestConcurrentWriteRelsSucceed(t *testing.T) {
func TestAnythingAfterCloseDoesNotPanic(t *testing.T) {
require := require.New(t)

ds, err := NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour)
ds, err := memdb.NewMemdbDatastore(0, 1*time.Hour, 1*time.Hour)
require.NoError(err)

lowestRevision, err := ds.HeadRevision(t.Context())
Expand All @@ -129,21 +131,21 @@ func TestAnythingAfterCloseDoesNotPanic(t *testing.T) {

select {
case err := <-errChan:
require.ErrorIs(err, ErrMemDBIsClosed)
require.ErrorIs(err, memdb.ErrMemDBIsClosed)
case <-time.After(time.Second):
require.Fail("expected an error but waited too long")
}

_, err = ds.Statistics(t.Context())
require.ErrorIs(err, ErrMemDBIsClosed)
require.ErrorIs(err, memdb.ErrMemDBIsClosed)

err = ds.CheckRevision(t.Context(), lowestRevision.Revision)
require.ErrorIs(err, ErrMemDBIsClosed)
require.ErrorIs(err, memdb.ErrMemDBIsClosed)

_, err = ds.OptimizedRevision(t.Context())
require.ErrorIs(err, ErrMemDBIsClosed)
require.ErrorIs(err, memdb.ErrMemDBIsClosed)

reader := ds.SnapshotReader(datastore.NoRevision)
_, err = reader.CountRelationships(t.Context(), "blah")
require.ErrorIs(err, ErrMemDBIsClosed)
require.ErrorIs(err, memdb.ErrMemDBIsClosed)
}
4 changes: 4 additions & 0 deletions internal/datastore/mysql/datastore.go
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,10 @@ func (mds *mysqlDatastore) MetricsID() (string, error) {
return common.MetricsIDFromURL(mds.url)
}

func (mds *mysqlDatastore) EngineName() string {
return Engine
}

func (mds *mysqlDatastore) SnapshotReader(rev datastore.Revision) datastore.Reader {
createTxFunc := func(ctx context.Context) (*sql.Tx, txCleanupFunc, error) {
tx, err := mds.db.BeginTx(ctx, mds.readTxOptions)
Expand Down
129 changes: 129 additions & 0 deletions internal/datastore/mysql/enginebuilder.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
package mysql

import (
"context"
"errors"
"fmt"

"github.com/ccoveille/go-safecast/v2"
sqlDriver "github.com/go-sql-driver/mysql"

"github.com/authzed/spicedb/internal/datastore/common"
"github.com/authzed/spicedb/internal/datastore/mysql/migrations"
"github.com/authzed/spicedb/internal/datastore/proxy"
log "github.com/authzed/spicedb/internal/logging"
datastorecfg "github.com/authzed/spicedb/pkg/cmd/datastore/dsconfig"
"github.com/authzed/spicedb/pkg/datastore"
"github.com/authzed/spicedb/pkg/datastore/migration"
)

func init() {
datastorecfg.RegisterEngine(Engine, newDatastoreFromConfig)
migration.RegisterMigratableEngine(Engine, migrations.Manager, newMigrationDriverFromConfig, "add_schema_tables")
}

func newMigrationDriverFromConfig(ctx context.Context, cfg *migration.Config) (*migrations.MySQLDriver, error) {
credentialsProvider, err := cfg.CredentialsProvider(ctx)
if err != nil {
return nil, err
}

// Do this outside NewMySQLDriverFromDSN to avoid races on MySQL datastore tests
if err := sqlDriver.SetLogger(&log.Logger); err != nil {
return nil, fmt.Errorf("unable to set logging to mysql driver: %w", err)
}

Check warning on line 34 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L33-L34

Added lines #L33 - L34 were not covered by tests

return migrations.NewMySQLDriverFromDSN(cfg.DatastoreURI, cfg.MySQLTablePrefix, credentialsProvider)
}

func newDatastoreFromConfig(ctx context.Context, opts datastorecfg.Config) (datastore.Datastore, error) {
primary, err := newPrimaryDatastoreFromConfig(ctx, opts)
if err != nil {
return nil, err
}

Check warning on line 43 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L42-L43

Added lines #L42 - L43 were not covered by tests

if len(opts.ReadReplicaURIs) > datastorecfg.MaxReplicaCount {
return nil, fmt.Errorf("too many read replicas, max is %d", datastorecfg.MaxReplicaCount)
}

Check warning on line 47 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L46-L47

Added lines #L46 - L47 were not covered by tests

replicas := make([]datastore.ReadOnlyDatastore, 0, len(opts.ReadReplicaURIs))
for index, replicaURI := range opts.ReadReplicaURIs {
uintIndex, err := safecast.Convert[uint32](index)
if err != nil {
return nil, errors.New("too many replicas")
}
replica, err := newReplicaDatastoreFromConfig(ctx, uintIndex, replicaURI, opts)
if err != nil {
return nil, err
}
replicas = append(replicas, replica)

Check warning on line 59 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L51-L59

Added lines #L51 - L59 were not covered by tests
}

return proxy.NewCheckingReplicatedDatastore(primary, replicas...)
}

func commonDatastoreOptionsFromConfig(opts datastorecfg.Config) ([]Option, error) {
maxRetries, err := safecast.Convert[uint8](opts.MaxRetries)
if err != nil {
return nil, errors.New("max-retries could not be cast to uint8")
}

Check warning on line 69 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L68-L69

Added lines #L68 - L69 were not covered by tests

watchChangeBufferMaximumSize, err := common.WatchBufferSize(opts.WatchChangeBufferMaximumSize)
if err != nil {
return nil, err
}

Check warning on line 74 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L73-L74

Added lines #L73 - L74 were not covered by tests

return []Option{
TablePrefix(opts.TablePrefix),
MaxRetries(maxRetries),
OverrideLockWaitTimeout(1),
WithEnablePrometheusStats(opts.EnableDatastoreMetrics),
WatchBufferLength(opts.WatchBufferLength),
WatchBufferWriteTimeout(opts.WatchBufferWriteTimeout),
WatchChangeBufferMaximumSize(watchChangeBufferMaximumSize),
MaxRevisionStalenessPercent(opts.MaxRevisionStalenessPercent),
RevisionQuantization(opts.RevisionQuantization),
FilterMaximumIDCount(opts.FilterMaximumIDCount),
AllowedMigrations(opts.AllowedMigrations),
WithColumnOptimization(opts.ExperimentalColumnOptimization),
}, nil
}

func newReplicaDatastoreFromConfig(ctx context.Context, replicaIndex uint32, replicaURI string, opts datastorecfg.Config) (datastore.ReadOnlyDatastore, error) {
mysqlOpts := []Option{ //nolint: prealloc // we're not concerned about perf here
MaxOpenConns(opts.ReadReplicaConnPool.MaxOpenConns),
ConnMaxIdleTime(opts.ReadReplicaConnPool.MaxIdleTime),
ConnMaxLifetime(opts.ReadReplicaConnPool.MaxLifetime),
CredentialsProviderName(opts.ReadReplicaCredentialsProviderName),
}

commonOptions, err := commonDatastoreOptionsFromConfig(opts)
if err != nil {
return nil, err
}
mysqlOpts = append(mysqlOpts, commonOptions...)
return NewReadOnlyMySQLDatastore(ctx, replicaURI, replicaIndex, mysqlOpts...)

Check warning on line 105 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L92-L105

Added lines #L92 - L105 were not covered by tests
}

func newPrimaryDatastoreFromConfig(ctx context.Context, opts datastorecfg.Config) (datastore.Datastore, error) {
mysqlOpts := []Option{ //nolint: prealloc // we're not concerned about perf here
GCInterval(opts.GCInterval),
GCWindow(opts.GCWindow),
GCInterval(opts.GCInterval),
GCEnabled(!opts.ReadOnly),
GCMaxOperationTime(opts.GCMaxOperationTime),
MaxOpenConns(opts.ReadConnPool.MaxOpenConns),
ConnMaxIdleTime(opts.ReadConnPool.MaxIdleTime),
ConnMaxLifetime(opts.ReadConnPool.MaxLifetime),
WithWatchDisabled(opts.DisableWatchSupport),
CredentialsProviderName(opts.CredentialsProviderName),
FollowerReadDelay(opts.FollowerReadDelay),
}

commonOptions, err := commonDatastoreOptionsFromConfig(opts)
if err != nil {
return nil, err
}

Check warning on line 126 in internal/datastore/mysql/enginebuilder.go

View check run for this annotation

Codecov / codecov/patch

internal/datastore/mysql/enginebuilder.go#L125-L126

Added lines #L125 - L126 were not covered by tests
mysqlOpts = append(mysqlOpts, commonOptions...)
return NewMySQLDatastore(ctx, opts.URI, mysqlOpts...)
}
Loading
Loading