mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 02:44:59 +00:00
Friday's big jam
This commit is contained in:
parent
5baee81eb5
commit
6f5b7e332e
24 changed files with 1020 additions and 667 deletions
|
|
@ -13,17 +13,12 @@ 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)
|
||||
|
||||
// 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)
|
||||
|
||||
// RemoveJobs removes jobs for the given database.
|
||||
RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error)
|
||||
|
||||
// BalanceDatabase forces a database balance. TODO(tlt): currently this is
|
||||
// only used in tests, so perhaps we can get rid of it.
|
||||
BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error)
|
||||
|
||||
// CurrentState returns the workers and jobs currently active for the given
|
||||
|
|
@ -90,9 +85,6 @@ 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) 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) {
|
||||
return []dax.WorkerDiff{}, nil
|
||||
}
|
||||
|
|
@ -136,4 +128,6 @@ func (b *NopBalancer) WorkerServiceProviders(tx dax.Transaction /*, future optio
|
|||
return nil, nil
|
||||
}
|
||||
|
||||
func (b *NopBalancer) AssignFreeServiceToDatabase(tx dax.Transaction, wspID dax.WorkerServiceProviderID, qdb *dax.QualifiedDatabase) (*dax.WorkerService, error)
|
||||
func (b *NopBalancer) AssignFreeServiceToDatabase(tx dax.Transaction, wspID dax.WorkerServiceProviderID, qdb *dax.QualifiedDatabase) (*dax.WorkerService, error) {
|
||||
return nil, nil
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,8 +4,6 @@ package balancer
|
|||
import (
|
||||
"log"
|
||||
"math"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
|
|
@ -18,6 +16,8 @@ import (
|
|||
// Ensure type implements interface.
|
||||
var _ controller.Balancer = (*Balancer)(nil)
|
||||
|
||||
// TODO: continue figuring out which things we removed need to come back and implement them, based on what's not implemented on the balancer interface
|
||||
|
||||
// Balancer is an implementation of the balancer.Balancer interface which
|
||||
// isolates workers and jobs by database. It helps manage the relationships
|
||||
// between workers and jobs. The logic it uses to balance jobs across workers is
|
||||
|
|
@ -25,9 +25,6 @@ var _ controller.Balancer = (*Balancer)(nil)
|
|||
// jobs. It does not take anything else (such as job size, worker capabilities,
|
||||
// etc) into consideration.
|
||||
type Balancer struct {
|
||||
// current represents the current state of worker/job assigments.
|
||||
current WorkerJobService
|
||||
|
||||
workerRegistry controller.WorkerRegistry
|
||||
|
||||
// freeJobs is the set of jobs which have yet to be assigned to a worker.
|
||||
|
|
@ -40,21 +37,21 @@ type Balancer struct {
|
|||
|
||||
schemar schemar.Schemar
|
||||
|
||||
wsp WorkerServiceProviderService
|
||||
store controller.Store
|
||||
|
||||
logger logger.Logger
|
||||
}
|
||||
|
||||
// New returns a new instance of Balancer.
|
||||
func New(wr controller.WorkerRegistry, fjs FreeJobService, wjs WorkerJobService, fws FreeWorkerService, schemar schemar.Schemar, wsp WorkerServiceProviderService, logger logger.Logger) *Balancer {
|
||||
func New(wr controller.WorkerRegistry, fjs FreeJobService, wjs WorkerJobService, fws FreeWorkerService, schemar schemar.Schemar, store controller.Store, logger logger.Logger) *Balancer {
|
||||
return &Balancer{
|
||||
current: wjs,
|
||||
workerRegistry: wr,
|
||||
freeJobs: fjs,
|
||||
freeWorkers: fws,
|
||||
schemar: schemar,
|
||||
wsp: wsp,
|
||||
logger: logger,
|
||||
store: store,
|
||||
|
||||
logger: logger,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -66,189 +63,38 @@ func New(wr controller.WorkerRegistry, fjs FreeJobService, wjs WorkerJobService,
|
|||
func (b *Balancer) AddWorker(tx dax.Transaction, node *dax.Node) ([]dax.WorkerDiff, error) {
|
||||
b.logger.Debugf("AddWorker(%s)", node.Address)
|
||||
|
||||
if err := b.workerRegistry.AddWorker(tx, node); err != nil {
|
||||
qdbidp, err := b.store.AddWorker(tx, node)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "creating node on node service: %s", node.Address)
|
||||
}
|
||||
|
||||
diffs := NewInternalDiffs()
|
||||
|
||||
// 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 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", node.Address)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
diffs, err := b.balanceDatabase(tx, *qdbidp, node.ServiceID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "balancing db")
|
||||
}
|
||||
|
||||
return diffs.Output(), 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)
|
||||
|
||||
// Find out how many free workers we have.
|
||||
freeWorkers, err := b.freeWorkers.ListWorkers(tx, qdbid, roleType)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting free worker list")
|
||||
}
|
||||
freeWorkerCount := len(freeWorkers)
|
||||
|
||||
// If there are no free workers, return early.
|
||||
if freeWorkerCount == 0 {
|
||||
b.logger.Debugf("No free workers for '%s'", roleType)
|
||||
return InternalDiffs{}, nil
|
||||
}
|
||||
|
||||
// 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 database")
|
||||
}
|
||||
|
||||
// Create a map[database]int where int is the number of workers required to
|
||||
// reach that database's minWorkerCount. This map will only contain database
|
||||
// which need more workers in order to reach their minimum.
|
||||
m := make(map[dax.QualifiedDatabaseID]int)
|
||||
|
||||
for _, qdb := range qdbs {
|
||||
qdbid := qdb.QualifiedID()
|
||||
|
||||
minWorkers := qdb.Options.WorkersMin
|
||||
if minWorkers == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
// If the database doesn't have any jobs, there's no need to go any
|
||||
// further. In other words, we don't want to assign a worker to a
|
||||
// database until it has at least one job.
|
||||
if hasJobs, err := b.databaseHasJobs(tx, roleType, qdbid); err != nil {
|
||||
return nil, errors.Wrapf(err, "checking has jobs: (%s) %s", roleType, qdbid)
|
||||
} else if !hasJobs {
|
||||
continue
|
||||
}
|
||||
|
||||
// Get the number of workers assigned to this database.
|
||||
workerCount, err := b.current.WorkerCount(tx, roleType, qdbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", roleType, qdbid)
|
||||
}
|
||||
|
||||
diff := minWorkers - workerCount
|
||||
|
||||
// If we have more workers than the min required, or if we have the
|
||||
// exact number of workers , don't do anything for that database.
|
||||
if diff <= 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
m[qdbid] = diff
|
||||
}
|
||||
|
||||
diffs := NewInternalDiffs()
|
||||
|
||||
// Create an ordered slice of map keys so that tests are predictable.
|
||||
qdbids := make([]dax.QualifiedDatabaseID, 0, len(m))
|
||||
for qdbid := range m {
|
||||
qdbids = append(qdbids, qdbid)
|
||||
}
|
||||
sort.Sort(dax.QualifiedDatabaseIDs(qdbids))
|
||||
|
||||
// For each database, if there are enough free workers to
|
||||
// satisfy its min, then pop that number of workers from the free list. If
|
||||
// not, continue to the next database until either reaching the end of the
|
||||
// database list or until there are no more free workers in the list,
|
||||
// whichever comes first.
|
||||
for _, qdbid := range qdbids {
|
||||
need := m[qdbid]
|
||||
|
||||
if freeWorkerCount == 0 {
|
||||
break
|
||||
}
|
||||
|
||||
if freeWorkerCount >= need {
|
||||
addrs, err := b.freeWorkers.PopWorkers(tx, roleType, need)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "popping free worker: (%s)", roleType)
|
||||
}
|
||||
|
||||
if diff, err := b.addDatabaseWorkers(tx, roleType, qdbid, addrs...); err != nil {
|
||||
return nil, errors.Wrapf(err, "adding database workers: (%s) %s, %v", roleType, qdbid, addrs)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return diffs, nil
|
||||
}
|
||||
|
||||
// addDatabaseWorkers adds workers from the free worker list to the pool of
|
||||
// workers for a specific database.
|
||||
func (b *Balancer) addDatabaseWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addrs ...dax.Address) (InternalDiffs, error) {
|
||||
for _, addr := range addrs {
|
||||
if err := b.current.CreateWorker(tx, roleType, qdbid, addr); err != nil {
|
||||
return nil, errors.Wrap(err, "creating worker")
|
||||
}
|
||||
}
|
||||
|
||||
// Process the freeJobs.
|
||||
return b.processFreeJobs(tx, roleType, qdbid)
|
||||
}
|
||||
|
||||
// databaseHasJobs returns true if the database has at least one job.
|
||||
func (b *Balancer) databaseHasJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (bool, error) {
|
||||
// Free jobs.
|
||||
if freeJobs, err := b.freeJobs.ListJobs(tx, roleType, qdbid); err != nil {
|
||||
return false, errors.Wrapf(err, "getting free jobs: (%s) %s", roleType, qdbid)
|
||||
} else if len(freeJobs) > 0 {
|
||||
return true, nil
|
||||
}
|
||||
|
||||
// Assigned jobs.
|
||||
if wis, err := b.current.WorkersJobs(tx, roleType, qdbid); err != nil {
|
||||
return false, errors.Wrapf(err, "getting free jobs: (%s) %s", roleType, qdbid)
|
||||
} else {
|
||||
for _, wi := range wis {
|
||||
if len(wi.Jobs) > 0 {
|
||||
return true, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// RemoveWorker removes a worker from the system. If the worker is currently
|
||||
// assigned to a database and has jobs, it will be removed and its jobs will
|
||||
// be either transferred to other workers or placed on the free job list.
|
||||
func (b *Balancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.WorkerDiff, error) {
|
||||
diffs := NewInternalDiffs()
|
||||
diffs := controller.NewInternalDiffs()
|
||||
|
||||
// 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.AllRoleTypes {
|
||||
if diff, err := b.removeDatabaseWorker(tx, rt, qdbid, addr); err != nil {
|
||||
return nil, errors.Wrapf(err, "removing worker: (%s) %s", rt, addr)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
}
|
||||
worker, err := b.store.WorkerForAddress(tx, addr)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting worker for address")
|
||||
}
|
||||
|
||||
// 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)
|
||||
// drop from workers table. This *should* also remove all associations of jobs to this worker.
|
||||
if err := b.store.RemoveWorker(tx, worker.ID); err != nil {
|
||||
return nil, errors.Wrap(err, "removing worker")
|
||||
}
|
||||
|
||||
if dbkey != "" {
|
||||
qdbid := dbkey.QualifiedDatabaseID()
|
||||
if worker.DatabaseID != nil {
|
||||
// Balance the affected database.
|
||||
if diff, err := b.balanceDatabase(tx, qdbid); err != nil {
|
||||
return nil, errors.Wrapf(err, "balancing database: %s", qdbid)
|
||||
if diff, err := b.balanceDatabase(tx, *worker.DatabaseID, worker.ServiceID); err != nil {
|
||||
return nil, errors.Wrapf(err, "balancing database: %s", worker.DatabaseID)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
|
|
@ -257,34 +103,6 @@ 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 {
|
||||
return nil, errors.Wrap(err, "listing jobs")
|
||||
}
|
||||
|
||||
// Before removing the worker, mark its jobs as free.
|
||||
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
|
||||
// case where the worker has died and no longer exists), return the diffs
|
||||
// which represent the removal of jobs from the worker.
|
||||
diff := NewInternalDiffs()
|
||||
for _, job := range jobs {
|
||||
diff.Removed(addr, job)
|
||||
}
|
||||
|
||||
return diff, nil
|
||||
}
|
||||
|
||||
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) {
|
||||
start := time.Now()
|
||||
defer func() {
|
||||
|
|
@ -306,9 +124,9 @@ func (b *Balancer) AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.Q
|
|||
// pass a table, we're still encoding the tableKey in the job. In theory, we
|
||||
// could exclude tableKey from the job coming into this method, and add it
|
||||
// here.
|
||||
qdbid := qtid.QualifiedDatabaseID
|
||||
dbid := qtid.QualifiedDatabaseID.DatabaseID
|
||||
|
||||
diff, err := b.addJobs(tx, roleType, qdbid, jobs...)
|
||||
diff, err := b.addJobs(tx, roleType, dbid, jobs...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "adding job")
|
||||
}
|
||||
|
|
@ -316,59 +134,46 @@ func (b *Balancer) AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.Q
|
|||
return diff.Output(), nil
|
||||
}
|
||||
|
||||
func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs ...dax.Job) (InternalDiffs, error) {
|
||||
diffs := NewInternalDiffs()
|
||||
func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, jobs ...dax.Job) (controller.InternalDiffs, error) {
|
||||
diffs := controller.NewInternalDiffs()
|
||||
|
||||
if len(jobs) == 0 {
|
||||
return diffs, nil
|
||||
}
|
||||
|
||||
if cnt, err := b.current.WorkerCount(tx, roleType, qdbid); err != nil {
|
||||
ws, err := b.store.WorkerService(tx, dbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting worker service")
|
||||
}
|
||||
|
||||
if cnt, err := b.store.WorkerCount(tx, roleType, ws.ID); err != nil {
|
||||
return nil, errors.Wrap(err, "getting worker count")
|
||||
} else if cnt == 0 {
|
||||
if err := b.freeJobs.CreateJobs(tx, roleType, qdbid, jobs...); err != nil {
|
||||
return nil, errors.Wrap(err, "creating free job")
|
||||
}
|
||||
|
||||
// Since we've just added free jobs to the database, try to add a worker
|
||||
// for the database. This case would happen because when a database is
|
||||
// first created, it is not assigned any workers. A database is not
|
||||
// assigned workers until it has at least one job (which this database
|
||||
// now has).
|
||||
if diff, err := b.balanceDatabaseForRole(tx, roleType, qdbid); err != nil {
|
||||
return nil, errors.Wrapf(err, "balancing database for role: (%s)", roleType)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
|
||||
// Now check, again, to see if the database has a worker.
|
||||
if cnt2, err := b.current.WorkerCount(tx, roleType, qdbid); err != nil {
|
||||
return nil, errors.Wrap(err, "getting worker count, again")
|
||||
} else if cnt2 == 0 {
|
||||
// TODO: we might want to inform the user that a job is in the free list
|
||||
// because there are no workers.
|
||||
return InternalDiffs{}, nil
|
||||
if err := b.store.CreateFreeJobs(tx, roleType, dbid, jobs...); err != nil {
|
||||
return nil, errors.Wrap(err, "creating free jobs")
|
||||
}
|
||||
}
|
||||
|
||||
diff, err := b.addDatabaseJobs(tx, roleType, qdbid, jobs...)
|
||||
diff, err := b.addDatabaseJobs(tx, roleType, dbid, ws.ID, jobs...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "adding database jobs: (%s) %s", roleType, qdbid)
|
||||
return nil, errors.Wrapf(err, "adding database jobs: (%s) %s", roleType, dbid)
|
||||
}
|
||||
diffs.Merge(diff)
|
||||
|
||||
return diffs, nil
|
||||
}
|
||||
|
||||
// addDatabaseJobs adds the job for the provided database.
|
||||
func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs ...dax.Job) (InternalDiffs, error) {
|
||||
workerJobs, err := b.current.WorkersJobs(tx, roleType, qdbid)
|
||||
// addDatabaseJobs adds the job for the provided database. TODO, bad doc
|
||||
func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID, jobs ...dax.Job) (controller.InternalDiffs, error) {
|
||||
workerJobs, err := b.store.WorkersJobs(tx, roleType, svcID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting workers jobs: %s", roleType)
|
||||
}
|
||||
|
||||
jset := dax.NewSet[dax.Job]()
|
||||
addrToID := make(map[dax.Address]string)
|
||||
for _, workerInfo := range workerJobs {
|
||||
addrToID[workerInfo.Address] = workerInfo.ID
|
||||
jset.Merge(dax.NewSet(workerInfo.Jobs...))
|
||||
}
|
||||
|
||||
|
|
@ -379,7 +184,7 @@ func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, qd
|
|||
jobCounts[v.Address] = len(v.Jobs)
|
||||
}
|
||||
|
||||
diffs := NewInternalDiffs()
|
||||
diffs := controller.NewInternalDiffs()
|
||||
|
||||
jobsToCreate := make(map[dax.Address][]dax.Job)
|
||||
|
||||
|
|
@ -409,7 +214,7 @@ func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, qd
|
|||
}
|
||||
|
||||
for addr, jobs := range jobsToCreate {
|
||||
if err := b.current.AssignWorkerToJobs(tx, roleType, qdbid, addr, jobs...); err != nil {
|
||||
if err := b.store.AssignWorkerToJobs(tx, roleType, dbid, addrToID[addr], jobs...); err != nil {
|
||||
return nil, errors.Wrap(err, "creating job")
|
||||
}
|
||||
for _, job := range jobs {
|
||||
|
|
@ -420,88 +225,48 @@ func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, qd
|
|||
return diffs, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) {
|
||||
qdbid := qtid.QualifiedDatabaseID
|
||||
|
||||
// If no jobs are provided, remove all jobs for table.
|
||||
if len(jobs) == 0 {
|
||||
diffs, err := b.removeJobsForTable(tx, roleType, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "removing jobs for table: (%s) %s", roleType, qtid)
|
||||
}
|
||||
return diffs.Output(), nil
|
||||
func (b *Balancer) CurrentState(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerInfo, error) {
|
||||
dbid := qdbid.DatabaseID
|
||||
ws, err := b.store.WorkerService(tx, dbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting worker service")
|
||||
}
|
||||
return b.store.WorkersJobs(tx, roleType, ws.ID)
|
||||
}
|
||||
|
||||
diffs := NewInternalDiffs()
|
||||
func (b *Balancer) WorkerState(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) (dax.WorkerInfo, error) {
|
||||
return b.store.WorkerJobs(tx, roleType, addr)
|
||||
}
|
||||
|
||||
for _, job := range jobs {
|
||||
if diff, err := b.removeJob(tx, roleType, qdbid, job); err != nil {
|
||||
return nil, errors.Wrapf(err, "removing job: (%s) %s, %s", roleType, qdbid, job)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
func (b *Balancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.WorkerDiff, error) {
|
||||
diffs, err := b.store.DeleteJobsForTable(tx, roleType, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "deleting jobs from storage")
|
||||
}
|
||||
|
||||
return diffs.Output(), nil
|
||||
}
|
||||
|
||||
func (b *Balancer) removeJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) (InternalDiffs, error) {
|
||||
idiffs, err := b.current.DeleteJobsForTable(tx, roleType, qtid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "deleting jobs for table: (%s) %s", roleType, qtid)
|
||||
}
|
||||
if err := b.freeJobs.DeleteJobsForTable(tx, roleType, qtid); err != nil {
|
||||
return nil, errors.Wrapf(err, "deleting free jobs for table: (%s) %s", roleType, qtid)
|
||||
}
|
||||
return idiffs, nil
|
||||
}
|
||||
|
||||
// Balance calls balanceDatabase on every database in the schemar.
|
||||
func (b *Balancer) Balance(tx dax.Transaction) ([]dax.WorkerDiff, error) {
|
||||
diffs, err := b.balance(tx)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "balancing all")
|
||||
}
|
||||
|
||||
return diffs.Output(), nil
|
||||
}
|
||||
|
||||
func (b *Balancer) balance(tx dax.Transaction) (InternalDiffs, error) {
|
||||
qdbs, err := b.schemar.Databases(tx, "")
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting all databases")
|
||||
}
|
||||
|
||||
diffs := NewInternalDiffs()
|
||||
|
||||
for _, qdb := range qdbs {
|
||||
qdbid := qdb.QualifiedID()
|
||||
if diff, err := b.balanceDatabase(tx, qdbid); err != nil {
|
||||
return nil, errors.Wrapf(err, "balancing database: %s", qdbid)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
}
|
||||
|
||||
return diffs, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) {
|
||||
diffs, err := b.balanceDatabase(tx, qdbid)
|
||||
dbid := qdbid.DatabaseID
|
||||
ws, err := b.store.WorkerService(tx, dbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "balancing database: %s", qdbid)
|
||||
return nil, errors.Wrap(err, "getting worker service")
|
||||
}
|
||||
|
||||
diffs, err := b.balanceDatabase(tx, qdbid.DatabaseID, ws.ID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "balancing")
|
||||
}
|
||||
return diffs.Output(), nil
|
||||
}
|
||||
|
||||
func (b *Balancer) balanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) (InternalDiffs, error) {
|
||||
diffs := NewInternalDiffs()
|
||||
func (b *Balancer) balanceDatabase(tx dax.Transaction, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) {
|
||||
diffs := controller.NewInternalDiffs()
|
||||
|
||||
for _, role := range dax.AllRoleTypes {
|
||||
diff, err := b.balanceDatabaseForRole(tx, role, qdbid)
|
||||
diff, err := b.balanceDatabaseForRole(tx, role, dbid, svcID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", role, qdbid)
|
||||
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", role, dbid)
|
||||
}
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
|
|
@ -509,34 +274,26 @@ func (b *Balancer) balanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDataba
|
|||
return diffs, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (InternalDiffs, error) {
|
||||
b.logger.Debugf("balancing database %s for role: %s\n", qdbid, roleType)
|
||||
diffs := NewInternalDiffs()
|
||||
|
||||
// Before balancing, make sure the database has its minimum number of
|
||||
// workers satisfied.
|
||||
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)
|
||||
}
|
||||
func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) {
|
||||
b.logger.Debugf("balancing database %s for role: %s\n", dbid, roleType)
|
||||
diffs := controller.NewInternalDiffs()
|
||||
|
||||
// If there are no workers, we can't properly balance.
|
||||
if cnt, err := b.current.WorkerCount(tx, roleType, qdbid); err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", roleType, qdbid)
|
||||
if cnt, err := b.store.WorkerCount(tx, roleType, svcID); err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", roleType, dbid)
|
||||
} else if cnt == 0 {
|
||||
return InternalDiffs{}, nil
|
||||
return controller.InternalDiffs{}, errors.Errorf("unexpected? no workers for database %v (TODO: this used to silently return nil)", dbid)
|
||||
}
|
||||
|
||||
// Process the freeJobs.
|
||||
if diff, err := b.processFreeJobs(tx, roleType, qdbid); err != nil {
|
||||
return nil, errors.Wrapf(err, "processing free jobs: (%s) %s", roleType, qdbid)
|
||||
if diff, err := b.processFreeJobs(tx, roleType, dbid, svcID); err != nil {
|
||||
return nil, errors.Wrapf(err, "processing free jobs: (%s) %s", roleType, dbid)
|
||||
} else {
|
||||
diffs.Merge(diff)
|
||||
}
|
||||
|
||||
// Balance the jobs among workers.
|
||||
diff, err := b.balanceDatabaseJobs(tx, roleType, qdbid, diffs)
|
||||
diff, err := b.balanceDatabaseJobs(tx, roleType, dbid, svcID, diffs)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "balancing jobs")
|
||||
}
|
||||
|
|
@ -544,129 +301,6 @@ func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleT
|
|||
return diff, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) CurrentState(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerInfo, error) {
|
||||
return b.current.WorkersJobs(tx, roleType, qdbid)
|
||||
}
|
||||
|
||||
func (b *Balancer) WorkerState(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) (dax.WorkerInfo, error) {
|
||||
info := dax.WorkerInfo{
|
||||
Address: addr,
|
||||
}
|
||||
|
||||
dbkey := b.current.DatabaseForWorker(tx, addr)
|
||||
if dbkey == "" {
|
||||
return info, nil
|
||||
}
|
||||
qdbid := dbkey.QualifiedDatabaseID()
|
||||
|
||||
jobs, err := b.current.ListJobs(tx, roleType, qdbid, addr)
|
||||
if err != nil {
|
||||
return dax.WorkerInfo{}, errors.Wrapf(err, "listing jobs: (%s) %s, %s", roleType, qdbid, addr)
|
||||
}
|
||||
info.Jobs = jobs
|
||||
|
||||
return info, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) WorkersForJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs ...dax.Job) ([]dax.WorkerInfo, error) {
|
||||
out := make(map[dax.Address]dax.Set[dax.Job])
|
||||
|
||||
workerJobs, err := b.current.WorkersJobs(tx, roleType, qdbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker jobs: (%s) %s", roleType, qdbid)
|
||||
}
|
||||
for _, workerInfo := range workerJobs {
|
||||
jset := dax.NewSet(workerInfo.Jobs...)
|
||||
|
||||
matches := dax.NewSet[dax.Job]()
|
||||
for _, job := range jobs {
|
||||
if jset.Contains(job) {
|
||||
matches.Add(job)
|
||||
}
|
||||
}
|
||||
|
||||
if len(matches) > 0 {
|
||||
out[workerInfo.Address] = matches
|
||||
}
|
||||
}
|
||||
|
||||
workers := make([]dax.WorkerInfo, 0, len(out))
|
||||
|
||||
for addr, jset := range out {
|
||||
workers = append(workers, dax.WorkerInfo{
|
||||
Address: addr,
|
||||
Jobs: jset.Sorted(),
|
||||
})
|
||||
}
|
||||
|
||||
sort.Sort(dax.WorkerInfos(workers))
|
||||
|
||||
return workers, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) WorkersForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.WorkerInfo, error) {
|
||||
out := make(map[dax.Address]dax.Set[dax.Job])
|
||||
|
||||
qdbid := qtid.QualifiedDatabaseID
|
||||
|
||||
prefix := string(qtid.Key())
|
||||
|
||||
workerJobs, err := b.current.WorkersJobs(tx, roleType, qdbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker jobs: (%s) %s", roleType, qdbid)
|
||||
}
|
||||
for _, workerInfo := range workerJobs {
|
||||
|
||||
matches := dax.NewSet[dax.Job]()
|
||||
for _, job := range workerInfo.Jobs {
|
||||
if strings.HasPrefix(string(job), prefix) {
|
||||
matches.Add(job)
|
||||
}
|
||||
}
|
||||
if len(matches) > 0 {
|
||||
out[workerInfo.Address] = matches
|
||||
}
|
||||
}
|
||||
|
||||
workers := make([]dax.WorkerInfo, 0, len(out))
|
||||
|
||||
for addr, jset := range out {
|
||||
workers = append(workers, dax.WorkerInfo{
|
||||
Address: addr,
|
||||
Jobs: jset.Sorted(),
|
||||
})
|
||||
}
|
||||
|
||||
sort.Sort(dax.WorkerInfos(workers))
|
||||
|
||||
return workers, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) removeJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job dax.Job) (InternalDiffs, error) {
|
||||
if addr, ok, err := b.workerForJob(tx, roleType, qdbid, job); err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker for job: %s", job)
|
||||
} else if ok {
|
||||
if err := b.current.DeleteJob(tx, roleType, qdbid, addr, job); err != nil {
|
||||
return nil, errors.Wrapf(err, "deleting job: (%s) %s, %s, %s", roleType, qdbid, addr, job)
|
||||
}
|
||||
|
||||
diffs := NewInternalDiffs()
|
||||
diffs.Removed(addr, job)
|
||||
|
||||
return diffs, nil
|
||||
}
|
||||
|
||||
// Just in case the job is in the free list (and wasn't assigned to a
|
||||
// worker), remove it; there's no need to provide a diff. There should never
|
||||
// be a case where the same job is both in the free list and assigned to a
|
||||
// worker.
|
||||
if err := b.freeJobs.DeleteJob(tx, roleType, qdbid, job); err != nil {
|
||||
return nil, errors.Wrapf(err, "deleting free job: (%s) %s, %s", roleType, qdbid, job)
|
||||
}
|
||||
|
||||
return InternalDiffs{}, nil
|
||||
}
|
||||
|
||||
// balanceDatabaseJobs moves jobs among workers with the goal of having an equal
|
||||
// number of jobs per worker. This method takes an `internalDiffs` as input for
|
||||
// cases where some action has preceeded this call which also resulted in
|
||||
|
|
@ -674,72 +308,83 @@ func (b *Balancer) removeJob(tx dax.Transaction, roleType dax.RoleType, qdbid da
|
|||
// the internalDiffs.merge() method, but we would need to modify that method to
|
||||
// be smarter about the order in which it applies the add/remove operations.
|
||||
// Until that's in place, we'll pass in a value here.
|
||||
func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, diffs InternalDiffs) (InternalDiffs, error) {
|
||||
numWorkers, err := b.current.WorkerCount(tx, roleType, qdbid)
|
||||
func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID, diffs controller.InternalDiffs) (controller.InternalDiffs, error) {
|
||||
workerInfos, err := b.store.WorkersJobs(tx, roleType, svcID)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", roleType, qdbid)
|
||||
return nil, errors.Wrapf(err, "getting current state: (%s) %s", roleType, dbid)
|
||||
}
|
||||
|
||||
numWorkers := len(workerInfos)
|
||||
if numWorkers == 0 {
|
||||
return nil, errors.Errorf("can't balance a database ID %v with no workers", dbid)
|
||||
}
|
||||
numJobs := 0
|
||||
if addrs, err := b.current.ListWorkers(tx, roleType, qdbid); err != nil {
|
||||
return nil, errors.Wrapf(err, "listing workers: (%s) %s", roleType, qdbid)
|
||||
} else {
|
||||
for _, addr := range addrs {
|
||||
jobCounts, err := b.current.JobCounts(tx, roleType, qdbid, addr)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting job count: (%s) %s, %s", roleType, qdbid, addr)
|
||||
}
|
||||
numJobs += jobCounts[addr]
|
||||
}
|
||||
for _, info := range workerInfos {
|
||||
numJobs += len(info.Jobs)
|
||||
}
|
||||
|
||||
minJobsPerWorker := numJobs / numWorkers
|
||||
numWorkersAboveMin := numJobs % numWorkers
|
||||
|
||||
// workerInfos is used now in order to guarantee a sort order.
|
||||
workerInfos, err := b.CurrentState(tx, roleType, qdbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting current state: (%s) %s", roleType, qdbid)
|
||||
}
|
||||
removedJobs := make(dax.Jobs, 0)
|
||||
addedJobs := make(map[string]dax.Jobs)
|
||||
|
||||
// Loop through each worker, and if the number of jobs for the worker
|
||||
// exceeds the target, then remove the job and add it back (which is
|
||||
// effectively how we rebalance a job).
|
||||
// remove jobs loop
|
||||
//
|
||||
// Find all the workers that have more jobs than the target and
|
||||
// add those jobs to the list of removedJobs.
|
||||
for i, workerInfo := range workerInfos {
|
||||
numTargetJobs := minJobsPerWorker
|
||||
if i < numWorkersAboveMin {
|
||||
numTargetJobs += 1
|
||||
}
|
||||
|
||||
jobCounts, err := b.current.JobCounts(tx, roleType, qdbid, workerInfo.Address)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "getting job count: (%s) %s, %s", roleType, qdbid, workerInfo.Address)
|
||||
}
|
||||
numCurrentJobs := jobCounts[workerInfo.Address]
|
||||
|
||||
numCurrentJobs := len(workerInfo.Jobs)
|
||||
// If we don't need to remove jobs from this worker, then just continue
|
||||
// on to the next worker.
|
||||
if numCurrentJobs <= numTargetJobs {
|
||||
continue
|
||||
}
|
||||
|
||||
sortedJobs, err := b.current.ListJobs(tx, roleType, qdbid, workerInfo.Address)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "listing jobs: (%s) %s, %s", roleType, qdbid, workerInfo.Address)
|
||||
for i := numTargetJobs; i < numCurrentJobs; i++ {
|
||||
diffs.Removed(workerInfo.Address, workerInfo.Jobs[i])
|
||||
}
|
||||
removedJobs = append(removedJobs, workerInfo.Jobs[numTargetJobs:]...)
|
||||
}
|
||||
|
||||
// add jobs loop
|
||||
//
|
||||
// Find all the workers that have fewer jobs than the target, and
|
||||
// add jobs to them from the removedJobs list. This should balance
|
||||
// perfectly such that removedJobs is empty at the end of this
|
||||
// loop.
|
||||
for i, workerInfo := range workerInfos {
|
||||
numTargetJobs := minJobsPerWorker
|
||||
if i < numWorkersAboveMin {
|
||||
numTargetJobs += 1
|
||||
}
|
||||
|
||||
// Remove the extra jobs from the end of the list, and add them back
|
||||
// again (which should place them on a worker with fewer jobs).
|
||||
for i := numCurrentJobs - 1; i >= numTargetJobs; i-- {
|
||||
if rj, err := b.removeJob(tx, roleType, qdbid, sortedJobs[i]); err != nil {
|
||||
return nil, errors.Wrapf(err, "removing job: %s", sortedJobs[i])
|
||||
} else {
|
||||
diffs.Merge(rj)
|
||||
}
|
||||
if aj, err := b.addJobs(tx, roleType, qdbid, sortedJobs[i]); err != nil {
|
||||
return nil, errors.Wrapf(err, "adding job: %s", sortedJobs[i])
|
||||
} else {
|
||||
diffs.Merge(aj)
|
||||
}
|
||||
numCurrentJobs := len(workerInfo.Jobs)
|
||||
// If current jobs is already at the target, we don't need to
|
||||
// add jobs to this worker, so continue.
|
||||
if numCurrentJobs >= numTargetJobs {
|
||||
continue
|
||||
}
|
||||
|
||||
for i := 0; i < numTargetJobs-numCurrentJobs; i++ {
|
||||
diffs.Added(workerInfo.Address, removedJobs[0])
|
||||
addedJobs[workerInfo.ID] = append(addedJobs[workerInfo.ID], removedJobs[0])
|
||||
removedJobs = removedJobs[1:]
|
||||
}
|
||||
}
|
||||
if len(removedJobs) != 0 {
|
||||
panic("everything is terrible")
|
||||
}
|
||||
|
||||
for wid, jobs := range addedJobs {
|
||||
err := b.store.AssignWorkerToJobs(tx, roleType, dbid, wid, jobs...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "assigning worker to jobs in DB")
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -747,14 +392,14 @@ func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType
|
|||
}
|
||||
|
||||
// processFreeJobs assigns all jobs in the free list to a worker.
|
||||
func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (InternalDiffs, error) {
|
||||
diffs := NewInternalDiffs()
|
||||
jobs, err := b.freeJobs.ListJobs(tx, roleType, qdbid)
|
||||
func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) {
|
||||
diffs := controller.NewInternalDiffs()
|
||||
jobs, err := b.store.ListFreeJobs(tx, roleType, dbid)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "listing free jobs: %s", roleType)
|
||||
}
|
||||
|
||||
if aj, err := b.addDatabaseJobs(tx, roleType, qdbid, jobs...); err != nil {
|
||||
if aj, err := b.addDatabaseJobs(tx, roleType, dbid, svcID, jobs...); err != nil {
|
||||
return nil, errors.Wrapf(err, "adding jobs: %s", jobs)
|
||||
} else {
|
||||
diffs.Merge(aj)
|
||||
|
|
@ -763,45 +408,24 @@ func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, qd
|
|||
return diffs, nil
|
||||
}
|
||||
|
||||
// workerForJob returns the worker currently assigned to the given job.
|
||||
func (b *Balancer) workerForJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job dax.Job) (dax.Address, bool, error) {
|
||||
workerJobs, err := b.current.WorkersJobs(tx, roleType, qdbid)
|
||||
if err != nil {
|
||||
return "", false, errors.Wrapf(err, "getting workers jobs: (%s) %s", roleType, qdbid)
|
||||
}
|
||||
for _, workerInfo := range workerJobs {
|
||||
jset := dax.NewSet(workerInfo.Jobs...)
|
||||
if jset.Contains(job) {
|
||||
return workerInfo.Address, true, nil
|
||||
}
|
||||
}
|
||||
return "", false, nil
|
||||
}
|
||||
|
||||
func (b *Balancer) ReadNode(tx dax.Transaction, addr dax.Address) (*dax.Node, error) {
|
||||
return b.workerRegistry.Worker(tx, addr)
|
||||
}
|
||||
|
||||
func (b *Balancer) Nodes(tx dax.Transaction) ([]*dax.Node, error) {
|
||||
return b.workerRegistry.Workers(tx)
|
||||
}
|
||||
|
||||
func (b *Balancer) CreateWorkerServiceProvider(tx dax.Transaction, sp dax.WorkerServiceProvider) error {
|
||||
return b.wsp.CreateWorkerServiceProvider(tx, sp)
|
||||
return b.store.CreateWorkerServiceProvider(tx, sp)
|
||||
}
|
||||
|
||||
func (b *Balancer) CreateWorkerService(tx dax.Transaction, srv dax.WorkerService) error {
|
||||
return b.wsp.CreateWorkerService(tx, srv)
|
||||
return b.store.CreateWorkerService(tx, srv)
|
||||
}
|
||||
|
||||
func (b *Balancer) WorkerServiceProviders(tx dax.Transaction /*, future optional filters */) (dax.WorkerServiceProviders, error) {
|
||||
return b.wsp.WorkerServiceProviders(tx)
|
||||
return b.store.WorkerServiceProviders(tx)
|
||||
}
|
||||
|
||||
func (b *Balancer) AssignFreeServiceToDatabase(tx dax.Transaction, wspID dax.WorkerServiceProviderID, qdb *dax.QualifiedDatabase) (*dax.WorkerService, error) {
|
||||
return b.wsp.AssignFreeServiceToDatabase(tx, wspID, qdb)
|
||||
return b.store.AssignFreeServiceToDatabase(tx, wspID, qdb)
|
||||
}
|
||||
|
||||
func (b *Balancer) WorkerServices(tx dax.Transaction, wsp dax.WorkerServiceProviderID) (dax.WorkerServices, error) {
|
||||
return b.wsp.WorkerServices(tx, wsp)
|
||||
return b.store.WorkerServices(tx, wsp)
|
||||
}
|
||||
|
||||
type WorkerJobService interface {
|
||||
|
|
@ -810,12 +434,9 @@ type WorkerJobService interface {
|
|||
WorkerCount(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (int, error)
|
||||
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
|
||||
ReleaseWorkers(tx dax.Transaction, addrs ...dax.Address) 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)
|
||||
DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) (controller.InternalDiffs, error)
|
||||
JobCounts(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr ...dax.Address) (map[dax.Address]int, error)
|
||||
ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) (dax.Jobs, error)
|
||||
|
||||
|
|
@ -832,7 +453,7 @@ type FreeJobService interface {
|
|||
|
||||
type FreeWorkerService interface {
|
||||
PopWorkers(tx dax.Transaction, roleType dax.RoleType, num int) ([]dax.Address, error)
|
||||
ListWorkers(tx dax.Transaction, qdbid dax.QualifiedDatabaseID, roleType dax.RoleType) (dax.Addresses, error)
|
||||
ListWorkers(tx dax.Transaction, roleType dax.RoleType) (dax.Addresses, error)
|
||||
}
|
||||
|
||||
type WorkerServiceProviderService interface {
|
||||
|
|
|
|||
|
|
@ -13,6 +13,8 @@ func TestMain(m *testing.M) {
|
|||
os.Exit(run(m))
|
||||
}
|
||||
|
||||
// TODO moved to sqldb package, remove from here once the individual service tests aren't needed
|
||||
|
||||
// SQLTransactor is a global connection to a SQL database which is
|
||||
// created, migrated, and destroyed for each test run. The database
|
||||
// gets a randomized name and is used by all the tests in this
|
||||
|
|
|
|||
|
|
@ -29,6 +29,24 @@ func TestWorkerJobService(t *testing.T) {
|
|||
}
|
||||
}()
|
||||
|
||||
wspSvc := sqldb.NewWorkerServiceProviderService(nil)
|
||||
err = wspSvc.CreateWorkerServiceProvider(tx, dax.WorkerServiceProvider{
|
||||
ID: "wspID",
|
||||
Roles: []dax.RoleType{"compute", "translate"},
|
||||
Address: "wspaddr",
|
||||
Description: "wspdesc",
|
||||
})
|
||||
require.NoError(t, err, "creating wsp")
|
||||
|
||||
wspSvc.CreateWorkerService(tx, dax.WorkerService{
|
||||
ID: orgID,
|
||||
Roles: []dax.RoleType{},
|
||||
WorkerServiceProviderID: "",
|
||||
DatabaseID: dbID,
|
||||
WorkersMin: 1,
|
||||
WorkersMax: 1,
|
||||
})
|
||||
|
||||
// must have a database to do workerjob stuff
|
||||
schemar := sqldb.NewSchemar(nil)
|
||||
err = schemar.CreateDatabase(tx,
|
||||
|
|
@ -42,6 +60,7 @@ func TestWorkerJobService(t *testing.T) {
|
|||
|
||||
node := &dax.Node{
|
||||
Address: nodeAddr,
|
||||
ServiceID: "mysvc1",
|
||||
RoleTypes: []dax.RoleType{role},
|
||||
}
|
||||
|
||||
|
|
@ -50,9 +69,6 @@ func TestWorkerJobService(t *testing.T) {
|
|||
err = workerReg.AddWorker(tx, node)
|
||||
require.NoError(t, err)
|
||||
|
||||
err = wjSvc.CreateWorker(tx, role, qdbid, nodeAddr)
|
||||
require.NoError(t, err)
|
||||
|
||||
// we create a qtid to prefix jobs so that we can then test the
|
||||
// "DeleteJobsForTable" method. Is it strange that a table is not
|
||||
// explicitly mentioned in the interface on the way in, but is in
|
||||
|
|
|
|||
|
|
@ -2,11 +2,8 @@
|
|||
package controller
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
|
|
@ -33,8 +30,9 @@ var _ dax.WorkerRegistry = (*Controller)(nil)
|
|||
type Controller struct {
|
||||
// Schemar is used by the controller to get table, and other schema,
|
||||
// information.
|
||||
Schemar schemar.Schemar
|
||||
Balancer Balancer
|
||||
Schemar schemar.Schemar
|
||||
Balancer Balancer
|
||||
WSPGetter func(addr dax.Address, log logger.Logger) WorkerServiceProvider
|
||||
|
||||
Transactor dax.Transactor
|
||||
|
||||
|
|
@ -656,14 +654,11 @@ func (c *Controller) CreateDatabase(ctx context.Context, qdb *dax.QualifiedDatab
|
|||
return errors.Wrap(err, "assigning free service to database")
|
||||
}
|
||||
|
||||
// Encode the request.
|
||||
postBody, err := json.Marshal(svc)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "marshalling post request")
|
||||
}
|
||||
buf := bytes.NewBuffer(postBody)
|
||||
client := c.WSPGetter(wsp.Address, c.Logger())
|
||||
|
||||
http.Post(fmt.Sprintf("%s/claim", wsp.Address), "application/json", buf)
|
||||
if err := client.ClaimService(ctx, svc); err != nil {
|
||||
return errors.Wrap(err, "claiming service")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
|
@ -697,9 +692,10 @@ func (c *Controller) DropDatabase(ctx context.Context, qdbid dax.QualifiedDataba
|
|||
for worker := range workerSet {
|
||||
addrs = append(addrs, worker)
|
||||
}
|
||||
if err := c.Balancer.ReleaseWorkers(tx, addrs...); err != nil {
|
||||
return errors.Wrap(err, "freeing workers")
|
||||
}
|
||||
// if err := c.Balancer.ReleaseWorkers(tx, addrs...); err != nil {
|
||||
// return errors.Wrap(err, "freeing workers")
|
||||
// }
|
||||
// TODO: this becomes balancer.DropWorkerService
|
||||
|
||||
// Drop the database record from the schema.
|
||||
if err := c.Schemar.DropDatabase(tx, qdbid); err != nil {
|
||||
|
|
@ -1014,53 +1010,53 @@ func (c *Controller) Tables(ctx context.Context, qdbid dax.QualifiedDatabaseID,
|
|||
return c.Schemar.Tables(tx, qdbid, ids...)
|
||||
}
|
||||
|
||||
// RemoveShards deregisters the table/shard combinations with the controller and
|
||||
// sends the necessary directives.
|
||||
func (c *Controller) RemoveShards(ctx context.Context, qtid dax.QualifiedTableID, shards ...dax.ShardNum) error {
|
||||
// // RemoveShards deregisters the table/shard combinations with the controller and
|
||||
// // sends the necessary directives.
|
||||
// func (c *Controller) RemoveShards(ctx context.Context, qtid dax.QualifiedTableID, shards ...dax.ShardNum) error {
|
||||
|
||||
var directives []*dax.Directive
|
||||
// var directives []*dax.Directive
|
||||
|
||||
fn := func(tx dax.Transaction, writable bool) error {
|
||||
// workerSet maintains the set of workers which have a job assignment change
|
||||
// and therefore need to be sent an updated Directive.
|
||||
workerSet := NewAddressSet()
|
||||
// fn := func(tx dax.Transaction, writable bool) error {
|
||||
// // workerSet maintains the set of workers which have a job assignment change
|
||||
// // and therefore need to be sent an updated Directive.
|
||||
// workerSet := NewAddressSet()
|
||||
|
||||
for _, s := range shards {
|
||||
// We don't currently use the returned diff, other than to determine
|
||||
// which worker was affected, because we send the full Directive every
|
||||
// time.
|
||||
diffs, err := c.Balancer.RemoveJobs(tx, dax.RoleTypeCompute, qtid, shard(qtid.Key(), s).Job())
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "removing job")
|
||||
}
|
||||
for _, diff := range diffs {
|
||||
workerSet.Add(dax.Address(diff.Address))
|
||||
}
|
||||
}
|
||||
// for _, s := range shards {
|
||||
// // We don't currently use the returned diff, other than to determine
|
||||
// // which worker was affected, because we send the full Directive every
|
||||
// // time.
|
||||
// diffs, err := c.Balancer.RemoveJobs(tx, dax.RoleTypeCompute, qtid, shard(qtid.Key(), s).Job())
|
||||
// if err != nil {
|
||||
// return errors.Wrap(err, "removing job")
|
||||
// }
|
||||
// for _, diff := range diffs {
|
||||
// workerSet.Add(dax.Address(diff.Address))
|
||||
// }
|
||||
// }
|
||||
|
||||
// Convert the slice of addresses into a slice of addressMethod containing
|
||||
// the appropriate method.
|
||||
addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodFull)
|
||||
// // Convert the slice of addresses into a slice of addressMethod containing
|
||||
// // the appropriate method.
|
||||
// addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodFull)
|
||||
|
||||
var err error
|
||||
directives, err = c.buildDirectives(ctx, tx, addrMethods)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "building directives")
|
||||
}
|
||||
// var err error
|
||||
// directives, err = c.buildDirectives(ctx, tx, addrMethods)
|
||||
// if err != nil {
|
||||
// return errors.Wrap(err, "building directives")
|
||||
// }
|
||||
|
||||
return nil
|
||||
}
|
||||
// return nil
|
||||
// }
|
||||
|
||||
if err := dax.RetryWithTx(ctx, c.Transactor, fn, true, txRetry); err != nil {
|
||||
return errors.Wrap(err, "retry with tx: write")
|
||||
}
|
||||
// if err := dax.RetryWithTx(ctx, c.Transactor, fn, true, txRetry); err != nil {
|
||||
// return errors.Wrap(err, "retry with tx: write")
|
||||
// }
|
||||
|
||||
if err := c.sendDirectives(ctx, directives); err != nil {
|
||||
return NewErrDirectiveSendFailure(err.Error())
|
||||
}
|
||||
// if err := c.sendDirectives(ctx, directives); err != nil {
|
||||
// return NewErrDirectiveSendFailure(err.Error())
|
||||
// }
|
||||
|
||||
return nil
|
||||
}
|
||||
// return nil
|
||||
// }
|
||||
|
||||
func (c *Controller) sendDirectives(ctx context.Context, directives []*dax.Directive) error {
|
||||
if len(directives) == 0 {
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ import (
|
|||
"github.com/featurebasedb/featurebase/v3/dax/controller"
|
||||
controllerhttp "github.com/featurebasedb/featurebase/v3/dax/controller/http"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller/sqldb"
|
||||
wsphttp "github.com/featurebasedb/featurebase/v3/dax/worker_service_provider/http"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
fbnet "github.com/featurebasedb/featurebase/v3/net"
|
||||
|
|
@ -31,6 +32,7 @@ func New(uri *fbnet.URI, cfg controller.Config) *controllerService {
|
|||
}
|
||||
|
||||
controller := controller.New(cfg)
|
||||
controller.WSPGetter = wsphttp.NewClient
|
||||
controllerSvc := &controllerService{
|
||||
uri: uri,
|
||||
controller: controller,
|
||||
|
|
|
|||
|
|
@ -12,7 +12,7 @@ func NewBalancer(log logger.Logger) *balancer.Balancer {
|
|||
wjs := NewWorkerJobService(log)
|
||||
fws := NewFreeWorkerService(log)
|
||||
ns := NewWorkerRegistry(log)
|
||||
wsp := NewWorkerServiceProviderService(log)
|
||||
store := NewStore(log)
|
||||
|
||||
return balancer.New(ns, fjs, wjs, fws, schemar, wsp, log)
|
||||
return balancer.New(ns, fjs, wjs, fws, schemar, store, log)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -49,7 +49,7 @@ func (fw *freeWorkerService) PopWorkers(tx dax.Transaction, roleType dax.RoleTyp
|
|||
return ret, nil
|
||||
}
|
||||
|
||||
func (fw *freeWorkerService) ListWorkers(tx dax.Transaction, qdbid dax.QualifiedDatabaseID, roleType dax.RoleType) (dax.Addresses, error) {
|
||||
func (fw *freeWorkerService) ListWorkers(tx dax.Transaction, roleType dax.RoleType) (dax.Addresses, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
|
|
|
|||
40
dax/controller/sqldb/sqldb_init_test.go
Normal file
40
dax/controller/sqldb/sqldb_init_test.go
Normal file
|
|
@ -0,0 +1,40 @@
|
|||
package sqldb_test
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller/sqldb"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
)
|
||||
|
||||
func TestMain(m *testing.M) {
|
||||
os.Exit(run(m))
|
||||
}
|
||||
|
||||
// SQLTransactor is a global connection to a SQL database which is
|
||||
// created, migrated, and destroyed for each test run. The database
|
||||
// gets a randomized name and is used by all the tests in this
|
||||
// package.
|
||||
var SQLTransactor sqldb.Transactor
|
||||
|
||||
// run is a separate function so that we can defer cleanups. (deferred functions won't run if os.Exit is called)
|
||||
func run(m *testing.M) int {
|
||||
// We connect to a randomized database, create it, and run migrations. Then we drop it when tests are done.
|
||||
conf := sqldb.GetTestConfigRandomDB("balancer_test")
|
||||
var err error
|
||||
SQLTransactor, err = sqldb.NewTransactor(conf, logger.StderrLogger)
|
||||
if err != nil {
|
||||
fmt.Printf("couldn't set up transactor: %v", err)
|
||||
return -1
|
||||
}
|
||||
if err := SQLTransactor.Start(); err != nil {
|
||||
fmt.Printf("couldn't start transactor: %v", err)
|
||||
return -1
|
||||
}
|
||||
|
||||
defer sqldb.DropDatabase(SQLTransactor)
|
||||
code := m.Run()
|
||||
return code
|
||||
}
|
||||
405
dax/controller/sqldb/store.go
Normal file
405
dax/controller/sqldb/store.go
Normal file
|
|
@ -0,0 +1,405 @@
|
|||
package sqldb
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller/balancer"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/models"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/gobuffalo/nulls"
|
||||
"github.com/gofrs/uuid"
|
||||
)
|
||||
|
||||
var _ controller.Store = (*store)(nil)
|
||||
|
||||
func NewStore(log logger.Logger) *store {
|
||||
if log == nil {
|
||||
log = logger.NopLogger
|
||||
}
|
||||
return &store{
|
||||
log: log,
|
||||
}
|
||||
}
|
||||
|
||||
type store struct {
|
||||
log logger.Logger
|
||||
}
|
||||
|
||||
// AddWorker adds the given node to the workers table. For
|
||||
// convenience, it returns the DatabaseID of the database
|
||||
// with which the worker's WorkerService is associated. If the Service
|
||||
// is not assigned to a database, it returns nil.
|
||||
func (s *store) AddWorker(tx dax.Transaction, node *dax.Node) (*dax.DatabaseID, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
worker := &models.Worker{
|
||||
Address: node.Address,
|
||||
ServiceID: node.ServiceID,
|
||||
}
|
||||
for _, roleType := range node.RoleTypes {
|
||||
if err := worker.SetRole(roleType); err != nil {
|
||||
return nil, errors.Wrapf(err, "setting role: %s", roleType)
|
||||
}
|
||||
}
|
||||
|
||||
if err := dt.C.Create(worker); err != nil {
|
||||
return nil, errors.Wrapf(err, "putting worker into database")
|
||||
}
|
||||
|
||||
workerSvc := &models.WorkerService{}
|
||||
if err := dt.C.Find(workerSvc, worker.ServiceID); err != nil {
|
||||
return nil, errors.Wrap(err, "getting worker service")
|
||||
}
|
||||
if !workerSvc.DatabaseID.Valid {
|
||||
return nil, nil
|
||||
}
|
||||
return (*dax.DatabaseID)(&workerSvc.DatabaseID.String), nil
|
||||
}
|
||||
|
||||
// WorkerCount returns the number of workers with the given service
|
||||
// ID.
|
||||
func (s *store) WorkerCount(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) (int, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
cnt, err := dt.C.Where(fmt.Sprintf("service_id = ? AND role_%s = true", role), svcID).Count(&models.Worker{})
|
||||
if err != nil {
|
||||
return 0, errors.Wrapf(err, "counting workers for service '%s'", svcID)
|
||||
}
|
||||
return cnt, nil
|
||||
}
|
||||
|
||||
// WorkerCount returns the number of workers with the given service
|
||||
// ID.
|
||||
func (s *store) WorkerCountDatabase(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (int, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
sql := `SELECT
|
||||
count(w.id) as cnt
|
||||
FROM
|
||||
worker_services ws
|
||||
INNER JOIN databases d ON ws.database_id = d.id
|
||||
INNER JOIN workers w ON w.service_id = ws.id
|
||||
WHERE
|
||||
d.id = ?`
|
||||
|
||||
var res struct {
|
||||
Count int `db:"cnt"`
|
||||
}
|
||||
if err := dt.C.RawQuery(sql, dbid).First(&res); err != nil {
|
||||
return 0, errors.Wrapf(err, "counting workers for database '%s'", dbid)
|
||||
}
|
||||
return res.Count, nil
|
||||
}
|
||||
|
||||
func (s *store) ListFreeJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (dax.Jobs, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
jobs := make(models.Jobs, 0)
|
||||
err := dt.C.Where("role = ? and database_id = ? and worker_id is NULL", role, dbid).Order("name asc").All(&jobs)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "querying for jobs")
|
||||
}
|
||||
|
||||
djs := make(dax.Jobs, len(jobs))
|
||||
for i, job := range jobs {
|
||||
djs[i] = job.Name
|
||||
}
|
||||
return djs, nil
|
||||
}
|
||||
|
||||
func (s *store) WorkerJobs(tx dax.Transaction, role dax.RoleType, addr dax.Address) (dax.WorkerInfo, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.WorkerInfo{}, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
sql := `SELECT
|
||||
jobs.name as name
|
||||
workers.id as id
|
||||
FROM
|
||||
jobs INNER JOIN workers on jobs.worker_id = worker.id
|
||||
WHERE
|
||||
workers.address = ?`
|
||||
|
||||
results := []struct {
|
||||
Name dax.Job `db:"name"`
|
||||
ID uuid.UUID `db:"id"`
|
||||
}{}
|
||||
if err := dt.C.RawQuery(sql, addr).First(&results); err != nil {
|
||||
return dax.WorkerInfo{}, errors.Wrap(err, "querying for jobs")
|
||||
}
|
||||
|
||||
if len(results) == 0 {
|
||||
return dax.WorkerInfo{Address: addr}, nil
|
||||
}
|
||||
wi := dax.WorkerInfo{
|
||||
Address: addr,
|
||||
ID: dax.WorkerID(results[0].ID),
|
||||
Jobs: make([]dax.Job, len(results)),
|
||||
}
|
||||
for i, res := range results {
|
||||
wi.Jobs[i] = res.Name
|
||||
}
|
||||
return wi, nil
|
||||
}
|
||||
|
||||
func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) ([]dax.WorkerInfo, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
results := []struct {
|
||||
Address dax.Address `db:"address"`
|
||||
WorkerID string `db:"worker_id"`
|
||||
JobName nulls.String `db:"job_name"`
|
||||
}{}
|
||||
|
||||
sql := `SELECT
|
||||
workers.address as address
|
||||
workers.id as worker_id
|
||||
jobs.name as job_name
|
||||
FROM
|
||||
workers LEFT JOIN jobs on jobs.worker_id = workers.id
|
||||
WHERE
|
||||
jobs.role = ?
|
||||
workers.service_id = ?
|
||||
ORDER BY
|
||||
address ASC
|
||||
job_name ASC`
|
||||
|
||||
err := dt.C.RawQuery(sql, role, svcID).All(&results)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "querying")
|
||||
}
|
||||
|
||||
infos := make([]dax.WorkerInfo, 0)
|
||||
if len(results) == 0 {
|
||||
return infos, nil
|
||||
}
|
||||
|
||||
// convert flat results list to []dax.WorkerInfo
|
||||
addr := dax.Address("")
|
||||
infosIDX := -1
|
||||
for _, res := range results {
|
||||
if res.Address != addr {
|
||||
addr = res.Address
|
||||
infosIDX++
|
||||
infos = append(infos, dax.WorkerInfo{
|
||||
Address: addr,
|
||||
ID: res.WorkerID,
|
||||
})
|
||||
}
|
||||
if res.JobName.Valid {
|
||||
infos[infosIDX].Jobs = append(infos[infosIDX].Jobs, dax.Job(res.JobName.String))
|
||||
}
|
||||
}
|
||||
return infos, nil
|
||||
}
|
||||
|
||||
func (s *store) CreateFreeJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, jobs ...dax.Job) error {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
modelJobs := make([]models.Job, len(jobs))
|
||||
for i, job := range jobs {
|
||||
modelJobs[i] = models.Job{
|
||||
Name: job,
|
||||
Role: role,
|
||||
DatabaseID: dbid,
|
||||
}
|
||||
}
|
||||
|
||||
err := dt.C.Create(&modelJobs)
|
||||
return errors.Wrap(err, "creating jobs")
|
||||
}
|
||||
|
||||
func (s *store) AssignWorkerToJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, workerID dax.WorkerID, jobs ...dax.Job) error {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
modelJobs := models.Jobs{}
|
||||
sql := `UPDATE
|
||||
jobs
|
||||
SET
|
||||
worker_id = ?
|
||||
WHERE
|
||||
role = ? AND
|
||||
database_id = ? AND
|
||||
name in (?)
|
||||
RETURNING
|
||||
jobs.id,
|
||||
jobs.name`
|
||||
|
||||
err := dt.C.RawQuery(sql, workerID, role, dbid, jobs).All(&modelJobs)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "updating jobs")
|
||||
}
|
||||
|
||||
// Create any jobs not in the returned "modelJobs", which means they
|
||||
// didn't exist already and wouldn't have been updated by the
|
||||
// UPDATE statement.
|
||||
toBeAssigned := sjobsNotAssigned(jobs, modelJobs, role, uuid.UUID(workerID), dbid)
|
||||
|
||||
if err := dt.C.Create(toBeAssigned); err != nil {
|
||||
return errors.Wrap(err, "creating jobs")
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *store) ListWorkers(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) (dax.Addresses, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
workers := models.Workers{}
|
||||
sql := fmt.Sprintf("role_%s = true and service_id = ?", role)
|
||||
err := dt.C.Select("address").Where(sql, svcID).Order("address asc").All(&workers)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting workers")
|
||||
}
|
||||
|
||||
ret := make(dax.Addresses, len(workers))
|
||||
for i, wrkr := range workers {
|
||||
ret[i] = wrkr.Address
|
||||
}
|
||||
|
||||
return ret, nil
|
||||
}
|
||||
|
||||
// TODO: rename once we delete jobsNotAssigned
|
||||
func sjobsNotAssigned(incomingJobs []dax.Job, assigned models.Jobs, roleType dax.RoleType, workerUUID uuid.UUID, dbid dax.DatabaseID) (toBeAssigned models.Jobs) {
|
||||
outer:
|
||||
for _, incJob := range incomingJobs {
|
||||
for _, assignedJob := range assigned {
|
||||
if assignedJob.Name == incJob {
|
||||
continue outer
|
||||
}
|
||||
}
|
||||
toBeAssigned = append(toBeAssigned,
|
||||
models.Job{
|
||||
Name: incJob,
|
||||
Role: roleType,
|
||||
DatabaseID: dbid,
|
||||
WorkerID: nulls.NewUUID(workerUUID),
|
||||
},
|
||||
)
|
||||
}
|
||||
return toBeAssigned
|
||||
}
|
||||
|
||||
func (s *store) WorkerForAddress(tx dax.Transaction, addr dax.Address) (dax.Node, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.Node{}, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
result := struct {
|
||||
ID dax.WorkerID `db:"id"`
|
||||
Address dax.Address `db:"address"`
|
||||
ServiceID dax.WorkerServiceID `db:"service_id"`
|
||||
DatabaseID nulls.String `db:"database_id"`
|
||||
}{}
|
||||
sql := `SELECT
|
||||
w.id AS id,
|
||||
w.address AS address,
|
||||
w.service_id AS service_id,
|
||||
worker_services.database_id AS database_id
|
||||
FROM
|
||||
workers w INNER JOIN worker_services ON w.service_id = worker_services.id
|
||||
WHERE
|
||||
w.address = ?`
|
||||
err := dt.C.RawQuery(sql, addr).First(&result)
|
||||
if err != nil {
|
||||
return dax.Node{}, errors.Wrap(err, "getting workers")
|
||||
}
|
||||
|
||||
var dbid *dax.DatabaseID
|
||||
if result.DatabaseID.Valid {
|
||||
tmp := dax.DatabaseID(result.DatabaseID.String)
|
||||
dbid = &tmp
|
||||
}
|
||||
|
||||
return dax.Node{
|
||||
ID: result.ID,
|
||||
Address: result.Address,
|
||||
ServiceID: result.ServiceID,
|
||||
DatabaseID: dbid,
|
||||
// RoleTypes: []dax.RoleType{}, do we need??
|
||||
}, nil
|
||||
|
||||
}
|
||||
|
||||
func (s *store) WorkerService(tx dax.Transaction, dbid dax.DatabaseID) (dax.WorkerService, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.WorkerService{}, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
ws := &models.WorkerService{}
|
||||
if err := dt.C.Where("database_id = ?", dbid).First(ws); err != nil {
|
||||
return dax.WorkerService{}, errors.Wrapf(err, "finding worker service for database: '%s'", dbid)
|
||||
}
|
||||
return toDaxWorkerService(*ws), nil
|
||||
}
|
||||
|
||||
func (s *store) RemoveWorker(tx dax.Transaction, id dax.WorkerID) error {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
err := dt.C.Destroy(&models.Worker{ID: uuid.UUID(id)})
|
||||
return errors.Wrap(err, "destroying")
|
||||
}
|
||||
|
||||
func (s *store) DeleteJobsForTable(tx dax.Transaction, role dax.RoleType, qtid dax.QualifiedTableID) (controller.InternalDiffs, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
results := []struct {
|
||||
ID uuid.UUID `db:"id"`
|
||||
Name dax.Job `db:"name"`
|
||||
Address dax.Address `db:"address"`
|
||||
}{}
|
||||
err := dt.C.RawQuery("select j.id, j.name, w.address from jobs j inner join workers w on j.worker_id = w.id where j.role = ? and j.database_id = ? and j.name LIKE ?", role, qtid.QualifiedDatabaseID.DatabaseID, fmt.Sprintf("%s%%", qtid.Key())).All(&results)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "querying for jobs")
|
||||
}
|
||||
|
||||
idiffs := make(balancer.InternalDiffs)
|
||||
ids := make([]uuid.UUID, 0, len(results))
|
||||
for _, job := range results {
|
||||
idiffs.Removed(job.Address, job.Name)
|
||||
ids = append(ids, job.ID)
|
||||
}
|
||||
|
||||
if len(ids) > 0 {
|
||||
err = dt.C.RawQuery("DELETE FROM jobs WHERE id in (?)", ids).Exec()
|
||||
}
|
||||
|
||||
return idiffs, errors.Wrap(err, "deleting jobs")
|
||||
}
|
||||
130
dax/controller/sqldb/store_test.go
Normal file
130
dax/controller/sqldb/store_test.go
Normal file
|
|
@ -0,0 +1,130 @@
|
|||
package sqldb_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller/sqldb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
func TestStore(t *testing.T) {
|
||||
tx, err := SQLTransactor.BeginTx(context.Background(), true)
|
||||
require.NoError(t, err, "getting transaction")
|
||||
|
||||
defer func() {
|
||||
err := tx.Rollback()
|
||||
if err != nil {
|
||||
t.Logf("rolling back: %v", err)
|
||||
}
|
||||
}()
|
||||
|
||||
var s controller.Store = sqldb.NewStore(nil)
|
||||
|
||||
wsp := dax.WorkerServiceProvider{
|
||||
ID: "wspID1",
|
||||
Roles: []dax.RoleType{"compute"},
|
||||
Address: "wsp1.example.com:8082",
|
||||
Description: "the description",
|
||||
}
|
||||
wsp2 := dax.WorkerServiceProvider{
|
||||
ID: "wspID2",
|
||||
Roles: []dax.RoleType{"compute", "translate"},
|
||||
Address: "wsp2.example.com:8082",
|
||||
Description: "the description2",
|
||||
}
|
||||
|
||||
t.Run("add a worker service provider", func(t *testing.T) {
|
||||
err := s.CreateWorkerServiceProvider(tx, wsp)
|
||||
require.NoError(t, err)
|
||||
})
|
||||
|
||||
t.Run("read the wsp back", func(t *testing.T) {
|
||||
wsps, err := s.WorkerServiceProviders(tx)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, dax.WorkerServiceProviders{wsp}, wsps)
|
||||
})
|
||||
|
||||
t.Run("add another worker service provider", func(t *testing.T) {
|
||||
err := s.CreateWorkerServiceProvider(tx, wsp2)
|
||||
require.NoError(t, err)
|
||||
})
|
||||
|
||||
t.Run("read both wsps back", func(t *testing.T) {
|
||||
wsps, err := s.WorkerServiceProviders(tx)
|
||||
require.NoError(t, err)
|
||||
require.ElementsMatch(t, dax.WorkerServiceProviders{wsp, wsp2}, wsps)
|
||||
})
|
||||
|
||||
// Three worker services, two from WSP1, 3rd from WSP2
|
||||
ws := dax.WorkerService{
|
||||
ID: "wsID1",
|
||||
Roles: []dax.RoleType{"compute"},
|
||||
WorkerServiceProviderID: "wspID1",
|
||||
DatabaseID: "",
|
||||
WorkersMin: 1,
|
||||
WorkersMax: 1,
|
||||
}
|
||||
|
||||
ws2 := ws
|
||||
ws2.ID = "wsID2"
|
||||
|
||||
ws3 := ws
|
||||
ws3.ID = "wsID3"
|
||||
ws3.Roles = []dax.RoleType{"compute", "translate"}
|
||||
ws3.WorkerServiceProviderID = "wspID2"
|
||||
|
||||
t.Run("create 3 worker services across 2 wsps", func(t *testing.T) {
|
||||
require.NoError(t, s.CreateWorkerService(tx, ws), "creating worker service 1")
|
||||
require.NoError(t, s.CreateWorkerService(tx, ws2), "creating worker service 2")
|
||||
require.NoError(t, s.CreateWorkerService(tx, ws3), "creating worker service 3")
|
||||
})
|
||||
|
||||
t.Run("query all worker services regardless of wsp", func(t *testing.T) {
|
||||
wsvcs, err := s.WorkerServices(tx, "")
|
||||
require.NoError(t, err, "getting all worker services")
|
||||
require.ElementsMatch(t, dax.WorkerServices{ws, ws2, ws3}, wsvcs)
|
||||
})
|
||||
|
||||
t.Run("query worker services for wspID1", func(t *testing.T) {
|
||||
wsvcs, err := s.WorkerServices(tx, "wspID1")
|
||||
require.NoError(t, err, "getting all worker services")
|
||||
require.ElementsMatch(t, dax.WorkerServices{ws, ws2}, wsvcs)
|
||||
})
|
||||
|
||||
node1 := &dax.Node{
|
||||
Address: "nodeaddr1",
|
||||
ServiceID: ws.ID,
|
||||
RoleTypes: ws.Roles,
|
||||
}
|
||||
|
||||
t.Run("register a worker in ws1", func(t *testing.T) {
|
||||
qdbidp, err := s.AddWorker(tx, node1)
|
||||
require.NoError(t, err)
|
||||
require.Nil(t, qdbidp)
|
||||
})
|
||||
|
||||
t.Run("test WorkerCount", func(t *testing.T) {
|
||||
cnt, err := s.WorkerCount(tx, dax.RoleTypeCompute, ws.ID)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, 1, cnt)
|
||||
})
|
||||
|
||||
t.Run("list workers", func(t *testing.T) {
|
||||
addrs, err := s.ListWorkers(tx, dax.RoleTypeCompute, ws.ID)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, dax.Addresses{node1.Address}, addrs)
|
||||
})
|
||||
|
||||
t.Run("worker for address", func(t *testing.T) {
|
||||
node, err := s.WorkerForAddress(tx, "nodeaddr1")
|
||||
var expectedDBID *dax.DatabaseID = nil
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, dax.Address("nodeaddr1"), node.Address)
|
||||
require.Equal(t, ws.ID, node.ServiceID)
|
||||
require.Equal(t, expectedDBID, node.DatabaseID)
|
||||
})
|
||||
|
||||
}
|
||||
|
|
@ -4,24 +4,10 @@ import (
|
|||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/models"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
"github.com/gobuffalo/nulls"
|
||||
)
|
||||
|
||||
func NewWorkerServiceProviderService(log logger.Logger) *workerServiceProviderService {
|
||||
if log == nil {
|
||||
log = logger.NopLogger
|
||||
}
|
||||
return &workerServiceProviderService{
|
||||
log: log,
|
||||
}
|
||||
}
|
||||
|
||||
type workerServiceProviderService struct {
|
||||
log logger.Logger
|
||||
}
|
||||
|
||||
func (w *workerServiceProviderService) CreateWorkerServiceProvider(tx dax.Transaction, sp dax.WorkerServiceProvider) error {
|
||||
func (w *store) CreateWorkerServiceProvider(tx dax.Transaction, sp dax.WorkerServiceProvider) error {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
|
|
@ -39,7 +25,7 @@ func (w *workerServiceProviderService) CreateWorkerServiceProvider(tx dax.Transa
|
|||
return errors.Wrap(err, "inserting to DB")
|
||||
}
|
||||
|
||||
func (w *workerServiceProviderService) CreateWorkerService(tx dax.Transaction, srv dax.WorkerService) error {
|
||||
func (w *store) CreateWorkerService(tx dax.Transaction, srv dax.WorkerService) error {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
|
|
@ -64,7 +50,7 @@ func (w *workerServiceProviderService) CreateWorkerService(tx dax.Transaction, s
|
|||
return errors.Wrap(err, "inserting to DB")
|
||||
}
|
||||
|
||||
func (w *workerServiceProviderService) WorkerServiceProviders(tx dax.Transaction /*, future optional filters */) (dax.WorkerServiceProviders, error) {
|
||||
func (w *store) WorkerServiceProviders(tx dax.Transaction /*, future optional filters */) (dax.WorkerServiceProviders, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
|
|
@ -98,7 +84,7 @@ func toDaxWSP(wsp models.WorkerServiceProvider) dax.WorkerServiceProvider {
|
|||
}
|
||||
|
||||
// AssignFreeServiceToDatabase finds a WorkerService with the given service provider ID and
|
||||
func (w *workerServiceProviderService) AssignFreeServiceToDatabase(tx dax.Transaction, wspID dax.WorkerServiceProviderID, qdb *dax.QualifiedDatabase) (*dax.WorkerService, error) {
|
||||
func (w *store) AssignFreeServiceToDatabase(tx dax.Transaction, wspID dax.WorkerServiceProviderID, qdb *dax.QualifiedDatabase) (*dax.WorkerService, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
|
|
@ -122,7 +108,7 @@ func (w *workerServiceProviderService) AssignFreeServiceToDatabase(tx dax.Transa
|
|||
|
||||
}
|
||||
|
||||
func (w *workerServiceProviderService) WorkerServices(tx dax.Transaction, wspID dax.WorkerServiceProviderID) (dax.WorkerServices, error) {
|
||||
func (w *store) WorkerServices(tx dax.Transaction, wspID dax.WorkerServiceProviderID) (dax.WorkerServices, error) {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
|
|
|
|||
|
|
@ -34,8 +34,10 @@ func (w *workerJobService) WorkersJobs(tx dax.Transaction, roleType dax.RoleType
|
|||
|
||||
// First, get all workers for the database.
|
||||
workers := models.Workers{}
|
||||
sql := fmt.Sprintf("role_%s = true and database_id = ?", roleType)
|
||||
err := dt.C.Where(sql, qdbid.DatabaseID).Order("address asc").All(&workers)
|
||||
q := dt.C.Q().Eager()
|
||||
q = q.InnerJoin("services", "services.id = workers.service_id")
|
||||
q = q.Where(fmt.Sprintf("services.role_%s = true and services.database_id = ?", roleType), qdbid.DatabaseID)
|
||||
err := q.Order("address asc").All(&workers)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "getting workers")
|
||||
}
|
||||
|
|
@ -74,10 +76,9 @@ func (w *workerJobService) WorkerCount(tx dax.Transaction, roleType dax.RoleType
|
|||
if !ok {
|
||||
return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
worker := &models.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")
|
||||
var count int
|
||||
err := dt.C.RawQuery("SELECT count(*) from workers inner join services on workers.service_id = services.id where services.database_id = ?", qdbid.DatabaseID).First(&count)
|
||||
return count, errors.Wrap(err, "getting count")
|
||||
}
|
||||
|
||||
func (w *workerJobService) ListWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Addresses, error) {
|
||||
|
|
@ -101,33 +102,6 @@ func (w *workerJobService) ListWorkers(tx dax.Transaction, roleType dax.RoleType
|
|||
return ret, nil
|
||||
}
|
||||
|
||||
func (w *workerJobService) CreateWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
worker := &models.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) ReleaseWorkers(tx dax.Transaction, addrs ...dax.Address) error {
|
||||
dt, ok := tx.(*DaxTransaction)
|
||||
if !ok {
|
||||
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
|
||||
}
|
||||
|
||||
if len(addrs) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
err := dt.C.RawQuery("UPDATE workers set database_id = NULL where address in (?)", addrs).Exec()
|
||||
return errors.Wrap(err, "updating workers")
|
||||
}
|
||||
|
||||
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 {
|
||||
|
|
@ -149,7 +123,7 @@ func (w *workerJobService) AssignWorkerToJobs(tx dax.Transaction, roleType dax.R
|
|||
|
||||
// Assign jobs not in "jobs", and therefore didn't get updated by the
|
||||
// previous sql statement.
|
||||
toBeAssigned := jobsNotAssigned(job, jobs, roleType, worker)
|
||||
toBeAssigned := jobsNotAssigned(job, jobs, roleType, worker, qdbid.DatabaseID)
|
||||
|
||||
if err := dt.C.Create(toBeAssigned); err != nil {
|
||||
return errors.Wrap(err, "creating jobs")
|
||||
|
|
@ -158,7 +132,7 @@ func (w *workerJobService) AssignWorkerToJobs(tx dax.Transaction, roleType dax.R
|
|||
return nil
|
||||
}
|
||||
|
||||
func jobsNotAssigned(incomingJobs []dax.Job, assigned models.Jobs, roleType dax.RoleType, worker *models.Worker) (toBeAssigned models.Jobs) {
|
||||
func jobsNotAssigned(incomingJobs []dax.Job, assigned models.Jobs, roleType dax.RoleType, worker *models.Worker, dbid dax.DatabaseID) (toBeAssigned models.Jobs) {
|
||||
outer:
|
||||
for _, incJob := range incomingJobs {
|
||||
for _, assignedJob := range assigned {
|
||||
|
|
@ -170,7 +144,7 @@ outer:
|
|||
models.Job{
|
||||
Name: incJob,
|
||||
Role: roleType,
|
||||
DatabaseID: dax.DatabaseID(worker.DatabaseID.String),
|
||||
DatabaseID: dbid,
|
||||
Worker: worker,
|
||||
},
|
||||
)
|
||||
|
|
|
|||
|
|
@ -5,7 +5,6 @@ import (
|
|||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/models"
|
||||
"github.com/gobuffalo/nulls"
|
||||
"github.com/gofrs/uuid"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
|
@ -24,8 +23,7 @@ func TestJobsNotUpdated(t *testing.T) {
|
|||
toCreate := jobsNotAssigned(incJobs, created, dax.RoleTypeCompute, &models.Worker{
|
||||
ID: u2,
|
||||
RoleCompute: true,
|
||||
DatabaseID: nulls.NewString("dbid"),
|
||||
})
|
||||
}, "dbid")
|
||||
|
||||
require.Equal(t, 3, len(toCreate))
|
||||
|
||||
|
|
@ -33,8 +31,7 @@ func TestJobsNotUpdated(t *testing.T) {
|
|||
toCreate = jobsNotAssigned(incJobs, models.Jobs{}, dax.RoleTypeCompute, &models.Worker{
|
||||
ID: u2,
|
||||
RoleCompute: true,
|
||||
DatabaseID: nulls.NewString("dbid"),
|
||||
})
|
||||
}, "dbid")
|
||||
|
||||
require.Equal(t, 4, len(toCreate))
|
||||
}
|
||||
|
|
|
|||
54
dax/controller/store.go
Normal file
54
dax/controller/store.go
Normal file
|
|
@ -0,0 +1,54 @@
|
|||
package controller
|
||||
|
||||
import (
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
)
|
||||
|
||||
// Store wraps up (and will replace??) all the controller interfaces
|
||||
// for accessing storage. Since all the storage is in one database
|
||||
// under the hood and is interrelated, breaking it into separate
|
||||
// "services" is confusing the heck outta me.
|
||||
type Store interface {
|
||||
// CreateWorkerServiceProvider adds a WorkerServiceProvider (WSP) to
|
||||
// storage. Any WorkerService which registers must have a
|
||||
// WorkerServiceProviderID that matches an existing WSP, and the
|
||||
// Controller will ask for which WSPs are available when asked to
|
||||
// create a database. It will then ask one of the WSPs for a
|
||||
// WorkerService to assign to the database.
|
||||
CreateWorkerServiceProvider(tx dax.Transaction, sp dax.WorkerServiceProvider) error
|
||||
|
||||
// CreateWorkerService adds a WorkerService to storage. Generally
|
||||
// a WSP can create WorkerServices (which can create Workers)
|
||||
// before they are requested. Because the workers will register
|
||||
// themselves as soon as they come up, and they must be associated
|
||||
// with a WorkerService, the WSP registers all Services with the
|
||||
// Controller which stores the knowledge of their existence by
|
||||
// calling this method.
|
||||
CreateWorkerService(tx dax.Transaction, srv dax.WorkerService) error
|
||||
|
||||
WorkerServiceProviders(tx dax.Transaction /*, future optional filters */) (dax.WorkerServiceProviders, error)
|
||||
|
||||
// WorkerServices returns all worker services which came from the
|
||||
// WorkerServiceProvider with the given ID. If that ID is empty,
|
||||
// then all WorkerServices are returned.
|
||||
WorkerServices(tx dax.Transaction, wspID dax.WorkerServiceProviderID) (dax.WorkerServices, error)
|
||||
|
||||
WorkerService(tx dax.Transaction, dbid dax.DatabaseID) (dax.WorkerService, error)
|
||||
|
||||
AddWorker(tx dax.Transaction, node *dax.Node) (*dax.DatabaseID, error)
|
||||
WorkerCount(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) (int, error)
|
||||
WorkerCountDatabase(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (int, error)
|
||||
ListFreeJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (dax.Jobs, error)
|
||||
WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) ([]dax.WorkerInfo, error)
|
||||
AssignWorkerToJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, workerID string, jobs ...dax.Job) error
|
||||
ListWorkers(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) (dax.Addresses, error)
|
||||
|
||||
AssignFreeServiceToDatabase(tx dax.Transaction, wspID dax.WorkerServiceProviderID, qdb *dax.QualifiedDatabase) (*dax.WorkerService, error)
|
||||
|
||||
WorkerForAddress(tx dax.Transaction, addr dax.Address) (dax.Node, error)
|
||||
RemoveWorker(tx dax.Transaction, id dax.WorkerID) error
|
||||
|
||||
CreateFreeJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, jobs ...dax.Job) error
|
||||
DeleteJobsForTable(tx dax.Transaction, role dax.RoleType, qtid dax.QualifiedTableID) (InternalDiffs, error)
|
||||
WorkerJobs(tx dax.Transaction, role dax.RoleType, addr dax.Address) (dax.WorkerInfo, error)
|
||||
}
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package balancer
|
||||
package controller
|
||||
|
||||
import (
|
||||
"sort"
|
||||
13
dax/controller/worker_service_provider.go
Normal file
13
dax/controller/worker_service_provider.go
Normal file
|
|
@ -0,0 +1,13 @@
|
|||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
)
|
||||
|
||||
type WorkerServiceProvider interface {
|
||||
ClaimService(ctx context.Context, svc *dax.WorkerService) error
|
||||
UpdateService(ctx context.Context, svc *dax.WorkerService) error
|
||||
DropService(ctx context.Context, svc *dax.WorkerService) error
|
||||
}
|
||||
|
|
@ -24,4 +24,5 @@ create_table("worker_services") {
|
|||
|
||||
add_column("workers", "service_id", "string", {})
|
||||
add_foreign_key("workers", "service_id", {"worker_services": ["id"]}, {"on_delete": "cascade"})
|
||||
drop_column("workers", "database_id")
|
||||
|
||||
|
|
|
|||
|
|
@ -6,7 +6,6 @@ import (
|
|||
|
||||
"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"
|
||||
"github.com/gobuffalo/validate/v3/validators"
|
||||
|
|
@ -23,7 +22,6 @@ type Worker struct {
|
|||
RoleCompute bool `json:"role_compute" db:"role_compute"`
|
||||
RoleTranslate bool `json:"role_translate" db:"role_translate"`
|
||||
RoleQuery bool `json:"role_query" db:"role_query"`
|
||||
DatabaseID nulls.String `json:"database_id" db:"database_id"` // this can be empty which means the worker is unassigned. Probably get rid of this now that every worker is associated w/ a Service and every Service has an ID
|
||||
CreatedAt time.Time `json:"created_at" db:"created_at"`
|
||||
UpdatedAt time.Time `json:"updated_at" db:"updated_at"`
|
||||
}
|
||||
|
|
|
|||
|
|
@ -6,11 +6,18 @@ import (
|
|||
"strings"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/gofrs/uuid"
|
||||
)
|
||||
|
||||
type WorkerID uuid.UUID
|
||||
|
||||
// Node is used in API requests, like RegisterNode (before being assigned
|
||||
// roles).
|
||||
type Node struct {
|
||||
// ID is a unique identifier assigned by the metadata storage
|
||||
// layer (controller.Store).
|
||||
ID WorkerID `json:"id"`
|
||||
|
||||
// Address is the node's network address
|
||||
Address Address `json:"address"`
|
||||
|
||||
|
|
@ -19,6 +26,10 @@ type Node struct {
|
|||
// assigned to.
|
||||
ServiceID WorkerServiceID `json:"service_id"`
|
||||
|
||||
// DatabaseID may be nil if this Worker is not part of a service
|
||||
// that's been assigned to a database.
|
||||
DatabaseID *DatabaseID `json:"database_id"`
|
||||
|
||||
// RoleTypes allows a registering node to specify which role type(s) it is
|
||||
// capable of filling. The controller will not assign a role to this node
|
||||
// with a type not included in RoleTypes.
|
||||
|
|
|
|||
109
dax/worker_service_provider/http/client.go
Normal file
109
dax/worker_service_provider/http/client.go
Normal file
|
|
@ -0,0 +1,109 @@
|
|||
package http
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
)
|
||||
|
||||
const (
|
||||
defaultScheme = "http"
|
||||
)
|
||||
|
||||
var _ controller.WorkerServiceProvider = (*Client)(nil)
|
||||
|
||||
// Client is an HTTP client that operates on the WorkerServiceProvider.
|
||||
type Client struct {
|
||||
address dax.Address
|
||||
httpClient *http.Client
|
||||
logger logger.Logger
|
||||
}
|
||||
|
||||
// NewClient returns a new instance of Client.
|
||||
func NewClient(address dax.Address, logger logger.Logger) controller.WorkerServiceProvider {
|
||||
return &Client{
|
||||
address: address,
|
||||
logger: logger,
|
||||
httpClient: &http.Client{
|
||||
Timeout: time.Second * 30,
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
// Health returns true if the client address returns status OK at its /health
|
||||
// endpoint.
|
||||
func (c *Client) Health() bool {
|
||||
url := fmt.Sprintf("%s/health", c.address.WithScheme(defaultScheme))
|
||||
|
||||
if resp, err := c.httpClient.Get(url); err != nil {
|
||||
return false
|
||||
} else if resp.StatusCode != http.StatusOK {
|
||||
defer resp.Body.Close()
|
||||
return false
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
func (c *Client) ClaimService(ctx context.Context, svc *dax.WorkerService) error {
|
||||
postBody, err := json.Marshal(svc)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "marshalling post request")
|
||||
}
|
||||
buf := bytes.NewBuffer(postBody)
|
||||
|
||||
resp, err := c.httpClient.Post(fmt.Sprintf("%s/claim", c.address), "application/json", buf)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "making post /claim request")
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return errors.Wrapf(errors.UnmarshalJSON(resp.Body), "status code: %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
}
|
||||
|
||||
func (c *Client) UpdateService(ctx context.Context, svc *dax.WorkerService) error {
|
||||
postBody, err := json.Marshal(svc)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "marshalling post request")
|
||||
}
|
||||
buf := bytes.NewBuffer(postBody)
|
||||
|
||||
resp, err := c.httpClient.Post(fmt.Sprintf("%s/update", c.address), "application/json", buf)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "making post /update request")
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return errors.Wrapf(errors.UnmarshalJSON(resp.Body), "status code: %d", resp.StatusCode)
|
||||
}
|
||||
return nil
|
||||
|
||||
}
|
||||
|
||||
func (c *Client) DropService(ctx context.Context, svc *dax.WorkerService) error {
|
||||
postBody, err := json.Marshal(svc)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "marshalling post request")
|
||||
}
|
||||
buf := bytes.NewBuffer(postBody)
|
||||
|
||||
resp, err := c.httpClient.Post(fmt.Sprintf("%s/drop", c.address), "application/json", buf)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "making post /drop request")
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return errors.Wrapf(errors.UnmarshalJSON(resp.Body), "status code: %d", resp.StatusCode)
|
||||
}
|
||||
return nil
|
||||
|
||||
}
|
||||
|
|
@ -47,8 +47,8 @@ func (s *server) postClaim(w http.ResponseWriter, r *http.Request) {
|
|||
body := r.Body
|
||||
defer body.Close()
|
||||
|
||||
req := dax.WorkerService{}
|
||||
if err := json.NewDecoder(body).Decode(&req); err != nil {
|
||||
req := &dax.WorkerService{}
|
||||
if err := json.NewDecoder(body).Decode(req); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
|
@ -64,8 +64,8 @@ func (s *server) postUpdate(w http.ResponseWriter, r *http.Request) {
|
|||
body := r.Body
|
||||
defer body.Close()
|
||||
|
||||
req := dax.WorkerService{}
|
||||
if err := json.NewDecoder(body).Decode(&req); err != nil {
|
||||
req := &dax.WorkerService{}
|
||||
if err := json.NewDecoder(body).Decode(req); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
|
@ -81,8 +81,8 @@ func (s *server) postDrop(w http.ResponseWriter, r *http.Request) {
|
|||
body := r.Body
|
||||
defer body.Close()
|
||||
|
||||
req := dax.WorkerService{}
|
||||
if err := json.NewDecoder(body).Decode(&req); err != nil {
|
||||
req := &dax.WorkerService{}
|
||||
if err := json.NewDecoder(body).Decode(req); err != nil {
|
||||
http.Error(w, err.Error(), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
|
|
|||
|
|
@ -7,6 +7,7 @@ import (
|
|||
|
||||
"github.com/featurebasedb/featurebase/v3/dax"
|
||||
computersvc "github.com/featurebasedb/featurebase/v3/dax/computer/service"
|
||||
"github.com/featurebasedb/featurebase/v3/dax/controller"
|
||||
controllerclient "github.com/featurebasedb/featurebase/v3/dax/controller/client"
|
||||
"github.com/featurebasedb/featurebase/v3/errors"
|
||||
"github.com/featurebasedb/featurebase/v3/logger"
|
||||
|
|
@ -21,6 +22,8 @@ func NewConfig() Config {
|
|||
}
|
||||
}
|
||||
|
||||
var _ controller.WorkerServiceProvider = (*WSP)(nil)
|
||||
|
||||
type Config struct {
|
||||
ID string `toml:"id"`
|
||||
Address *fbnet.URI `toml:"-"`
|
||||
|
|
@ -149,7 +152,7 @@ func (w *WSP) addService() error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (w *WSP) ClaimService(ctx context.Context, svc dax.WorkerService) error {
|
||||
func (w *WSP) ClaimService(ctx context.Context, svc *dax.WorkerService) error {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
|
||||
|
|
@ -166,7 +169,7 @@ func (w *WSP) ClaimService(ctx context.Context, svc dax.WorkerService) error {
|
|||
}
|
||||
isClaim := mysvc.DatabaseID == ""
|
||||
|
||||
mysvc.WorkerService = svc
|
||||
mysvc.WorkerService = *svc
|
||||
w.services[svc.ID] = mysvc
|
||||
|
||||
// claiming a service might update min/max workers, so we scale.
|
||||
|
|
@ -235,11 +238,11 @@ func (w *WSP) removeWorker(mysvc *workerService) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (w *WSP) UpdateService(ctx context.Context, svc dax.WorkerService) error {
|
||||
func (w *WSP) UpdateService(ctx context.Context, svc *dax.WorkerService) error {
|
||||
return w.ClaimService(ctx, svc)
|
||||
}
|
||||
|
||||
func (w *WSP) DropService(ctx context.Context, svc dax.WorkerService) error {
|
||||
func (w *WSP) DropService(ctx context.Context, svc *dax.WorkerService) error {
|
||||
w.mu.Lock()
|
||||
defer w.mu.Unlock()
|
||||
|
||||
|
|
@ -24,6 +24,7 @@ type Jobber interface {
|
|||
// WorkerInfo represents a Worker and the Jobs to which it has been assigned.
|
||||
type WorkerInfo struct {
|
||||
Address Address
|
||||
ID WorkerID // internal identifier used in SQL database, used internally as an optimization to avoid extra queries
|
||||
Jobs []Job
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue