From 09dd745bdd7644ca71ae89681c96b05f3236dffc Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Sat, 8 Apr 2023 08:16:55 -0500 Subject: [PATCH] it compiles! --- ctl/dax.go | 1 - dax/controller/balancer.go | 4 +- dax/controller/balancer/balancer.go | 157 +++++++---- dax/controller/balancer/free_job_test.go | 86 ------ dax/controller/balancer/free_worker_test.go | 51 ---- dax/controller/balancer/node_test.go | 65 ----- dax/controller/balancer/worker_job_test.go | 140 ---------- dax/controller/sqldb/balancer.go | 7 +- dax/controller/sqldb/freejob.go | 115 -------- dax/controller/sqldb/freeworker.go | 72 ----- dax/controller/sqldb/store.go | 85 ++++-- dax/controller/sqldb/store_test.go | 7 - dax/controller/sqldb/worker.go | 114 -------- dax/controller/sqldb/workerjob.go | 285 -------------------- dax/controller/sqldb/workerjob_test.go | 37 --- dax/controller/store.go | 9 +- dax/test/dax/dax_test.go | 2 - 17 files changed, 174 insertions(+), 1063 deletions(-) delete mode 100644 dax/controller/balancer/free_job_test.go delete mode 100644 dax/controller/balancer/free_worker_test.go delete mode 100644 dax/controller/balancer/node_test.go delete mode 100644 dax/controller/balancer/worker_job_test.go delete mode 100644 dax/controller/sqldb/freejob.go delete mode 100644 dax/controller/sqldb/freeworker.go delete mode 100644 dax/controller/sqldb/worker.go delete mode 100644 dax/controller/sqldb/workerjob.go delete mode 100644 dax/controller/sqldb/workerjob_test.go diff --git a/ctl/dax.go b/ctl/dax.go index 244dedb8a..7cb41529a 100644 --- a/ctl/dax.go +++ b/ctl/dax.go @@ -42,7 +42,6 @@ func BuildDAXFlags(cmd *cobra.Command, srv *server.Command) { flags.StringVar(&srv.Config.WorkerServiceProvider.Config.ControllerAddress, "wsp.config.controller-address", srv.Config.WorkerServiceProvider.Config.ControllerAddress, "Address of remote Controller process.") // Computer flags.BoolVar(&srv.Config.Computer.Run, "computer.run", srv.Config.Computer.Run, "Run the Computer service in process.") - flags.IntVar(&srv.Config.Computer.N, "computer.n", srv.Config.Computer.N, "The number of Computer services to run in process.") flags.StringVar(&srv.Config.Computer.WorkerServiceID, "computer.worker_service_id", srv.Config.Computer.WorkerServiceID, "ID of WorkerService which spawned this computer.") flags.AddFlagSet(serverFlagSet(&srv.Config.Computer.Config, "computer.config")) } diff --git a/dax/controller/balancer.go b/dax/controller/balancer.go index c91afd3c0..e46acebf2 100644 --- a/dax/controller/balancer.go +++ b/dax/controller/balancer.go @@ -17,7 +17,7 @@ type Balancer interface { AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) // RemoveJobs removes jobs for the given database. - RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) + RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.WorkerDiff, error) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) @@ -88,7 +88,7 @@ func (b *NopBalancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax. func (b *NopBalancer) AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) { return []dax.WorkerDiff{}, nil } -func (b *NopBalancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) { +func (b *NopBalancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.WorkerDiff, error) { return []dax.WorkerDiff{}, nil } func (b *NopBalancer) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) { diff --git a/dax/controller/balancer/balancer.go b/dax/controller/balancer/balancer.go index 6189359ec..01d579891 100644 --- a/dax/controller/balancer/balancer.go +++ b/dax/controller/balancer/balancer.go @@ -4,11 +4,12 @@ package balancer import ( "log" "math" + "sort" + "strings" "time" "github.com/featurebasedb/featurebase/v3/dax" "github.com/featurebasedb/featurebase/v3/dax/controller" - "github.com/featurebasedb/featurebase/v3/dax/controller/schemar" "github.com/featurebasedb/featurebase/v3/errors" "github.com/featurebasedb/featurebase/v3/logger" ) @@ -25,31 +26,15 @@ var _ controller.Balancer = (*Balancer)(nil) // jobs. It does not take anything else (such as job size, worker capabilities, // etc) into consideration. type Balancer struct { - workerRegistry controller.WorkerRegistry - - // freeJobs is the set of jobs which have yet to be assigned to a worker. - // This could be because there are no available workers, or because a worker - // has been removed and the jobs for which it was responsible have yet to be - // reassigned. - freeJobs FreeJobService - - freeWorkers FreeWorkerService - - schemar schemar.Schemar - store controller.Store logger logger.Logger } // New returns a new instance of Balancer. -func New(wr controller.WorkerRegistry, fjs FreeJobService, wjs WorkerJobService, fws FreeWorkerService, schemar schemar.Schemar, store controller.Store, logger logger.Logger) *Balancer { +func New(store controller.Store, logger logger.Logger) *Balancer { return &Balancer{ - workerRegistry: wr, - freeJobs: fjs, - freeWorkers: fws, - schemar: schemar, - store: store, + store: store, logger: logger, } @@ -68,7 +53,7 @@ func (b *Balancer) AddWorker(tx dax.Transaction, node *dax.Node) ([]dax.WorkerDi return nil, errors.Wrapf(err, "creating node on node service: %s", node.Address) } - diffs, err := b.balanceDatabase(tx, *qdbidp, node.ServiceID) + diffs, err := b.balanceDatabase(tx, *qdbidp) if err != nil { return nil, errors.Wrap(err, "balancing db") } @@ -93,7 +78,7 @@ func (b *Balancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.Wor if worker.DatabaseID != nil { // Balance the affected database. - if diff, err := b.balanceDatabase(tx, *worker.DatabaseID, worker.ServiceID); err != nil { + if diff, err := b.balanceDatabase(tx, *worker.DatabaseID); err != nil { return nil, errors.Wrapf(err, "balancing database: %s", worker.DatabaseID) } else { diffs.Merge(diff) @@ -141,12 +126,7 @@ func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.D return diffs, nil } - ws, err := b.store.WorkerService(tx, dbid) - if err != nil { - return nil, errors.Wrap(err, "getting worker service") - } - - if cnt, err := b.store.WorkerCount(tx, roleType, ws.ID); err != nil { + if cnt, err := b.store.WorkerCount(tx, roleType, dbid); err != nil { return nil, errors.Wrap(err, "getting worker count") } else if cnt == 0 { if err := b.store.CreateFreeJobs(tx, roleType, dbid, jobs...); err != nil { @@ -154,7 +134,7 @@ func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.D } } - diff, err := b.addDatabaseJobs(tx, roleType, dbid, ws.ID, jobs...) + diff, err := b.addDatabaseJobs(tx, roleType, dbid, jobs...) if err != nil { return nil, errors.Wrapf(err, "adding database jobs: (%s) %s", roleType, dbid) } @@ -164,14 +144,14 @@ func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.D } // addDatabaseJobs adds the job for the provided database. TODO, bad doc -func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID, jobs ...dax.Job) (controller.InternalDiffs, error) { - workerJobs, err := b.store.WorkersJobs(tx, roleType, svcID) +func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, jobs ...dax.Job) (controller.InternalDiffs, error) { + workerJobs, err := b.store.WorkersJobs(tx, roleType, dbid) if err != nil { return nil, errors.Wrapf(err, "getting workers jobs: %s", roleType) } jset := dax.NewSet[dax.Job]() - addrToID := make(map[dax.Address]string) + addrToID := make(map[dax.Address]dax.WorkerID) for _, workerInfo := range workerJobs { addrToID[workerInfo.Address] = workerInfo.ID jset.Merge(dax.NewSet(workerInfo.Jobs...)) @@ -227,11 +207,7 @@ func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, db func (b *Balancer) CurrentState(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerInfo, error) { dbid := qdbid.DatabaseID - ws, err := b.store.WorkerService(tx, dbid) - if err != nil { - return nil, errors.Wrap(err, "getting worker service") - } - return b.store.WorkersJobs(tx, roleType, ws.ID) + return b.store.WorkersJobs(tx, roleType, dbid) } func (b *Balancer) WorkerState(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) (dax.WorkerInfo, error) { @@ -246,25 +222,95 @@ func (b *Balancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid da return diffs.Output(), nil } -func (b *Balancer) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) { - dbid := qdbid.DatabaseID - ws, err := b.store.WorkerService(tx, dbid) +// TODO can probably simplify this with better SQL +func (b *Balancer) WorkersForJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs ...dax.Job) ([]dax.WorkerInfo, error) { + out := make(map[dax.Address]dax.Set[dax.Job]) + + workerJobs, err := b.store.WorkersJobs(tx, roleType, qdbid.DatabaseID) if err != nil { - return nil, errors.Wrap(err, "getting worker service") + return nil, errors.Wrapf(err, "getting worker jobs: (%s) %s", roleType, qdbid) + } + for _, workerInfo := range workerJobs { + jset := dax.NewSet(workerInfo.Jobs...) + + matches := dax.NewSet[dax.Job]() + for _, job := range jobs { + if jset.Contains(job) { + matches.Add(job) + } + } + + if len(matches) > 0 { + out[workerInfo.Address] = matches + } } - diffs, err := b.balanceDatabase(tx, qdbid.DatabaseID, ws.ID) + workers := make([]dax.WorkerInfo, 0, len(out)) + + for addr, jset := range out { + workers = append(workers, dax.WorkerInfo{ + Address: addr, + Jobs: jset.Sorted(), + }) + } + + sort.Sort(dax.WorkerInfos(workers)) + + return workers, nil +} + +// TODO can probably simplify this with better SQL +func (b *Balancer) WorkersForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.WorkerInfo, error) { + out := make(map[dax.Address]dax.Set[dax.Job]) + + dbid := qtid.QualifiedDatabaseID.DatabaseID + + prefix := string(qtid.Key()) + + workerJobs, err := b.store.WorkersJobs(tx, roleType, dbid) + if err != nil { + return nil, errors.Wrapf(err, "getting worker jobs: (%s) %s", roleType, dbid) + } + for _, workerInfo := range workerJobs { + + matches := dax.NewSet[dax.Job]() + for _, job := range workerInfo.Jobs { + if strings.HasPrefix(string(job), prefix) { + matches.Add(job) + } + } + if len(matches) > 0 { + out[workerInfo.Address] = matches + } + } + + workers := make([]dax.WorkerInfo, 0, len(out)) + + for addr, jset := range out { + workers = append(workers, dax.WorkerInfo{ + Address: addr, + Jobs: jset.Sorted(), + }) + } + + sort.Sort(dax.WorkerInfos(workers)) + + return workers, nil +} + +func (b *Balancer) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) { + diffs, err := b.balanceDatabase(tx, qdbid.DatabaseID) if err != nil { return nil, errors.Wrap(err, "balancing") } return diffs.Output(), nil } -func (b *Balancer) balanceDatabase(tx dax.Transaction, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) { +func (b *Balancer) balanceDatabase(tx dax.Transaction, dbid dax.DatabaseID) (controller.InternalDiffs, error) { diffs := controller.NewInternalDiffs() for _, role := range dax.AllRoleTypes { - diff, err := b.balanceDatabaseForRole(tx, role, dbid, svcID) + diff, err := b.balanceDatabaseForRole(tx, role, dbid) if err != nil { return nil, errors.Wrapf(err, "getting worker count: (%s) %s", role, dbid) } @@ -274,26 +320,26 @@ func (b *Balancer) balanceDatabase(tx dax.Transaction, dbid dax.DatabaseID, svcI return diffs, nil } -func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) { +func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID) (controller.InternalDiffs, error) { b.logger.Debugf("balancing database %s for role: %s\n", dbid, roleType) diffs := controller.NewInternalDiffs() // If there are no workers, we can't properly balance. - if cnt, err := b.store.WorkerCount(tx, roleType, svcID); err != nil { + if cnt, err := b.store.WorkerCount(tx, roleType, dbid); err != nil { return nil, errors.Wrapf(err, "getting worker count: (%s) %s", roleType, dbid) } else if cnt == 0 { return controller.InternalDiffs{}, errors.Errorf("unexpected? no workers for database %v (TODO: this used to silently return nil)", dbid) } // Process the freeJobs. - if diff, err := b.processFreeJobs(tx, roleType, dbid, svcID); err != nil { + if diff, err := b.processFreeJobs(tx, roleType, dbid); err != nil { return nil, errors.Wrapf(err, "processing free jobs: (%s) %s", roleType, dbid) } else { diffs.Merge(diff) } // Balance the jobs among workers. - diff, err := b.balanceDatabaseJobs(tx, roleType, dbid, svcID, diffs) + diff, err := b.balanceDatabaseJobs(tx, roleType, dbid, diffs) if err != nil { return nil, errors.Wrap(err, "balancing jobs") } @@ -308,8 +354,8 @@ func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleT // the internalDiffs.merge() method, but we would need to modify that method to // be smarter about the order in which it applies the add/remove operations. // Until that's in place, we'll pass in a value here. -func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID, diffs controller.InternalDiffs) (controller.InternalDiffs, error) { - workerInfos, err := b.store.WorkersJobs(tx, roleType, svcID) +func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, diffs controller.InternalDiffs) (controller.InternalDiffs, error) { + workerInfos, err := b.store.WorkersJobs(tx, roleType, dbid) if err != nil { return nil, errors.Wrapf(err, "getting current state: (%s) %s", roleType, dbid) } @@ -327,7 +373,7 @@ func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType numWorkersAboveMin := numJobs % numWorkers removedJobs := make(dax.Jobs, 0) - addedJobs := make(map[string]dax.Jobs) + addedJobs := make(map[dax.WorkerID]dax.Jobs) // remove jobs loop // @@ -392,14 +438,14 @@ func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType } // processFreeJobs assigns all jobs in the free list to a worker. -func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) { +func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID) (controller.InternalDiffs, error) { diffs := controller.NewInternalDiffs() jobs, err := b.store.ListFreeJobs(tx, roleType, dbid) if err != nil { return nil, errors.Wrapf(err, "listing free jobs: %s", roleType) } - if aj, err := b.addDatabaseJobs(tx, roleType, dbid, svcID, jobs...); err != nil { + if aj, err := b.addDatabaseJobs(tx, roleType, dbid, jobs...); err != nil { return nil, errors.Wrapf(err, "adding jobs: %s", jobs) } else { diffs.Merge(aj) @@ -408,6 +454,13 @@ func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, db return diffs, nil } +func (b *Balancer) ReadNode(tx dax.Transaction, addr dax.Address) (*dax.Node, error) { + return b.store.Worker(tx, addr) +} +func (b *Balancer) Nodes(tx dax.Transaction) ([]*dax.Node, error) { + return b.store.Workers(tx) +} + func (b *Balancer) CreateWorkerServiceProvider(tx dax.Transaction, sp dax.WorkerServiceProvider) error { return b.store.CreateWorkerServiceProvider(tx, sp) } diff --git a/dax/controller/balancer/free_job_test.go b/dax/controller/balancer/free_job_test.go deleted file mode 100644 index 31477940c..000000000 --- a/dax/controller/balancer/free_job_test.go +++ /dev/null @@ -1,86 +0,0 @@ -package balancer_test - -import ( - "context" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/sqldb" - "github.com/stretchr/testify/require" -) - -func TestFreeJobService(t *testing.T) { - tx, err := SQLTransactor.BeginTx(context.Background(), true) - require.NoError(t, err, "getting transaction") - - defer func() { - err := tx.Rollback() - if err != nil { - t.Logf("rolling back: %v", err) - } - }() - - // must have a database to do job stuff - schemar := sqldb.NewSchemar(nil) - err = schemar.CreateDatabase(tx, - &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ID: dbID, Name: dbName}}) - require.NoError(t, err) - - fjSvc := sqldb.NewFreeJobService(nil) - qdbid := dax.QualifiedDatabaseID{OrganizationID: orgID, DatabaseID: dbID} - qtid := dax.QualifiedTableID{ - QualifiedDatabaseID: qdbid, - Name: tableName, - ID: tableID, - } - job1 := dax.Job(qtid.Key() + "job1") - job2 := dax.Job(qtid.Key() + "job2") - job3 := dax.Job(qtid.Key() + "job3") - - node := &dax.Node{ - Address: nodeAddr, - RoleTypes: []dax.RoleType{role}, - } - - err = fjSvc.CreateJobs(tx, role, qdbid, job1, job2, job3) - require.NoError(t, err) - - err = fjSvc.DeleteJob(tx, role, qdbid, job2) - require.NoError(t, err) - - jobs, err := fjSvc.ListJobs(tx, role, qdbid) - require.NoError(t, err) - require.ElementsMatch(t, dax.Jobs{job1, job3}, jobs) - - workerReg := sqldb.NewWorkerRegistry(nil) - err = workerReg.AddWorker(tx, node) - require.NoError(t, err) - - wjSvc := sqldb.NewWorkerJobService(nil) - err = wjSvc.CreateWorker(tx, role, qdbid, nodeAddr) - require.NoError(t, err) - - err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job1) - require.NoError(t, err) - - jobs, err = fjSvc.ListJobs(tx, role, qdbid) - require.NoError(t, err) - require.ElementsMatch(t, dax.Jobs{job3}, jobs) - - err = fjSvc.MarkJobsAsFree(tx, role, qdbid, dax.Jobs{job1}) - require.NoError(t, err) - - jobs, err = fjSvc.ListJobs(tx, role, qdbid) - require.NoError(t, err) - require.ElementsMatch(t, dax.Jobs{job1, job3}, jobs) - - err = fjSvc.DeleteJobsForTable(tx, role, qtid) - require.NoError(t, err) - - jobs, err = fjSvc.ListJobs(tx, role, qdbid) - require.NoError(t, err) - require.Empty(t, jobs) - -} diff --git a/dax/controller/balancer/free_worker_test.go b/dax/controller/balancer/free_worker_test.go deleted file mode 100644 index 7c0f11141..000000000 --- a/dax/controller/balancer/free_worker_test.go +++ /dev/null @@ -1,51 +0,0 @@ -package balancer_test - -import ( - "context" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/sqldb" - "github.com/stretchr/testify/require" -) - -func TestFreeWorkerService(t *testing.T) { - tx, err := SQLTransactor.BeginTx(context.Background(), true) - require.NoError(t, err, "getting transaction") - - defer func() { - err := tx.Rollback() - if err != nil { - t.Logf("rolling back: %v", err) - } - }() - - node1 := &dax.Node{Address: nodeAddr, RoleTypes: dax.AllRoleTypes} - node2 := &dax.Node{Address: nodeAddr2, RoleTypes: dax.AllRoleTypes} - node3 := &dax.Node{Address: nodeAddr3, RoleTypes: dax.AllRoleTypes} - node4 := &dax.Node{Address: nodeAddr4, RoleTypes: dax.AllRoleTypes} - node5 := &dax.Node{Address: nodeAddr5, RoleTypes: dax.AllRoleTypes} - - workerReg := sqldb.NewWorkerRegistry(nil) - - // Add some workers. - require.NoError(t, workerReg.AddWorker(tx, node1)) - require.NoError(t, workerReg.AddWorker(tx, node2)) - require.NoError(t, workerReg.AddWorker(tx, node3)) - require.NoError(t, workerReg.AddWorker(tx, node4)) - require.NoError(t, workerReg.AddWorker(tx, node5)) - - // Remove one of the workers. - require.NoError(t, workerReg.RemoveWorker(tx, node2.Address)) - - fwSvc := sqldb.NewFreeWorkerService(nil) - - addrs, err := fwSvc.ListWorkers(tx, role) - require.NoError(t, err) - require.ElementsMatch(t, dax.Addresses{nodeAddr, nodeAddr3, nodeAddr4, nodeAddr5}, addrs) - - addrs, err = fwSvc.PopWorkers(tx, role, 2) - require.NoError(t, err) - require.Equal(t, 2, len(addrs)) - require.NotEqual(t, addrs[0], addrs[1]) -} diff --git a/dax/controller/balancer/node_test.go b/dax/controller/balancer/node_test.go deleted file mode 100644 index cacd8d869..000000000 --- a/dax/controller/balancer/node_test.go +++ /dev/null @@ -1,65 +0,0 @@ -package balancer_test - -import ( - "context" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/sqldb" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -const ( - nodeAddr = "myaddress" - nodeAddr2 = "myaddress2" - nodeAddr3 = "myaddress3" - nodeAddr4 = "myaddress4" - nodeAddr5 = "myaddress5" -) - -func TestWorkerRegistry(t *testing.T) { - tx, err := SQLTransactor.BeginTx(context.Background(), true) - require.NoError(t, err, "getting transaction") - - defer func() { - err := tx.Rollback() - if err != nil { - t.Logf("rolling back: %v", err) - } - }() - - workerReg := sqldb.NewWorkerRegistry(nil) - - err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr, RoleTypes: []dax.RoleType{dax.RoleTypeCompute}}) - require.NoError(t, err) - - node, err := workerReg.Worker(tx, nodeAddr) - require.NoError(t, err) - require.EqualValues(t, nodeAddr, node.Address) - require.EqualValues(t, 1, len(node.RoleTypes)) - require.EqualValues(t, "compute", node.RoleTypes[0]) - - err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr2, RoleTypes: []dax.RoleType{dax.RoleTypeTranslate, dax.RoleTypeCompute}}) - require.NoError(t, err, "create node 2") - - err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr3, RoleTypes: []dax.RoleType{dax.RoleTypeCompute}}) - require.NoError(t, err, "create node 3") - - nodes, err := workerReg.Workers(tx) - require.NoError(t, err) - assert.EqualValues(t, 3, len(nodes)) - for _, node := range nodes { - assert.Contains(t, node.RoleTypes, dax.RoleType("compute"), "node should have compute role but is: %+v", node) - } - - err = workerReg.RemoveWorker(tx, nodeAddr2) - require.NoError(t, err, "deleting node") - - nodes, err = workerReg.Workers(tx) - require.NoError(t, err) - require.EqualValues(t, 2, len(nodes)) - for _, node := range nodes { - require.Contains(t, node.RoleTypes, dax.RoleType("compute"), "node should have compute role but is: %+v", node) - } -} diff --git a/dax/controller/balancer/worker_job_test.go b/dax/controller/balancer/worker_job_test.go deleted file mode 100644 index 2a43779c8..000000000 --- a/dax/controller/balancer/worker_job_test.go +++ /dev/null @@ -1,140 +0,0 @@ -package balancer_test - -import ( - "context" - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/sqldb" - "github.com/stretchr/testify/require" -) - -const ( - orgID = "orgid" - dbID = "blah" - dbName = "nameofdb" - role = "compute" - tableName = "tbl" - tableID = "tblid" -) - -func TestWorkerJobService(t *testing.T) { - tx, err := SQLTransactor.BeginTx(context.Background(), true) - require.NoError(t, err, "getting transaction") - - defer func() { - err := tx.Rollback() - if err != nil { - t.Logf("rolling back: %v", err) - } - }() - - wspSvc := sqldb.NewWorkerServiceProviderService(nil) - err = wspSvc.CreateWorkerServiceProvider(tx, dax.WorkerServiceProvider{ - ID: "wspID", - Roles: []dax.RoleType{"compute", "translate"}, - Address: "wspaddr", - Description: "wspdesc", - }) - require.NoError(t, err, "creating wsp") - - wspSvc.CreateWorkerService(tx, dax.WorkerService{ - ID: orgID, - Roles: []dax.RoleType{}, - WorkerServiceProviderID: "", - DatabaseID: dbID, - WorkersMin: 1, - WorkersMax: 1, - }) - - // must have a database to do workerjob stuff - schemar := sqldb.NewSchemar(nil) - err = schemar.CreateDatabase(tx, - &dax.QualifiedDatabase{ - OrganizationID: orgID, - Database: dax.Database{ID: dbID, Name: dbName}}) - require.NoError(t, err) - - wjSvc := sqldb.NewWorkerJobService(nil) - qdbid := dax.QualifiedDatabaseID{OrganizationID: orgID, DatabaseID: dbID} - - node := &dax.Node{ - Address: nodeAddr, - ServiceID: "mysvc1", - RoleTypes: []dax.RoleType{role}, - } - - // have to create a free worker before you can create a worker job worker - workerReg := sqldb.NewWorkerRegistry(nil) - err = workerReg.AddWorker(tx, node) - require.NoError(t, err) - - // we create a qtid to prefix jobs so that we can then test the - // "DeleteJobsForTable" method. Is it strange that a table is not - // explicitly mentioned in the interface on the way in, but is in - // the Delete method? Why yes, yes it is... thank you for asking. #TODO - qtid := dax.QualifiedTableID{ - QualifiedDatabaseID: qdbid, - Name: tableName, - ID: tableID, - } - job1 := dax.Job(qtid.Key() + "job1") - job2 := dax.Job(qtid.Key() + "job2") - job3 := dax.Job(qtid.Key() + "job3") - - fjSvc := sqldb.NewFreeJobService(nil) - err = fjSvc.CreateJobs(tx, role, qdbid, job1, job2, job3) - require.NoError(t, err) - - err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job1, job2) - require.NoError(t, err) - - jobs, err := wjSvc.ListJobs(tx, role, qdbid, nodeAddr) - require.NoError(t, err) - require.ElementsMatch(t, dax.Jobs{job1, job2}, jobs) - - workerInfos, err := wjSvc.WorkersJobs(tx, role, qdbid) - require.NoError(t, err) - require.Equal(t, 1, len(workerInfos)) - require.ElementsMatch(t, []dax.Job{job1, job2}, workerInfos[0].Jobs) - - cnt, err := wjSvc.WorkerCount(tx, role, qdbid) - require.NoError(t, err) - require.Equal(t, 1, cnt) - - addrs, err := wjSvc.ListWorkers(tx, role, qdbid) - require.NoError(t, err) - require.ElementsMatch(t, dax.Addresses{nodeAddr}, addrs) - - err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job3) - require.NoError(t, err) - - jcs, err := wjSvc.JobCounts(tx, role, qdbid, nodeAddr) - require.NoError(t, err) - require.Equal(t, 3, jcs[nodeAddr]) - - err = wjSvc.DeleteJob(tx, role, qdbid, nodeAddr, job3) - require.NoError(t, err) - - idiffs, err := wjSvc.DeleteJobsForTable(tx, role, qtid) - require.NoError(t, err) - workerDiffs := idiffs.Output() - require.Equal(t, 1, len(workerDiffs)) - require.EqualValues(t, nodeAddr, workerDiffs[0].Address) - require.Empty(t, workerDiffs[0].AddedJobs) - require.ElementsMatch(t, []dax.Job{job1, job2}, workerDiffs[0].RemovedJobs) - - jobs, err = wjSvc.ListJobs(tx, role, qdbid, nodeAddr) - require.NoError(t, err) - require.Empty(t, jobs) - - dk := wjSvc.DatabaseForWorker(tx, nodeAddr) - require.EqualValues(t, "db__orgid__blah", dk) - - err = wjSvc.ReleaseWorkers(tx, nodeAddr) - require.NoError(t, err) - - addrs, err = wjSvc.ListWorkers(tx, role, qdbid) - require.NoError(t, err) - require.Empty(t, addrs) -} diff --git a/dax/controller/sqldb/balancer.go b/dax/controller/sqldb/balancer.go index 4c8aee363..81fa1ff6b 100644 --- a/dax/controller/sqldb/balancer.go +++ b/dax/controller/sqldb/balancer.go @@ -7,12 +7,7 @@ import ( // NewBalancer returns a new instance of controller.Balancer. func NewBalancer(log logger.Logger) *balancer.Balancer { - schemar := NewSchemar(log) - fjs := NewFreeJobService(log) - wjs := NewWorkerJobService(log) - fws := NewFreeWorkerService(log) - ns := NewWorkerRegistry(log) store := NewStore(log) - return balancer.New(ns, fjs, wjs, fws, schemar, store, log) + return balancer.New(store, log) } diff --git a/dax/controller/sqldb/freejob.go b/dax/controller/sqldb/freejob.go deleted file mode 100644 index 505703bc0..000000000 --- a/dax/controller/sqldb/freejob.go +++ /dev/null @@ -1,115 +0,0 @@ -package sqldb - -import ( - "fmt" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/balancer" - "github.com/featurebasedb/featurebase/v3/dax/models" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/pkg/errors" -) - -func NewFreeJobService(log logger.Logger) balancer.FreeJobService { - if log == nil { - log = logger.NopLogger - } - return &freeJobService{ - log: log, - } -} - -type freeJobService struct { - log logger.Logger -} - -func (fj *freeJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job ...dax.Job) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - // jobNames is used as input to the "name in (...)" query. - jobNames := make([]interface{}, 0, len(job)) - for i := range job { - jobNames = append(jobNames, job[i].Job()) - } - - // existing will contain the list of jobs which already exist. - existing := &models.Jobs{} - if err := dt.C.Where("name in (?)", jobNames...).All(existing); err != nil { - return errors.Wrap(err, "getting existing jobs") - } - - jobs := make(models.Jobs, 0, len(job)) - for _, j := range job { - // Check to be sure this job doesn't already exist. - if existing.Contains(j) { - continue - } - jobs = append(jobs, models.Job{ - Name: j, - Role: roleType, - DatabaseID: qdbid.DatabaseID, - }) - } - - if len(jobs) == 0 { - return nil - } - - err := dt.C.Create(jobs) - return errors.Wrap(err, "creating free jobs") -} - -func (fj *freeJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job dax.Job) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - err := dt.C.RawQuery("DELETE from jobs where role = ? and database_id = ? and name = ? and worker_id is NULL", roleType, qdbid.DatabaseID, job).Exec() - return errors.Wrap(err, "deleting") -} - -func (fj *freeJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - err := dt.C.RawQuery("DELETE from jobs where role = ? and database_id = ? and name LIKE ? and worker_id is NULL", - roleType, qtid.DatabaseID, fmt.Sprintf("%s%%", qtid.Key())).Exec() - return errors.Wrap(err, "deleting") -} - -func (fj *freeJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Jobs, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - jobs := make(models.Jobs, 0) - err := dt.C.Where("role = ? and database_id = ? and worker_id is NULL", roleType, qdbid.DatabaseID).Order("name asc").All(&jobs) - if err != nil { - return nil, errors.Wrap(err, "querying for jobs") - } - - djs := make(dax.Jobs, len(jobs)) - for i, job := range jobs { - djs[i] = job.Name - } - return djs, nil -} - -// MarkJobsAsFree disassociates any worker that was previously assigned to this -// job. -func (fj *freeJobService) MarkJobsAsFree(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs dax.Jobs) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - err := dt.C.RawQuery("UPDATE jobs SET worker_id = NULL WHERE role = ? and database_id = ?", roleType, qdbid.DatabaseID).Exec() - return errors.Wrap(err, "marking jobs free") -} diff --git a/dax/controller/sqldb/freeworker.go b/dax/controller/sqldb/freeworker.go deleted file mode 100644 index 4ddc0a2e8..000000000 --- a/dax/controller/sqldb/freeworker.go +++ /dev/null @@ -1,72 +0,0 @@ -package sqldb - -import ( - "fmt" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/balancer" - "github.com/featurebasedb/featurebase/v3/dax/models" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/pkg/errors" -) - -func NewFreeWorkerService(log logger.Logger) balancer.FreeWorkerService { - if log == nil { - log = logger.NopLogger - } - return &freeWorkerService{ - log: log, - } -} - -type freeWorkerService struct { - log logger.Logger -} - -func (fw *freeWorkerService) PopWorkers(tx dax.Transaction, roleType dax.RoleType, num int) ([]dax.Address, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - results := make([]struct { - Address dax.Address `db:"address"` - }, 0, num) - sel := fmt.Sprintf("select address from workers where role_%s = true and database_id is NULL limit ?", roleType) - err := dt.C.RawQuery(sel, num).All(&results) - if err != nil { - return nil, errors.Wrap(err, "querying") - } - if len(results) < num { - return nil, errors.Errorf("not enough free workers to get: wanted %d, have: %d", num, len(results)) - } - - ret := make([]dax.Address, num) - for i, res := range results { - ret[i] = res.Address - } - - return ret, nil -} - -func (fw *freeWorkerService) ListWorkers(tx dax.Transaction, roleType dax.RoleType) (dax.Addresses, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - workers := make(models.Workers, 0) - where := fmt.Sprintf("role_%s = true and database_id is NULL", roleType) - err := dt.C.Select("address").Where(where).Order("address asc").All(&workers) - if err != nil { - return nil, errors.Wrap(err, "querying for free workers") - } - - ret := make(dax.Addresses, len(workers)) - for i, w := range workers { - ret[i] = w.Address - } - - return ret, nil - -} diff --git a/dax/controller/sqldb/store.go b/dax/controller/sqldb/store.go index 1b929fca8..47f288852 100644 --- a/dax/controller/sqldb/store.go +++ b/dax/controller/sqldb/store.go @@ -5,7 +5,6 @@ import ( "github.com/featurebasedb/featurebase/v3/dax" "github.com/featurebasedb/featurebase/v3/dax/controller" - "github.com/featurebasedb/featurebase/v3/dax/controller/balancer" "github.com/featurebasedb/featurebase/v3/dax/models" "github.com/featurebasedb/featurebase/v3/errors" "github.com/featurebasedb/featurebase/v3/logger" @@ -62,24 +61,8 @@ func (s *store) AddWorker(tx dax.Transaction, node *dax.Node) (*dax.DatabaseID, return (*dax.DatabaseID)(&workerSvc.DatabaseID.String), nil } -// WorkerCount returns the number of workers with the given service -// ID. -func (s *store) WorkerCount(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) (int, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - cnt, err := dt.C.Where(fmt.Sprintf("service_id = ? AND role_%s = true", role), svcID).Count(&models.Worker{}) - if err != nil { - return 0, errors.Wrapf(err, "counting workers for service '%s'", svcID) - } - return cnt, nil -} - -// WorkerCount returns the number of workers with the given service -// ID. -func (s *store) WorkerCountDatabase(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (int, error) { +// WorkerCount returns the number of workers associated with the given database. +func (s *store) WorkerCount(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (int, error) { dt, ok := tx.(*DaxTransaction) if !ok { return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") @@ -158,7 +141,7 @@ func (s *store) WorkerJobs(tx dax.Transaction, role dax.RoleType, addr dax.Addre return wi, nil } -func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) ([]dax.WorkerInfo, error) { +func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) ([]dax.WorkerInfo, error) { dt, ok := tx.(*DaxTransaction) if !ok { return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") @@ -166,7 +149,7 @@ func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.Wor results := []struct { Address dax.Address `db:"address"` - WorkerID string `db:"worker_id"` + WorkerID dax.WorkerID `db:"worker_id"` JobName nulls.String `db:"job_name"` }{} @@ -176,14 +159,15 @@ func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.Wor jobs.name as job_name FROM workers LEFT JOIN jobs on jobs.worker_id = workers.id + INNER JOIN worker_service ws ON ws.id = workers.service_id WHERE jobs.role = ? - workers.service_id = ? + ws.database_id = ? ORDER BY address ASC job_name ASC` - err := dt.C.RawQuery(sql, role, svcID).All(&results) + err := dt.C.RawQuery(sql, role, dbid).All(&results) if err != nil { return nil, errors.Wrap(err, "querying") } @@ -390,7 +374,7 @@ func (s *store) DeleteJobsForTable(tx dax.Transaction, role dax.RoleType, qtid d return nil, errors.Wrap(err, "querying for jobs") } - idiffs := make(balancer.InternalDiffs) + idiffs := make(controller.InternalDiffs) ids := make([]uuid.UUID, 0, len(results)) for _, job := range results { idiffs.Removed(job.Address, job.Name) @@ -403,3 +387,56 @@ func (s *store) DeleteJobsForTable(tx dax.Transaction, role dax.RoleType, qtid d return idiffs, errors.Wrap(err, "deleting jobs") } + +func (s *store) Workers(tx dax.Transaction) ([]*dax.Node, error) { + dt, ok := tx.(*DaxTransaction) + if !ok { + return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") + } + + workers := []*models.Worker{} + dt.C.Eager().Order("address asc").All(&workers) + + ret := make([]*dax.Node, len(workers)) + for i, worker := range workers { + ret[i] = &dax.Node{ + Address: worker.Address, + RoleTypes: workerRoleTypes(worker), + } + } + + return ret, nil +} + +func (w *store) Worker(tx dax.Transaction, addr dax.Address) (*dax.Node, error) { + dt, ok := tx.(*DaxTransaction) + if !ok { + return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") + } + + worker := &models.Worker{} + err := dt.C.Eager().Where("address = ?", addr).First(worker) + if err != nil { + return nil, errors.Wrapf(err, "getting worker: %s", addr) + } + + return &dax.Node{ + Address: worker.Address, + RoleTypes: workerRoleTypes(worker), + }, nil +} + +func workerRoleTypes(worker *models.Worker) []dax.RoleType { + roleTypes := make([]dax.RoleType, 0) + if worker.RoleCompute { + roleTypes = append(roleTypes, dax.RoleTypeCompute) + } + if worker.RoleTranslate { + roleTypes = append(roleTypes, dax.RoleTypeTranslate) + } + if worker.RoleQuery { + roleTypes = append(roleTypes, dax.RoleTypeQuery) + } + + return roleTypes +} diff --git a/dax/controller/sqldb/store_test.go b/dax/controller/sqldb/store_test.go index 370af80b1..d43e41164 100644 --- a/dax/controller/sqldb/store_test.go +++ b/dax/controller/sqldb/store_test.go @@ -106,12 +106,6 @@ func TestStore(t *testing.T) { require.Nil(t, qdbidp) }) - t.Run("test WorkerCount", func(t *testing.T) { - cnt, err := s.WorkerCount(tx, dax.RoleTypeCompute, ws.ID) - require.NoError(t, err) - require.Equal(t, 1, cnt) - }) - t.Run("list workers", func(t *testing.T) { addrs, err := s.ListWorkers(tx, dax.RoleTypeCompute, ws.ID) require.NoError(t, err) @@ -126,5 +120,4 @@ func TestStore(t *testing.T) { require.Equal(t, ws.ID, node.ServiceID) require.Equal(t, expectedDBID, node.DatabaseID) }) - } diff --git a/dax/controller/sqldb/worker.go b/dax/controller/sqldb/worker.go deleted file mode 100644 index 4e1c8692c..000000000 --- a/dax/controller/sqldb/worker.go +++ /dev/null @@ -1,114 +0,0 @@ -package sqldb - -import ( - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller" - "github.com/featurebasedb/featurebase/v3/dax/models" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/pkg/errors" -) - -var _ controller.WorkerRegistry = (*workerRegistry)(nil) - -func NewWorkerRegistry(log logger.Logger) *workerRegistry { - if log == nil { - log = logger.NopLogger - } - return &workerRegistry{ - log: log, - } -} - -type workerRegistry struct { - log logger.Logger -} - -func (w *workerRegistry) AddWorker(tx dax.Transaction, node *dax.Node) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - worker := &models.Worker{ - Address: node.Address, - ServiceID: node.ServiceID, - } - for _, roleType := range node.RoleTypes { - if err := worker.SetRole(roleType); err != nil { - return errors.Wrapf(err, "setting role: %s", roleType) - } - } - - return dt.C.Create(worker) -} - -func (w *workerRegistry) Worker(tx dax.Transaction, addr dax.Address) (*dax.Node, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - worker := &models.Worker{} - err := dt.C.Eager().Where("address = ?", addr).First(worker) - if err != nil { - return nil, errors.Wrapf(err, "getting worker: %s", addr) - } - - return &dax.Node{ - Address: worker.Address, - RoleTypes: workerRoleTypes(worker), - }, nil -} - -func (w *workerRegistry) RemoveWorker(tx dax.Transaction, addr dax.Address) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - worker := &models.Worker{} - err := dt.C.Eager().Where("address = ?", addr).First(worker) - if isNoRowsError(err) { - return nil - } else if err != nil { - return errors.Wrapf(err, "finding worker: %s", addr) - } - - err = dt.C.Destroy(worker) - return errors.Wrap(err, "destroying worker") -} - -func (w *workerRegistry) Workers(tx dax.Transaction) ([]*dax.Node, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - workers := []*models.Worker{} - dt.C.Eager().Order("address asc").All(&workers) - - ret := make([]*dax.Node, len(workers)) - for i, worker := range workers { - ret[i] = &dax.Node{ - Address: worker.Address, - RoleTypes: workerRoleTypes(worker), - } - } - - return ret, nil -} - -func workerRoleTypes(worker *models.Worker) []dax.RoleType { - roleTypes := make([]dax.RoleType, 0) - if worker.RoleCompute { - roleTypes = append(roleTypes, dax.RoleTypeCompute) - } - if worker.RoleTranslate { - roleTypes = append(roleTypes, dax.RoleTypeTranslate) - } - if worker.RoleQuery { - roleTypes = append(roleTypes, dax.RoleTypeQuery) - } - - return roleTypes -} diff --git a/dax/controller/sqldb/workerjob.go b/dax/controller/sqldb/workerjob.go deleted file mode 100644 index 2fdd8cf58..000000000 --- a/dax/controller/sqldb/workerjob.go +++ /dev/null @@ -1,285 +0,0 @@ -package sqldb - -import ( - "fmt" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/controller/balancer" - "github.com/featurebasedb/featurebase/v3/dax/models" - "github.com/featurebasedb/featurebase/v3/logger" - "github.com/gofrs/uuid" - "github.com/pkg/errors" -) - -func NewWorkerJobService(log logger.Logger) balancer.WorkerJobService { - if log == nil { - log = logger.NopLogger - } - return &workerJobService{ - log: log, - } -} - -type workerJobService struct { - log logger.Logger -} - -// WorkersJobs returns all the workers for the database along with the jobs -// associated to each worker, even if the number of jobs is 0. -func (w *workerJobService) WorkersJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerInfo, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - // First, get all workers for the database. - workers := models.Workers{} - q := dt.C.Q().Eager() - q = q.InnerJoin("services", "services.id = workers.service_id") - q = q.Where(fmt.Sprintf("services.role_%s = true and services.database_id = ?", roleType), qdbid.DatabaseID) - err := q.Order("address asc").All(&workers) - if err != nil { - return nil, errors.Wrap(err, "getting workers") - } - - // Then, get the jobs for each worker. Ideally, we would do this in a single - // sql query, but it wasn't clear how to do an Eager() LeftJoin() where - // there is a where clause condition on the right side of the join (in this - // case, `jobs.role = ?`). - ret := make([]dax.WorkerInfo, len(workers)) - for i, worker := range workers { - ret[i].Address = worker.Address - jobs, err := jobsForWorker(dt, &worker, roleType) - if err != nil { - return nil, errors.Wrap(err, "getting jobs for worker") - } - ret[i].Jobs = jobs - } - - return ret, nil -} - -func jobsForWorker(dt *DaxTransaction, worker *models.Worker, roleType dax.RoleType) ([]dax.Job, error) { - jobs := models.Jobs{} - if err := dt.C.Where("worker_id = ? and role = ?", worker.ID, roleType).Order("name asc").All(&jobs); err != nil { - return nil, errors.Wrapf(err, "getting jobs for worker: %s", worker.ID) - } - ret := make([]dax.Job, len(jobs)) - for i := range jobs { - ret[i] = jobs[i].Name - } - return ret, nil -} - -func (w *workerJobService) WorkerCount(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (int, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - var count int - err := dt.C.RawQuery("SELECT count(*) from workers inner join services on workers.service_id = services.id where services.database_id = ?", qdbid.DatabaseID).First(&count) - return count, errors.Wrap(err, "getting count") -} - -func (w *workerJobService) ListWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Addresses, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - workers := models.Workers{} - sql := fmt.Sprintf("role_%s = true and database_id = ?", roleType) - err := dt.C.Select("address").Where(sql, qdbid.DatabaseID).Order("address asc").All(&workers) - if err != nil { - return nil, errors.Wrap(err, "getting workers") - } - - ret := make(dax.Addresses, len(workers)) - for i, wrkr := range workers { - ret[i] = wrkr.Address - } - - return ret, nil -} - -func (w *workerJobService) AssignWorkerToJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job ...dax.Job) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - worker := &models.Worker{} - sql := fmt.Sprintf("address = ? and role_%s = true", roleType) - err := dt.C.Where(sql, addr).First(worker) - if err != nil { - return errors.Wrapf(err, "getting worker: (%s) %s", roleType, addr) - } - - jobs := models.Jobs{} - err = dt.C.RawQuery("UPDATE jobs SET worker_id = ? WHERE role = ? and name in (?) RETURNING jobs.ID, jobs.Name", worker.ID, roleType, job).All(&jobs) - if err != nil { - return errors.Wrap(err, "updating jobs") - } - - // Assign jobs not in "jobs", and therefore didn't get updated by the - // previous sql statement. - toBeAssigned := jobsNotAssigned(job, jobs, roleType, worker, qdbid.DatabaseID) - - if err := dt.C.Create(toBeAssigned); err != nil { - return errors.Wrap(err, "creating jobs") - } - - return nil -} - -func jobsNotAssigned(incomingJobs []dax.Job, assigned models.Jobs, roleType dax.RoleType, worker *models.Worker, dbid dax.DatabaseID) (toBeAssigned models.Jobs) { -outer: - for _, incJob := range incomingJobs { - for _, assignedJob := range assigned { - if assignedJob.Name == incJob { - continue outer - } - } - toBeAssigned = append(toBeAssigned, - models.Job{ - Name: incJob, - Role: roleType, - DatabaseID: dbid, - Worker: worker, - }, - ) - } - return toBeAssigned -} - -func (w *workerJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job dax.Job) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - worker := &models.Worker{} - sql := fmt.Sprintf("role_%s = true and database_id = ? and address = ?", roleType) - err := dt.C.Select("id").Where(sql, qdbid.DatabaseID, addr).First(worker) - if err != nil { - return errors.Wrap(err, "getting worker") - } - - jerb := &models.Job{} - dt.C.Select("id").Where("role = ? and worker_id = ? and name = ?", roleType, worker.ID, job).First(jerb) - if err != nil { - return errors.Wrap(err, "getting job") - } - - err = dt.C.Destroy(jerb) - if err != nil { - return errors.Wrap(err, "destroying job") - } - - return nil -} - -func (w *workerJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) (balancer.InternalDiffs, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - results := []struct { - ID uuid.UUID `db:"id"` - Name dax.Job `db:"name"` - Address dax.Address `db:"address"` - }{} - err := dt.C.RawQuery("select j.id, j.name, w.address from jobs j inner join workers w on j.worker_id = w.id where j.role = ? and j.database_id = ? and j.name LIKE ?", roleType, qtid.QualifiedDatabaseID.DatabaseID, fmt.Sprintf("%s%%", qtid.Key())).All(&results) - if err != nil { - return nil, errors.Wrap(err, "querying for jobs") - } - - idiffs := make(balancer.InternalDiffs) - ids := make([]uuid.UUID, 0, len(results)) - for _, job := range results { - idiffs.Removed(job.Address, job.Name) - ids = append(ids, job.ID) - } - - if len(ids) > 0 { - err = dt.C.RawQuery("DELETE FROM jobs WHERE id in (?)", ids).Exec() - } - - return idiffs, errors.Wrap(err, "deleting jobs") -} - -func (w *workerJobService) JobCounts(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addrs ...dax.Address) (map[dax.Address]int, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - results := []struct { - Address dax.Address `db:"address"` - Count int `db:"count"` - }{} - var err error - if len(addrs) == 0 { - qstring := `select address, count(*) as count - from workers w inner join jobs j on j.worker_id = w.id - where w.database_id = ? and w.role_%s = true - and j.role = ? - group by w.address` - sql := fmt.Sprintf(qstring, roleType) - err = dt.C.RawQuery(sql, qdbid.DatabaseID, roleType).All(&results) - } else { - qstring := `select address, count(*) as count - from workers w inner join jobs j on j.worker_id = w.id - where w.database_id = ? and w.role_%s = true - and j.role = ? - and w.address in (?) - group by w.address` - sql := fmt.Sprintf(qstring, roleType) - err = dt.C.RawQuery(sql, qdbid.DatabaseID, roleType, addrs).All(&results) - } - if err != nil { - return nil, errors.Wrap(err, "querying for jobs") - } - ret := make(map[dax.Address]int) - for _, res := range results { - ret[res.Address] = res.Count - } - - return ret, nil -} - -func (w *workerJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) (dax.Jobs, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - worker := &models.Worker{} - sql := fmt.Sprintf("role_%s = true and database_id = ? and address = ?", roleType) - err := dt.C.Where(sql, qdbid.DatabaseID, addr).First(worker) - if isNoRowsError(err) { - return nil, nil - } else if err != nil { - return nil, errors.Wrap(err, "getting worker") - } - - return jobsForWorker(dt, worker, roleType) -} - -func (w *workerJobService) DatabaseForWorker(tx dax.Transaction, addr dax.Address) dax.DatabaseKey { - dt, ok := tx.(*DaxTransaction) - if !ok { - panic("wrong transaction type passed to sqldb DatabaseForWorker") - } - - db := &models.Database{} - err := dt.C.RawQuery("select d.ID, d.organization_id from databases d inner join workers w on d.id = w.database_id where w.address = ?", addr).First(db) - if isNoRowsError(err) { - return "" - } else if err != nil { - panic(err) - } - - return dax.QualifiedDatabase{OrganizationID: dax.OrganizationID(db.OrganizationID), Database: dax.Database{ID: dax.DatabaseID(db.ID)}}.Key() -} diff --git a/dax/controller/sqldb/workerjob_test.go b/dax/controller/sqldb/workerjob_test.go deleted file mode 100644 index a07e5f4dd..000000000 --- a/dax/controller/sqldb/workerjob_test.go +++ /dev/null @@ -1,37 +0,0 @@ -package sqldb - -import ( - "testing" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/featurebasedb/featurebase/v3/dax/models" - "github.com/gofrs/uuid" - "github.com/stretchr/testify/require" -) - -func TestJobsNotUpdated(t *testing.T) { - u, _ := uuid.NewV4() - u2, _ := uuid.NewV4() - incJobs := []dax.Job{"job1", "job2", "job3", "job4"} - created := models.Jobs{ - models.Job{ - ID: u, - Name: "job2", - }, - } - - toCreate := jobsNotAssigned(incJobs, created, dax.RoleTypeCompute, &models.Worker{ - ID: u2, - RoleCompute: true, - }, "dbid") - - require.Equal(t, 3, len(toCreate)) - - // Test when 0 jobs are updated - toCreate = jobsNotAssigned(incJobs, models.Jobs{}, dax.RoleTypeCompute, &models.Worker{ - ID: u2, - RoleCompute: true, - }, "dbid") - - require.Equal(t, 4, len(toCreate)) -} diff --git a/dax/controller/store.go b/dax/controller/store.go index 52f12934c..ccadd22d3 100644 --- a/dax/controller/store.go +++ b/dax/controller/store.go @@ -36,11 +36,10 @@ type Store interface { WorkerService(tx dax.Transaction, dbid dax.DatabaseID) (dax.WorkerService, error) AddWorker(tx dax.Transaction, node *dax.Node) (*dax.DatabaseID, error) - WorkerCount(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) (int, error) - WorkerCountDatabase(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (int, error) + WorkerCount(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (int, error) ListFreeJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (dax.Jobs, error) - WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) ([]dax.WorkerInfo, error) - AssignWorkerToJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, workerID string, jobs ...dax.Job) error + WorkersJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) ([]dax.WorkerInfo, error) + AssignWorkerToJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, workerID dax.WorkerID, jobs ...dax.Job) error ListWorkers(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) (dax.Addresses, error) AssignFreeServiceToDatabase(tx dax.Transaction, wspID dax.WorkerServiceProviderID, qdb *dax.QualifiedDatabase) (*dax.WorkerService, error) @@ -51,4 +50,6 @@ type Store interface { CreateFreeJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, jobs ...dax.Job) error DeleteJobsForTable(tx dax.Transaction, role dax.RoleType, qtid dax.QualifiedTableID) (InternalDiffs, error) WorkerJobs(tx dax.Transaction, role dax.RoleType, addr dax.Address) (dax.WorkerInfo, error) + Workers(tx dax.Transaction) ([]*dax.Node, error) + Worker(tx dax.Transaction, addr dax.Address) (*dax.Node, error) } diff --git a/dax/test/dax/dax_test.go b/dax/test/dax/dax_test.go index 6d79e1c13..5ae82e878 100644 --- a/dax/test/dax/dax_test.go +++ b/dax/test/dax/dax_test.go @@ -224,7 +224,6 @@ func TestDAXIntegration(t *testing.T) { t.Run("Poller", func(t *testing.T) { cfg := test.DefaultConfig() - cfg.Computer.N = 2 opt := server.OptCommandConfig(cfg) mc := test.MustRunManagedCommand(t, opt) @@ -690,7 +689,6 @@ func TestDAXIntegration(t *testing.T) { t.Run("DatabaseOptions", func(t *testing.T) { cfg := test.DefaultConfig() - cfg.Computer.N = 4 opt := server.OptCommandConfig(cfg) mc := test.MustRunManagedCommand(t, opt) defer mc.Close()