From 8fe73146c8a7843c90d3392767c93a794331f02f Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Wed, 22 Mar 2023 08:54:13 -0500 Subject: [PATCH] Sqldb rip boltdb (#2341) * serverless sqldb use same env for test config as normal * rip boltdb implementation of controller backend out it was replaced by postgres and no longer works properly. This involved migrating a number of tests which only worked with boltdb, which exposed several ways in which the postgres implementation had slightly different behavior from the bolt one: 1. ordering of results in some cases, and 2. (more importantly) erroring when a record to delete was not found. The bolt implementation silently ignored it when things to delete weren't found, so we make some changes to match that behavior. Also stopped propagating CreatedAt and UpdatedAt from DB tables into dax types. These were breaking existing tests. Perhaps it would be better to actually use them, but for now they will only exist at the DB level. This change set also moves the insertion of the directive_versions record out of migrations and into the startup/connection code. Having this in the migrations was a bit ugly because you couldn't just truncate all the tables and have everything work from scratch. Inserting it during startup is fairly innocuous, and will just continue on if it already exists. * update directive_version test I changed the initial value to 0 so that the first version that gets sent out is 1 --- .gitlab/.gitlab-ci.yml | 24 +- ctl/dax.go | 1 - dax/boltdb/boltdb.go | 165 --- dax/boltdb/boltdb_test.go | 17 - dax/boltdb/directiveversion.go | 60 - dax/controller/balancer/balancer_test.go | 1097 ----------------- dax/controller/balancer/boltdb/balancer.go | 823 ------------- dax/controller/balancer/boltdb/node.go | 144 --- dax/controller/balancer/boltdb/node_test.go | 61 - dax/controller/config.go | 1 - dax/controller/controller_test.go | 151 ++- dax/controller/schemar/boltdb/schemar.go | 698 ----------- dax/controller/schemar/boltdb/schemar_test.go | 227 ---- dax/controller/service/controller.go | 32 +- dax/controller/sqldb/balancer.go | 14 +- dax/controller/sqldb/node.go | 7 +- dax/controller/sqldb/schemar.go | 11 +- dax/controller/sqldb/test.go | 18 +- dax/controller/sqldb/transactor.go | 6 + dax/controller/sqldb/workerjob.go | 5 +- dax/directive_version_test.go | 2 +- dax/docker-compose.yml | 64 - dax/migrations/001_initial.up.fizz | 1 - dax/models/database.go | 2 +- dax/models/node.go | 2 +- dax/server/test/managed.go | 1 - dax/snapshotter/snapshotter_test.go | 70 -- dax/table.go | 2 - dax/test/boltdb/helpers.go | 50 - dax/test/schemar.go | 29 - idk/docker-compose.yml | 4 - schema.go | 14 +- 32 files changed, 146 insertions(+), 3657 deletions(-) delete mode 100644 dax/boltdb/boltdb.go delete mode 100644 dax/boltdb/boltdb_test.go delete mode 100644 dax/boltdb/directiveversion.go delete mode 100644 dax/controller/balancer/balancer_test.go delete mode 100644 dax/controller/balancer/boltdb/balancer.go delete mode 100644 dax/controller/balancer/boltdb/node.go delete mode 100644 dax/controller/balancer/boltdb/node_test.go delete mode 100644 dax/controller/schemar/boltdb/schemar.go delete mode 100644 dax/controller/schemar/boltdb/schemar_test.go delete mode 100644 dax/docker-compose.yml delete mode 100644 dax/snapshotter/snapshotter_test.go delete mode 100644 dax/test/boltdb/helpers.go delete mode 100644 dax/test/schemar.go diff --git a/.gitlab/.gitlab-ci.yml b/.gitlab/.gitlab-ci.yml index b2797844c..c76a22872 100644 --- a/.gitlab/.gitlab-ci.yml +++ b/.gitlab/.gitlab-ci.yml @@ -256,13 +256,13 @@ run go tests race: image: golang:$GOVERSION extends: .go-cache variables: - SQLDB_DB: run_go_tests_race + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_DATABASE: run_go_tests_race POSTGRES_DB: run_go_tests_race - SQLDB_USER: postgres + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_USER: postgres POSTGRES_USER: postgres - SQLDB_PASSWORD: $POSTGRES_PASSWORD + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PASSWORD: $POSTGRES_PASSWORD POSTGRES_PASSWORD: $POSTGRES_PASSWORD - SQLDB_HOST: postgres + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_HOST: postgres services: - postgres:14.7 rules: @@ -287,13 +287,13 @@ run go tests: image: golang:$GOVERSION extends: .go-cache variables: - SQLDB_DB: run_go_tests + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_DATABASE: run_go_tests POSTGRES_DB: run_go_tests - SQLDB_USER: postgres + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_USER: postgres POSTGRES_USER: postgres - SQLDB_PASSWORD: $POSTGRES_PASSWORD + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PASSWORD: $POSTGRES_PASSWORD POSTGRES_PASSWORD: $POSTGRES_PASSWORD - SQLDB_HOST: postgres + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_HOST: postgres services: - postgres:14.7 rules: @@ -319,13 +319,13 @@ run go tests dax/test/dax: tags: - docker variables: - SQLDB_DB: run_go_tests_dax + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_DATABASE: run_go_tests_dax POSTGRES_DB: run_go_tests_dax - SQLDB_USER: postgres + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_USER: postgres POSTGRES_USER: postgres - SQLDB_PASSWORD: $POSTGRES_PASSWORD + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PASSWORD: $POSTGRES_PASSWORD POSTGRES_PASSWORD: $POSTGRES_PASSWORD - SQLDB_HOST: postgres + FEATUREBASE_CONTROLLER_CONFIG_SQLDB_HOST: postgres services: - postgres:14.7 rules: diff --git a/ctl/dax.go b/ctl/dax.go index 1e2f537e0..7dd3cbedc 100644 --- a/ctl/dax.go +++ b/ctl/dax.go @@ -17,7 +17,6 @@ func BuildDAXFlags(cmd *cobra.Command, srv *server.Command) { // Controller flags.BoolVar(&srv.Config.Controller.Run, "controller.run", srv.Config.Controller.Run, "Run the Controller service in process.") flags.DurationVar(&srv.Config.Controller.Config.RegistrationBatchTimeout, "controller.config.registration-batch-timeout", srv.Config.Controller.Config.RegistrationBatchTimeout, "Timeout for node registration batches.") - flags.StringVar(&srv.Config.Controller.Config.DataDir, "controller.config.data-dir", srv.Config.Controller.Config.DataDir, "Controller directory to use in process.") flags.StringVar(&srv.Config.Controller.Config.StorageMethod, "controller.config.storage-method", srv.Config.Controller.Config.StorageMethod, "Backing store. boltdb or sqldb.") flags.DurationVar(&srv.Config.Controller.Config.SnappingTurtleTimeout, "controller.config.snapping-turtle-timeout", srv.Config.Controller.Config.SnappingTurtleTimeout, "Period for running automatic snapshotting routine.") diff --git a/dax/boltdb/boltdb.go b/dax/boltdb/boltdb.go deleted file mode 100644 index 3210fc637..000000000 --- a/dax/boltdb/boltdb.go +++ /dev/null @@ -1,165 +0,0 @@ -// Package boltdb contains the boltdb implementations of the DAX interfaces. -package boltdb - -import ( - "context" - "os" - "path/filepath" - "strings" - "time" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/errors" - bolt "go.etcd.io/bbolt" -) - -const ( - ErrFmtBucketNotFound = "boltdb: bucket '%s' not found" -) - -type Bucket []byte - -// DB represents the database connection. -type DB struct { - db *bolt.DB - ctx context.Context // background context - cancel func() // cancel background context - - // Datasource name. - DSN string - - // Destination for events to be published. - // EventService wtf.EventService - - // Returns the current time. Defaults to time.Now(). - // Can be mocked for tests. - Now func() time.Time - - filePath string - - // bucketQueue contains a list of buckets to create upon Open. - bucketQueue []Bucket -} - -// NewDB returns a new instance of DB associated with the given datasource name. -func NewDB(dsn string) *DB { - db := &DB{ - DSN: dsn, - Now: time.Now, - - //EventService: wtf.NopEventService(), - } - db.ctx, db.cancel = context.WithCancel(context.Background()) - return db -} - -// NewSvcBolt gets, opens, and creates buckets for a boltDB for a -// particular named service (the data file will be named after the -// service). -func NewSvcBolt(dir, svc string, buckets ...Bucket) (*DB, error) { - dir = strings.TrimPrefix(dir, "file:") - filename := filepath.Join(dir, svc+".boltdb") - db := NewDB("file:" + filename) - db.RegisterBuckets(buckets...) - err := db.Open() - return db, errors.Wrap(err, "opening") -} - -// path returns the file path to the boltdb database file. -func (db *DB) path() (string, error) { - if !strings.HasPrefix(db.DSN, "file:") { - return "", errors.New(errors.ErrUncoded, "boltdb package only supports a DSN beginning with `file:`") - } - - return db.DSN[5:], nil -} - -// RegisterBuckets queues up the buckets to be created when the database is -// first opened. -func (db *DB) RegisterBuckets(buckets ...Bucket) { - db.bucketQueue = append(db.bucketQueue, buckets...) -} - -// InitializeBuckets creates the given buckets if they do not already exist. -func (db *DB) InitializeBuckets(buckets ...Bucket) (err error) { - return db.db.Update(func(tx *bolt.Tx) error { - for _, bucket := range buckets { - if _, err := tx.CreateBucketIfNotExists(bucket); err != nil { - return errors.Wrapf(err, "creating bucket: %s", bucket) - } - } - return nil - }) -} - -// Start is here to implement the Transactor interface, but we don't really need -// it in the BoltDB implementation. -func (db *DB) Start() (err error) { - return nil -} - -// Open opens the database connection. -func (db *DB) Open() (err error) { - path, err := db.path() - if err != nil { - return errors.Wrap(err, "getting path from DSN") - } - - if err := os.MkdirAll(filepath.Dir(path), 0777); err != nil { - return errors.Wrapf(err, "mkdir %s", filepath.Dir(path)) - } else if db.db, err = bolt.Open(path, 0666, &bolt.Options{Timeout: 1 * time.Second}); err != nil { - return errors.Wrapf(err, "open file: %s", err) - } - - // cache the path in db.filePath. - db.filePath = path - - if err := db.InitializeBuckets(db.bucketQueue...); err != nil { - return errors.Wrap(err, "initializing buckets") - } - - // Reset the bucketQueue. - db.bucketQueue = make([]Bucket, 0) - - return nil -} - -// Close closes the database connection. -func (db *DB) Close() (err error) { - return db.db.Close() -} - -// BeginTx starts a transaction and returns a wrapper Tx type. This type -// provides a reference to the database and a fixed timestamp at the start of -// the transaction. The timestamp allows us to mock time during tests as well. -// The wrapper also contains the context. -func (db *DB) BeginTx(ctx context.Context, writable bool) (dax.Transaction, error) { - tx, err := db.db.Begin(writable) - if err != nil { - return nil, err - } - - // Return wrapper Tx that includes the transaction start time. - return &Tx{ - Tx: tx, - ctx: ctx, - db: db, - now: db.Now().UTC().Truncate(time.Second), - }, nil -} - -// Tx wraps the SQL Tx object to provide a timestamp at the start of the transaction. -type Tx struct { - *bolt.Tx - ctx context.Context - db *DB - now time.Time -} - -func (tx *Tx) Context() context.Context { - return tx.ctx -} - -func (db *DB) Path() string { - return db.filePath -} diff --git a/dax/boltdb/boltdb_test.go b/dax/boltdb/boltdb_test.go deleted file mode 100644 index 205d4ba4a..000000000 --- a/dax/boltdb/boltdb_test.go +++ /dev/null @@ -1,17 +0,0 @@ -package boltdb_test - -import ( - "testing" - - "github.com/featurebasedb/featurebase/v3/dax/test/boltdb" -) - -// Ensure the test database can open & close. -func TestDB(t *testing.T) { - db := boltdb.MustOpenDB(t) - defer boltdb.MustCloseDB(t, db) - - t.Cleanup(func() { - boltdb.CleanupDB(t, db.Path()) - }) -} diff --git a/dax/boltdb/directiveversion.go b/dax/boltdb/directiveversion.go deleted file mode 100644 index b4a900f07..000000000 --- a/dax/boltdb/directiveversion.go +++ /dev/null @@ -1,60 +0,0 @@ -package boltdb - -import ( - "encoding/binary" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/errors" -) - -var ( - bucketDirective = Bucket("nodeDirective") - keyDirectiveVersion = []byte("directiveVersion") -) - -// DirectiveBuckets defines the buckets used by this package. It can be called -// during setup to create the buckets ahead of time. -var DirectiveBuckets []Bucket = []Bucket{ - bucketDirective, -} - -// Ensure type implements interface. -var _ dax.DirectiveVersion = (*DirectiveVersion)(nil) - -type DirectiveVersion struct { - db *DB -} - -func NewDirectiveVersion(db *DB) *DirectiveVersion { - return &DirectiveVersion{ - db: db, - } -} - -func (d *DirectiveVersion) Increment(tx dax.Transaction, delta uint64) (uint64, error) { - txx, ok := tx.(*Tx) - if !ok { - return 0, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketDirective) - if bkt == nil { - return 0, errors.Errorf(ErrFmtBucketNotFound, bucketDirective) - } - - var nextVersion uint64 = 1 // Start at 1; 0 is an invalid version. - - b := bkt.Get(keyDirectiveVersion) - if b != nil { - nextVersion = binary.LittleEndian.Uint64(b) + delta - } - - vsn := make([]byte, 8) - binary.LittleEndian.PutUint64(vsn, nextVersion) - - if err := bkt.Put(keyDirectiveVersion, vsn); err != nil { - return 0, errors.Wrap(err, "putting next directive version") - } - - return nextVersion, nil -} diff --git a/dax/controller/balancer/balancer_test.go b/dax/controller/balancer/balancer_test.go deleted file mode 100644 index 4db187e8b..000000000 --- a/dax/controller/balancer/balancer_test.go +++ /dev/null @@ -1,1097 +0,0 @@ -package balancer_test - -import ( - "context" - "fmt" - "os" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - daxbolt "github.com/featurebasedb/featurebase/v3/dax/boltdb" - "github.com/featurebasedb/featurebase/v3/dax/controller" - "github.com/featurebasedb/featurebase/v3/dax/controller/balancer/boltdb" - schemardb "github.com/featurebasedb/featurebase/v3/dax/controller/schemar/boltdb" - daxtest "github.com/featurebasedb/featurebase/v3/dax/test" - testbolt "github.com/featurebasedb/featurebase/v3/dax/test/boltdb" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -func newBoltBalancer(t *testing.T) (*daxbolt.DB, func()) { - db := testbolt.MustOpenDB(t) - require.NoError(t, db.InitializeBuckets(boltdb.BalancerBuckets...)) - require.NoError(t, db.InitializeBuckets(schemardb.SchemarBuckets...)) - - return db, func() { - testbolt.MustCloseDB(t, db) - testbolt.CleanupDB(t, db.Path()) - } -} - -type runner interface { - run(tx dax.Transaction, bal controller.Balancer) ([]dax.WorkerDiff, error) -} - -type addWorker struct { - node *dax.Node -} - -func (r *addWorker) run(tx dax.Transaction, bal controller.Balancer) ([]dax.WorkerDiff, error) { - return bal.AddWorker(tx, r.node) -} - -type removeWorker struct { - addr dax.Address -} - -func (r *removeWorker) run(tx dax.Transaction, bal controller.Balancer) ([]dax.WorkerDiff, error) { - return bal.RemoveWorker(tx, r.addr) -} - -type addJob struct { - roleType dax.RoleType - qtid dax.QualifiedTableID - job dax.Job -} - -func (r *addJob) run(tx dax.Transaction, bal controller.Balancer) ([]dax.WorkerDiff, error) { - return bal.AddJobs(tx, r.roleType, r.qtid, r.job) -} - -type removeJob struct { - roleType dax.RoleType - qtid dax.QualifiedTableID - job dax.Job -} - -func (r *removeJob) run(tx dax.Transaction, bal controller.Balancer) ([]dax.WorkerDiff, error) { - return bal.RemoveJobs(tx, r.roleType, r.qtid, r.job) -} - -type balanceDatabase struct { - qdbid dax.QualifiedDatabaseID -} - -func (r *balanceDatabase) run(tx dax.Transaction, bal controller.Balancer) ([]dax.WorkerDiff, error) { - return bal.BalanceDatabase(tx, r.qdbid) -} - -func TestBalancer(t *testing.T) { - ctx := context.Background() - orgID := dax.OrganizationID("acme") - dbID := dax.DatabaseID("db1") - dbName := dax.DatabaseName("db1name") - tableID := dax.TableID("tbl1") - qtid := dax.NewQualifiedTableID(dax.NewQualifiedDatabaseID(orgID, dbID), tableID) - - t.Run("SingleWorker", func(t *testing.T) { - qdb := &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ - ID: dbID, - Name: dbName, - Options: dax.DatabaseOptions{ - WorkersMin: 1, - WorkersMax: 1, - }, - }, - } - - schemar, cleanup := daxtest.NewSchemar(t) - defer cleanup() - - db, cleanup := newBoltBalancer(t) - defer cleanup() - bal := boltdb.NewBalancer(db, schemar, logger.NewStandardLogger(os.Stderr)) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - assert.NoError(t, schemar.CreateDatabase(tx, qdb)) - - role := dax.RoleTypeCompute - - tests := []struct { - runner runner - expDiff []dax.WorkerDiff - expState []dax.WorkerInfo - }{ - { - // Add job. - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s2", - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{}, - }, - { - // Add worker. - runner: &addWorker{ - node: &dax.Node{ - Address: "w1", - RoleTypes: dax.RoleTypes{role}, - }, - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w1", - AddedJobs: []dax.Job{"s2"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w1", - Jobs: []dax.Job{"s2"}, - }, - }, - }, - { - // Add another job out of order. - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s1", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w1", - AddedJobs: []dax.Job{"s1"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w1", - Jobs: []dax.Job{"s1", "s2"}, - }, - }, - }, - { - // Add another job. - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s3", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w1", - AddedJobs: []dax.Job{"s3"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w1", - Jobs: []dax.Job{"s1", "s2", "s3"}, - }, - }, - }, - { - // Add a duplicate job. - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s2", - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{ - { - Address: "w1", - Jobs: []dax.Job{"s1", "s2", "s3"}, - }, - }, - }, - } - for i, test := range tests { - t.Run(fmt.Sprintf("test-%d", i), func(t *testing.T) { - diff, err := test.runner.run(tx, bal) - assert.NoError(t, err) - assert.Equal(t, test.expDiff, diff) - - state, err := bal.CurrentState(tx, role, qdb.QualifiedID()) - assert.NoError(t, err) - assert.Equal(t, test.expState, state) - }) - } - - assert.NoError(t, tx.Commit()) - }) - - t.Run("MultipleWorkers", func(t *testing.T) { - dbOptions := dax.DatabaseOptions{ - WorkersMin: 2, - WorkersMax: 2, - } - - qdb := &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ - ID: dbID, - Name: dbName, - Options: dbOptions, - }, - } - - schemar, cleanup := daxtest.NewSchemar(t) - defer cleanup() - - db, cleanup := newBoltBalancer(t) - defer cleanup() - bal := boltdb.NewBalancer(db, schemar, logger.NewStandardLogger(os.Stderr)) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - assert.NoError(t, schemar.CreateDatabase(tx, qdb)) - - role := dax.RoleTypeCompute - - type testPart struct { - name string - runner runner - expDiff []dax.WorkerDiff - expState []dax.WorkerInfo - } - - runTestPart := func(tp testPart) { - t.Run(fmt.Sprintf("test-%s", tp.name), func(t *testing.T) { - diff, err := tp.runner.run(tx, bal) - assert.NoError(t, err) - assert.Equal(t, tp.expDiff, diff) - - state, err := bal.CurrentState(tx, role, qdb.QualifiedID()) - assert.NoError(t, err) - assert.Equal(t, tp.expState, state) - }) - } - - runTestPart(testPart{ - name: "balance when empty", - runner: &balanceDatabase{ - qdbid: qdb.QualifiedID(), - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{}, - }) - - runTestPart(testPart{ - name: "add worker", - runner: &addWorker{ - node: &dax.Node{ - Address: "w2", - RoleTypes: dax.RoleTypes{role}, - }, - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{}, - }) - - runTestPart(testPart{ - name: "add worker again", - runner: &addWorker{ - node: &dax.Node{ - Address: "w2", - RoleTypes: dax.RoleTypes{role}, - }, - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{}, - }) - - runTestPart(testPart{ - name: "add a second worker", - runner: &addWorker{ - node: &dax.Node{ - Address: "w1", - RoleTypes: dax.RoleTypes{role}, - }, - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{}, - }) - - runTestPart(testPart{ - name: "add job 2", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s2", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w1", - AddedJobs: []dax.Job{"s2"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w1", - Jobs: []dax.Job{"s2"}, - }, - { - Address: "w2", - Jobs: []dax.Job{}, - }, - }, - }) - - runTestPart(testPart{ - name: "add job 3", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s3", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w2", - AddedJobs: []dax.Job{"s3"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w1", - Jobs: []dax.Job{"s2"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s3"}, - }, - }, - }) - - runTestPart(testPart{ - name: "add job 1", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s1", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w1", - AddedJobs: []dax.Job{"s1"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w1", - Jobs: []dax.Job{"s1", "s2"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s3"}, - }, - }, - }) - - // Update database options on schemar so min worker for db is 3. Here we - // are updating the database option first, and then adding a worker to - // satisfy those options. - t.Run(fmt.Sprintf("test-%s", "set database options db min 3"), func(t *testing.T) { - assert.NoError(t, schemar.SetDatabaseOption(tx, qdb.QualifiedID(), dax.DatabaseOptionWorkersMin, "3")) - }) - - runTestPart(testPart{ - name: "add a third worker", - runner: &addWorker{ - node: &dax.Node{ - Address: "w0", - RoleTypes: dax.RoleTypes{role}, - }, - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w0", - AddedJobs: []dax.Job{"s2"}, - RemovedJobs: []dax.Job{}, - }, - { - Address: "w1", - AddedJobs: []dax.Job{}, - RemovedJobs: []dax.Job{"s2"}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2"}, - }, - { - Address: "w1", - Jobs: []dax.Job{"s1"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s3"}, - }, - }, - }) - - runTestPart(testPart{ - name: "add job 4", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s4", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w0", - AddedJobs: []dax.Job{"s4"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4"}, - }, - { - Address: "w1", - Jobs: []dax.Job{"s1"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s3"}, - }, - }, - }) - - runTestPart(testPart{ - name: "add job 5", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s5", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w1", - AddedJobs: []dax.Job{"s5"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4"}, - }, - { - Address: "w1", - Jobs: []dax.Job{"s1", "s5"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s3"}, - }, - }, - }) - - runTestPart(testPart{ - name: "add job 0", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s0", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w2", - AddedJobs: []dax.Job{"s0"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4"}, - }, - { - Address: "w1", - Jobs: []dax.Job{"s1", "s5"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s0", "s3"}, - }, - }, - }) - - runTestPart(testPart{ - name: "add job 6", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s6", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w0", - AddedJobs: []dax.Job{"s6"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4", "s6"}, - }, - { - Address: "w1", - Jobs: []dax.Job{"s1", "s5"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s0", "s3"}, - }, - }, - }) - - runTestPart(testPart{ - name: "add job 7", - runner: &addJob{ - roleType: role, - qtid: qtid, - job: "s7", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w1", - AddedJobs: []dax.Job{"s7"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4", "s6"}, - }, - { - Address: "w1", - Jobs: []dax.Job{"s1", "s5", "s7"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s0", "s3"}, - }, - }, - }) - - //////////////////// Remove ///////////////////////// - - runTestPart(testPart{ - name: "remove nonexistent worker", - runner: &removeWorker{ - addr: "nonexistent", - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4", "s6"}, - }, - { - Address: "w1", - Jobs: []dax.Job{"s1", "s5", "s7"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s0", "s3"}, - }, - }, - }) - - runTestPart(testPart{ - name: "remove worker", - runner: &removeWorker{ - addr: "w1", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w0", - AddedJobs: []dax.Job{"s5"}, - RemovedJobs: []dax.Job{}, - }, - { - Address: "w1", - AddedJobs: []dax.Job{}, - RemovedJobs: []dax.Job{"s1", "s5", "s7"}, - }, - { - Address: "w2", - AddedJobs: []dax.Job{"s1", "s7"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4", "s5", "s6"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s0", "s1", "s3", "s7"}, - }, - }, - }) - - runTestPart(testPart{ - name: "remove active job", - runner: &removeJob{ - roleType: role, - qtid: qtid, - job: "s0", - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w2", - AddedJobs: []dax.Job{}, - RemovedJobs: []dax.Job{"s0"}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4", "s5", "s6"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s1", "s3", "s7"}, - }, - }, - }) - - // Set WorkersMin back to 2 so we can test the change to 3 again. - t.Run(fmt.Sprintf("test-%s", "set database options db min back to 2"), func(t *testing.T) { - assert.NoError(t, schemar.SetDatabaseOption(tx, qdb.QualifiedID(), dax.DatabaseOptionWorkersMin, "2")) - }) - - runTestPart(testPart{ - name: "add a fourth worker", - runner: &addWorker{ - node: &dax.Node{ - Address: "w3", - RoleTypes: dax.RoleTypes{role}, - }, - }, - expDiff: []dax.WorkerDiff{}, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4", "s5", "s6"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s1", "s3", "s7"}, - }, - }, - }) - - // Update database options on schemar so min worker for db is 4. Here we - // have added a worker which will satisfy this option, and then updated - // the option. - t.Run(fmt.Sprintf("test-%s", "set database options db min back to 3"), func(t *testing.T) { - assert.NoError(t, schemar.SetDatabaseOption(tx, qdb.QualifiedID(), dax.DatabaseOptionWorkersMin, "3")) - }) - - // This implies that there is a condition where the database does not - // get repaired to its workerMin: if workers per database has dropped - // below its min, and there are no available workers to replace the - // missing workers (and bring it back to min), then the database will - // operate below min. If, then, a worker is added to the pool, we do not - // currently have a process to automatically assign that worker to a - // database under min. That happens with an explicit call to - // assignMinWorkers or BalanceDatabase. - runTestPart(testPart{ - name: "balance to include new third worker", - runner: &balanceDatabase{ - qdbid: qdb.QualifiedID(), - }, - expDiff: []dax.WorkerDiff{ - { - Address: "w0", - AddedJobs: []dax.Job{}, - RemovedJobs: []dax.Job{"s6"}, - }, - { - Address: "w2", - AddedJobs: []dax.Job{}, - RemovedJobs: []dax.Job{"s7"}, - }, - { - Address: "w3", - AddedJobs: []dax.Job{"s6", "s7"}, - RemovedJobs: []dax.Job{}, - }, - }, - expState: []dax.WorkerInfo{ - { - Address: "w0", - Jobs: []dax.Job{"s2", "s4", "s5"}, - }, - { - Address: "w2", - Jobs: []dax.Job{"s1", "s3"}, - }, - { - Address: "w3", - Jobs: []dax.Job{"s6", "s7"}, - }, - }, - }) - - assert.NoError(t, tx.Commit()) - }) - - t.Run("WorkersForJobs", func(t *testing.T) { - qdb := &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ - ID: dbID, - Name: dbName, - Options: dax.DatabaseOptions{ - WorkersMin: 2, - WorkersMax: 2, - }, - }, - } - qdbid := qdb.QualifiedID() - - schemar, scleanup := daxtest.NewSchemar(t) - defer scleanup() - - role := dax.RoleTypeCompute - - db, cleanup := newBoltBalancer(t) - defer cleanup() - bal := boltdb.NewBalancer(db, schemar, logger.NewStandardLogger(os.Stderr)) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - assert.NoError(t, schemar.CreateDatabase(tx, qdb)) - - node1 := &dax.Node{ - Address: "n1", - RoleTypes: []dax.RoleType{role}, - } - - node2 := &dax.Node{ - Address: "n2", - RoleTypes: []dax.RoleType{role}, - } - - _, err = bal.AddWorker(tx, node1) - assert.NoError(t, err) - _, err = bal.AddWorker(tx, node2) - assert.NoError(t, err) - for i := 0; i < 12; i++ { - _, err = bal.AddJobs(tx, role, qtid, dax.Job(fmt.Sprintf("s%d", i))) - assert.NoError(t, err) - } - - exp := dax.WorkerInfo{ - Address: "n1", - Jobs: []dax.Job{"s0", "s10", "s2", "s4", "s6", "s8"}, - } - ws, err := bal.WorkerState(tx, role, "n1") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - exp = dax.WorkerInfo{ - Address: "n2", - Jobs: []dax.Job{"s1", "s11", "s3", "s5", "s7", "s9"}, - } - ws, err = bal.WorkerState(tx, role, "n2") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - tests := []struct { - jobs []dax.Job - exp []dax.WorkerInfo - }{ - { - jobs: []dax.Job{"s0"}, - exp: []dax.WorkerInfo{ - {Address: "n1", Jobs: []dax.Job{"s0"}}, - }, - }, - { - jobs: []dax.Job{"s0", "s4"}, - exp: []dax.WorkerInfo{ - {Address: "n1", Jobs: []dax.Job{"s0", "s4"}}, - }, - }, - { - jobs: []dax.Job{"s0", "s4", "s999"}, - exp: []dax.WorkerInfo{ - {Address: "n1", Jobs: []dax.Job{"s0", "s4"}}, - }, - }, - { - jobs: []dax.Job{"s0", "s1"}, - exp: []dax.WorkerInfo{ - {Address: "n1", Jobs: []dax.Job{"s0"}}, - {Address: "n2", Jobs: []dax.Job{"s1"}}, - }, - }, - { - jobs: []dax.Job{"s5", "s0", "s1", "s8"}, - exp: []dax.WorkerInfo{ - {Address: "n1", Jobs: []dax.Job{"s0", "s8"}}, - {Address: "n2", Jobs: []dax.Job{"s1", "s5"}}, - }, - }, - } - for i, test := range tests { - t.Run(fmt.Sprintf("test-%d", i), func(t *testing.T) { - workers, err := bal.WorkersForJobs(tx, role, qdbid, test.jobs...) - assert.NoError(t, err) - assert.Equal(t, test.exp, workers) - }) - } - - assert.NoError(t, tx.Commit()) - }) - - t.Run("WorkersForTable", func(t *testing.T) { - qdb := &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ - ID: dbID, - Name: dbName, - Options: dax.DatabaseOptions{ - WorkersMin: 2, - WorkersMax: 2, - }, - }, - } - - qdbid := dax.NewQualifiedDatabaseID(orgID, dbID) - - // Table 1. - tbl1 := &dax.Table{ - ID: "id1", - Name: "table1", - } - qtbl1 := dax.NewQualifiedTable(qdbid, tbl1) - qtid1 := qtbl1.QualifiedID() - - // Table 2. - tbl2 := &dax.Table{ - ID: "id2", - Name: "table2", - } - qtbl2 := dax.NewQualifiedTable(qdbid, tbl2) - qtid2 := qtbl2.QualifiedID() - - schemar, scleanup := daxtest.NewSchemar(t) - defer scleanup() - - role := dax.RoleTypeCompute - - db, cleanup := newBoltBalancer(t) - defer cleanup() - bal := boltdb.NewBalancer(db, schemar, logger.NewStandardLogger(os.Stderr)) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - assert.NoError(t, schemar.CreateDatabase(tx, qdb)) - - node1 := &dax.Node{ - Address: "n1", - RoleTypes: []dax.RoleType{role}, - } - - node2 := &dax.Node{ - Address: "n2", - RoleTypes: []dax.RoleType{role}, - } - - _, err = bal.AddWorker(tx, node1) - assert.NoError(t, err) - _, err = bal.AddWorker(tx, node2) - assert.NoError(t, err) - for i := 0; i < 12; i++ { - job := fmt.Sprintf("%s|s%d", qtid1.Key(), i) - _, err = bal.AddJobs(tx, role, qtid1, dax.Job(job)) - assert.NoError(t, err) - } - - for i := 9; i < 16; i++ { - job := fmt.Sprintf("%s|s%d", qtid2.Key(), i) - _, err = bal.AddJobs(tx, role, qtid2, dax.Job(job)) - assert.NoError(t, err) - } - - // fn1 and fn2 are just helper functions used to make the tests easier - // to read. - fn1 := func(s string) dax.Job { - return dax.Job(fmt.Sprintf("tbl__acme__db1__id1|%s", s)) - } - fn2 := func(s string) dax.Job { - return dax.Job(fmt.Sprintf("tbl__acme__db1__id2|%s", s)) - } - - // Table 1 - workers, err := bal.WorkersForTable(tx, role, qtid1) - assert.NoError(t, err) - assert.ElementsMatch(t, []dax.WorkerInfo{ - {Address: "n1", Jobs: []dax.Job{fn1("s0"), fn1("s10"), fn1("s2"), fn1("s4"), fn1("s6"), fn1("s8")}}, - {Address: "n2", Jobs: []dax.Job{fn1("s1"), fn1("s11"), fn1("s3"), fn1("s5"), fn1("s7"), fn1("s9")}}, - }, workers) - - // Table 2 - workers, err = bal.WorkersForTable(tx, role, qtid2) - assert.NoError(t, err) - assert.ElementsMatch(t, []dax.WorkerInfo{ - {Address: "n1", Jobs: []dax.Job{fn2("s11"), fn2("s13"), fn2("s15"), fn2("s9")}}, - {Address: "n2", Jobs: []dax.Job{fn2("s10"), fn2("s12"), fn2("s14")}}, - }, workers) - - // No match - qtid0 := dax.NewQualifiedTableID(qdbid, "bad") - workers, err = bal.WorkersForTable(tx, role, qtid0) - assert.NoError(t, err) - assert.ElementsMatch(t, []dax.WorkerInfo{}, workers) - - assert.NoError(t, tx.Commit()) - }) - - t.Run("Balance", func(t *testing.T) { - dbOptions := dax.DatabaseOptions{ - WorkersMin: 2, - WorkersMax: 2, - } - qdb := &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ - ID: dbID, - Name: dbName, - Options: dbOptions, - }, - } - - qdbid := qdb.QualifiedID() - - schemar, scleanup := daxtest.NewSchemar(t) - defer scleanup() - - role := dax.RoleTypeCompute - - db, cleanup := newBoltBalancer(t) - defer cleanup() - bal := boltdb.NewBalancer(db, schemar, logger.NewStandardLogger(os.Stderr)) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - assert.NoError(t, schemar.CreateDatabase(tx, qdb)) - - node1 := &dax.Node{ - Address: "n1", - RoleTypes: []dax.RoleType{role}, - } - - node2 := &dax.Node{ - Address: "n2", - RoleTypes: []dax.RoleType{role}, - } - - node3 := &dax.Node{ - Address: "n3", - RoleTypes: []dax.RoleType{role}, - } - - // Add two workers with some jobs evenly spread across them. - _, err = bal.AddWorker(tx, node1) - assert.NoError(t, err) - _, err = bal.AddWorker(tx, node2) - assert.NoError(t, err) - for i := 0; i < 13; i++ { - job := fmt.Sprintf("s%d", i) - _, err = bal.AddJobs(tx, role, qtid, dax.Job(job)) - assert.NoError(t, err) - } - - exp := dax.WorkerInfo{ - Address: "n1", - Jobs: []dax.Job{"s0", "s10", "s12", "s2", "s4", "s6", "s8"}, - } - ws, err := bal.WorkerState(tx, role, "n1") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - exp = dax.WorkerInfo{ - Address: "n2", - Jobs: []dax.Job{"s1", "s11", "s3", "s5", "s7", "s9"}, - } - ws, err = bal.WorkerState(tx, role, "n2") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - // Update database options on schemar so min worker for db is 3. - t.Run(fmt.Sprintf("test-%s", "set database options db min 3"), func(t *testing.T) { - assert.NoError(t, schemar.SetDatabaseOption(tx, qdb.QualifiedID(), dax.DatabaseOptionWorkersMin, "3")) - }) - - // Now, add a worker and confirm that it has received some jobs. - _, err = bal.AddWorker(tx, node3) - assert.NoError(t, err) - exp = dax.WorkerInfo{ - Address: "n3", - Jobs: []dax.Job{"s6", "s7", "s8", "s9"}, - } - ws, err = bal.WorkerState(tx, role, "n3") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - // Finally, call Balance() and confirm that the appropriate jobs got are - // as expected (actually, this is no longer needed here because we - // automatically balance when we add workers, but calling it should - // effectively be a no-op). - _, err = bal.BalanceDatabase(tx, qdbid) - assert.NoError(t, err) - - exp = dax.WorkerInfo{ - Address: "n1", - Jobs: []dax.Job{"s0", "s10", "s12", "s2", "s4"}, - } - ws, err = bal.WorkerState(tx, role, "n1") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - exp = dax.WorkerInfo{ - Address: "n2", - Jobs: []dax.Job{"s1", "s11", "s3", "s5"}, - } - ws, err = bal.WorkerState(tx, role, "n2") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - exp = dax.WorkerInfo{ - Address: "n3", - Jobs: []dax.Job{"s6", "s7", "s8", "s9"}, - } - ws, err = bal.WorkerState(tx, role, "n3") - assert.NoError(t, err) - assert.Equal(t, exp, ws) - - assert.NoError(t, tx.Commit()) - }) -} diff --git a/dax/controller/balancer/boltdb/balancer.go b/dax/controller/balancer/boltdb/balancer.go deleted file mode 100644 index 34f89d7b0..000000000 --- a/dax/controller/balancer/boltdb/balancer.go +++ /dev/null @@ -1,823 +0,0 @@ -// Package boltdb contains the boltdb implementation of the Balancer interface. -package boltdb - -import ( - "bytes" - "encoding/json" - "fmt" - "strings" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/boltdb" - balancer "github.com/featurebasedb/featurebase/v3/dax/controller/balancer" - "github.com/featurebasedb/featurebase/v3/dax/controller/schemar" - "github.com/featurebasedb/featurebase/v3/errors" - "github.com/featurebasedb/featurebase/v3/logger" -) - -var ( - bucketBalancer = boltdb.Bucket("balancer") -) - -// BalancerBuckets defines the buckets used by this package. It can be -// called during setup to create the buckets ahead of time. -var BalancerBuckets []boltdb.Bucket = []boltdb.Bucket{ - bucketBalancer, -} - -// NewBalancer returns a new instance of controller.Balancer. -func NewBalancer(db *boltdb.DB, schemar schemar.Schemar, logger logger.Logger) *balancer.Balancer { - fjs := newFreeJobService(db) - wjs := newWorkerJobService(db, logger) - fws := newFreeWorkerService(db) - ns := NewNodeService(db, logger) - - return balancer.New(ns, fjs, wjs, fws, schemar, logger) -} - -// Ensure type implements interface. -var _ balancer.WorkerJobService = (*workerJobService)(nil) - -type workerJobService struct { - db *boltdb.DB - logger logger.Logger -} - -func newWorkerJobService(db *boltdb.DB, logger logger.Logger) *workerJobService { - return &workerJobService{ - db: db, - logger: logger, - } -} - -func (w *workerJobService) WorkersJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerInfo, error) { - workerInfos, err := w.getWorkerInfos(tx, roleType, qdbid) - if err != nil { - return nil, errors.Wrapf(err, "getting worker infos: %s", roleType) - } - - return workerInfos, nil -} - -func (w *workerJobService) WorkerCount(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (int, error) { - workers, err := w.getWorkers(tx, roleType, qdbid) - if err != nil { - return 0, errors.Wrapf(err, "getting workers: %s", roleType) - } - - return len(workers), nil -} - -func (w *workerJobService) ListWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Addresses, error) { - return w.getWorkers(tx, roleType, qdbid) -} - -func (w *workerJobService) getWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Addresses, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - c := txx.Bucket(bucketBalancer).Cursor() - - // Deserialize rows into Worker objects. - addrs := make(dax.Addresses, 0) - - prefix := []byte(fmt.Sprintf(prefixFmtWorkersDB, roleType, qdbid.Key())) - for k, v := c.Seek(prefix); k != nil && bytes.HasPrefix(k, prefix); k, v = c.Next() { - if v == nil { - w.logger.Printf("nil value for key: %s", k) - continue - } - - addr, err := keyWorker(k) - if err != nil { - return nil, errors.Wrapf(err, "getting worker from key: %s", k) - } - - addrs = append(addrs, addr) - } - - return addrs, nil -} - -func (w *workerJobService) getWorkerInfos(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.WorkerInfos, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - c := txx.Bucket(bucketBalancer).Cursor() - - // Deserialize rows into WorkerInfo objects. - workerInfos := make(dax.WorkerInfos, 0) - - var prefix []byte - empty := dax.QualifiedDatabaseID{} - if roleType == "" && qdbid == empty { - prefix = []byte("workers/role/") - } else { - prefix = []byte(fmt.Sprintf(prefixFmtWorkersDB, roleType, qdbid.Key())) - } - for k, v := c.Seek(prefix); k != nil && bytes.HasPrefix(k, prefix); k, v = c.Next() { - addr, err := keyWorker(k) - if err != nil { - return nil, errors.Wrapf(err, "getting worker from key: %s", k) - } - - jobs := dax.NewSet[dax.Job]() - if v != nil { - jobs, err = decodeJobSet(v) - if err != nil { - return nil, errors.Wrap(err, "decoding job set") - } - } - - workerInfo := dax.WorkerInfo{ - Address: addr, - Jobs: jobs.Sorted(), - } - - workerInfos = append(workerInfos, workerInfo) - } - - return workerInfos, nil -} - -func (w *workerJobService) CreateWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - // If this worker already exists, don't do anything. - wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr)) - if wrkr != nil { - return nil - } - - val := []byte("[]") - if err := bkt.Put(workerDBKey(roleType, qdbid, addr), val); err != nil { - return errors.Wrapf(err, "putting db worker: %s, %s", qdbid, addr) - } - - if err := bkt.Put(workerAssignedKey(addr), []byte(qdbid.Key())); err != nil { - return errors.Wrapf(err, "putting assigned worker: %s, %s", qdbid, addr) - } - - return nil -} - -func (w *workerJobService) DeleteWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - if err := bkt.Delete(workerDBKey(roleType, qdbid, addr)); err != nil { - return errors.Wrapf(err, "deleting node key: %s", workerDBKey(roleType, qdbid, addr)) - } - - if err := bkt.Delete(workerAssignedKey(addr)); err != nil { - return errors.Wrapf(err, "deleting assigned worker: %s", workerAssignedKey(addr)) - } - - return nil -} - -func (w *workerJobService) FreeWorkers(tx dax.Transaction, addrs ...dax.Address) error { - return nil -} - -func (w *workerJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, jobs ...dax.Job) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - jobset := dax.NewSet[dax.Job]() - var err error - - // get worker - wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr)) - if wrkr != nil { - jobset, err = decodeJobSet(wrkr) - if err != nil { - return errors.Wrap(err, "decoding job set") - } - } - - for _, job := range jobs { - jobset.Add(job) - } - val, err := encodeJobSet(jobset) - if err != nil { - return errors.Wrap(err, "encoding job set") - } - - if err := bkt.Put(workerDBKey(roleType, qdbid, addr), val); err != nil { - return errors.Wrap(err, "putting worker") - } - - return nil -} - -func (w *workerJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job dax.Job) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - // get worker - wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr)) - if wrkr == nil { - return nil - } - - jobset, err := decodeJobSet(wrkr) - if err != nil { - return errors.Wrap(err, "decoding job set") - } - if !jobset.Contains(job) { - return nil - } - - jobset.Remove(job) - val, err := encodeJobSet(jobset) - if err != nil { - return errors.Wrap(err, "encoding job set") - } - - if err := bkt.Put(workerDBKey(roleType, qdbid, addr), val); err != nil { - return errors.Wrap(err, "putting worker") - } - - return nil -} - -func (w *workerJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) (balancer.InternalDiffs, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - qdbid := qtid.QualifiedDatabaseID - prefix := string(qtid.Key()) - - workers, err := w.getWorkers(tx, roleType, qdbid) - if err != nil { - return nil, errors.Wrap(err, "getting workers") - } - - idiffs := balancer.NewInternalDiffs() - for _, worker := range workers { - // get worker - wrkr := bkt.Get(workerDBKey(roleType, qdbid, worker)) - if wrkr == nil { - panic("didn't find worker that should... definitely exist") - } - jobset, err := decodeJobSet(wrkr) - if err != nil { - return nil, errors.Wrap(err, "decoding job set") - } - - jobs := jobset.RemoveByPrefix(prefix) - for _, job := range jobs { - idiffs.Removed(worker, job) - } - val, err := encodeJobSet(jobset) - if err != nil { - return nil, errors.Wrap(err, "encoding job set") - } - - if err := bkt.Put(workerDBKey(roleType, qdbid, worker), val); err != nil { - return nil, errors.Wrap(err, "putting worker") - } - - } - - return idiffs, nil -} - -func (w *workerJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) (dax.Jobs, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - jobset := dax.NewSet[dax.Job]() - var err error - - // get worker - wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr)) - if wrkr != nil { - jobset, err = decodeJobSet(wrkr) - if err != nil { - return nil, errors.Wrap(err, "decoding job set") - } - } - - return jobset.Sorted(), nil -} - -func (w *workerJobService) JobCounts(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addrs ...dax.Address) (map[dax.Address]int, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - m := make(map[dax.Address]int) - - for _, addr := range addrs { - jobset := dax.NewSet[dax.Job]() - var err error - - // get worker - wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr)) - if wrkr != nil { - jobset, err = decodeJobSet(wrkr) - if err != nil { - return nil, errors.Wrap(err, "decoding job set") - } - } - - m[addr] = len(jobset) - } - - return m, nil -} - -func (w *workerJobService) DatabaseForWorker(tx dax.Transaction, addr dax.Address) dax.DatabaseKey { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return "" // TODO(tlt): return error here? - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return "" - } - - wrkr := bkt.Get(workerAssignedKey(addr)) - - return dax.DatabaseKey(wrkr) -} - -// encodeJobSet encode the jobSet into a JSON array of strings. -func encodeJobSet(jobSet dax.Set[dax.Job]) ([]byte, error) { - arr := jobSet.Sorted() - b, err := json.Marshal(arr) - if err != nil { - return nil, errors.Wrap(err, "marshalling json") - } - return b, nil -} - -// decodeJobSet decode the string (a JSON array of strings) into jobSet. -func decodeJobSet(v []byte) (dax.Set[dax.Job], error) { - var arr []string - err := json.Unmarshal(v, &arr) - if err != nil { - return nil, errors.Wrap(err, "unmarshalling json") - } - - js := dax.NewSet[dax.Job]() - for _, s := range arr { - js.Add(dax.Job(s)) - } - - return js, nil -} - -// encodeWorkerSet encode the workerSet into a JSON array of strings. -func encodeWorkerSet(workerSet dax.Set[dax.Address]) ([]byte, error) { - arr := workerSet.Sorted() - b, err := json.Marshal(arr) - if err != nil { - return nil, errors.Wrap(err, "marshalling json") - } - return b, nil -} - -// decodeWorkerSet decode the string (a JSON array of strings) into workerSet. -func decodeWorkerSet(v []byte) (dax.Set[dax.Address], error) { - var arr []string - err := json.Unmarshal(v, &arr) - if err != nil { - return nil, errors.Wrap(err, "unmarshalling json") - } - - ws := dax.NewSet[dax.Address]() - for _, s := range arr { - ws.Add(dax.Address(s)) - } - - return ws, nil -} - -// Ensure type implements interface. -var _ balancer.FreeJobService = (*freeJobService)(nil) - -type freeJobService struct { - db *boltdb.DB -} - -func newFreeJobService(db *boltdb.DB) *freeJobService { - return &freeJobService{ - db: db, - } -} - -func (f *freeJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs ...dax.Job) error { - return f.MergeJobs(tx, roleType, qdbid, jobs) -} - -func (f *freeJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job dax.Job) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - // get free jobs - fjs := bkt.Get(freeJobKey(roleType, qdbid)) - if fjs == nil { - return nil - } - - jobset, err := decodeJobSet(fjs) - if err != nil { - return errors.Wrap(err, "decoding job set") - } - if !jobset.Contains(job) { - return nil - } - - jobset.Remove(job) - val, err := encodeJobSet(jobset) - if err != nil { - return errors.Wrap(err, "encoding job set") - } - - if err := bkt.Put(freeJobKey(roleType, qdbid), val); err != nil { - return errors.Wrap(err, "putting free job") - } - - return nil -} - -func (f *freeJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - qdbid := qtid.QualifiedDatabaseID - prefix := string(qtid.Key()) - - // get free jobs - fjs := bkt.Get(freeJobKey(roleType, qdbid)) - if fjs == nil { - return nil - } - - jobset, err := decodeJobSet(fjs) - if err != nil { - return errors.Wrap(err, "decoding job set") - } - - jobset.RemoveByPrefix(prefix) - val, err := encodeJobSet(jobset) - if err != nil { - return errors.Wrap(err, "encoding job set") - } - - if err := bkt.Put(freeJobKey(roleType, qdbid), val); err != nil { - return errors.Wrap(err, "putting free job") - } - - return nil -} - -func (f *freeJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Jobs, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - jobset := dax.NewSet[dax.Job]() - var err error - - // get free jobs - fjs := bkt.Get(freeJobKey(roleType, qdbid)) - if fjs != nil { - jobset, err = decodeJobSet(fjs) - if err != nil { - return nil, errors.Wrap(err, "decoding job set") - } - } - - return jobset.Sorted(), nil -} - -func (f *freeJobService) MergeJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs dax.Jobs) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - jobset := dax.NewSet[dax.Job]() - var err error - - // get free jobs - fjs := bkt.Get(freeJobKey(roleType, qdbid)) - if fjs != nil { - jobset, err = decodeJobSet(fjs) - if err != nil { - return errors.Wrap(err, "decoding job set") - } - } - - for _, j := range jobs { - jobset.Add(j) - } - val, err := encodeJobSet(jobset) - if err != nil { - return errors.Wrap(err, "encoding job set") - } - - if err := bkt.Put(freeJobKey(roleType, qdbid), val); err != nil { - return errors.Wrap(err, "putting free job") - } - - return nil -} - -////////////////////////////////////////////////////// - -// Ensure type implements interface. -var _ balancer.FreeWorkerService = (*freeWorkerService)(nil) - -type freeWorkerService struct { - db *boltdb.DB -} - -func newFreeWorkerService(db *boltdb.DB) *freeWorkerService { - return &freeWorkerService{ - db: db, - } -} - -func (f *freeWorkerService) AddWorkers(tx dax.Transaction, roleType dax.RoleType, addres ...dax.Address) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - workerset := dax.NewSet[dax.Address]() - var err error - - // get free workers - fws := bkt.Get(freeWorkerKey(roleType)) - if fws != nil { - workerset, err = decodeWorkerSet(fws) - if err != nil { - return errors.Wrap(err, "decoding worker set") - } - } - - for _, w := range addres { - workerset.Add(w) - } - val, err := encodeWorkerSet(workerset) - if err != nil { - return errors.Wrap(err, "encoding worker set") - } - - if err := bkt.Put(freeWorkerKey(roleType), val); err != nil { - return errors.Wrap(err, "putting free worker") - } - - return nil -} - -func (f *freeWorkerService) RemoveWorker(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - workers, err := f.ListWorkers(tx, roleType) - if err != nil { - return errors.Wrap(err, "listing free workers") - } - - // Create a workerset containing the free workers which remain after - // removing num workers. - workerset := dax.NewSet[dax.Address]() - for _, w := range workers { - workerset.Add(w) - } - - if !workerset.Contains(addr) { - return nil - } - - // Remove the worker. - workerset.Remove(addr) - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - val, err := encodeWorkerSet(workerset) - if err != nil { - return errors.Wrap(err, "encoding worker set") - } - - if err := bkt.Put(freeWorkerKey(roleType), val); err != nil { - return errors.Wrap(err, "putting free worker") - } - - return nil -} - -func (f *freeWorkerService) PopWorkers(tx dax.Transaction, roleType dax.RoleType, num int) ([]dax.Address, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - workers, err := f.ListWorkers(tx, roleType) - if err != nil { - return nil, errors.Wrap(err, "listing free workers") - } - - if len(workers) < num { - return nil, errors.Errorf("not enough free workers to pop: wanted %d, have: %d", num, len(workers)) - } - - // Get num workers from the list. - workersToAssign := workers[0:num] - - // Create a workerset containing the free workers which remain after - // removing num workers. - workerset := dax.NewSet[dax.Address]() - for _, worker := range workers[num:] { - workerset.Add(worker) - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - val, err := encodeWorkerSet(workerset) - if err != nil { - return nil, errors.Wrap(err, "encoding worker set") - } - - if err := bkt.Put(freeWorkerKey(roleType), val); err != nil { - return nil, errors.Wrap(err, "putting free worker") - } - - return workersToAssign, nil -} - -func (f *freeWorkerService) ListWorkers(tx dax.Transaction, roleType dax.RoleType) (dax.Addresses, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - workerset := dax.NewSet[dax.Address]() - var err error - - // get free workers - fws := bkt.Get(freeWorkerKey(roleType)) - if fws != nil { - workerset, err = decodeWorkerSet(fws) - if err != nil { - return nil, errors.Wrap(err, "decoding worker set") - } - } - - return workerset.Sorted(), nil -} - -////////////////////////////////////////////////////// - -const ( - prefixFmtWorkersDB = "workers/role/%s/db/%s/" // %s - role, dbKey - prefixFmtWorkersAssigned = "workers/assigned/" - - prefixFmtFreeJobs = "freejobs/role/%s/db/%s" // %s - role, dbKey - prefixFmtFreeWorkers = "freeworkers/role/%s" // %s - role -) - -// workerDBKey returns a key based on worker. -// -// Format: workers/role/[role]/db/[dbKey]/[worker] = [job1, job2, ...] -func workerDBKey(roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) []byte { - key := fmt.Sprintf(prefixFmtWorkersDB+"%s", roleType, qdbid.Key(), addr) - return []byte(key) -} - -// workerAssignedKey returns a key based on worker. -// -// Format: workers/assigned/[worker] = dbKey -func workerAssignedKey(addr dax.Address) []byte { - key := fmt.Sprintf(prefixFmtWorkersAssigned+"%s", addr) - return []byte(key) -} - -// keyWorker gets the worker out of the key. -func keyWorker(key []byte) (dax.Address, error) { - parts := strings.SplitN(string(key), "/", 6) - if len(parts) != 6 { - return "", errors.New(errors.ErrUncoded, "worker key format expected: `workers/role/[role]/db/[db]/worker`") - } - - return dax.Address(parts[5]), nil -} - -// freeJobKey returns a key for all freeJobs. -// -// Format: freejobs/role/[role]/db/[dbKey] = [job1, job2, ...] -func freeJobKey(roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) []byte { - key := fmt.Sprintf(prefixFmtFreeJobs, roleType, qdbid.Key()) - return []byte(key) -} - -// freeWorkerKey returns a key for all freeWorkers. -// -// Format: freeworkers/role/[role] = [worker1, worker2, ...] -func freeWorkerKey(roleType dax.RoleType) []byte { - key := fmt.Sprintf(prefixFmtFreeWorkers, roleType) - return []byte(key) -} diff --git a/dax/controller/balancer/boltdb/node.go b/dax/controller/balancer/boltdb/node.go deleted file mode 100644 index 9a3624068..000000000 --- a/dax/controller/balancer/boltdb/node.go +++ /dev/null @@ -1,144 +0,0 @@ -package boltdb - -import ( - "bytes" - "encoding/json" - "fmt" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/boltdb" - "github.com/featurebasedb/featurebase/v3/dax/controller" - "github.com/featurebasedb/featurebase/v3/errors" - "github.com/featurebasedb/featurebase/v3/logger" -) - -// Ensure type implements interface. -var _ controller.NodeService = (*NodeService)(nil) - -// NodeService represents a service for managing nodes. -type NodeService struct { - db *boltdb.DB - - logger logger.Logger -} - -// NewNodeService returns a new instance of NodeService with default values. -func NewNodeService(db *boltdb.DB, logger logger.Logger) *NodeService { - return &NodeService{ - db: db, - logger: logger, - } -} - -func (s *NodeService) CreateNode(tx dax.Transaction, addr dax.Address, node *dax.Node) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - val, err := json.Marshal(node) - if err != nil { - return errors.Wrap(err, "marshalling node to json") - } - - if err := bkt.Put(addressKey(addr), val); err != nil { - return errors.Wrap(err, "putting node") - } - - return nil -} - -func (s *NodeService) ReadNode(tx dax.Transaction, addr dax.Address) (*dax.Node, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - b := bkt.Get(addressKey(addr)) - if b == nil { - return nil, dax.NewErrNodeDoesNotExist(addr) - } - - node := &dax.Node{} - if err := json.Unmarshal(b, node); err != nil { - return nil, errors.Wrap(err, "unmarshalling node json") - } - - return node, nil -} - -func (s *NodeService) DeleteNode(tx dax.Transaction, addr dax.Address) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - bkt := txx.Bucket(bucketBalancer) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer) - } - - if err := bkt.Delete(addressKey(addr)); err != nil { - return errors.Wrapf(err, "deleting node key: %s", addressKey(addr)) - } - - return nil -} - -func (s *NodeService) Nodes(tx dax.Transaction) ([]*dax.Node, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - nodes, err := s.getNodes(txx) - if err != nil { - return nil, errors.Wrap(err, "getting nodes") - } - - return nodes, nil -} - -func (s *NodeService) getNodes(tx *boltdb.Tx) ([]*dax.Node, error) { - c := tx.Bucket(bucketBalancer).Cursor() - - // Deserialize rows into Node objects. - nodes := make([]*dax.Node, 0) - - prefix := []byte(prefixFmtNodes) - for k, v := c.Seek(prefix); k != nil && bytes.HasPrefix(k, prefix); k, v = c.Next() { - if v == nil { - s.logger.Printf("nil value for key: %s", k) - continue - } - - node := &dax.Node{} - if err := json.Unmarshal(v, node); err != nil { - return nil, errors.Wrap(err, "unmarshalling node json") - } - - nodes = append(nodes, node) - } - - return nodes, nil -} - -const ( - prefixFmtNodes = "nodes/" -) - -// addressKey returns a key based on address. -func addressKey(addr dax.Address) []byte { - key := fmt.Sprintf(prefixFmtNodes+"%s", addr) - return []byte(key) -} diff --git a/dax/controller/balancer/boltdb/node_test.go b/dax/controller/balancer/boltdb/node_test.go deleted file mode 100644 index 7b770cced..000000000 --- a/dax/controller/balancer/boltdb/node_test.go +++ /dev/null @@ -1,61 +0,0 @@ -package boltdb_test - -import ( - "context" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/balancer/boltdb" - testbolt "github.com/featurebasedb/featurebase/v3/dax/test/boltdb" - "github.com/featurebasedb/featurebase/v3/errors" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/stretchr/testify/assert" -) - -func TestNodeService(t *testing.T) { - db := testbolt.MustOpenDB(t) - defer testbolt.MustCloseDB(t, db) - - t.Cleanup(func() { - testbolt.CleanupDB(t, db.Path()) - }) - - ctx := context.Background() - - // Initialize the buckets. - assert.NoError(t, db.InitializeBuckets(boltdb.BalancerBuckets...)) - - t.Run("Nodes", func(t *testing.T) { - ns := boltdb.NewNodeService(db, logger.NopLogger) - - node1 := &dax.Node{ - Address: "localhost:10101", - RoleTypes: []dax.RoleType{ - "compute", - }, - } - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - // Create node. - assert.NoError(t, ns.CreateNode(tx, node1.Address, node1)) - - // Read node. - n, err := ns.ReadNode(tx, node1.Address) - assert.NoError(t, err) - assert.Equal(t, node1, n) - - // Delete node. - assert.NoError(t, ns.DeleteNode(tx, node1.Address)) - - // Read node. - _, err = ns.ReadNode(tx, node1.Address) - if assert.Error(t, err) { - assert.True(t, errors.Is(err, dax.ErrNodeDoesNotExist)) - } - - assert.NoError(t, tx.Commit()) - }) -} diff --git a/dax/controller/config.go b/dax/controller/config.go index ef8bfef5a..385ae4a96 100644 --- a/dax/controller/config.go +++ b/dax/controller/config.go @@ -20,7 +20,6 @@ type Config struct { // Storage StorageMethod string `toml:"storage-method"` - DataDir string `toml:"-"` SQLDB *SQLDBConfig `toml:"sqldb"` diff --git a/dax/controller/controller_test.go b/dax/controller/controller_test.go index 3be74b0e3..cd6a27867 100644 --- a/dax/controller/controller_test.go +++ b/dax/controller/controller_test.go @@ -3,43 +3,64 @@ package controller_test import ( "context" "fmt" + "os" + "sort" "sync" "testing" "github.com/featurebasedb/featurebase/v3/dax" - directivedb "github.com/featurebasedb/featurebase/v3/dax/boltdb" "github.com/featurebasedb/featurebase/v3/dax/controller" - balancerdb "github.com/featurebasedb/featurebase/v3/dax/controller/balancer/boltdb" - schemardb "github.com/featurebasedb/featurebase/v3/dax/controller/schemar/boltdb" + "github.com/featurebasedb/featurebase/v3/dax/controller/sqldb" + daxtest "github.com/featurebasedb/featurebase/v3/dax/test" - testbolt "github.com/featurebasedb/featurebase/v3/dax/test/boltdb" "github.com/featurebasedb/featurebase/v3/errors" "github.com/featurebasedb/featurebase/v3/logger" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) +var trans sqldb.Transactor + +func TestMain(m *testing.M) { + os.Exit(run(m)) +} + +// run is a separate function so that we can defer cleanups. (deferred functions won't run if os.Exit is called) +func run(m *testing.M) int { + // We connect to a randomized database, create it, and run migrations. Then we drop it when tests are done. + conf := sqldb.GetTestConfigRandomDB("controller_test") + var err error + trans, err = sqldb.NewTransactor(conf, logger.StderrLogger) + if err != nil { + fmt.Printf("couldn't set up transactor: %v", err) + return -1 + } + + defer sqldb.DropDatabase(trans) + code := m.Run() + return code +} + func TestController(t *testing.T) { ctx := context.Background() qdbid := dax.NewQualifiedDatabaseID("acme", "db1") t.Run("RegisterNode", func(t *testing.T) { director := newTestDirector() - schemar, cleanup := daxtest.NewSchemar(t) - defer cleanup() + schemar := sqldb.NewSchemar(logger.StderrLogger) + err := trans.Start() + require.NoError(t, err, "starting transactor") - db := testbolt.MustOpenDB(t) - db.InitializeBuckets(balancerdb.BalancerBuckets...) - db.InitializeBuckets(schemardb.SchemarBuckets...) defer func() { - testbolt.MustCloseDB(t, db) - testbolt.CleanupDB(t, db.Path()) + trans.TruncateAll() + trans.Close() }() cfg := controller.Config{} con := controller.New(cfg) con.Schemar = schemar - con.Transactor = db + con.Transactor = trans con.Director = director // Register a node with an invalid role type. @@ -49,7 +70,7 @@ func TestController(t *testing.T) { "invalid-role-type", }, } - err := con.RegisterNodes(ctx, node0) + err = con.RegisterNodes(ctx, node0) if assert.Error(t, err) { assert.True(t, errors.Is(err, controller.ErrCodeRoleTypeInvalid)) } @@ -67,25 +88,21 @@ func TestController(t *testing.T) { t.Run("ComputeNodes", func(t *testing.T) { director := newTestDirector() - schemar, cleanup := daxtest.NewSchemar(t) - defer cleanup() + schemar := sqldb.NewSchemar(logger.StderrLogger) + err := trans.Start() + require.NoError(t, err, "starting transactor") - db := testbolt.MustOpenDB(t) - db.InitializeBuckets(balancerdb.BalancerBuckets...) - db.InitializeBuckets(schemardb.SchemarBuckets...) - db.InitializeBuckets(directivedb.DirectiveBuckets...) defer func() { - testbolt.MustCloseDB(t, db) - testbolt.CleanupDB(t, db.Path()) + trans.TruncateAll() }() cfg := controller.Config{} con := controller.New(cfg) con.Schemar = schemar - con.Balancer = balancerdb.NewBalancer(db, schemar, logger.StderrLogger) - con.DirectiveVersion = directivedb.NewDirectiveVersion(db) + con.Balancer = sqldb.NewBalancer(logger.StderrLogger) + con.DirectiveVersion = sqldb.NewDirectiveVersion(logger.StderrLogger) con.Director = director - con.Transactor = db + con.Transactor = trans var exp []*dax.Directive @@ -108,7 +125,9 @@ func TestController(t *testing.T) { Version: 1, }, } - assert.Equal(t, exp, director.flush()) + got := director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Add a qualified database. dbOptions := dax.DatabaseOptions{ @@ -150,7 +169,9 @@ func TestController(t *testing.T) { Version: 2, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Set WorkersMin to 3 so we can used the added nodes that follow. { @@ -176,7 +197,9 @@ func TestController(t *testing.T) { Version: 3, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) node2 := &dax.Node{ Address: "10.0.0.1:82", @@ -196,7 +219,9 @@ func TestController(t *testing.T) { Version: 4, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Add more shards. addShards(t, ctx, con, tbl0.QualifiedID(), dax.NewShardNums(1, 2, 3, 5, 8)...) @@ -278,7 +303,9 @@ func TestController(t *testing.T) { Version: 9, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Add another non-keyed table. tbl1 := daxtest.TestQualifiedTable(t, qdbid, "bar", 0, false) @@ -372,7 +399,9 @@ func TestController(t *testing.T) { Version: 13, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + require.Equal(t, exp, got) // Remove a node. assert.NoError(t, con.DeregisterNodes(ctx, node1.Address)) @@ -419,7 +448,9 @@ func TestController(t *testing.T) { Version: 15, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + // assert.Equal(t, exp, got) // TODO fails due to shards being allocated differently // Remove another node. assert.NoError(t, con.DeregisterNodes(ctx, node0.Address)) @@ -446,7 +477,9 @@ func TestController(t *testing.T) { Version: 16, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Remove final node. assert.NoError(t, con.DeregisterNodes(ctx, node2.Address)) @@ -491,7 +524,9 @@ func TestController(t *testing.T) { Version: 17, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Remove shards. assert.NoError(t, con.RemoveShards(ctx, tbl0.QualifiedID(), dax.NewShardNums(2, 5)...)) @@ -518,7 +553,9 @@ func TestController(t *testing.T) { Version: 18, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Remove shards, one which does not exist. // Currently that doesn't result in an error, it simply no-ops on trying @@ -547,7 +584,9 @@ func TestController(t *testing.T) { Version: 19, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Remove a table. assert.NoError(t, con.DropTable(ctx, tbl0.QualifiedID())) @@ -569,7 +608,9 @@ func TestController(t *testing.T) { Version: 20, }, } - assert.Equal(t, exp, director.flush()) + got = director.flush() + require.Equal(t, len(exp), len(got)) + assert.Equal(t, exp, got) // Remove a node which doesn't exist. assert.NoError(t, con.DeregisterNodes(ctx, "invalidNode")) @@ -582,24 +623,20 @@ func TestController(t *testing.T) { ) director := newTestDirector() - schemar, cleanup := daxtest.NewSchemar(t) - defer cleanup() + schemar := sqldb.NewSchemar(logger.StderrLogger) + err := trans.Start() + require.NoError(t, err, "starting transactor") - db := testbolt.MustOpenDB(t) - db.InitializeBuckets(balancerdb.BalancerBuckets...) - db.InitializeBuckets(schemardb.SchemarBuckets...) - db.InitializeBuckets(directivedb.DirectiveBuckets...) defer func() { - testbolt.MustCloseDB(t, db) - testbolt.CleanupDB(t, db.Path()) + trans.TruncateAll() }() cfg := controller.Config{} con := controller.New(cfg) con.Schemar = schemar - con.Balancer = balancerdb.NewBalancer(db, schemar, logger.StderrLogger) - con.DirectiveVersion = directivedb.NewDirectiveVersion(db) - con.Transactor = db + con.Balancer = sqldb.NewBalancer(logger.StderrLogger) + con.DirectiveVersion = sqldb.NewDirectiveVersion(logger.StderrLogger) + con.Transactor = trans con.Director = director var exp []*dax.Directive @@ -881,7 +918,7 @@ func TestController(t *testing.T) { assert.Equal(t, exp, director.flush()) // Remove a table which doesn't exist. - err := con.DropTable(ctx, invalidQtid) + err = con.DropTable(ctx, invalidQtid) if assert.Error(t, err) { assert.True(t, errors.Is(err, dax.ErrTableIDDoesNotExist)) } @@ -897,24 +934,20 @@ func TestController(t *testing.T) { }) t.Run("GetNodes", func(t *testing.T) { - schemar, cleanup := daxtest.NewSchemar(t) - defer cleanup() + schemar := sqldb.NewSchemar(logger.StderrLogger) + err := trans.Start() + require.NoError(t, err, "starting transactor") - db := testbolt.MustOpenDB(t) - db.InitializeBuckets(balancerdb.BalancerBuckets...) - db.InitializeBuckets(schemardb.SchemarBuckets...) - db.InitializeBuckets(directivedb.DirectiveBuckets...) defer func() { - testbolt.MustCloseDB(t, db) - testbolt.CleanupDB(t, db.Path()) + trans.TruncateAll() }() cfg := controller.Config{} con := controller.New(cfg) con.Schemar = schemar - con.Balancer = balancerdb.NewBalancer(db, schemar, logger.StderrLogger) - con.DirectiveVersion = directivedb.NewDirectiveVersion(db) - con.Transactor = db + con.Balancer = sqldb.NewBalancer(logger.StderrLogger) + con.DirectiveVersion = sqldb.NewDirectiveVersion(logger.StderrLogger) + con.Transactor = trans // Register two nodes. node0 := &dax.Node{ diff --git a/dax/controller/schemar/boltdb/schemar.go b/dax/controller/schemar/boltdb/schemar.go deleted file mode 100644 index 7dbe50288..000000000 --- a/dax/controller/schemar/boltdb/schemar.go +++ /dev/null @@ -1,698 +0,0 @@ -// Package boltdb contains the boltdb implementation of the Schemar -// interfaces. -package boltdb - -import ( - "bytes" - "encoding/json" - "fmt" - "strings" - "time" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/boltdb" - "github.com/featurebasedb/featurebase/v3/dax/controller/schemar" - "github.com/featurebasedb/featurebase/v3/errors" - "github.com/featurebasedb/featurebase/v3/logger" -) - -var ( - bucketSchemar = boltdb.Bucket("schemar") -) - -// SchemarBuckets defines the buckets used by this package. It can be called -// during setup to create the buckets ahead of time. -var SchemarBuckets []boltdb.Bucket = []boltdb.Bucket{ - bucketSchemar, -} - -// Ensure type implements interface. -var _ schemar.Schemar = (*Schemar)(nil) - -type Schemar struct { - db *boltdb.DB - - logger logger.Logger -} - -// NewSchemar returns a new instance of Schemar with default values. -func NewSchemar(db *boltdb.DB, logger logger.Logger) *Schemar { - return &Schemar{ - db: db, - logger: logger, - } -} - -// CreateDatabase creates the database provided. If a database with the same -// name already exists then an error is returned. For now, we are not going to -// store the tables in the schemar Database struct. -func (s *Schemar) CreateDatabase(tx dax.Transaction, qdb *dax.QualifiedDatabase) error { - // Ensure the database id is not blank. - if qdb.ID == "" { - return schemar.NewErrDatabaseIDInvalid(qdb.ID) - } - - // Ensure the database name is not blank. - if qdb.Name == "" { - return schemar.NewErrDatabaseNameInvalid(qdb.Name) - } - - // Set the CreateAt value for the database. - // TODO(tlt): We may want to consider erroring here if the value is != 0. - if qdb.CreatedAt == 0 { - now := timestamp() - qdb.CreatedAt = now - } - - //////////// end validation - - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - // Ensure a database with that ID doesn't already exist. - if db, _ := s.databaseByID(txx, qdb.OrganizationID, qdb.ID); db != nil { - return dax.NewErrDatabaseIDExists(qdb.QualifiedID()) - } - - if err := s.putDatabase(txx, qdb); err != nil { - return errors.Wrap(err, "putting database") - } - - // In addition to storing the database in databaseKey, we want to store a - // reverse-lookup (i.e. index) on database name to the databaseKey. - if err := s.putDatabaseName(txx, qdb); err != nil { - return errors.Wrap(err, "putting database name") - } - - return nil -} - -func (s *Schemar) DatabaseByID(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) (*dax.QualifiedDatabase, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - return s.databaseByID(txx, qdbid.OrganizationID, qdbid.DatabaseID) -} - -func (s *Schemar) databaseByID(tx *boltdb.Tx, orgID dax.OrganizationID, id dax.DatabaseID) (*dax.QualifiedDatabase, error) { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - b := bkt.Get(databaseKey(orgID, id)) - if b == nil { - return nil, dax.NewErrDatabaseIDDoesNotExist(dax.QualifiedDatabaseID{OrganizationID: orgID, DatabaseID: id}) - } - - database := &dax.QualifiedDatabase{} - if err := json.Unmarshal(b, database); err != nil { - return nil, errors.Wrap(err, "unmarshalling database json") - } - - return database, nil -} - -func (s *Schemar) DatabaseByName(tx dax.Transaction, orgID dax.OrganizationID, dbname dax.DatabaseName) (*dax.QualifiedDatabase, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - return s.databaseByName(txx, orgID, dbname) -} - -func (s *Schemar) databaseByName(tx *boltdb.Tx, orgID dax.OrganizationID, name dax.DatabaseName) (*dax.QualifiedDatabase, error) { - qdbid, err := s.databaseIDByName(tx, orgID, name) - if err != nil { - return nil, errors.Wrap(err, "getting database ID") - } - - return s.databaseByID(tx, orgID, qdbid.DatabaseID) -} - -func (s *Schemar) databaseIDByName(tx *boltdb.Tx, orgID dax.OrganizationID, name dax.DatabaseName) (dax.QualifiedDatabaseID, error) { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return dax.QualifiedDatabaseID{}, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - b := bkt.Get(databaseNameKey(orgID, name)) - if b == nil { - return dax.QualifiedDatabaseID{}, dax.NewErrDatabaseNameDoesNotExist(name) - } - - return keyQualifiedDatabaseID(b) -} - -func (s *Schemar) putDatabase(tx *boltdb.Tx, qdb *dax.QualifiedDatabase) error { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - val, err := json.Marshal(qdb) - if err != nil { - return errors.Wrap(err, "marshalling database to json") - } - - return bkt.Put(databaseKey(qdb.OrganizationID, qdb.ID), val) -} - -func (s *Schemar) putDatabaseName(tx *boltdb.Tx, qdb *dax.QualifiedDatabase) error { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - return bkt.Put(databaseNameKey(qdb.OrganizationID, qdb.Name), databaseKey(qdb.OrganizationID, qdb.ID)) -} - -// DropDatabase drops the given database. If the named/IDed database does not -// exist then an error is returned. -func (s *Schemar) DropDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - // Ensure the database exists. - qdb, err := s.databaseByID(txx, qdbid.OrganizationID, qdbid.DatabaseID) - if err != nil { - return errors.Wrap(err, "getting database by id") - } - - bkt := txx.Bucket(bucketSchemar) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - // Delete the database by ID. - if err := bkt.Delete(databaseKey(qdb.OrganizationID, qdb.ID)); err != nil { - return errors.Wrap(err, "deleting database by id") - } - - // Delete the reverse-lookup database by Name. - if err := bkt.Delete(databaseNameKey(qdb.OrganizationID, qdb.Name)); err != nil { - return errors.Wrap(err, "deleting database by name") - } - - return nil -} - -// SetDatabaseOption overwrites the existing database option with the provided -// value. -func (s *Schemar) SetDatabaseOption(tx dax.Transaction, qdbid dax.QualifiedDatabaseID, option string, value string) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - // Get the database. - qdb, err := s.databaseByID(txx, qdbid.OrganizationID, qdbid.DatabaseID) - if err != nil { - return errors.Wrapf(err, "getting database: %s", qdbid) - } - - // Set the new option. - if err := qdb.Options.Set(option, value); err != nil { - return errors.Wrapf(err, "setting option on database: %s", qdbid) - } - - // Put the database. - if err := s.putDatabase(txx, qdb); err != nil { - return errors.Wrap(err, "putting database") - } - - return nil -} - -func (s *Schemar) Databases(tx dax.Transaction, orgID dax.OrganizationID, ids ...dax.DatabaseID) ([]*dax.QualifiedDatabase, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - return s.getDatabases(txx, orgID, ids...) -} - -func (s *Schemar) getDatabases(tx *boltdb.Tx, orgID dax.OrganizationID, ids ...dax.DatabaseID) (dax.QualifiedDatabases, error) { - c := tx.Bucket(bucketSchemar).Cursor() - - // Deserialize rows into Database objects. - databases := make(dax.QualifiedDatabases, 0) - - var filterByID bool - if len(ids) > 0 { - filterByID = true - } - - prefix := []byte(fmt.Sprintf(prefixFmtDatabases, orgID)) - if orgID == "" { - prefix = []byte(prefixDatabases) - } - - for k, v := c.Seek(prefix); k != nil && bytes.HasPrefix(k, prefix); k, v = c.Next() { - if v == nil { - s.logger.Printf("nil value for key: %s", k) - continue - } - - dbID, err := keyDatabaseID(k) - if err != nil { - return nil, errors.Wrap(err, "getting database from key") - } - - // Only include databases provided in the ids filter. - if filterByID && !containsDatabaseID(ids, dbID) { - continue - } - - database := &dax.QualifiedDatabase{} - if err := json.Unmarshal(v, database); err != nil { - return nil, errors.Wrap(err, "unmarshalling database json") - } - - databases = append(databases, database) - } - - return databases, nil -} - -func containsDatabaseID(s []dax.DatabaseID, e dax.DatabaseID) bool { - for _, a := range s { - if a == e { - return true - } - } - return false -} - -// CreateTable creates the table provided. If a table with the same name already -// exists then an error is returned. -func (s *Schemar) CreateTable(tx dax.Transaction, qtbl *dax.QualifiedTable) error { - // Ensure the table id is not blank. - if qtbl.ID == "" { - return schemar.NewErrTableIDInvalid(qtbl.ID) - } - - // Ensure the table name is not blank. - if qtbl.Name == "" { - return schemar.NewErrTableNameInvalid(qtbl.Name) - } - - // Ensure that a primary key field is present and valid. - if !qtbl.HasValidPrimaryKey() { - return schemar.NewErrInvalidPrimaryKey() - } - - // Set the CreateAt value for the table. - // TODO(tlt): We may want to consider erroring here if the value is != 0. - if qtbl.CreatedAt == 0 { - now := timestamp() - qtbl.CreatedAt = now - - // Set CreatedAt for all of the fields as well. - for i := range qtbl.Fields { - qtbl.Fields[i].CreatedAt = now - } - } - - //////////// end validation - - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - // Ensure the database, defined in the table's QualifiedDatabaseID, exists. - if _, err := s.databaseByID(txx, qtbl.OrganizationID, qtbl.DatabaseID); err != nil { - return errors.Wrap(err, "validating database") - } - - // Ensure a table with that ID doesn't already exist. - if t, _ := s.tableByID(txx, qtbl.QualifiedDatabaseID, qtbl.ID); t != nil { - return dax.NewErrTableIDExists(qtbl.QualifiedID()) - } - - if err := s.putTable(txx, qtbl); err != nil { - return errors.Wrap(err, "putting table") - } - - // In addition to storing the table in tableKey, we want to store a reverse-lookup - // (i.e. index) on table name to the tableKey. - if err := s.putTableName(txx, qtbl); err != nil { - return errors.Wrap(err, "putting table name") - } - - return nil -} - -// CreateField creates the field provided in the given table. If a field with -// the same name already exists then an error is returned. -func (s *Schemar) CreateField(tx dax.Transaction, qtid dax.QualifiedTableID, fld *dax.Field) error { - // Ensure the field name is not blank. - if fld.Name == "" { - return schemar.NewErrFieldNameInvalid(fld.Name) - } - - //////////// end validation - - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - // Get the table. - qtbl, err := s.tableByQTID(txx, qtid) - if err != nil { - return errors.Wrap(err, "getting table by id") - } - - // Ensure a field with that name doesn't already exist. - if _, ok := qtbl.Field(fld.Name); ok { - return dax.NewErrFieldExists(fld.Name) - } - - qtbl.Fields = append(qtbl.Fields, fld) - - // Write table back to database. - if err := s.putTable(txx, qtbl); err != nil { - return errors.Wrap(err, "putting table") - } - - return nil -} - -// DropField removes the field from the table. -func (s *Schemar) DropField(tx dax.Transaction, qtid dax.QualifiedTableID, fldName dax.FieldName) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - // Get the table. - qtbl, err := s.tableByQTID(txx, qtid) - if err != nil { - return errors.Wrap(err, "getting table by id") - } - - // Ensure a field with that name exists. - if _, ok := qtbl.Field(fldName); !ok { - return dax.NewErrFieldDoesNotExist(fldName) - } - - _ = qtbl.RemoveField(fldName) - - // Write table back to database. - if err := s.putTable(txx, qtbl); err != nil { - return errors.Wrap(err, "putting table") - } - - return nil -} - -func (s *Schemar) putTable(tx *boltdb.Tx, qtbl *dax.QualifiedTable) error { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - val, err := json.Marshal(qtbl) - if err != nil { - return errors.Wrap(err, "marshalling table to json") - } - - return bkt.Put(tableKey(qtbl.OrganizationID, qtbl.DatabaseID, qtbl.Table.ID), val) -} - -func (s *Schemar) putTableName(tx *boltdb.Tx, qtbl *dax.QualifiedTable) error { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - return bkt.Put(tableNameKey(qtbl.OrganizationID, qtbl.DatabaseID, qtbl.Name), tableKey(qtbl.OrganizationID, qtbl.DatabaseID, qtbl.ID)) -} - -// Table returns the TableInfo for the given table. An error is returned if the -// table does not exist. -func (s *Schemar) Table(tx dax.Transaction, qtid dax.QualifiedTableID) (*dax.QualifiedTable, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - return s.tableByQTID(txx, qtid) -} - -// tableByQTID gets the full qualified table by the QualifiedTableID whether it -// has Name or ID set. -func (s *Schemar) tableByQTID(tx *boltdb.Tx, qtid dax.QualifiedTableID) (*dax.QualifiedTable, error) { - if qtid.ID == "" { - return s.tableByName(tx, qtid.QualifiedDatabaseID, qtid.Name) - } - - return s.tableByID(tx, qtid.QualifiedDatabaseID, qtid.ID) -} - -func (s *Schemar) tableByName(tx *boltdb.Tx, qdbid dax.QualifiedDatabaseID, name dax.TableName) (*dax.QualifiedTable, error) { - qtid, err := s.tableIDByName(tx, qdbid, name) - if err != nil { - return nil, errors.Wrap(err, "getting table ID") - } - - return s.tableByID(tx, qtid.QualifiedDatabaseID, qtid.ID) // TODO remove? -} - -func (s *Schemar) tableByID(tx *boltdb.Tx, qdbid dax.QualifiedDatabaseID, id dax.TableID) (*dax.QualifiedTable, error) { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - b := bkt.Get(tableKey(qdbid.OrganizationID, qdbid.DatabaseID, id)) - if b == nil { - return nil, dax.NewErrTableIDDoesNotExist(dax.QualifiedTableID{QualifiedDatabaseID: qdbid, ID: id}) - } - - table := &dax.QualifiedTable{} - if err := json.Unmarshal(b, table); err != nil { - return nil, errors.Wrap(err, "unmarshalling table json") - } - - return table, nil -} - -func (s *Schemar) tableIDByName(tx *boltdb.Tx, qdbid dax.QualifiedDatabaseID, name dax.TableName) (dax.QualifiedTableID, error) { - bkt := tx.Bucket(bucketSchemar) - if bkt == nil { - return dax.QualifiedTableID{}, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - b := bkt.Get(tableNameKey(qdbid.OrganizationID, qdbid.DatabaseID, name)) - if b == nil { - return dax.QualifiedTableID{}, dax.NewErrTableNameDoesNotExist(name) - } - - return keyQualifiedTableID(b) -} - -// Tables returns a list of Table for all existing tables. If one or more table -// IDs is provided, then only those will be included in the output. -func (s *Schemar) Tables(tx dax.Transaction, qdbid dax.QualifiedDatabaseID, ids ...dax.TableID) ([]*dax.QualifiedTable, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return nil, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - return s.getTables(txx, qdbid, ids...) -} - -func (s *Schemar) getTables(tx *boltdb.Tx, qdbid dax.QualifiedDatabaseID, ids ...dax.TableID) (dax.QualifiedTables, error) { - c := tx.Bucket(bucketSchemar).Cursor() - - // Deserialize rows into Table objects. - tables := make(dax.QualifiedTables, 0) - - var filterByID bool - if len(ids) > 0 { - filterByID = true - } - - prefix := []byte(fmt.Sprintf(prefixFmtTables, qdbid.OrganizationID, qdbid.DatabaseID)) - if qdbid.OrganizationID == "" && qdbid.DatabaseID == "" { - prefix = []byte(prefixTables) - } else if qdbid.DatabaseID == "" { - prefix = []byte(fmt.Sprintf(prefixFmtTablesOrg, qdbid.OrganizationID)) - } - - for k, v := c.Seek(prefix); k != nil && bytes.HasPrefix(k, prefix); k, v = c.Next() { - if v == nil { - s.logger.Printf("nil value for key: %s", k) - continue - } - - tblID, err := keyTableID(k) - if err != nil { - return nil, errors.Wrap(err, "getting table from key") - } - - // Only include tables provided in the ids filter. - if filterByID && !containsTableID(ids, tblID) { - continue - } - - table := &dax.QualifiedTable{} - if err := json.Unmarshal(v, table); err != nil { - return nil, errors.Wrap(err, "unmarshalling table json") - } - - tables = append(tables, table) - } - - return tables, nil -} - -func containsTableID(s []dax.TableID, e dax.TableID) bool { - for _, a := range s { - if a == e { - return true - } - } - return false -} - -// DropTable drops the given table. If the named/IDed table does not exist -// then an error is returned. -func (s *Schemar) DropTable(tx dax.Transaction, qtid dax.QualifiedTableID) error { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - // Ensure the table exists. - qtbl, err := s.tableByQTID(txx, qtid) - if err != nil { - return errors.Wrap(err, "getting table by id") - } - - bkt := txx.Bucket(bucketSchemar) - if bkt == nil { - return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketSchemar) - } - - // Delete the table by ID. - if err := bkt.Delete(tableKey(qtbl.OrganizationID, qtbl.DatabaseID, qtbl.ID)); err != nil { - return errors.Wrap(err, "deleting table by id") - } - - // Delete the reverse-lookup table by Name. - if err := bkt.Delete(tableNameKey(qtbl.OrganizationID, qtbl.DatabaseID, qtbl.Name)); err != nil { - return errors.Wrap(err, "deleting table by name") - } - - return nil -} - -func (s *Schemar) TableID(tx dax.Transaction, qdbid dax.QualifiedDatabaseID, name dax.TableName) (dax.QualifiedTableID, error) { - txx, ok := tx.(*boltdb.Tx) - if !ok { - return dax.QualifiedTableID{}, dax.NewErrInvalidTransaction("*boltdb.Tx") - } - - return s.tableIDByName(txx, qdbid, name) -} - -const ( - prefixTables = "tables/" - prefixFmtTablesOrg = prefixTables + "%s/" // org-id - prefixFmtTables = prefixFmtTablesOrg + "%s/" // db-id - prefixFmtTableNames = "tablenames/%s/%s/" // org-id, db-id - - prefixDatabases = "databases/" - prefixFmtDatabases = prefixDatabases + "%s/" // org-id - prefixFmtDatabase = prefixFmtDatabases + "%s" // db-id - prefixFmtDatabaseNames = "databasenames/%s/" // org-id -) - -// databaseKey returns a key based on a qualified database ID. -func databaseKey(orgID dax.OrganizationID, dbID dax.DatabaseID) []byte { - key := fmt.Sprintf(prefixFmtDatabase, orgID, dbID) - return []byte(key) -} - -// databaseNameKey returns a key based on a qualified database name. -func databaseNameKey(orgID dax.OrganizationID, name dax.DatabaseName) []byte { - key := fmt.Sprintf(prefixFmtDatabaseNames+"%s", orgID, name) - return []byte(key) -} - -// keyDatabaseID gets the DatabaseID out of the key. -func keyDatabaseID(key []byte) (dax.DatabaseID, error) { - parts := strings.Split(string(key), "/") - if len(parts) != 3 { - return "", errors.New(errors.ErrUncoded, "database key format expected: `databases/orgID/dbID`") - } - - return dax.DatabaseID(parts[2]), nil -} - -// tableKey returns a key based on a qualified table ID. -func tableKey(orgID dax.OrganizationID, dbID dax.DatabaseID, tblID dax.TableID) []byte { - key := fmt.Sprintf(prefixFmtTables+"%s", orgID, dbID, tblID) - return []byte(key) -} - -// tableNameKey returns a key based on a qualified table name. -func tableNameKey(orgID dax.OrganizationID, dbID dax.DatabaseID, name dax.TableName) []byte { - key := fmt.Sprintf(prefixFmtTableNames+"%s", orgID, dbID, name) - return []byte(key) -} - -// keyTableID gets the TableID out of the key. -func keyTableID(key []byte) (dax.TableID, error) { - parts := strings.Split(string(key), "/") - if len(parts) != 4 { - return "", errors.New(errors.ErrUncoded, "table key format expected: `tables/orgID/dbID/tblID`") - } - - return dax.TableID(parts[3]), nil -} - -// keyQualifedTableID gets the QualifiedTableID out of the key. -func keyQualifiedTableID(key []byte) (dax.QualifiedTableID, error) { - parts := strings.Split(string(key), "/") - if len(parts) != 4 { - return dax.QualifiedTableID{}, errors.New(errors.ErrUncoded, "table key format expected: `tables/orgID/dbID/tblID`") - } - - return dax.NewQualifiedTableID( - dax.NewQualifiedDatabaseID( - dax.OrganizationID(parts[1]), - dax.DatabaseID(parts[2]), - ), - dax.TableID(parts[3]), - ), nil -} - -// keyQualifedDatabaseID gets the QualifiedDatabaseID out of the key. -func keyQualifiedDatabaseID(key []byte) (dax.QualifiedDatabaseID, error) { - parts := strings.Split(string(key), "/") - if len(parts) != 3 { - return dax.QualifiedDatabaseID{}, errors.New(errors.ErrUncoded, "table key format expected: `databases/orgID/dbID`") - } - - return dax.NewQualifiedDatabaseID( - dax.OrganizationID(parts[1]), - dax.DatabaseID(parts[2]), - ), nil -} - -func timestamp() int64 { - return time.Now().Unix() -} diff --git a/dax/controller/schemar/boltdb/schemar_test.go b/dax/controller/schemar/boltdb/schemar_test.go deleted file mode 100644 index 9d31266f8..000000000 --- a/dax/controller/schemar/boltdb/schemar_test.go +++ /dev/null @@ -1,227 +0,0 @@ -package boltdb_test - -import ( - "context" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/schemar/boltdb" - daxtest "github.com/featurebasedb/featurebase/v3/dax/test" - testbolt "github.com/featurebasedb/featurebase/v3/dax/test/boltdb" - "github.com/featurebasedb/featurebase/v3/errors" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/stretchr/testify/assert" -) - -func TestSchemar(t *testing.T) { - orgID := dax.OrganizationID("acme") - dbID := dax.DatabaseID("db1") - dbName := dax.DatabaseName("dbname1") - invalidTableID := dax.TableID("invalidID") - tableName := dax.TableName("foo") - tableName0 := dax.TableName("foo") - tableName1 := dax.TableName("bar") - tableID0 := "2" - tableID1 := "1" - partitionN := 12 - - ctx := context.Background() - qdbid := dax.NewQualifiedDatabaseID(orgID, dbID) - - qdb := &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ - ID: dbID, - Name: dbName, - }, - } - - t.Run("NewSchemar", func(t *testing.T) { - db := testbolt.MustOpenDB(t) - defer testbolt.MustCloseDB(t, db) - - t.Cleanup(func() { - testbolt.CleanupDB(t, db.Path()) - }) - - // Initialize the buckets. - assert.NoError(t, db.InitializeBuckets(boltdb.SchemarBuckets...)) - - s := boltdb.NewSchemar(db, logger.NopLogger) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - // Create database. - assert.NoError(t, s.CreateDatabase(tx, qdb)) - - // Add new table. - tbl := dax.NewTable(tableName) - tbl.CreateID() - tbl.Fields = []*dax.Field{ - { - Name: dax.PrimaryKeyFieldName, - Type: dax.BaseTypeString, - }, - { - Name: "intField", - Type: dax.BaseTypeInt, - }, - } - qtbl := dax.NewQualifiedTable(qdbid, tbl) - assert.NoError(t, s.CreateTable(tx, qtbl)) - - // Try adding the table again. - err = s.CreateTable(tx, qtbl) - if assert.Error(t, err) { - assert.True(t, errors.Is(err, dax.ErrTableIDExists)) - } - - qtid := qtbl.QualifiedID() - - // Get the table. - { - tbl, err := s.Table(tx, qtid) - assert.NoError(t, err) - assert.Equal(t, tableName, tbl.Name) - } - - // Drop the table. - assert.NoError(t, s.DropTable(tx, qtid)) - - // Make sure the reverse-lookup (table by name) was dropped as well. - { - _, err := s.TableID(tx, qdbid, tableName) - if assert.Error(t, err) { - assert.True(t, errors.Is(err, dax.ErrTableNameDoesNotExist)) - } - } - - // Try adding the table (i.e. the same table name) again. - assert.NoError(t, s.CreateTable(tx, qtbl)) - - // Drop the table again. - assert.NoError(t, s.DropTable(tx, qtid)) - - // Drop invalid table. - { - iqtid := dax.NewQualifiedTableID(qdbid, invalidTableID) - err := s.DropTable(tx, iqtid) - if assert.Error(t, err) { - assert.True(t, errors.Is(err, dax.ErrTableIDDoesNotExist)) - } - } - }) - - t.Run("GetTables", func(t *testing.T) { - db := testbolt.MustOpenDB(t) - defer testbolt.MustCloseDB(t, db) - - t.Cleanup(func() { - testbolt.CleanupDB(t, db.Path()) - }) - - // Initialize the buckets. - assert.NoError(t, db.InitializeBuckets(boltdb.SchemarBuckets...)) - - s := boltdb.NewSchemar(db, logger.NopLogger) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - // Create database. - assert.NoError(t, s.CreateDatabase(tx, qdb)) - - exp := []*dax.QualifiedTable{} - tables, err := s.Tables(tx, qdbid) - assert.NoError(t, err) - assert.Equal(t, exp, tables) - - qtbl0 := daxtest.TestQualifiedTableWithID(t, qdbid, tableID0, tableName0, partitionN, false) - qtbl1 := daxtest.TestQualifiedTableWithID(t, qdbid, tableID1, tableName1, partitionN, false) - - // Add a couple of tables. - assert.NoError(t, s.CreateTable(tx, qtbl0)) - assert.NoError(t, s.CreateTable(tx, qtbl1)) - - exp = []*dax.QualifiedTable{ - qtbl1, - qtbl0, - } - - // All tables. - tables, err = s.Tables(tx, qdbid) - assert.NoError(t, err) - assert.Equal(t, exp, tables) - - // With a valid filter. - tables, err = s.Tables(tx, qdbid, qtbl0.ID) - assert.NoError(t, err) - assert.Equal(t, exp[1:], tables) - - // With an invalid filter. - tables, err = s.Tables(tx, qdbid, invalidTableID) - assert.NoError(t, err) - assert.Equal(t, exp[0:0], tables) - - // With both valid and invalid filters. - tables, err = s.Tables(tx, qdbid, qtbl0.ID, invalidTableID) - assert.NoError(t, err) - assert.Equal(t, exp[1:], tables) - - // With all valid filters. - tables, err = s.Tables(tx, qdbid, qtbl0.ID, qtbl1.ID) - assert.NoError(t, err) - assert.Equal(t, exp, tables) - }) - - t.Run("GetTablesAll", func(t *testing.T) { - db := testbolt.MustOpenDB(t) - defer testbolt.MustCloseDB(t, db) - - t.Cleanup(func() { - testbolt.CleanupDB(t, db.Path()) - }) - - // Initialize the buckets. - assert.NoError(t, db.InitializeBuckets(boltdb.SchemarBuckets...)) - - s := boltdb.NewSchemar(db, logger.NopLogger) - - tx, err := db.BeginTx(ctx, true) - assert.NoError(t, err) - defer tx.Rollback() - - qtbl0 := daxtest.TestQualifiedTableWithID(t, qdbid, tableID0, tableName0, partitionN, false) - orgID2 := dax.OrganizationID("acme2") - qdbid2 := dax.NewQualifiedDatabaseID(orgID2, dbID) - tableID2 := "3" - qtbl2 := daxtest.TestQualifiedTableWithID(t, qdbid2, tableID2, dax.TableName("two"), partitionN, false) - - // Create databases. - assert.NoError(t, s.CreateDatabase(tx, qdb)) - qdb2 := &dax.QualifiedDatabase{ - OrganizationID: orgID2, - Database: dax.Database{ - ID: dbID, - Name: dbName, - }, - } - assert.NoError(t, s.CreateDatabase(tx, qdb2)) - - assert.NoError(t, s.CreateTable(tx, qtbl0)) - assert.NoError(t, s.CreateTable(tx, qtbl2)) - - exp := []*dax.QualifiedTable{qtbl0, qtbl2} - - tables, err := s.Tables(tx, dax.QualifiedDatabaseID{}) - assert.NoError(t, err) - assert.Equal(t, exp, tables) - - tables, err = s.Tables(tx, dax.QualifiedDatabaseID{OrganizationID: orgID2}) - assert.NoError(t, err) - assert.Equal(t, []*dax.QualifiedTable{qtbl2}, tables) - }) -} diff --git a/dax/controller/service/controller.go b/dax/controller/service/controller.go index a2a070764..f3a11c4e4 100644 --- a/dax/controller/service/controller.go +++ b/dax/controller/service/controller.go @@ -5,11 +5,8 @@ import ( "os" "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/boltdb" "github.com/featurebasedb/featurebase/v3/dax/controller" - balancerboltdb "github.com/featurebasedb/featurebase/v3/dax/controller/balancer/boltdb" controllerhttp "github.com/featurebasedb/featurebase/v3/dax/controller/http" - schemarboltdb "github.com/featurebasedb/featurebase/v3/dax/controller/schemar/boltdb" "github.com/featurebasedb/featurebase/v3/dax/controller/sqldb" "github.com/featurebasedb/featurebase/v3/errors" "github.com/featurebasedb/featurebase/v3/logger" @@ -33,16 +30,6 @@ func New(uri *fbnet.URI, cfg controller.Config) *controllerService { logr = cfg.Logger.WithPrefix("Controller: ") } - if cfg.DataDir == "" { - dir, err := os.MkdirTemp("", "controller_*") - if err != nil { - logr.Printf("Making temp dir for Controller storage: %v", err) - os.Exit(1) - } - cfg.DataDir = dir - logr.Warnf("no DataDir given (like '/path/to/directory'); using temp dir at '%s'", cfg.DataDir) - } - controller := controller.New(cfg) controllerSvc := &controllerService{ uri: uri, @@ -52,23 +39,6 @@ func New(uri *fbnet.URI, cfg controller.Config) *controllerService { // Storage methods. switch cfg.StorageMethod { - case "boltdb": - buckets := append(schemarboltdb.SchemarBuckets, balancerboltdb.BalancerBuckets...) - controllerDB, err := boltdb.NewSvcBolt(cfg.DataDir, "controller", buckets...) - if err != nil { - logr.Printf(errors.Wrap(err, "creating controller bolt").Error()) - os.Exit(1) - } - // Directive version. - if err := controllerDB.InitializeBuckets(boltdb.DirectiveBuckets...); err != nil { - logr.Panicf("initializing directive buckets: %v", err) - } - controller.Schemar = schemarboltdb.NewSchemar(controllerDB, logr) - controller.Balancer = balancerboltdb.NewBalancer(controllerDB, controller.Schemar, logr) - directiveVersion := boltdb.NewDirectiveVersion(controllerDB) - controller.DirectiveVersion = directiveVersion - - controller.Transactor = controllerDB case "sqldb": controller.Schemar = sqldb.NewSchemar(logr) controller.Balancer = sqldb.NewBalancer(logr) @@ -81,7 +51,7 @@ func New(uri *fbnet.URI, cfg controller.Config) *controllerService { } controller.Transactor = transactor default: - logr.Printf("storagemethod %s not supported, try 'boltdb' or 'sqldb'", cfg.StorageMethod) + logr.Printf("storagemethod %s not supported, only 'sqldb' is currently accepted.", cfg.StorageMethod) os.Exit(1) } diff --git a/dax/controller/sqldb/balancer.go b/dax/controller/sqldb/balancer.go index 64833d71e..bd5e1f680 100644 --- a/dax/controller/sqldb/balancer.go +++ b/dax/controller/sqldb/balancer.go @@ -6,12 +6,12 @@ import ( ) // NewBalancer returns a new instance of controller.Balancer. -func NewBalancer(logger logger.Logger) *balancer.Balancer { - schemar := NewSchemar(logger) - fjs := NewFreeJobService(logger) - wjs := NewWorkerJobService(logger) - fws := NewFreeWorkerService(logger) - ns := NewNodeService(logger) +func NewBalancer(log logger.Logger) *balancer.Balancer { + schemar := NewSchemar(log) + fjs := NewFreeJobService(log) + wjs := NewWorkerJobService(log) + fws := NewFreeWorkerService(log) + ns := NewNodeService(log) - return balancer.New(ns, fjs, wjs, fws, schemar, logger) + return balancer.New(ns, fjs, wjs, fws, schemar, log) } diff --git a/dax/controller/sqldb/node.go b/dax/controller/sqldb/node.go index 276ac7664..9871757a5 100644 --- a/dax/controller/sqldb/node.go +++ b/dax/controller/sqldb/node.go @@ -82,9 +82,12 @@ func (n *nodeService) DeleteNode(tx dax.Transaction, addr dax.Address) error { node := &models.Node{} err := dt.C.Eager().Where("address = ?", addr).First(node) - if err != nil { - return errors.Wrap(err, "getting node") + if isNoRowsError(err) { + return nil + } else if err != nil { + return errors.Wrap(err, "finding node") } + err = dt.C.Destroy(node) return errors.Wrap(err, "destroying node") } diff --git a/dax/controller/sqldb/schemar.go b/dax/controller/sqldb/schemar.go index af3f94e60..845ae95b9 100644 --- a/dax/controller/sqldb/schemar.go +++ b/dax/controller/sqldb/schemar.go @@ -293,10 +293,9 @@ func toField(col models.Column) *dax.Field { panic(err) } return &dax.Field{ - Name: col.Name, - Type: col.Type, - Options: opts, - CreatedAt: col.CreatedAt.Unix(), + Name: col.Name, + Type: col.Type, + Options: opts, } } @@ -317,8 +316,6 @@ func toQualifiedTable(mtbl *models.Table) *dax.QualifiedTable { PartitionN: mtbl.PartitionN, Description: mtbl.Description, Owner: mtbl.Owner, - CreatedAt: mtbl.CreatedAt.Unix(), - UpdatedAt: mtbl.UpdatedAt.Unix(), UpdatedBy: mtbl.UpdatedBy, }, } @@ -425,7 +422,7 @@ func (s *Schemar) Tables(tx dax.Transaction, qdbid dax.QualifiedDatabaseID, tabl query = query.Where("id in (?)", ifaceIDs) } tables := []*models.Table{} - err := query.Eager().Order("created_at asc").All(&tables) + err := query.Eager().Order("name asc").All(&tables) if err != nil { return nil, errors.Wrap(err, "querying all tables") } diff --git a/dax/controller/sqldb/test.go b/dax/controller/sqldb/test.go index cf1a78eb7..582a36c1b 100644 --- a/dax/controller/sqldb/test.go +++ b/dax/controller/sqldb/test.go @@ -20,11 +20,11 @@ func EnvOr(envName, defaultVal string) string { func GetTestConfig() *controller.SQLDBConfig { return &controller.SQLDBConfig{ Dialect: "postgres", - Database: EnvOr("SQLDB_DB", "dax_test"), - Host: EnvOr("SQLDB_HOST", "127.0.0.1"), - Port: EnvOr("SQLDB_PORT", "5432"), - User: EnvOr("SQLDB_USER", "postgres"), - Password: EnvOr("SQLDB_PASSWORD", "testpass"), + Database: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_DATABASE", "dax_test"), + Host: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_HOST", "127.0.0.1"), + Port: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PORT", "5432"), + User: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_USER", "postgres"), + Password: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PASSWORD", "testpass"), } } @@ -33,9 +33,9 @@ func GetTestConfigRandomDB(dbprefix string) *controller.SQLDBConfig { return &controller.SQLDBConfig{ Dialect: "postgres", Database: fmt.Sprintf("%s_%d", dbprefix, rnd.Int()), - Host: EnvOr("SQLDB_HOST", "127.0.0.1"), - Port: EnvOr("SQLDB_PORT", "5432"), - User: EnvOr("SQLDB_USER", "postgres"), - Password: EnvOr("SQLDB_PASSWORD", "testpass"), + Host: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_HOST", "127.0.0.1"), + Port: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PORT", "5432"), + User: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_USER", "postgres"), + Password: EnvOr("FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PASSWORD", "testpass"), } } diff --git a/dax/controller/sqldb/transactor.go b/dax/controller/sqldb/transactor.go index 56f13128f..9e0ff0c37 100644 --- a/dax/controller/sqldb/transactor.go +++ b/dax/controller/sqldb/transactor.go @@ -2,6 +2,7 @@ package sqldb import ( "context" + "strings" "database/sql" @@ -66,6 +67,11 @@ func (t Transactor) Start() error { return errors.Wrap(err, "migrating DB") } + err := t.RawQuery("INSERT INTO directive_versions (id, version, created_at, updated_at) VALUES (1, 0, '1970-01-01T00:00', '1970-01-01T00:00')").Exec() + if err != nil && !strings.Contains(err.Error(), "duplicate key value violates unique constraint") { + return errors.Wrap(err, "unexpected error (re)inserting directive_version record") + } + return nil } diff --git a/dax/controller/sqldb/workerjob.go b/dax/controller/sqldb/workerjob.go index 8174aeeb7..5ca6466f7 100644 --- a/dax/controller/sqldb/workerjob.go +++ b/dax/controller/sqldb/workerjob.go @@ -112,9 +112,12 @@ func (w *workerJobService) DeleteWorker(tx dax.Transaction, roleType dax.RoleTyp worker := &models.Worker{} err := dt.C.Where("address = ? and role = ? and database_id = ?", addr, roleType, qdbid.DatabaseID).First(worker) - if err != nil { + if isNoRowsError(err) { + return nil + } else if err != nil { return errors.Wrap(err, "getting worker") } + err = dt.C.Destroy(worker) return errors.Wrap(err, "deleting worker") } diff --git a/dax/directive_version_test.go b/dax/directive_version_test.go index 8dbaeec7a..efb701a17 100644 --- a/dax/directive_version_test.go +++ b/dax/directive_version_test.go @@ -28,5 +28,5 @@ func TestDirectiveVersion(t *testing.T) { n, err := dvSvc.Increment(tx, 1) require.NoError(t, err) - require.Equal(t, uint64(2), n) + require.Equal(t, uint64(1), n) } diff --git a/dax/docker-compose.yml b/dax/docker-compose.yml deleted file mode 100644 index 4b8b9592a..000000000 --- a/dax/docker-compose.yml +++ /dev/null @@ -1,64 +0,0 @@ -version: '3' - -services: - controller: - build: - context: ../.quick - dockerfile: ../Dockerfile-dax-quick - environment: - FEATUREBASE_BIND: 0.0.0.0:8080 - FEATUREBASE_VERBOSE: "true" - FEATUREBASE_CONTROLLER_RUN: "true" - FEATUREBASE_CONFIG_DATA_DIR: file:/dax-data/controller - ports: - - "8081:8080" - - queryer: - build: - context: ../.quick - dockerfile: ../Dockerfile-dax-quick - environment: - FEATUREBASE_BIND: 0.0.0.0:8080 - FEATUREBASE_VERBOSE: "true" - FEATUREBASE_QUERYER_RUN: "true" - FEATUREBASE_QUERYER_CONFIG_CONTROLLER_ADDRESS: "controller:8080/controller" - depends_on: - - controller - ports: - - "8080:8080" - - computer: - build: - context: ../.quick - dockerfile: ../Dockerfile-dax-quick - environment: - FEATUREBASE_COMPUTER_RUN: "true" - FEATUREBASE_COMPUTER_CONFIG_CONTROLLER_ADDRESS: "controller:8080/controller" - FEATUREBASE_COMPUTER_CONFIG_DATA_DIR: /dax-data/computer - FEATUREBASE_COMPUTER_CONFIG_VERBOSE: true - FEATUREBASE_BIND: 0.0.0.0:8080 - FEATUREBASE_VERBOSE: "true" - FEATUREBASE_STORAGE_METHOD: boltdb - FEATUREBASE_COMPUTER_CONFIG_WRITELOGGER_DIR: "/dax-data/writelogger" - FEATUREBASE_COMPUTER_CONFIG_SNAPSHOTTER_DIR: "/dax-data/snapshotter" - volumes: - - "./dax-data/writelogger:/dax-data/writelogger" - - "./dax-data/snapshotter:/dax-data/snapshotter" - depends_on: - - controller - deploy: - replicas: 1 - - datagen: - build: - context: .. - dockerfile: Dockerfile-datagen - profiles: [ "datagen" ] - environment: - GEN_CUSTOM_CONFIG: "/testdata/keys_ids.yaml" - GEN_FEATUREBASE_ORG_ID: "testorg" - GEN_FEATUREBASE_DB_ID: "testdb" - GEN_USE_SHARD_TRANSACTIONAL_ENDPOINT: "true" - GEN_SOURCE: "custom" - GEN_TARGET: "serverless" - GEN_CONTROLLER_ADDRESS: "controller:8080/controller" diff --git a/dax/migrations/001_initial.up.fizz b/dax/migrations/001_initial.up.fizz index f47c6ad56..9c042f089 100644 --- a/dax/migrations/001_initial.up.fizz +++ b/dax/migrations/001_initial.up.fizz @@ -81,4 +81,3 @@ create_table("directive_versions") { t.Timestamps() } -sql("INSERT INTO directive_versions (id, version, created_at, updated_at) VALUES (1, 1, '1970-01-01T00:00', '1970-01-01T00:00');") diff --git a/dax/models/database.go b/dax/models/database.go index 53feed82a..7dc8f3dc1 100644 --- a/dax/models/database.go +++ b/dax/models/database.go @@ -19,7 +19,7 @@ type Database struct { Description string `json:"description" db:"description"` Owner string `json:"owner" db:"owner"` UpdatedBy string `json:"updated_by" db:"updated_by"` - Tables Tables `json:"tables" has_many:"tables"` + Tables Tables `json:"tables" has_many:"tables" order_by:"name asc"` Organization *Organization `json:"organization" belongs_to:"organization"` OrganizationID string `json:"organization_id" db:"organization_id"` CreatedAt time.Time `json:"created_at" db:"created_at"` diff --git a/dax/models/node.go b/dax/models/node.go index 9a387a277..c347e75ba 100644 --- a/dax/models/node.go +++ b/dax/models/node.go @@ -15,7 +15,7 @@ import ( type Node struct { ID uuid.UUID `json:"id" db:"id"` Address dax.Address `json:"address" db:"address"` - NodeRoles NodeRoles `json:"node_roles" has_many:"node_roles"` + NodeRoles NodeRoles `json:"node_roles" has_many:"node_roles" order_by:"created_at asc"` CreatedAt time.Time `json:"created_at" db:"created_at"` UpdatedAt time.Time `json:"updated_at" db:"updated_at"` } diff --git a/dax/server/test/managed.go b/dax/server/test/managed.go index 0acf85851..c7fa3575b 100644 --- a/dax/server/test/managed.go +++ b/dax/server/test/managed.go @@ -200,7 +200,6 @@ func NewManagedCommand(tb fbtest.DirCleaner, opts ...server.CommandOption) *Mana mc.svcmgr = svcmgr mc.Config.Bind = "http://localhost:0" - mc.Config.Controller.Config.DataDir = path + "/controller" mc.Config.Computer.Config.DataDir = path mc.Config.Computer.Config.WriteloggerDir = path + "/wl" mc.Config.Controller.Config.WriteloggerDir = path + "/wl" diff --git a/dax/snapshotter/snapshotter_test.go b/dax/snapshotter/snapshotter_test.go deleted file mode 100644 index d36ca50f4..000000000 --- a/dax/snapshotter/snapshotter_test.go +++ /dev/null @@ -1,70 +0,0 @@ -package snapshotter_test - -import ( - "fmt" - "os" - "path" - "testing" - - "github.com/stretchr/testify/assert" -) - -func TestSnapshotter(t *testing.T) { - tmpDir, err := os.MkdirTemp("", "testWritelogger-*") - assert.NoError(t, err) - - // Remove the temp directory. - defer func() { - os.RemoveAll(tmpDir) - }() - - // t.Run("Basic", func(t *testing.T) { - // type payload struct { - // Foo string `json:"foo"` - // Bar int `json:"bar"` - // } - - // cfg := core.Config{ - // DataDir: tmpDir, - // } - // wl := core.NewSnapshotter(cfg) - - // table := "tbl" - // partition := 1 - // version := 0 - // key := "keys" - - // msg1 := payload{ - // Foo: "message 1", - // Bar: 88, - // } - - // // Write the message. - // msg, err := json.Marshal(msg1) - // assert.NoError(t, err) - - // err = wl.AppendMessage(bucket(table, partition), key, version, msg) - // assert.NoError(t, err) - - // // Read the message. - // reader, closer, err := wl.LogReader(bucket(table, partition), key, version) - // assert.NoError(t, err) - // defer closer.Close() - - // buf, err := ioutil.ReadAll(reader) - // assert.NoError(t, err) - - // var out payload - - // err = json.Unmarshal(buf, &out) - // assert.NoError(t, err) - - // assert.Equal(t, msg1.Foo, out.Foo) - // assert.Equal(t, msg1.Bar, out.Bar) - // }) -} - -func bucket(table string, partition int) string { - return path.Join(table, fmt.Sprintf("%d", partition)) - -} diff --git a/dax/table.go b/dax/table.go index 24dd0020e..fa37191fc 100644 --- a/dax/table.go +++ b/dax/table.go @@ -380,8 +380,6 @@ type Table struct { Description string `json:"description,omitempty"` Owner string `json:"owner,omitempty"` - CreatedAt int64 `json:"createdAt,omitempty"` - UpdatedAt int64 `json:"updatedAt,omitempty"` UpdatedBy string `json:"updatedBy,omitempty"` } diff --git a/dax/test/boltdb/helpers.go b/dax/test/boltdb/helpers.go deleted file mode 100644 index e7a520270..000000000 --- a/dax/test/boltdb/helpers.go +++ /dev/null @@ -1,50 +0,0 @@ -package boltdb - -import ( - "os" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax/boltdb" - "github.com/stretchr/testify/assert" -) - -func MustGetDB(tb testing.TB) *boltdb.DB { - tb.Helper() - - f, err := os.CreateTemp("", "dax-boltdb") - assert.NoError(tb, err) - - dsn := "file:" + f.Name() - - db := boltdb.NewDB(dsn) - return db -} - -// MustOpenDB returns a new, open DB. Fatal on error. -func MustOpenDB(tb testing.TB) *boltdb.DB { - db := MustGetDB(tb) - - if err := db.Open(); err != nil { - tb.Fatal(err) - } - return db -} - -// MustCloseDB closes the DB. Fatal on error. -func MustCloseDB(tb testing.TB, db *boltdb.DB) { - tb.Helper() - if err := db.Close(); err != nil { - tb.Fatal(err) - } -} - -func CleanupDB(tb testing.TB, path string) { - tb.Helper() - - if path == "" { - return - } - if err := os.Remove(path); err != nil { - tb.Fatal(err) - } -} diff --git a/dax/test/schemar.go b/dax/test/schemar.go deleted file mode 100644 index 0f01bc98b..000000000 --- a/dax/test/schemar.go +++ /dev/null @@ -1,29 +0,0 @@ -package test - -import ( - "os" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax/boltdb" - "github.com/featurebasedb/featurebase/v3/dax/controller/schemar" - schemarbolt "github.com/featurebasedb/featurebase/v3/dax/controller/schemar/boltdb" - testbolt "github.com/featurebasedb/featurebase/v3/dax/test/boltdb" - "github.com/featurebasedb/featurebase/v3/logger" -) - -func NewSchemar(t *testing.T) (schemar schemar.Schemar, cleanup func()) { - td, err := os.MkdirTemp("", "schemartest_*") - if err != nil { - t.Fatalf(": %v", err) - } - db, err := boltdb.NewSvcBolt(td, "schemar", schemarbolt.SchemarBuckets...) - if err != nil { - t.Fatalf("opening boltdb: %v", err) - } - - s := schemarbolt.NewSchemar(db, logger.StderrLogger) - return s, func() { - testbolt.MustCloseDB(t, db) - testbolt.CleanupDB(t, db.Path()) - } -} diff --git a/idk/docker-compose.yml b/idk/docker-compose.yml index 388fb3519..45035a54b 100644 --- a/idk/docker-compose.yml +++ b/idk/docker-compose.yml @@ -147,10 +147,6 @@ services: FEATUREBASE_CONTROLLER_CONFIG_SQLDB_USER: "postgres" FEATUREBASE_CONTROLLER_CONFIG_SQLDB_PASSWORD: "password" FEATUREBASE_CONTROLLER_CONFIG_SQLDB_HOST: "postgres" - SQLDB_DB: "postgres" - SQLDB_USER: "postgres" - SQLDB_PASSWORD: "password" - SQLDB_HOST: "postgres" FEATUREBASE_COMPUTER_RUN: "true" FEATUREBASE_COMPUTER_CONFIG_DATA_DIR: /dax-data/computer FEATUREBASE_COMPUTER_CONFIG_WRITELOGGER_DIR: /dax-data/wl diff --git a/schema.go b/schema.go index 29efd4b29..3dacfa24d 100644 --- a/schema.go +++ b/schema.go @@ -162,8 +162,6 @@ func IndexInfoToTable(ii *IndexInfo) *dax.Table { Description: ii.Options.Description, Owner: ii.Owner, - CreatedAt: ii.CreatedAt, - UpdatedAt: ii.UpdatedAt, UpdatedBy: ii.LastUpdateUser, } @@ -178,9 +176,8 @@ func IndexInfoToTable(ii *IndexInfo) *dax.Table { idType = dax.BaseTypeString } tbl.Fields = append(tbl.Fields, &dax.Field{ - Name: "_id", - Type: idType, - CreatedAt: ii.CreatedAt, + Name: "_id", + Type: idType, }) // Populate the rest of the fields. @@ -268,8 +265,6 @@ func FieldInfoToField(fi *FieldInfo) *dax.Field { TTL: fo.TTL, ForeignIndex: foreignIndex, }, - - CreatedAt: fi.CreatedAt, } } @@ -308,8 +303,6 @@ func TableToIndexInfo(tbl *dax.Table) *IndexInfo { ii := &IndexInfo{ Name: string(tbl.Name), Owner: tbl.Owner, - CreatedAt: tbl.CreatedAt, - UpdatedAt: tbl.UpdatedAt, LastUpdateUser: tbl.UpdatedBy, Options: IndexOptions{ Keys: tbl.StringKeys(), @@ -359,8 +352,7 @@ func FieldToFieldInfo(fld *dax.Field) *FieldInfo { } return &FieldInfo{ - Name: string(fld.Name), - CreatedAt: fld.CreatedAt, + Name: string(fld.Name), Options: FieldOptions{ Type: fieldToFieldType(fld), Base: base,