diff --git a/dax/controller/balancer.go b/dax/controller/balancer.go index e22a266e5..c91afd3c0 100644 --- a/dax/controller/balancer.go +++ b/dax/controller/balancer.go @@ -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 +} diff --git a/dax/controller/balancer/balancer.go b/dax/controller/balancer/balancer.go index 77341608d..6189359ec 100644 --- a/dax/controller/balancer/balancer.go +++ b/dax/controller/balancer/balancer.go @@ -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 { diff --git a/dax/controller/balancer/sqldb_init_test.go b/dax/controller/balancer/sqldb_init_test.go index 0f72ec05f..ec4683977 100644 --- a/dax/controller/balancer/sqldb_init_test.go +++ b/dax/controller/balancer/sqldb_init_test.go @@ -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 diff --git a/dax/controller/balancer/worker_job_test.go b/dax/controller/balancer/worker_job_test.go index 8e62d6854..2a43779c8 100644 --- a/dax/controller/balancer/worker_job_test.go +++ b/dax/controller/balancer/worker_job_test.go @@ -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 diff --git a/dax/controller/controller.go b/dax/controller/controller.go index d22b6efd6..b0db83a7b 100644 --- a/dax/controller/controller.go +++ b/dax/controller/controller.go @@ -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 { diff --git a/dax/controller/service/controller.go b/dax/controller/service/controller.go index f3a11c4e4..ea262add4 100644 --- a/dax/controller/service/controller.go +++ b/dax/controller/service/controller.go @@ -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, diff --git a/dax/controller/sqldb/balancer.go b/dax/controller/sqldb/balancer.go index ccf1a04b4..4c8aee363 100644 --- a/dax/controller/sqldb/balancer.go +++ b/dax/controller/sqldb/balancer.go @@ -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) } diff --git a/dax/controller/sqldb/freeworker.go b/dax/controller/sqldb/freeworker.go index 219d6b3d9..4ddc0a2e8 100644 --- a/dax/controller/sqldb/freeworker.go +++ b/dax/controller/sqldb/freeworker.go @@ -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") diff --git a/dax/controller/sqldb/sqldb_init_test.go b/dax/controller/sqldb/sqldb_init_test.go new file mode 100644 index 000000000..e368b9f36 --- /dev/null +++ b/dax/controller/sqldb/sqldb_init_test.go @@ -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 +} diff --git a/dax/controller/sqldb/store.go b/dax/controller/sqldb/store.go new file mode 100644 index 000000000..1b929fca8 --- /dev/null +++ b/dax/controller/sqldb/store.go @@ -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") +} diff --git a/dax/controller/sqldb/store_test.go b/dax/controller/sqldb/store_test.go new file mode 100644 index 000000000..370af80b1 --- /dev/null +++ b/dax/controller/sqldb/store_test.go @@ -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) + }) + +} diff --git a/dax/controller/sqldb/worker_service_provider.go b/dax/controller/sqldb/worker_service_provider.go index bebed93c7..94873b138 100644 --- a/dax/controller/sqldb/worker_service_provider.go +++ b/dax/controller/sqldb/worker_service_provider.go @@ -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") diff --git a/dax/controller/sqldb/workerjob.go b/dax/controller/sqldb/workerjob.go index 5e4eb15e4..2fdd8cf58 100644 --- a/dax/controller/sqldb/workerjob.go +++ b/dax/controller/sqldb/workerjob.go @@ -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, }, ) diff --git a/dax/controller/sqldb/workerjob_test.go b/dax/controller/sqldb/workerjob_test.go index f583d7d4c..a07e5f4dd 100644 --- a/dax/controller/sqldb/workerjob_test.go +++ b/dax/controller/sqldb/workerjob_test.go @@ -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)) } diff --git a/dax/controller/store.go b/dax/controller/store.go new file mode 100644 index 000000000..52f12934c --- /dev/null +++ b/dax/controller/store.go @@ -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) +} diff --git a/dax/controller/balancer/types.go b/dax/controller/types.go similarity index 98% rename from dax/controller/balancer/types.go rename to dax/controller/types.go index 90fcafa6e..c527f1ac0 100644 --- a/dax/controller/balancer/types.go +++ b/dax/controller/types.go @@ -1,4 +1,4 @@ -package balancer +package controller import ( "sort" diff --git a/dax/controller/worker_service_provider.go b/dax/controller/worker_service_provider.go new file mode 100644 index 000000000..267d90498 --- /dev/null +++ b/dax/controller/worker_service_provider.go @@ -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 +} diff --git a/dax/migrations/004_serviceprovider.up.fizz b/dax/migrations/004_serviceprovider.up.fizz index b0e7fedc2..8cd049a83 100644 --- a/dax/migrations/004_serviceprovider.up.fizz +++ b/dax/migrations/004_serviceprovider.up.fizz @@ -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") diff --git a/dax/models/worker.go b/dax/models/worker.go index 556b78fae..5f1009646 100644 --- a/dax/models/worker.go +++ b/dax/models/worker.go @@ -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"` } diff --git a/dax/worker.go b/dax/worker.go index bf4b14445..e7f49daa5 100644 --- a/dax/worker.go +++ b/dax/worker.go @@ -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. diff --git a/dax/worker_service_provider/http/client.go b/dax/worker_service_provider/http/client.go new file mode 100644 index 000000000..50bb42b15 --- /dev/null +++ b/dax/worker_service_provider/http/client.go @@ -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 + +} diff --git a/dax/worker_service_provider/http/handler.go b/dax/worker_service_provider/http/handler.go index 3399c6e98..fe7324024 100644 --- a/dax/worker_service_provider/http/handler.go +++ b/dax/worker_service_provider/http/handler.go @@ -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 } diff --git a/dax/worker_service_provider/worker_service_provider.go b/dax/worker_service_provider/wsp.go similarity index 95% rename from dax/worker_service_provider/worker_service_provider.go rename to dax/worker_service_provider/wsp.go index 5da751a18..74ce3d7df 100644 --- a/dax/worker_service_provider/worker_service_provider.go +++ b/dax/worker_service_provider/wsp.go @@ -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() diff --git a/dax/workerjob.go b/dax/workerjob.go index dae2f0486..6761162ea 100644 --- a/dax/workerjob.go +++ b/dax/workerjob.go @@ -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 }