diff --git a/dax/controller/balancer.go b/dax/controller/balancer.go index 87edb3cfd..63a57fd26 100644 --- a/dax/controller/balancer.go +++ b/dax/controller/balancer.go @@ -13,8 +13,8 @@ type Balancer interface { // be either transferred to other workers or placed on the free job list. RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.WorkerDiff, error) - // FreeWorkers dissociates the given workers from a database. - FreeWorkers(tx dax.Transaction, addrs ...dax.Address) error + // ReleaseWorkers dissociates the given workers from a database. + ReleaseWorkers(tx dax.Transaction, addrs ...dax.Address) error // AddJobs adds new jobs for the given database. AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) @@ -64,7 +64,7 @@ func (b *NopBalancer) AddWorker(tx dax.Transaction, node *dax.Node) ([]dax.Worke func (b *NopBalancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.WorkerDiff, error) { return []dax.WorkerDiff{}, nil } -func (b *NopBalancer) FreeWorkers(tx dax.Transaction, addrs ...dax.Address) error { +func (b *NopBalancer) ReleaseWorkers(tx dax.Transaction, addrs ...dax.Address) error { return nil } func (b *NopBalancer) AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) { diff --git a/dax/controller/balancer/balancer.go b/dax/controller/balancer/balancer.go index 6e7686f23..f93dad12f 100644 --- a/dax/controller/balancer/balancer.go +++ b/dax/controller/balancer/balancer.go @@ -28,7 +28,7 @@ type Balancer struct { // current represents the current state of worker/job assigments. current WorkerJobService - nodeService controller.NodeService + 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 @@ -44,47 +44,37 @@ type Balancer struct { } // New returns a new instance of Balancer. -func New(ns controller.NodeService, fjs FreeJobService, wjs WorkerJobService, fws FreeWorkerService, schemar schemar.Schemar, logger logger.Logger) *Balancer { +func New(wr controller.WorkerRegistry, fjs FreeJobService, wjs WorkerJobService, fws FreeWorkerService, schemar schemar.Schemar, logger logger.Logger) *Balancer { return &Balancer{ - current: wjs, - nodeService: ns, - freeJobs: fjs, - freeWorkers: fws, - schemar: schemar, - logger: logger, + current: wjs, + workerRegistry: wr, + freeJobs: fjs, + freeWorkers: fws, + schemar: schemar, + logger: logger, } } -// AddWorker adds the given Node to the Balancer's available worker pool. -// TODO(tlt): this method takes a Node (as opposed to a Worker) because in the -// future we may want to maintain separate worker pools based on RoleType -// (compute, translate, etc.). +// AddWorker adds the given Node to the Balancer's available worker pool. Note +// that a node is used for ALL of the role types specified. In other words, +// specifying roleTypes = {compute, translate}, does not mean that the node can +// be used as either a compute worker or a translate worker. It means that it +// will be used as both. func (b *Balancer) AddWorker(tx dax.Transaction, node *dax.Node) ([]dax.WorkerDiff, error) { - addr := node.Address - b.logger.Debugf("AddWorker(%s)", addr) + b.logger.Debugf("AddWorker(%s)", node.Address) - if err := b.nodeService.CreateNode(tx, addr, node); err != nil { - return nil, errors.Wrapf(err, "creating node on node service: %s", addr) + if err := b.workerRegistry.AddWorker(tx, node); err != nil { + return nil, errors.Wrapf(err, "creating node on node service: %s", node.Address) } diffs := NewInternalDiffs() - // This logic means that a node is used for ALL of the role types specified. - // In other words, specifying roleTypes = {compute, translate}, does not - // mean that the node can be used as either a compute worker or a translate - // worker. It means that it will be used as both. - for _, rt := range node.RoleTypes { - if err := b.addWorker(tx, rt, addr); err != nil { - return nil, errors.Wrapf(err, "adding worker: (%s) %s", rt, addr) - } - } - - // Process the freeWorkers. + // Process the newly added workers. // TODO(tlt): this is a little heavy-handed. I'm sure we'll need to be more - // intentional about knowing which databases needs workers, as opposed to + // intentional about knowing which databases need workers, as opposed to // this brute force loop over all databases every time. if diff, err := b.balance(tx); err != nil { - return nil, errors.Wrapf(err, "balancing new worker: %s", addr) + return nil, errors.Wrapf(err, "balancing new worker: %s", node.Address) } else { diffs.Merge(diff) } @@ -92,23 +82,9 @@ func (b *Balancer) AddWorker(tx dax.Transaction, node *dax.Node) ([]dax.WorkerDi return diffs.Output(), nil } -// addWorker adds a worker to the free worker list. From there, it can be used -// by any database which needs a worker. -func (b *Balancer) addWorker(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) error { - // If this worker already exists, don't do anything. - if dbkey := b.current.DatabaseForWorker(tx, addr); dbkey != "" { - return nil - } +func (b *Balancer) assignMinWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (InternalDiffs, error) { + b.logger.Debugf("assigning min workers for '%s', '%s'", roleType, qdbid) - if err := b.freeWorkers.AddWorkers(tx, roleType, addr); err != nil { - return errors.Wrap(err, "adding free worker") - } - - return nil -} - -func (b *Balancer) assignMinWorkers(tx dax.Transaction, roleType dax.RoleType) (InternalDiffs, error) { - b.logger.Debugf("assigning min workers for '%s'", roleType) // Find out how many free workers we have. freeWorkers, err := b.freeWorkers.ListWorkers(tx, roleType) if err != nil { @@ -122,10 +98,12 @@ func (b *Balancer) assignMinWorkers(tx dax.Transaction, roleType dax.RoleType) ( return InternalDiffs{}, nil } - // Get all database and their minWorkerCount (Database.Options.WorkersMin). - qdbs, err := b.schemar.Databases(tx, "") + // Get database and its minWorkerCount (Database.Options.WorkersMin). This + // used to get all databases, but now this method is specific to a single + // database. That's why we just get the one here. + qdbs, err := b.schemar.Databases(tx, qdbid.OrganizationID, qdbid.DatabaseID) if err != nil { - return nil, errors.Wrap(err, "getting all database") + return nil, errors.Wrap(err, "getting database") } // Create a map[database]int where int is the number of workers required to @@ -244,13 +222,12 @@ func (b *Balancer) databaseHasJobs(tx dax.Transaction, roleType dax.RoleType, qd func (b *Balancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.WorkerDiff, error) { diffs := NewInternalDiffs() - ////// The rest is database specific. //////////// - - // See if the worker is assigned to a database. If it's not, return early. + // See if the worker is assigned to a database. If it is, disassociate the + // worker from all of its jobs for the database. dbkey := b.current.DatabaseForWorker(tx, addr) if dbkey != "" { qdbid := dbkey.QualifiedDatabaseID() - for _, rt := range []dax.RoleType{dax.RoleTypeCompute, dax.RoleTypeTranslate} { + for _, rt := range dax.AllRoleTypes { if diff, err := b.removeDatabaseWorker(tx, rt, qdbid, addr); err != nil { return nil, errors.Wrapf(err, "removing worker: (%s) %s", rt, addr) } else { @@ -259,15 +236,8 @@ func (b *Balancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.Wor } } - // Remove the worker from the free worker list (if it's there). - for _, rt := range []dax.RoleType{dax.RoleTypeCompute, dax.RoleTypeTranslate} { - if err := b.freeWorkers.RemoveWorker(tx, rt, addr); err != nil { - return nil, errors.Wrapf(err, "removing worker from free list: (%s) %s", rt, addr) - } - } - - // Remove the worker (i.e. Node) from the node service. - if err := b.nodeService.DeleteNode(tx, addr); err != nil { + // Remove the worker from the worker registry. + if err := b.workerRegistry.RemoveWorker(tx, addr); err != nil { return nil, errors.Wrapf(err, "deleting node from node service: %s", addr) } @@ -284,6 +254,8 @@ func (b *Balancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.Wor return diffs.Output(), nil } +// removeDatabaseWorker is used to remove a worker that has been associated with +// a database. The worker here is determined by address. func (b *Balancer) removeDatabaseWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) (InternalDiffs, error) { jobs, err := b.current.ListJobs(tx, roleType, qdbid, addr) if err != nil { @@ -291,13 +263,8 @@ func (b *Balancer) removeDatabaseWorker(tx dax.Transaction, roleType dax.RoleTyp } // Before removing the worker, mark its jobs as free. - if err := b.freeJobs.MergeJobs(tx, roleType, qdbid, jobs); err != nil { - return nil, errors.Wrap(err, "merging free jobs") - } - - // Remove the worker. - if err := b.current.DeleteWorker(tx, roleType, qdbid, addr); err != nil { - return nil, errors.Wrap(err, "deleting worker") + if err := b.freeJobs.MarkJobsAsFree(tx, roleType, qdbid, jobs); err != nil { + return nil, errors.Wrap(err, "marking jobs as free") } // Even though this may not be useful to the caller (for example, in the @@ -311,8 +278,8 @@ func (b *Balancer) removeDatabaseWorker(tx dax.Transaction, roleType dax.RoleTyp return diff, nil } -func (b *Balancer) FreeWorkers(tx dax.Transaction, addrs ...dax.Address) error { - return errors.Wrap(b.current.FreeWorkers(tx, addrs...), "freeing workers") +func (b *Balancer) ReleaseWorkers(tx dax.Transaction, addrs ...dax.Address) error { + return errors.Wrap(b.current.ReleaseWorkers(tx, addrs...), "freeing workers") } func (b *Balancer) AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) { @@ -396,6 +363,7 @@ func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, qd if err != nil { return nil, errors.Wrapf(err, "getting workers jobs: %s", roleType) } + jset := dax.NewSet[dax.Job]() for _, workerInfo := range workerJobs { jset.Merge(dax.NewSet(workerInfo.Jobs...)) @@ -437,12 +405,12 @@ func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, qd jobCounts[lowWorker]++ } - for worker, jobs := range jobsToCreate { - if err := b.current.CreateJobs(tx, roleType, qdbid, worker, jobs...); err != nil { + for addr, jobs := range jobsToCreate { + if err := b.current.AssignWorkerToJobs(tx, roleType, qdbid, addr, jobs...); err != nil { return nil, errors.Wrap(err, "creating job") } for _, job := range jobs { - diffs.Added(worker, job) + diffs.Added(addr, job) } } @@ -527,7 +495,7 @@ func (b *Balancer) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDataba func (b *Balancer) balanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) (InternalDiffs, error) { diffs := NewInternalDiffs() - for _, role := range []dax.RoleType{dax.RoleTypeCompute, dax.RoleTypeTranslate} { + for _, role := range dax.AllRoleTypes { diff, err := b.balanceDatabaseForRole(tx, role, qdbid) if err != nil { return nil, errors.Wrapf(err, "getting worker count: (%s) %s", role, qdbid) @@ -544,8 +512,7 @@ func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleT // Before balancing, make sure the database has its minimum number of // workers satisfied. - // TODO(tlt): make assignMinWorkers database specific. - if diff, err := b.assignMinWorkers(tx, roleType); err != nil { + if diff, err := b.assignMinWorkers(tx, roleType, qdbid); err != nil { return nil, errors.Wrapf(err, "assigning min workers: (%s) %s", roleType, qdbid) } else { diffs.Merge(diff) @@ -809,11 +776,11 @@ func (b *Balancer) workerForJob(tx dax.Transaction, roleType dax.RoleType, qdbid } func (b *Balancer) ReadNode(tx dax.Transaction, addr dax.Address) (*dax.Node, error) { - return b.nodeService.ReadNode(tx, addr) + return b.workerRegistry.Worker(tx, addr) } func (b *Balancer) Nodes(tx dax.Transaction) ([]*dax.Node, error) { - return b.nodeService.Nodes(tx) + return b.workerRegistry.Workers(tx) } type WorkerJobService interface { @@ -823,10 +790,9 @@ type WorkerJobService interface { ListWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Addresses, error) CreateWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error - DeleteWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error - FreeWorkers(tx dax.Transaction, addrs ...dax.Address) error + ReleaseWorkers(tx dax.Transaction, addrs ...dax.Address) error - CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job ...dax.Job) error + AssignWorkerToJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job ...dax.Job) error DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job dax.Job) error DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) (InternalDiffs, error) JobCounts(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr ...dax.Address) (map[dax.Address]int, error) @@ -840,12 +806,10 @@ type FreeJobService interface { DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job dax.Job) error DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) error ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Jobs, error) - MergeJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs dax.Jobs) error + MarkJobsAsFree(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs dax.Jobs) error } type FreeWorkerService interface { - AddWorkers(tx dax.Transaction, roleType dax.RoleType, addrs ...dax.Address) error - RemoveWorker(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) error PopWorkers(tx dax.Transaction, roleType dax.RoleType, num int) ([]dax.Address, error) ListWorkers(tx dax.Transaction, roleType dax.RoleType) (dax.Addresses, error) } diff --git a/dax/controller/balancer/free_job_test.go b/dax/controller/balancer/free_job_test.go index dd72fb848..31477940c 100644 --- a/dax/controller/balancer/free_job_test.go +++ b/dax/controller/balancer/free_job_test.go @@ -39,6 +39,11 @@ func TestFreeJobService(t *testing.T) { 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) @@ -49,22 +54,22 @@ func TestFreeJobService(t *testing.T) { require.NoError(t, err) require.ElementsMatch(t, dax.Jobs{job1, job3}, jobs) - fwSvc := sqldb.NewFreeWorkerService(nil) - err = fwSvc.AddWorkers(tx, role, nodeAddr) + 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.CreateJobs(tx, role, qdbid, nodeAddr, job1) + 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.MergeJobs(tx, role, qdbid, dax.Jobs{job1}) + err = fjSvc.MarkJobsAsFree(tx, role, qdbid, dax.Jobs{job1}) require.NoError(t, err) jobs, err = fjSvc.ListJobs(tx, role, qdbid) diff --git a/dax/controller/balancer/free_worker_test.go b/dax/controller/balancer/free_worker_test.go index 56dff94d6..7c0f11141 100644 --- a/dax/controller/balancer/free_worker_test.go +++ b/dax/controller/balancer/free_worker_test.go @@ -20,12 +20,25 @@ func TestFreeWorkerService(t *testing.T) { } }() - fwSvc := sqldb.NewFreeWorkerService(nil) - err = fwSvc.AddWorkers(tx, role, nodeAddr, nodeAddr2, nodeAddr3, nodeAddr4, nodeAddr5) - require.NoError(t, 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} - err = fwSvc.RemoveWorker(tx, role, nodeAddr2) - require.NoError(t, err) + 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) diff --git a/dax/controller/balancer/node_test.go b/dax/controller/balancer/node_test.go index 186b31bb3..cacd8d869 100644 --- a/dax/controller/balancer/node_test.go +++ b/dax/controller/balancer/node_test.go @@ -18,7 +18,7 @@ const ( nodeAddr5 = "myaddress5" ) -func TestNodeService(t *testing.T) { +func TestWorkerRegistry(t *testing.T) { tx, err := SQLTransactor.BeginTx(context.Background(), true) require.NoError(t, err, "getting transaction") @@ -29,34 +29,34 @@ func TestNodeService(t *testing.T) { } }() - nodeSvc := sqldb.NewNodeService(nil) + workerReg := sqldb.NewWorkerRegistry(nil) - err = nodeSvc.CreateNode(tx, dax.Address(""), &dax.Node{Address: nodeAddr, RoleTypes: []dax.RoleType{"compute"}}) + err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr, RoleTypes: []dax.RoleType{dax.RoleTypeCompute}}) require.NoError(t, err) - node, err := nodeSvc.ReadNode(tx, nodeAddr) + 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 = nodeSvc.CreateNode(tx, dax.Address(""), &dax.Node{Address: nodeAddr2, RoleTypes: []dax.RoleType{"translate", "compute"}}) + err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr2, RoleTypes: []dax.RoleType{dax.RoleTypeTranslate, dax.RoleTypeCompute}}) require.NoError(t, err, "create node 2") - err = nodeSvc.CreateNode(tx, dax.Address(""), &dax.Node{Address: nodeAddr3, RoleTypes: []dax.RoleType{"compute"}}) + err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr3, RoleTypes: []dax.RoleType{dax.RoleTypeCompute}}) require.NoError(t, err, "create node 3") - nodes, err := nodeSvc.Nodes(tx) + 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 = nodeSvc.DeleteNode(tx, nodeAddr2) + err = workerReg.RemoveWorker(tx, nodeAddr2) require.NoError(t, err, "deleting node") - nodes, err = nodeSvc.Nodes(tx) + nodes, err = workerReg.Workers(tx) require.NoError(t, err) require.EqualValues(t, 2, len(nodes)) for _, node := range nodes { diff --git a/dax/controller/balancer/worker_job_test.go b/dax/controller/balancer/worker_job_test.go index f5ada4ab4..8e62d6854 100644 --- a/dax/controller/balancer/worker_job_test.go +++ b/dax/controller/balancer/worker_job_test.go @@ -40,9 +40,14 @@ func TestWorkerJobService(t *testing.T) { wjSvc := sqldb.NewWorkerJobService(nil) qdbid := dax.QualifiedDatabaseID{OrganizationID: orgID, DatabaseID: dbID} + node := &dax.Node{ + Address: nodeAddr, + RoleTypes: []dax.RoleType{role}, + } + // have to create a free worker before you can create a worker job worker - fwSvc := sqldb.NewFreeWorkerService(nil) - err = fwSvc.AddWorkers(tx, role, nodeAddr) + workerReg := sqldb.NewWorkerRegistry(nil) + err = workerReg.AddWorker(tx, node) require.NoError(t, err) err = wjSvc.CreateWorker(tx, role, qdbid, nodeAddr) @@ -65,7 +70,7 @@ func TestWorkerJobService(t *testing.T) { err = fjSvc.CreateJobs(tx, role, qdbid, job1, job2, job3) require.NoError(t, err) - err = wjSvc.CreateJobs(tx, role, qdbid, nodeAddr, job1, job2) + err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job1, job2) require.NoError(t, err) jobs, err := wjSvc.ListJobs(tx, role, qdbid, nodeAddr) @@ -85,7 +90,7 @@ func TestWorkerJobService(t *testing.T) { require.NoError(t, err) require.ElementsMatch(t, dax.Addresses{nodeAddr}, addrs) - err = wjSvc.CreateJobs(tx, role, qdbid, nodeAddr, job3) + err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job3) require.NoError(t, err) jcs, err := wjSvc.JobCounts(tx, role, qdbid, nodeAddr) @@ -110,7 +115,7 @@ func TestWorkerJobService(t *testing.T) { dk := wjSvc.DatabaseForWorker(tx, nodeAddr) require.EqualValues(t, "db__orgid__blah", dk) - err = wjSvc.DeleteWorker(tx, role, qdbid, nodeAddr) + err = wjSvc.ReleaseWorkers(tx, nodeAddr) require.NoError(t, err) addrs, err = wjSvc.ListWorkers(tx, role, qdbid) diff --git a/dax/controller/controller.go b/dax/controller/controller.go index e5afd553a..6ee2b74f0 100644 --- a/dax/controller/controller.go +++ b/dax/controller/controller.go @@ -25,7 +25,7 @@ const ( // Ensure type implements interface. var _ computer.Registrar = (*Controller)(nil) var _ dax.Schemar = (*Controller)(nil) -var _ dax.NodeService = (*Controller)(nil) +var _ dax.WorkerRegistry = (*Controller)(nil) type Controller struct { // Schemar is used by the controller to get table, and other schema, @@ -87,7 +87,7 @@ func New(cfg Config) *Controller { // Poller. pollerCfg := poller.Config{ AddressManager: c, - NodeService: c, + WorkerRegistry: c, NodePoller: poller.NewHTTPNodePoller(logr), PollInterval: cfg.PollInterval, Logger: logr, @@ -665,7 +665,7 @@ func (c *Controller) DropDatabase(ctx context.Context, qdbid dax.QualifiedDataba for worker := range workerSet { addrs = append(addrs, worker) } - if err := c.Balancer.FreeWorkers(tx, addrs...); err != nil { + if err := c.Balancer.ReleaseWorkers(tx, addrs...); err != nil { return errors.Wrap(err, "freeing workers") } @@ -1713,7 +1713,7 @@ func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID if err != nil { return nil, errors.Wrap(err, "getting nodes for table") } - computeNodes, err := c.assignedToComputeNodes(assignedNodes) + computeNodes, err := assignedToComputeNodes(assignedNodes) if err != nil { return nil, errors.Wrap(err, "converting assigned to compute nodes") } @@ -1725,13 +1725,13 @@ func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID return nil, errors.Wrap(err, "getting compute nodes read or write") } - return c.assignedToComputeNodes(assignedNodes) + return assignedToComputeNodes(assignedNodes) } // assignedToComputeNodes converts the provided []dax.AssignedNode to // []dax.ComputeNode. If any of the assigned nodes are not for RoleType // "compute", an error will be returned. -func (c *Controller) assignedToComputeNodes(nodes []dax.AssignedNode) ([]dax.ComputeNode, error) { +func assignedToComputeNodes(nodes []dax.AssignedNode) ([]dax.ComputeNode, error) { computeNodes := make([]dax.ComputeNode, 0) for _, node := range nodes { @@ -1775,7 +1775,7 @@ func (c *Controller) TranslateNodes(ctx context.Context, qtid dax.QualifiedTable if err != nil { return nil, errors.Wrap(err, "getting nodes for table") } - translateNodes, err := c.assignedToTranslateNodes(assignedNodes) + translateNodes, err := assignedToTranslateNodes(assignedNodes) if err != nil { return nil, errors.Wrap(err, "converting assigned to translate nodes") } @@ -1787,7 +1787,7 @@ func (c *Controller) TranslateNodes(ctx context.Context, qtid dax.QualifiedTable return nil, errors.Wrap(err, "getting translate nodes read or write") } - return c.assignedToTranslateNodes(assignedNodes) + return assignedToTranslateNodes(assignedNodes) } func (c *Controller) nodesForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.AssignedNode, error) { @@ -1808,7 +1808,7 @@ func (c *Controller) nodesForTable(tx dax.Transaction, roleType dax.RoleType, qt // assignedToTranslateNodes converts the provided []dax.AssignedNode to // []dax.TranslateNode. If any of the assigned nodes are not for RoleType // "translate", an error will be returned. -func (c *Controller) assignedToTranslateNodes(nodes []dax.AssignedNode) ([]dax.TranslateNode, error) { +func assignedToTranslateNodes(nodes []dax.AssignedNode) ([]dax.TranslateNode, error) { translateNodes := make([]dax.TranslateNode, 0) for _, node := range nodes { @@ -1872,7 +1872,7 @@ func (c *Controller) IngestPartition(ctx context.Context, qtid dax.QualifiedTabl } } - translateNodes, err := c.assignedToTranslateNodes(nodes) + translateNodes, err := assignedToTranslateNodes(nodes) if err != nil { return "", errors.Wrap(err, "converting assigned to translate nodes") } @@ -1954,7 +1954,7 @@ func (c *Controller) IngestShard(ctx context.Context, qtid dax.QualifiedTableID, retryAsWrite = true } - computeNodes, err := c.assignedToComputeNodes(nodes) + computeNodes, err := assignedToComputeNodes(nodes) if err != nil { return "", errors.Wrap(err, "converting assigned to compute nodes") } @@ -2201,18 +2201,18 @@ func (c *Controller) TableID(ctx context.Context, qdbid dax.QualifiedDatabaseID, return c.Schemar.TableID(tx, qdbid, name) } -// NodeService +// WorkerRegistry -func (c *Controller) CreateNode(context.Context, dax.Address, *dax.Node) error { - return errors.Errorf("Controller.CreateNode() not implemented") +func (c *Controller) AddWorker(context.Context, dax.Address, *dax.Node) error { + return errors.Errorf("Controller.AddWorker() not implemented") } -func (c *Controller) ReadNode(context.Context, dax.Address) (*dax.Node, error) { - return nil, errors.Errorf("Controller.ReadNode() not implemented") +func (c *Controller) Worker(context.Context, dax.Address) (*dax.Node, error) { + return nil, errors.Errorf("Controller.Worker() not implemented") } -func (c *Controller) DeleteNode(context.Context, dax.Address) error { +func (c *Controller) RemoveWorker(context.Context, dax.Address) error { return errors.Errorf("Controller.DeleteNode() not implemented") } -func (c *Controller) Nodes(ctx context.Context) ([]*dax.Node, error) { +func (c *Controller) Workers(ctx context.Context) ([]*dax.Node, error) { tx, err := c.Transactor.BeginTx(ctx, false) if err != nil { return nil, errors.Wrap(err, "beginning tx") diff --git a/dax/controller/node.go b/dax/controller/node.go deleted file mode 100644 index db64fe438..000000000 --- a/dax/controller/node.go +++ /dev/null @@ -1,36 +0,0 @@ -package controller - -import ( - "github.com/featurebasedb/featurebase/v3/dax" -) - -// NodeService represents a service for managing Nodes. -type NodeService interface { - CreateNode(dax.Transaction, dax.Address, *dax.Node) error - ReadNode(dax.Transaction, dax.Address) (*dax.Node, error) - DeleteNode(dax.Transaction, dax.Address) error - Nodes(dax.Transaction) ([]*dax.Node, error) -} - -// Ensure type implements interface. -var _ NodeService = &nopNodeService{} - -// nopNoder is a no-op implementation of the Noder interface. -type nopNodeService struct{} - -func NewNopNodeService() *nopNodeService { - return &nopNodeService{} -} - -func (n *nopNodeService) CreateNode(dax.Transaction, dax.Address, *dax.Node) error { - return nil -} -func (n *nopNodeService) ReadNode(dax.Transaction, dax.Address) (*dax.Node, error) { - return nil, nil -} -func (n *nopNodeService) DeleteNode(dax.Transaction, dax.Address) error { - return nil -} -func (n *nopNodeService) Nodes(dax.Transaction) ([]*dax.Node, error) { - return []*dax.Node{}, nil -} diff --git a/dax/controller/poller/config.go b/dax/controller/poller/config.go index 3afdc4cd9..a837e7fc8 100644 --- a/dax/controller/poller/config.go +++ b/dax/controller/poller/config.go @@ -9,7 +9,7 @@ import ( type Config struct { AddressManager dax.AddressManager - NodeService dax.NodeService + WorkerRegistry dax.WorkerRegistry NodePoller NodePoller PollInterval time.Duration Logger logger.Logger diff --git a/dax/controller/poller/poller.go b/dax/controller/poller/poller.go index 7725d86cb..8cb83fb39 100644 --- a/dax/controller/poller/poller.go +++ b/dax/controller/poller/poller.go @@ -16,7 +16,7 @@ type Poller struct { addressManager dax.AddressManager - nodeService dax.NodeService + workerRegistry dax.WorkerRegistry nodePoller NodePoller pollInterval time.Duration @@ -30,7 +30,7 @@ type Poller struct { func New(cfg Config) *Poller { p := &Poller{ addressManager: dax.NewNopAddressManager(), - nodeService: dax.NewNopNodeService(), + workerRegistry: dax.NewNopWorkerRegistry(), nodePoller: NewNopNodePoller(), pollInterval: time.Second, logger: logger.NopLogger, @@ -40,8 +40,8 @@ func New(cfg Config) *Poller { if cfg.AddressManager != nil { p.addressManager = cfg.AddressManager } - if cfg.NodeService != nil { - p.nodeService = cfg.NodeService + if cfg.WorkerRegistry != nil { + p.workerRegistry = cfg.WorkerRegistry } if cfg.NodePoller != nil { p.nodePoller = cfg.NodePoller @@ -57,7 +57,7 @@ func New(cfg Config) *Poller { } func (p *Poller) Addresses() []dax.Address { - nodes, err := p.nodeService.Nodes(context.Background()) + nodes, err := p.workerRegistry.Workers(context.Background()) if err != nil { p.logger.Errorf("POLLER: unable to get nodes from node service: %v", err) } diff --git a/dax/controller/poller/poller_test.go b/dax/controller/poller/poller_test.go index 6c3688a74..ca910275d 100644 --- a/dax/controller/poller/poller_test.go +++ b/dax/controller/poller/poller_test.go @@ -25,7 +25,7 @@ import ( func TestPoller(t *testing.T) { ctx := context.Background() - nodeService := newMemNodeService() + workerRegistry := newMemWorkerRegistry() // node 1 node1 := newMockNode(t, "health", 0) @@ -44,7 +44,7 @@ func TestPoller(t *testing.T) { } // manager - manager := newMockManager(t, ctx, "deregister-nodes", nodeService) + manager := newMockManager(t, ctx, "deregister-nodes", workerRegistry) defer manager.Close() managerAddr := dax.Address(manager.URL()) @@ -52,7 +52,7 @@ func TestPoller(t *testing.T) { cfg := poller.Config{ AddressManager: controllerhttp.NewAddressManager(managerAddr), NodePoller: poller.NewHTTPNodePoller(logger.NopLogger), - NodeService: nodeService, + WorkerRegistry: workerRegistry, } p := poller.New(cfg) @@ -62,9 +62,9 @@ func TestPoller(t *testing.T) { close(done) }() - // Add nodes to nodeService so they are available to the poller. - nodeService.CreateNode(ctx, addr1, daxNode1) - nodeService.CreateNode(ctx, addr2, daxNode2) + // Add workers to workerRegistry so they are available to the poller. + workerRegistry.AddWorker(ctx, addr1, daxNode1) + workerRegistry.AddWorker(ctx, addr2, daxNode2) go p.Run() defer p.Stop() @@ -83,13 +83,13 @@ type mockManager struct { t *testing.T server *httptest.Server - nodeService dax.NodeService + workerRegistry dax.WorkerRegistry } -func newMockManager(t *testing.T, ctx context.Context, deregisterPath string, nodeService dax.NodeService) *mockManager { +func newMockManager(t *testing.T, ctx context.Context, deregisterPath string, wr dax.WorkerRegistry) *mockManager { mm := &mockManager{ - t: t, - nodeService: nodeService, + t: t, + workerRegistry: wr, } // deregister is a function used in this mock to remove the address from the @@ -97,7 +97,7 @@ func newMockManager(t *testing.T, ctx context.Context, deregisterPath string, no // the Poller. deregister := func(addrs ...dax.Address) { for _, addr := range addrs { - mm.nodeService.DeleteNode(context.Background(), addr) + mm.workerRegistry.RemoveWorker(context.Background(), addr) } } @@ -176,25 +176,25 @@ func (m *mockNode) Close() { } } -type memNodeService struct { +type memWorkerRegistry struct { mu sync.RWMutex addresses map[dax.Address]*dax.Node } -func newMemNodeService() *memNodeService { - return &memNodeService{ +func newMemWorkerRegistry() *memWorkerRegistry { + return &memWorkerRegistry{ addresses: make(map[dax.Address]*dax.Node), } } -func (m *memNodeService) CreateNode(ctx context.Context, addr dax.Address, node *dax.Node) error { +func (m *memWorkerRegistry) AddWorker(ctx context.Context, addr dax.Address, node *dax.Node) error { m.mu.Lock() defer m.mu.Unlock() m.addresses[addr] = node return nil } -func (m *memNodeService) ReadNode(ctx context.Context, addr dax.Address) (*dax.Node, error) { +func (m *memWorkerRegistry) Worker(ctx context.Context, addr dax.Address) (*dax.Node, error) { m.mu.RLock() defer m.mu.RUnlock() node, ok := m.addresses[addr] @@ -204,14 +204,14 @@ func (m *memNodeService) ReadNode(ctx context.Context, addr dax.Address) (*dax.N return node, nil } -func (m *memNodeService) DeleteNode(ctx context.Context, addr dax.Address) error { +func (m *memWorkerRegistry) RemoveWorker(ctx context.Context, addr dax.Address) error { m.mu.Lock() defer m.mu.Unlock() delete(m.addresses, addr) return nil } -func (m *memNodeService) Nodes(ctx context.Context) ([]*dax.Node, error) { +func (m *memWorkerRegistry) Workers(ctx context.Context) ([]*dax.Node, error) { m.mu.RLock() defer m.mu.RUnlock() diff --git a/dax/controller/sqldb/balancer.go b/dax/controller/sqldb/balancer.go index bd5e1f680..9b7816da8 100644 --- a/dax/controller/sqldb/balancer.go +++ b/dax/controller/sqldb/balancer.go @@ -11,7 +11,7 @@ func NewBalancer(log logger.Logger) *balancer.Balancer { fjs := NewFreeJobService(log) wjs := NewWorkerJobService(log) fws := NewFreeWorkerService(log) - ns := NewNodeService(log) + ns := NewWorkerRegistry(log) return balancer.New(ns, fjs, wjs, fws, schemar, log) } diff --git a/dax/controller/sqldb/freejob.go b/dax/controller/sqldb/freejob.go index d578111e1..505703bc0 100644 --- a/dax/controller/sqldb/freejob.go +++ b/dax/controller/sqldb/freejob.go @@ -102,8 +102,9 @@ func (fj *freeJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qd return djs, nil } -// MergeJobs - AFAICT this means "mark these jobs as free" -func (fj *freeJobService) MergeJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs dax.Jobs) error { +// 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") diff --git a/dax/controller/sqldb/freeworker.go b/dax/controller/sqldb/freeworker.go index 3ba26ddf3..4ddc0a2e8 100644 --- a/dax/controller/sqldb/freeworker.go +++ b/dax/controller/sqldb/freeworker.go @@ -1,6 +1,8 @@ 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" @@ -21,34 +23,6 @@ type freeWorkerService struct { log logger.Logger } -func (fw *freeWorkerService) AddWorkers(tx dax.Transaction, roleType dax.RoleType, addrs ...dax.Address) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - workers := make(models.Workers, len(addrs)) - for i, addr := range addrs { - workers[i] = models.Worker{ - Address: addr, - Role: roleType, - } - } - - err := dt.C.Create(workers) - return errors.Wrap(err, "creating workers") -} - -func (fw *freeWorkerService) RemoveWorker(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - err := dt.C.RawQuery("DELETE from workers where database_id is null and role = ? and address = ?", roleType, addr).Exec() - return errors.Wrap(err, "deleting") -} - func (fw *freeWorkerService) PopWorkers(tx dax.Transaction, roleType dax.RoleType, num int) ([]dax.Address, error) { dt, ok := tx.(*DaxTransaction) if !ok { @@ -58,7 +32,8 @@ func (fw *freeWorkerService) PopWorkers(tx dax.Transaction, roleType dax.RoleTyp results := make([]struct { Address dax.Address `db:"address"` }, 0, num) - err := dt.C.RawQuery("select address from workers where role = ? and database_id is NULL limit ?", roleType, num).All(&results) + 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") } @@ -81,9 +56,10 @@ func (fw *freeWorkerService) ListWorkers(tx dax.Transaction, roleType dax.RoleTy } workers := make(models.Workers, 0) - err := dt.C.Select("address").Where("role = ? and database_id is NULL", roleType).Order("address asc").All(&workers) + 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 workers") + return nil, errors.Wrap(err, "querying for free workers") } ret := make(dax.Addresses, len(workers)) diff --git a/dax/controller/sqldb/node.go b/dax/controller/sqldb/node.go deleted file mode 100644 index 9871757a5..000000000 --- a/dax/controller/sqldb/node.go +++ /dev/null @@ -1,116 +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.NodeService = (*nodeService)(nil) - -func NewNodeService(log logger.Logger) *nodeService { - if log == nil { - log = logger.NopLogger - } - return &nodeService{ - log: log, - } -} - -type nodeService struct { - log logger.Logger -} - -func (n *nodeService) CreateNode(tx dax.Transaction, addr dax.Address, node *dax.Node) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - mnode := &models.Node{Address: node.Address} - err := dt.C.Create(mnode) - if err != nil { - return errors.Wrap(err, "creating node") - } - - nodeRoles := make(models.NodeRoles, len(node.RoleTypes)) - mnode.NodeRoles = nodeRoles - for i, rt := range node.RoleTypes { - mnode.NodeRoles[i] = models.NodeRole{ - NodeID: mnode.ID, - Role: rt, - } - } - err = dt.C.Create(&(mnode.NodeRoles)) - if err != nil { - return errors.Wrap(err, "creating node roles") - } - - return nil -} - -func (n *nodeService) ReadNode(tx dax.Transaction, addr dax.Address) (*dax.Node, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - node := &models.Node{} - err := dt.C.Eager().Where("address = ?", addr).First(node) - if err != nil { - return nil, errors.Wrap(err, "getting node") - } - - roleTypes := make([]dax.RoleType, len(node.NodeRoles)) - for i, nr := range node.NodeRoles { - roleTypes[i] = nr.Role - } - - return &dax.Node{ - Address: node.Address, - RoleTypes: roleTypes, - }, nil -} - -func (n *nodeService) DeleteNode(tx dax.Transaction, addr dax.Address) error { - dt, ok := tx.(*DaxTransaction) - if !ok { - return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - node := &models.Node{} - err := dt.C.Eager().Where("address = ?", addr).First(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") -} - -func (n *nodeService) Nodes(tx dax.Transaction) ([]*dax.Node, error) { - dt, ok := tx.(*DaxTransaction) - if !ok { - return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") - } - - nodes := []*models.Node{} - dt.C.Eager().Order("address asc").All(&nodes) - - ret := make([]*dax.Node, len(nodes)) - for i, node := range nodes { - ret[i] = &dax.Node{ - Address: node.Address, - RoleTypes: make([]dax.RoleType, len(node.NodeRoles)), - } - for j, nr := range node.NodeRoles { - ret[i].RoleTypes[j] = nr.Role - } - } - - return ret, nil -} diff --git a/dax/controller/sqldb/worker.go b/dax/controller/sqldb/worker.go new file mode 100644 index 000000000..99bd13bde --- /dev/null +++ b/dax/controller/sqldb/worker.go @@ -0,0 +1,139 @@ +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") + } + + workers := models.Workers{} + + // Determine if a worker for this address already exists. We use `All()` + // here instead of `First()` because `First()` returns an error if there's + // no match. + if err := dt.C.Where("address = ?", node.Address).All(&workers); err != nil { + return errors.Wrapf(err, "getting workers by address: %s", node.Address) + } + + switch len(workers) { + case 0: + // Continue on to create. + case 1: + // Since a worker for this address already exists, just update it and + // return. + worker := workers[0] + for _, roleType := range node.RoleTypes { + if err := worker.SetRole(roleType); err != nil { + return errors.Wrapf(err, "setting role: %s", roleType) + } + } + return dt.C.Update(worker) + default: + return errors.Errorf("found more than one worker for address: %s", node.Address) + } + + worker := &models.Worker{ + Address: node.Address, + } + 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 index 5ca6466f7..5e4eb15e4 100644 --- a/dax/controller/sqldb/workerjob.go +++ b/dax/controller/sqldb/workerjob.go @@ -24,37 +24,59 @@ 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{} - err := dt.C.Eager().Where("role = ? and database_id = ?", roleType, qdbid.DatabaseID).Order("address asc").All(&workers) + sql := fmt.Sprintf("role_%s = true and database_id = ?", roleType) + err := dt.C.Where(sql, qdbid.DatabaseID).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 - ret[i].Jobs = make([]dax.Job, len(worker.Jobs)) - for j, job := range worker.Jobs { - ret[i].Jobs[j] = job.Name + 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") } worker := &models.Worker{} - cnt, err := dt.C.Where("role = ? and database_id = ?", roleType, qdbid.DatabaseID).Count(worker) + sql := fmt.Sprintf("role_%s = true and database_id = ?", roleType) + cnt, err := dt.C.Where(sql, qdbid.DatabaseID).Count(worker) return cnt, errors.Wrap(err, "getting count") } @@ -65,7 +87,8 @@ func (w *workerJobService) ListWorkers(tx dax.Transaction, roleType dax.RoleType } workers := models.Workers{} - err := dt.C.Select("address").Where("role = ? and database_id = ?", roleType, qdbid.DatabaseID).Order("address asc").All(&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") } @@ -85,12 +108,13 @@ func (w *workerJobService) CreateWorker(tx dax.Transaction, roleType dax.RoleTyp } worker := &models.Worker{} - err := dt.C.RawQuery("UPDATE workers SET database_id = ? WHERE role = ? and address = ? RETURNING workers.ID", qdbid.DatabaseID, roleType, addr).First(worker) + sql := fmt.Sprintf("UPDATE workers SET database_id = ? WHERE role_%s = true and address = ? RETURNING workers.ID", roleType) + err := dt.C.RawQuery(sql, qdbid.DatabaseID, addr).First(worker) return errors.Wrap(err, "associating worker to database") } -func (w *workerJobService) FreeWorkers(tx dax.Transaction, addrs ...dax.Address) error { +func (w *workerJobService) ReleaseWorkers(tx dax.Transaction, addrs ...dax.Address) error { dt, ok := tx.(*DaxTransaction) if !ok { return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction") @@ -104,34 +128,17 @@ func (w *workerJobService) FreeWorkers(tx dax.Transaction, addrs ...dax.Address) return errors.Wrap(err, "updating workers") } -func (w *workerJobService) DeleteWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error { +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{} - err := dt.C.Where("address = ? and role = ? and database_id = ?", addr, roleType, qdbid.DatabaseID).First(worker) - 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") -} - -func (w *workerJobService) CreateJobs(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{} - err := dt.C.Where("address = ? and role = ?", addr, roleType).First(worker) + sql := fmt.Sprintf("address = ? and role_%s = true", roleType) + err := dt.C.Where(sql, addr).First(worker) if err != nil { - return errors.Wrap(err, "getting worker") + return errors.Wrapf(err, "getting worker: (%s) %s", roleType, addr) } jobs := models.Jobs{} @@ -140,35 +147,35 @@ func (w *workerJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, return errors.Wrap(err, "updating jobs") } - // create jobs not in "jobs" - toCreate := jobsNotUpdated(job, jobs, worker) + // Assign jobs not in "jobs", and therefore didn't get updated by the + // previous sql statement. + toBeAssigned := jobsNotAssigned(job, jobs, roleType, worker) - err = dt.C.Create(toCreate) - if err != nil { + if err := dt.C.Create(toBeAssigned); err != nil { return errors.Wrap(err, "creating jobs") } - return errors.Wrap(err, "creating jobs") + return nil } -func jobsNotUpdated(incomingJobs []dax.Job, created models.Jobs, worker *models.Worker) (toCreate models.Jobs) { +func jobsNotAssigned(incomingJobs []dax.Job, assigned models.Jobs, roleType dax.RoleType, worker *models.Worker) (toBeAssigned models.Jobs) { outer: for _, incJob := range incomingJobs { - for _, createdJob := range created { - if createdJob.Name == incJob { + for _, assignedJob := range assigned { + if assignedJob.Name == incJob { continue outer } } - toCreate = append(toCreate, + toBeAssigned = append(toBeAssigned, models.Job{ Name: incJob, - Role: worker.Role, + Role: roleType, DatabaseID: dax.DatabaseID(worker.DatabaseID.String), Worker: worker, }, ) } - return toCreate + return toBeAssigned } func (w *workerJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job dax.Job) error { @@ -178,7 +185,8 @@ func (w *workerJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, } worker := &models.Worker{} - err := dt.C.Select("id").Where("role = ? and database_id = ? and address = ?", roleType, qdbid.DatabaseID, addr).First(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") } @@ -227,7 +235,7 @@ func (w *workerJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.R return idiffs, errors.Wrap(err, "deleting jobs") } -func (w *workerJobService) JobCounts(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr ...dax.Address) (map[dax.Address]int, error) { +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") @@ -238,18 +246,23 @@ func (w *workerJobService) JobCounts(tx dax.Transaction, roleType dax.RoleType, Count int `db:"count"` }{} var err error - if len(addr) == 0 { + 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 = ? + where w.database_id = ? and w.role_%s = true + and j.role = ? group by w.address` - err = dt.C.RawQuery(qstring, qdbid.DatabaseID, roleType).All(&results) + 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.address in (?) and w.database_id = ? and w.role = ? + where w.database_id = ? and w.role_%s = true + and j.role = ? + and w.address in (?) group by w.address` - err = dt.C.RawQuery(qstring, addr, qdbid.DatabaseID, roleType).All(&results) + 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") @@ -269,20 +282,15 @@ func (w *workerJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, q } worker := &models.Worker{} - err := dt.C.Eager().Where("role = ? and database_id = ? and address = ?", roleType, qdbid.DatabaseID, addr).First(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") } - ret := make(dax.Jobs, len(worker.Jobs)) - // jobs are ordered by "name asc" defined on the worker model. - for i, job := range worker.Jobs { - ret[i] = job.Name - } - - return ret, nil + return jobsForWorker(dt, worker, roleType) } func (w *workerJobService) DatabaseForWorker(tx dax.Transaction, addr dax.Address) dax.DatabaseKey { diff --git a/dax/controller/sqldb/workerjob_test.go b/dax/controller/sqldb/workerjob_test.go index dd270ccb0..f583d7d4c 100644 --- a/dax/controller/sqldb/workerjob_test.go +++ b/dax/controller/sqldb/workerjob_test.go @@ -21,19 +21,19 @@ func TestJobsNotUpdated(t *testing.T) { }, } - toCreate := jobsNotUpdated(incJobs, created, &models.Worker{ - ID: u2, - Role: "compute", - DatabaseID: nulls.NewString("dbid"), + toCreate := jobsNotAssigned(incJobs, created, dax.RoleTypeCompute, &models.Worker{ + ID: u2, + RoleCompute: true, + DatabaseID: nulls.NewString("dbid"), }) require.Equal(t, 3, len(toCreate)) // Test when 0 jobs are updated - toCreate = jobsNotUpdated(incJobs, models.Jobs{}, &models.Worker{ - ID: u2, - Role: "compute", - DatabaseID: nulls.NewString("dbid"), + toCreate = jobsNotAssigned(incJobs, models.Jobs{}, dax.RoleTypeCompute, &models.Worker{ + ID: u2, + RoleCompute: true, + DatabaseID: nulls.NewString("dbid"), }) require.Equal(t, 4, len(toCreate)) diff --git a/dax/controller/worker.go b/dax/controller/worker.go new file mode 100644 index 000000000..a2d2ac8c0 --- /dev/null +++ b/dax/controller/worker.go @@ -0,0 +1,40 @@ +package controller + +import ( + "github.com/featurebasedb/featurebase/v3/dax" +) + +// WorkerRegistry represents a service for managing Nodes. Note that this +// interface mirrors the dax.WorkerRegistry interface, but its methods take +// dax.Transactions rather than Contexts. That's because the dax version of this +// interface is meant to be a the API boundary, where this is an interface for +// use within the Controller. +type WorkerRegistry interface { + AddWorker(dax.Transaction, *dax.Node) error + Worker(dax.Transaction, dax.Address) (*dax.Node, error) + RemoveWorker(dax.Transaction, dax.Address) error + Workers(dax.Transaction) ([]*dax.Node, error) +} + +// Ensure type implements interface. +var _ WorkerRegistry = &nopWorkerRegistry{} + +// nopWorkerRegistry is a no-op implementation of the WorkerRegistry interface. +type nopWorkerRegistry struct{} + +func NewNopWorkerRegistry() *nopWorkerRegistry { + return &nopWorkerRegistry{} +} + +func (n *nopWorkerRegistry) AddWorker(dax.Transaction, *dax.Node) error { + return nil +} +func (n *nopWorkerRegistry) Worker(dax.Transaction, dax.Address) (*dax.Node, error) { + return nil, nil +} +func (n *nopWorkerRegistry) RemoveWorker(dax.Transaction, dax.Address) error { + return nil +} +func (n *nopWorkerRegistry) Workers(dax.Transaction) ([]*dax.Node, error) { + return []*dax.Node{}, nil +} diff --git a/dax/migrations/003_node_to_worker.down.fizz b/dax/migrations/003_node_to_worker.down.fizz new file mode 100644 index 000000000..e69de29bb diff --git a/dax/migrations/003_node_to_worker.up.fizz b/dax/migrations/003_node_to_worker.up.fizz new file mode 100644 index 000000000..ee3d636c9 --- /dev/null +++ b/dax/migrations/003_node_to_worker.up.fizz @@ -0,0 +1,16 @@ +add_column("workers", "role_compute", "bool", {"default": false}) +add_column("workers", "role_translate", "bool", {"default": false}) +add_column("workers", "role_query", "bool", {"default": false}) + +sql("update jobs set worker_id = wc.id from workers wc inner join workers wt on wc.address = wt.address and wc.database_id = wt.database_id and wc.role = 'compute' and wt.role = 'translate' where worker_id = wt.id;") + +sql("update workers set role_compute = true where role = 'compute';") + +sql("update workers set role_translate = true from workers wt where workers.address = wt.address and workers.role = 'compute' and wt.role = 'translate';") + +sql("delete from workers where role = 'translate';") + +drop_column("workers", "role") + +drop_table("node_roles") +drop_table("nodes") \ No newline at end of file diff --git a/dax/models/node.go b/dax/models/node.go deleted file mode 100644 index c347e75ba..000000000 --- a/dax/models/node.go +++ /dev/null @@ -1,56 +0,0 @@ -package models - -import ( - "encoding/json" - "time" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/gobuffalo/pop/v6" - "github.com/gobuffalo/validate/v3" - "github.com/gobuffalo/validate/v3/validators" - "github.com/gofrs/uuid" -) - -// Node represents a host or server that is available to work on jobs. -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" order_by:"created_at asc"` - CreatedAt time.Time `json:"created_at" db:"created_at"` - UpdatedAt time.Time `json:"updated_at" db:"updated_at"` -} - -// String is not required by pop and may be deleted -func (t *Node) String() string { - jt, _ := json.MarshalIndent(t, " ", " ") //nolint:errchkjson - return string(jt) -} - -// Nodes is not required by pop and may be deleted -type Nodes []*Node - -// String is not required by pop and may be deleted -func (t Nodes) String() string { - jt, _ := json.MarshalIndent(t, " ", " ") //nolint:errchkjson - return string(jt) -} - -// Validate gets run every time you call a "pop.Validate*" (pop.ValidateAndSave, pop.ValidateAndCreate, pop.ValidateAndUpdate) method. -// This method is not required and may be deleted. -func (t *Node) Validate(tx *pop.Connection) (*validate.Errors, error) { - return validate.Validate( - &validators.StringIsPresent{Field: string(t.Address), Name: "Address"}, - ), nil -} - -// ValidateCreate gets run every time you call "pop.ValidateAndCreate" method. -// This method is not required and may be deleted. -func (t *Node) ValidateCreate(tx *pop.Connection) (*validate.Errors, error) { - return validate.NewErrors(), nil -} - -// ValidateUpdate gets run every time you call "pop.ValidateAndUpdate" method. -// This method is not required and may be deleted. -func (t *Node) ValidateUpdate(tx *pop.Connection) (*validate.Errors, error) { - return validate.NewErrors(), nil -} diff --git a/dax/models/node_roles.go b/dax/models/node_roles.go deleted file mode 100644 index 558e2940e..000000000 --- a/dax/models/node_roles.go +++ /dev/null @@ -1,31 +0,0 @@ -package models - -import ( - "encoding/json" - "time" - - "github.com/featurebasedb/featurebase/v3/dax" - "github.com/gofrs/uuid" -) - -// NodeRole holds information about what types of jobs (roles) each node can perform. -type NodeRole struct { - ID uuid.UUID `json:"id" db:"id"` - NodeID uuid.UUID `json:"node_id" db:"node_id"` - Role dax.RoleType `json:"role" db:"role"` - CreatedAt time.Time `json:"created_at" db:"created_at"` - UpdatedAt time.Time `json:"updated_at" db:"updated_at"` -} - -// String is not required by pop and may be deleted -func (t *NodeRole) String() string { - jt, _ := json.MarshalIndent(t, " ", " ") //nolint:errchkjson - return string(jt) -} - -type NodeRoles []NodeRole - -func (t NodeRoles) String() string { - jt, _ := json.MarshalIndent(t, " ", " ") //nolint:errchkjson - return string(jt) -} diff --git a/dax/models/worker.go b/dax/models/worker.go index c2ddd25f7..dd4b2f249 100644 --- a/dax/models/worker.go +++ b/dax/models/worker.go @@ -5,6 +5,7 @@ import ( "time" "github.com/featurebasedb/featurebase/v3/dax" + "github.com/featurebasedb/featurebase/v3/errors" "github.com/gobuffalo/nulls" "github.com/gobuffalo/pop/v6" "github.com/gobuffalo/validate/v3" @@ -15,13 +16,15 @@ import ( // Worker is a node plus a role that gets assigned to a database and // can be assigned jobs for that database. type Worker struct { - ID uuid.UUID `json:"id" db:"id"` - Address dax.Address `json:"address" db:"address"` - Role dax.RoleType `json:"role" db:"role"` - DatabaseID nulls.String `json:"database_id" db:"database_id"` // this can be empty which means the worker is unassigned - CreatedAt time.Time `json:"created_at" db:"created_at"` - UpdatedAt time.Time `json:"updated_at" db:"updated_at"` - Jobs Jobs `json:"jobs" has_many:"jobs" order_by:"name asc"` + ID uuid.UUID `json:"id" db:"id"` + Address dax.Address `json:"address" db:"address"` + DatabaseID nulls.String `json:"database_id" db:"database_id"` // this can be empty which means the worker is unassigned + CreatedAt time.Time `json:"created_at" db:"created_at"` + UpdatedAt time.Time `json:"updated_at" db:"updated_at"` + Jobs Jobs `json:"jobs" has_many:"jobs" order_by:"name asc"` + RoleCompute bool `json:"role_compute" db:"role_compute"` + RoleTranslate bool `json:"role_translate" db:"role_translate"` + RoleQuery bool `json:"role_query" db:"role_query"` } // String is not required by pop and may be deleted @@ -30,6 +33,23 @@ func (t *Worker) String() string { return string(jt) } +// SetRole applies a dax.RoleType to one of the boolean fields on the Worker +// model. It returns an error if the model does not support that role type. +func (t *Worker) SetRole(role dax.RoleType) error { + switch role { + case dax.RoleTypeCompute: + t.RoleCompute = true + case dax.RoleTypeTranslate: + t.RoleTranslate = true + case dax.RoleTypeQuery: + t.RoleQuery = true + default: + errors.Errorf("invalid role type for worker: %s", role) + } + + return nil +} + // Workers is not required by pop and may be deleted type Workers []Worker diff --git a/dax/role.go b/dax/role.go index 165356611..e38dce7f9 100644 --- a/dax/role.go +++ b/dax/role.go @@ -6,6 +6,11 @@ type RoleType string const ( RoleTypeCompute RoleType = "compute" RoleTypeTranslate RoleType = "translate" + RoleTypeQuery RoleType = "query" +) + +var ( + AllRoleTypes = []RoleType{RoleTypeCompute, RoleTypeTranslate, RoleTypeQuery} ) // RoleTypes is a list of RoleType, used primarily to introduce helper methods diff --git a/dax/test/dax/dax_test.go b/dax/test/dax/dax_test.go index 637345bc8..286020dc1 100644 --- a/dax/test/dax/dax_test.go +++ b/dax/test/dax/dax_test.go @@ -275,7 +275,7 @@ func TestDAXIntegration(t *testing.T) { // TODO: implement this without a sleep. time.Sleep(5 * time.Second) - // ensure paritions are still covered + // ensure partitions are still covered nodes, err = controllerClient.TranslateNodes(context.Background(), qtid, append(partitions0, partitions1...)...) assert.NoError(t, err) if assert.Len(t, nodes, 1) { @@ -465,7 +465,7 @@ func TestDAXIntegration(t *testing.T) { assert.NoError(t, svcmgr.ControllerStart()) assert.True(t, mc.Healthy(controllerKey)) - // ensure paritions are still covered + // ensure partitions are still covered nodes, err = controllerClient.TranslateNodes(context.Background(), qtid, partitions...) assert.NoError(t, err) if assert.Len(t, nodes, 1) { diff --git a/dax/node.go b/dax/worker.go similarity index 85% rename from dax/node.go rename to dax/worker.go index 62060af7b..e4ad397f4 100644 --- a/dax/node.go +++ b/dax/worker.go @@ -53,34 +53,34 @@ type AssignedNode struct { Role Role `json:"role"` } -// NodeService represents a service for managing Nodes. -type NodeService interface { - CreateNode(context.Context, Address, *Node) error - ReadNode(context.Context, Address) (*Node, error) - DeleteNode(context.Context, Address) error - Nodes(context.Context) ([]*Node, error) +// WorkerRegistry represents a service for managing Workers. +type WorkerRegistry interface { + AddWorker(context.Context, Address, *Node) error + Worker(context.Context, Address) (*Node, error) + RemoveWorker(context.Context, Address) error + Workers(context.Context) ([]*Node, error) } // Ensure type implements interface. -var _ NodeService = &nopNodeService{} +var _ WorkerRegistry = &nopWorkerRegistry{} -// nopNoder is a no-op implementation of the Noder interface. -type nopNodeService struct{} +// nopWorkerRegistry is a no-op implementation of the WorkerRegistry interface. +type nopWorkerRegistry struct{} -func NewNopNodeService() *nopNodeService { - return &nopNodeService{} +func NewNopWorkerRegistry() *nopWorkerRegistry { + return &nopWorkerRegistry{} } -func (n *nopNodeService) CreateNode(context.Context, Address, *Node) error { +func (n *nopWorkerRegistry) AddWorker(context.Context, Address, *Node) error { return nil } -func (n *nopNodeService) ReadNode(context.Context, Address) (*Node, error) { +func (n *nopWorkerRegistry) Worker(context.Context, Address) (*Node, error) { return nil, nil } -func (n *nopNodeService) DeleteNode(context.Context, Address) error { +func (n *nopWorkerRegistry) RemoveWorker(context.Context, Address) error { return nil } -func (n *nopNodeService) Nodes(context.Context) ([]*Node, error) { +func (n *nopWorkerRegistry) Workers(context.Context) ([]*Node, error) { return []*Node{}, nil }