Remove Node from data model; standardize on Worker (#2366)

* Remove Node from data model; standardize on Worker

This commit does a lot of things, but in general it attempts to simplify
the data model by getting rid of the Node and NodeRole models. Instead,
these will use the Worker model, which itself has individual boolean
fields for role types.

Get rid of roleType in some FreeWorker methods

rename NodeService to WorkerRegistry

simplify the freeworker interface

fix the tests

* Remove DeleteWorker method from workerJobService
This commit is contained in:
Travis Turner 2023-04-04 20:20:53 -05:00 committed by GitHub
parent 284f62dcb9
commit ea72396b4d
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
27 changed files with 465 additions and 512 deletions

View file

@ -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) {

View file

@ -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)
}

View file

@ -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)

View file

@ -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)

View file

@ -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 {

View file

@ -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)

View file

@ -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")

View file

@ -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
}

View file

@ -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

View file

@ -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)
}

View file

@ -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()

View file

@ -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)
}

View file

@ -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")

View file

@ -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))

View file

@ -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
}

View file

@ -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
}

View file

@ -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 {

View file

@ -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))

40
dax/controller/worker.go Normal file
View file

@ -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
}

View file

@ -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")

View file

@ -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
}

View file

@ -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)
}

View file

@ -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

View file

@ -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

View file

@ -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) {

View file

@ -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
}