mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
819 lines
20 KiB
Go
819 lines
20 KiB
Go
// Package boltdb contains the boltdb implementation of the Balancer interface.
|
|
package boltdb
|
|
|
|
import (
|
|
"bytes"
|
|
"encoding/json"
|
|
"fmt"
|
|
"strings"
|
|
|
|
"github.com/featurebasedb/featurebase/v3/dax"
|
|
"github.com/featurebasedb/featurebase/v3/dax/boltdb"
|
|
balancer "github.com/featurebasedb/featurebase/v3/dax/controller/balancer"
|
|
"github.com/featurebasedb/featurebase/v3/dax/controller/schemar"
|
|
"github.com/featurebasedb/featurebase/v3/errors"
|
|
"github.com/featurebasedb/featurebase/v3/logger"
|
|
)
|
|
|
|
var (
|
|
bucketBalancer = boltdb.Bucket("balancer")
|
|
)
|
|
|
|
// BalancerBuckets defines the buckets used by this package. It can be
|
|
// called during setup to create the buckets ahead of time.
|
|
var BalancerBuckets []boltdb.Bucket = []boltdb.Bucket{
|
|
bucketBalancer,
|
|
}
|
|
|
|
// NewBalancer returns a new instance of controller.Balancer.
|
|
func NewBalancer(db *boltdb.DB, schemar schemar.Schemar, logger logger.Logger) *balancer.Balancer {
|
|
fjs := newFreeJobService(db)
|
|
wjs := newWorkerJobService(db, logger)
|
|
fws := newFreeWorkerService(db)
|
|
ns := NewNodeService(db, logger)
|
|
|
|
return balancer.New(ns, fjs, wjs, fws, schemar, logger)
|
|
}
|
|
|
|
// Ensure type implements interface.
|
|
var _ balancer.WorkerJobService = (*workerJobService)(nil)
|
|
|
|
type workerJobService struct {
|
|
db *boltdb.DB
|
|
logger logger.Logger
|
|
}
|
|
|
|
func newWorkerJobService(db *boltdb.DB, logger logger.Logger) *workerJobService {
|
|
return &workerJobService{
|
|
db: db,
|
|
logger: logger,
|
|
}
|
|
}
|
|
|
|
func (w *workerJobService) WorkersJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerInfo, error) {
|
|
workerInfos, err := w.getWorkerInfos(tx, roleType, qdbid)
|
|
if err != nil {
|
|
return nil, errors.Wrapf(err, "getting worker infos: %s", roleType)
|
|
}
|
|
|
|
return workerInfos, nil
|
|
}
|
|
|
|
func (w *workerJobService) WorkerCount(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (int, error) {
|
|
workers, err := w.getWorkers(tx, roleType, qdbid)
|
|
if err != nil {
|
|
return 0, errors.Wrapf(err, "getting workers: %s", roleType)
|
|
}
|
|
|
|
return len(workers), nil
|
|
}
|
|
|
|
func (w *workerJobService) ListWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Addresses, error) {
|
|
return w.getWorkers(tx, roleType, qdbid)
|
|
}
|
|
|
|
func (w *workerJobService) getWorkers(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Addresses, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
c := txx.Bucket(bucketBalancer).Cursor()
|
|
|
|
// Deserialize rows into Worker objects.
|
|
addrs := make(dax.Addresses, 0)
|
|
|
|
prefix := []byte(fmt.Sprintf(prefixFmtWorkersDB, roleType, qdbid.Key()))
|
|
for k, v := c.Seek(prefix); k != nil && bytes.HasPrefix(k, prefix); k, v = c.Next() {
|
|
if v == nil {
|
|
w.logger.Printf("nil value for key: %s", k)
|
|
continue
|
|
}
|
|
|
|
addr, err := keyWorker(k)
|
|
if err != nil {
|
|
return nil, errors.Wrapf(err, "getting worker from key: %s", k)
|
|
}
|
|
|
|
addrs = append(addrs, addr)
|
|
}
|
|
|
|
return addrs, nil
|
|
}
|
|
|
|
func (w *workerJobService) getWorkerInfos(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.WorkerInfos, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
c := txx.Bucket(bucketBalancer).Cursor()
|
|
|
|
// Deserialize rows into WorkerInfo objects.
|
|
workerInfos := make(dax.WorkerInfos, 0)
|
|
|
|
var prefix []byte
|
|
empty := dax.QualifiedDatabaseID{}
|
|
if roleType == "" && qdbid == empty {
|
|
prefix = []byte("workers/role/")
|
|
} else {
|
|
prefix = []byte(fmt.Sprintf(prefixFmtWorkersDB, roleType, qdbid.Key()))
|
|
}
|
|
for k, v := c.Seek(prefix); k != nil && bytes.HasPrefix(k, prefix); k, v = c.Next() {
|
|
addr, err := keyWorker(k)
|
|
if err != nil {
|
|
return nil, errors.Wrapf(err, "getting worker from key: %s", k)
|
|
}
|
|
|
|
jobs := dax.NewSet[dax.Job]()
|
|
if v != nil {
|
|
jobs, err = decodeJobSet(v)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "decoding job set")
|
|
}
|
|
}
|
|
|
|
workerInfo := dax.WorkerInfo{
|
|
Address: addr,
|
|
Jobs: jobs.Sorted(),
|
|
}
|
|
|
|
workerInfos = append(workerInfos, workerInfo)
|
|
}
|
|
|
|
return workerInfos, nil
|
|
}
|
|
|
|
func (w *workerJobService) CreateWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
// If this worker already exists, don't do anything.
|
|
wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr))
|
|
if wrkr != nil {
|
|
return nil
|
|
}
|
|
|
|
val := []byte("[]")
|
|
if err := bkt.Put(workerDBKey(roleType, qdbid, addr), val); err != nil {
|
|
return errors.Wrapf(err, "putting db worker: %s, %s", qdbid, addr)
|
|
}
|
|
|
|
if err := bkt.Put(workerAssignedKey(addr), []byte(qdbid.Key())); err != nil {
|
|
return errors.Wrapf(err, "putting assigned worker: %s, %s", qdbid, addr)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (w *workerJobService) DeleteWorker(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
if err := bkt.Delete(workerDBKey(roleType, qdbid, addr)); err != nil {
|
|
return errors.Wrapf(err, "deleting node key: %s", workerDBKey(roleType, qdbid, addr))
|
|
}
|
|
|
|
if err := bkt.Delete(workerAssignedKey(addr)); err != nil {
|
|
return errors.Wrapf(err, "deleting assigned worker: %s", workerAssignedKey(addr))
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (w *workerJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, jobs ...dax.Job) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
jobset := dax.NewSet[dax.Job]()
|
|
var err error
|
|
|
|
// get worker
|
|
wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr))
|
|
if wrkr != nil {
|
|
jobset, err = decodeJobSet(wrkr)
|
|
if err != nil {
|
|
return errors.Wrap(err, "decoding job set")
|
|
}
|
|
}
|
|
|
|
for _, job := range jobs {
|
|
jobset.Add(job)
|
|
}
|
|
val, err := encodeJobSet(jobset)
|
|
if err != nil {
|
|
return errors.Wrap(err, "encoding job set")
|
|
}
|
|
|
|
if err := bkt.Put(workerDBKey(roleType, qdbid, addr), val); err != nil {
|
|
return errors.Wrap(err, "putting worker")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (w *workerJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job dax.Job) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
// get worker
|
|
wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr))
|
|
if wrkr == nil {
|
|
return nil
|
|
}
|
|
|
|
jobset, err := decodeJobSet(wrkr)
|
|
if err != nil {
|
|
return errors.Wrap(err, "decoding job set")
|
|
}
|
|
if !jobset.Contains(job) {
|
|
return nil
|
|
}
|
|
|
|
jobset.Remove(job)
|
|
val, err := encodeJobSet(jobset)
|
|
if err != nil {
|
|
return errors.Wrap(err, "encoding job set")
|
|
}
|
|
|
|
if err := bkt.Put(workerDBKey(roleType, qdbid, addr), val); err != nil {
|
|
return errors.Wrap(err, "putting worker")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (w *workerJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) (balancer.InternalDiffs, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
qdbid := qtid.QualifiedDatabaseID
|
|
prefix := string(qtid.Key())
|
|
|
|
workers, err := w.getWorkers(tx, roleType, qdbid)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "getting workers")
|
|
}
|
|
|
|
idiffs := balancer.NewInternalDiffs()
|
|
for _, worker := range workers {
|
|
// get worker
|
|
wrkr := bkt.Get(workerDBKey(roleType, qdbid, worker))
|
|
if wrkr == nil {
|
|
panic("didn't find worker that should... definitely exist")
|
|
}
|
|
jobset, err := decodeJobSet(wrkr)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "decoding job set")
|
|
}
|
|
|
|
jobs := jobset.RemoveByPrefix(prefix)
|
|
for _, job := range jobs {
|
|
idiffs.Removed(worker, job)
|
|
}
|
|
val, err := encodeJobSet(jobset)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "encoding job set")
|
|
}
|
|
|
|
if err := bkt.Put(workerDBKey(roleType, qdbid, worker), val); err != nil {
|
|
return nil, errors.Wrap(err, "putting worker")
|
|
}
|
|
|
|
}
|
|
|
|
return idiffs, nil
|
|
}
|
|
|
|
func (w *workerJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) (dax.Jobs, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
jobset := dax.NewSet[dax.Job]()
|
|
var err error
|
|
|
|
// get worker
|
|
wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr))
|
|
if wrkr != nil {
|
|
jobset, err = decodeJobSet(wrkr)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "decoding job set")
|
|
}
|
|
}
|
|
|
|
return jobset.Sorted(), nil
|
|
}
|
|
|
|
func (w *workerJobService) JobCounts(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addrs ...dax.Address) (map[dax.Address]int, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
m := make(map[dax.Address]int)
|
|
|
|
for _, addr := range addrs {
|
|
jobset := dax.NewSet[dax.Job]()
|
|
var err error
|
|
|
|
// get worker
|
|
wrkr := bkt.Get(workerDBKey(roleType, qdbid, addr))
|
|
if wrkr != nil {
|
|
jobset, err = decodeJobSet(wrkr)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "decoding job set")
|
|
}
|
|
}
|
|
|
|
m[addr] = len(jobset)
|
|
}
|
|
|
|
return m, nil
|
|
}
|
|
|
|
func (w *workerJobService) DatabaseForWorker(tx dax.Transaction, addr dax.Address) dax.DatabaseKey {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return "" // TODO(tlt): return error here?
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return ""
|
|
}
|
|
|
|
wrkr := bkt.Get(workerAssignedKey(addr))
|
|
|
|
return dax.DatabaseKey(wrkr)
|
|
}
|
|
|
|
// encodeJobSet encode the jobSet into a JSON array of strings.
|
|
func encodeJobSet(jobSet dax.Set[dax.Job]) ([]byte, error) {
|
|
arr := jobSet.Sorted()
|
|
b, err := json.Marshal(arr)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "marshalling json")
|
|
}
|
|
return b, nil
|
|
}
|
|
|
|
// decodeJobSet decode the string (a JSON array of strings) into jobSet.
|
|
func decodeJobSet(v []byte) (dax.Set[dax.Job], error) {
|
|
var arr []string
|
|
err := json.Unmarshal(v, &arr)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "unmarshalling json")
|
|
}
|
|
|
|
js := dax.NewSet[dax.Job]()
|
|
for _, s := range arr {
|
|
js.Add(dax.Job(s))
|
|
}
|
|
|
|
return js, nil
|
|
}
|
|
|
|
// encodeWorkerSet encode the workerSet into a JSON array of strings.
|
|
func encodeWorkerSet(workerSet dax.Set[dax.Address]) ([]byte, error) {
|
|
arr := workerSet.Sorted()
|
|
b, err := json.Marshal(arr)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "marshalling json")
|
|
}
|
|
return b, nil
|
|
}
|
|
|
|
// decodeWorkerSet decode the string (a JSON array of strings) into workerSet.
|
|
func decodeWorkerSet(v []byte) (dax.Set[dax.Address], error) {
|
|
var arr []string
|
|
err := json.Unmarshal(v, &arr)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "unmarshalling json")
|
|
}
|
|
|
|
ws := dax.NewSet[dax.Address]()
|
|
for _, s := range arr {
|
|
ws.Add(dax.Address(s))
|
|
}
|
|
|
|
return ws, nil
|
|
}
|
|
|
|
// Ensure type implements interface.
|
|
var _ balancer.FreeJobService = (*freeJobService)(nil)
|
|
|
|
type freeJobService struct {
|
|
db *boltdb.DB
|
|
}
|
|
|
|
func newFreeJobService(db *boltdb.DB) *freeJobService {
|
|
return &freeJobService{
|
|
db: db,
|
|
}
|
|
}
|
|
|
|
func (f *freeJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs ...dax.Job) error {
|
|
return f.MergeJobs(tx, roleType, qdbid, jobs)
|
|
}
|
|
|
|
func (f *freeJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job dax.Job) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
// get free jobs
|
|
fjs := bkt.Get(freeJobKey(roleType, qdbid))
|
|
if fjs == nil {
|
|
return nil
|
|
}
|
|
|
|
jobset, err := decodeJobSet(fjs)
|
|
if err != nil {
|
|
return errors.Wrap(err, "decoding job set")
|
|
}
|
|
if !jobset.Contains(job) {
|
|
return nil
|
|
}
|
|
|
|
jobset.Remove(job)
|
|
val, err := encodeJobSet(jobset)
|
|
if err != nil {
|
|
return errors.Wrap(err, "encoding job set")
|
|
}
|
|
|
|
if err := bkt.Put(freeJobKey(roleType, qdbid), val); err != nil {
|
|
return errors.Wrap(err, "putting free job")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (f *freeJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
qdbid := qtid.QualifiedDatabaseID
|
|
prefix := string(qtid.Key())
|
|
|
|
// get free jobs
|
|
fjs := bkt.Get(freeJobKey(roleType, qdbid))
|
|
if fjs == nil {
|
|
return nil
|
|
}
|
|
|
|
jobset, err := decodeJobSet(fjs)
|
|
if err != nil {
|
|
return errors.Wrap(err, "decoding job set")
|
|
}
|
|
|
|
jobset.RemoveByPrefix(prefix)
|
|
val, err := encodeJobSet(jobset)
|
|
if err != nil {
|
|
return errors.Wrap(err, "encoding job set")
|
|
}
|
|
|
|
if err := bkt.Put(freeJobKey(roleType, qdbid), val); err != nil {
|
|
return errors.Wrap(err, "putting free job")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (f *freeJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (dax.Jobs, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
jobset := dax.NewSet[dax.Job]()
|
|
var err error
|
|
|
|
// get free jobs
|
|
fjs := bkt.Get(freeJobKey(roleType, qdbid))
|
|
if fjs != nil {
|
|
jobset, err = decodeJobSet(fjs)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "decoding job set")
|
|
}
|
|
}
|
|
|
|
return jobset.Sorted(), nil
|
|
}
|
|
|
|
func (f *freeJobService) MergeJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs dax.Jobs) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
jobset := dax.NewSet[dax.Job]()
|
|
var err error
|
|
|
|
// get free jobs
|
|
fjs := bkt.Get(freeJobKey(roleType, qdbid))
|
|
if fjs != nil {
|
|
jobset, err = decodeJobSet(fjs)
|
|
if err != nil {
|
|
return errors.Wrap(err, "decoding job set")
|
|
}
|
|
}
|
|
|
|
for _, j := range jobs {
|
|
jobset.Add(j)
|
|
}
|
|
val, err := encodeJobSet(jobset)
|
|
if err != nil {
|
|
return errors.Wrap(err, "encoding job set")
|
|
}
|
|
|
|
if err := bkt.Put(freeJobKey(roleType, qdbid), val); err != nil {
|
|
return errors.Wrap(err, "putting free job")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
//////////////////////////////////////////////////////
|
|
|
|
// Ensure type implements interface.
|
|
var _ balancer.FreeWorkerService = (*freeWorkerService)(nil)
|
|
|
|
type freeWorkerService struct {
|
|
db *boltdb.DB
|
|
}
|
|
|
|
func newFreeWorkerService(db *boltdb.DB) *freeWorkerService {
|
|
return &freeWorkerService{
|
|
db: db,
|
|
}
|
|
}
|
|
|
|
func (f *freeWorkerService) AddWorkers(tx dax.Transaction, roleType dax.RoleType, addres ...dax.Address) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
workerset := dax.NewSet[dax.Address]()
|
|
var err error
|
|
|
|
// get free workers
|
|
fws := bkt.Get(freeWorkerKey(roleType))
|
|
if fws != nil {
|
|
workerset, err = decodeWorkerSet(fws)
|
|
if err != nil {
|
|
return errors.Wrap(err, "decoding worker set")
|
|
}
|
|
}
|
|
|
|
for _, w := range addres {
|
|
workerset.Add(w)
|
|
}
|
|
val, err := encodeWorkerSet(workerset)
|
|
if err != nil {
|
|
return errors.Wrap(err, "encoding worker set")
|
|
}
|
|
|
|
if err := bkt.Put(freeWorkerKey(roleType), val); err != nil {
|
|
return errors.Wrap(err, "putting free worker")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (f *freeWorkerService) RemoveWorker(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) error {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
workers, err := f.ListWorkers(tx, roleType)
|
|
if err != nil {
|
|
return errors.Wrap(err, "listing free workers")
|
|
}
|
|
|
|
// Create a workerset containing the free workers which remain after
|
|
// removing num workers.
|
|
workerset := dax.NewSet[dax.Address]()
|
|
for _, w := range workers {
|
|
workerset.Add(w)
|
|
}
|
|
|
|
if !workerset.Contains(addr) {
|
|
return nil
|
|
}
|
|
|
|
// Remove the worker.
|
|
workerset.Remove(addr)
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
val, err := encodeWorkerSet(workerset)
|
|
if err != nil {
|
|
return errors.Wrap(err, "encoding worker set")
|
|
}
|
|
|
|
if err := bkt.Put(freeWorkerKey(roleType), val); err != nil {
|
|
return errors.Wrap(err, "putting free worker")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (f *freeWorkerService) PopWorkers(tx dax.Transaction, roleType dax.RoleType, num int) ([]dax.Address, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
workers, err := f.ListWorkers(tx, roleType)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "listing free workers")
|
|
}
|
|
|
|
if len(workers) < num {
|
|
return nil, errors.Errorf("not enough free workers to pop: wanted %d, have: %d", num, len(workers))
|
|
}
|
|
|
|
// Get num workers from the list.
|
|
workersToAssign := workers[0:num]
|
|
|
|
// Create a workerset containing the free workers which remain after
|
|
// removing num workers.
|
|
workerset := dax.NewSet[dax.Address]()
|
|
for _, worker := range workers[num:] {
|
|
workerset.Add(worker)
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
val, err := encodeWorkerSet(workerset)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "encoding worker set")
|
|
}
|
|
|
|
if err := bkt.Put(freeWorkerKey(roleType), val); err != nil {
|
|
return nil, errors.Wrap(err, "putting free worker")
|
|
}
|
|
|
|
return workersToAssign, nil
|
|
}
|
|
|
|
func (f *freeWorkerService) ListWorkers(tx dax.Transaction, roleType dax.RoleType) (dax.Addresses, error) {
|
|
txx, ok := tx.(*boltdb.Tx)
|
|
if !ok {
|
|
return nil, dax.NewErrInvalidTransaction()
|
|
}
|
|
|
|
bkt := txx.Bucket(bucketBalancer)
|
|
if bkt == nil {
|
|
return nil, errors.Errorf(boltdb.ErrFmtBucketNotFound, bucketBalancer)
|
|
}
|
|
|
|
workerset := dax.NewSet[dax.Address]()
|
|
var err error
|
|
|
|
// get free workers
|
|
fws := bkt.Get(freeWorkerKey(roleType))
|
|
if fws != nil {
|
|
workerset, err = decodeWorkerSet(fws)
|
|
if err != nil {
|
|
return nil, errors.Wrap(err, "decoding worker set")
|
|
}
|
|
}
|
|
|
|
return workerset.Sorted(), nil
|
|
}
|
|
|
|
//////////////////////////////////////////////////////
|
|
|
|
const (
|
|
prefixFmtWorkersDB = "workers/role/%s/db/%s/" // %s - role, dbKey
|
|
prefixFmtWorkersAssigned = "workers/assigned/"
|
|
|
|
prefixFmtFreeJobs = "freejobs/role/%s/db/%s" // %s - role, dbKey
|
|
prefixFmtFreeWorkers = "freeworkers/role/%s" // %s - role
|
|
)
|
|
|
|
// workerDBKey returns a key based on worker.
|
|
//
|
|
// Format: workers/role/[role]/db/[dbKey]/[worker] = [job1, job2, ...]
|
|
func workerDBKey(roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) []byte {
|
|
key := fmt.Sprintf(prefixFmtWorkersDB+"%s", roleType, qdbid.Key(), addr)
|
|
return []byte(key)
|
|
}
|
|
|
|
// workerAssignedKey returns a key based on worker.
|
|
//
|
|
// Format: workers/assigned/[worker] = dbKey
|
|
func workerAssignedKey(addr dax.Address) []byte {
|
|
key := fmt.Sprintf(prefixFmtWorkersAssigned+"%s", addr)
|
|
return []byte(key)
|
|
}
|
|
|
|
// keyWorker gets the worker out of the key.
|
|
func keyWorker(key []byte) (dax.Address, error) {
|
|
parts := strings.SplitN(string(key), "/", 6)
|
|
if len(parts) != 6 {
|
|
return "", errors.New(errors.ErrUncoded, "worker key format expected: `workers/role/[role]/db/[db]/worker`")
|
|
}
|
|
|
|
return dax.Address(parts[5]), nil
|
|
}
|
|
|
|
// freeJobKey returns a key for all freeJobs.
|
|
//
|
|
// Format: freejobs/role/[role]/db/[dbKey] = [job1, job2, ...]
|
|
func freeJobKey(roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) []byte {
|
|
key := fmt.Sprintf(prefixFmtFreeJobs, roleType, qdbid.Key())
|
|
return []byte(key)
|
|
}
|
|
|
|
// freeWorkerKey returns a key for all freeWorkers.
|
|
//
|
|
// Format: freeworkers/role/[role] = [worker1, worker2, ...]
|
|
func freeWorkerKey(roleType dax.RoleType) []byte {
|
|
key := fmt.Sprintf(prefixFmtFreeWorkers, roleType)
|
|
return []byte(key)
|
|
}
|