it compiles!

This commit is contained in:
Matthew Jaffee 2023-04-08 08:16:55 -05:00
parent 6f5b7e332e
commit 09dd745bdd
17 changed files with 174 additions and 1063 deletions

View file

@ -42,7 +42,6 @@ func BuildDAXFlags(cmd *cobra.Command, srv *server.Command) {
flags.StringVar(&srv.Config.WorkerServiceProvider.Config.ControllerAddress, "wsp.config.controller-address", srv.Config.WorkerServiceProvider.Config.ControllerAddress, "Address of remote Controller process.")
// Computer
flags.BoolVar(&srv.Config.Computer.Run, "computer.run", srv.Config.Computer.Run, "Run the Computer service in process.")
flags.IntVar(&srv.Config.Computer.N, "computer.n", srv.Config.Computer.N, "The number of Computer services to run in process.")
flags.StringVar(&srv.Config.Computer.WorkerServiceID, "computer.worker_service_id", srv.Config.Computer.WorkerServiceID, "ID of WorkerService which spawned this computer.")
flags.AddFlagSet(serverFlagSet(&srv.Config.Computer.Config, "computer.config"))
}

View file

@ -17,7 +17,7 @@ type Balancer interface {
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)
RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.WorkerDiff, error)
BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error)
@ -88,7 +88,7 @@ func (b *NopBalancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.
func (b *NopBalancer) AddJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) {
return []dax.WorkerDiff{}, nil
}
func (b *NopBalancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID, jobs ...dax.Job) ([]dax.WorkerDiff, error) {
func (b *NopBalancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) ([]dax.WorkerDiff, error) {
return []dax.WorkerDiff{}, nil
}
func (b *NopBalancer) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) {

View file

@ -4,11 +4,12 @@ package balancer
import (
"log"
"math"
"sort"
"strings"
"time"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller"
"github.com/featurebasedb/featurebase/v3/dax/controller/schemar"
"github.com/featurebasedb/featurebase/v3/errors"
"github.com/featurebasedb/featurebase/v3/logger"
)
@ -25,31 +26,15 @@ var _ controller.Balancer = (*Balancer)(nil)
// jobs. It does not take anything else (such as job size, worker capabilities,
// etc) into consideration.
type Balancer struct {
workerRegistry controller.WorkerRegistry
// freeJobs is the set of jobs which have yet to be assigned to a worker.
// This could be because there are no available workers, or because a worker
// has been removed and the jobs for which it was responsible have yet to be
// reassigned.
freeJobs FreeJobService
freeWorkers FreeWorkerService
schemar schemar.Schemar
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, store controller.Store, logger logger.Logger) *Balancer {
func New(store controller.Store, logger logger.Logger) *Balancer {
return &Balancer{
workerRegistry: wr,
freeJobs: fjs,
freeWorkers: fws,
schemar: schemar,
store: store,
store: store,
logger: logger,
}
@ -68,7 +53,7 @@ func (b *Balancer) AddWorker(tx dax.Transaction, node *dax.Node) ([]dax.WorkerDi
return nil, errors.Wrapf(err, "creating node on node service: %s", node.Address)
}
diffs, err := b.balanceDatabase(tx, *qdbidp, node.ServiceID)
diffs, err := b.balanceDatabase(tx, *qdbidp)
if err != nil {
return nil, errors.Wrap(err, "balancing db")
}
@ -93,7 +78,7 @@ func (b *Balancer) RemoveWorker(tx dax.Transaction, addr dax.Address) ([]dax.Wor
if worker.DatabaseID != nil {
// Balance the affected database.
if diff, err := b.balanceDatabase(tx, *worker.DatabaseID, worker.ServiceID); err != nil {
if diff, err := b.balanceDatabase(tx, *worker.DatabaseID); err != nil {
return nil, errors.Wrapf(err, "balancing database: %s", worker.DatabaseID)
} else {
diffs.Merge(diff)
@ -141,12 +126,7 @@ func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.D
return diffs, 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 {
if cnt, err := b.store.WorkerCount(tx, roleType, dbid); err != nil {
return nil, errors.Wrap(err, "getting worker count")
} else if cnt == 0 {
if err := b.store.CreateFreeJobs(tx, roleType, dbid, jobs...); err != nil {
@ -154,7 +134,7 @@ func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.D
}
}
diff, err := b.addDatabaseJobs(tx, roleType, dbid, ws.ID, jobs...)
diff, err := b.addDatabaseJobs(tx, roleType, dbid, jobs...)
if err != nil {
return nil, errors.Wrapf(err, "adding database jobs: (%s) %s", roleType, dbid)
}
@ -164,14 +144,14 @@ func (b *Balancer) addJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.D
}
// 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)
func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, jobs ...dax.Job) (controller.InternalDiffs, error) {
workerJobs, err := b.store.WorkersJobs(tx, roleType, dbid)
if err != nil {
return nil, errors.Wrapf(err, "getting workers jobs: %s", roleType)
}
jset := dax.NewSet[dax.Job]()
addrToID := make(map[dax.Address]string)
addrToID := make(map[dax.Address]dax.WorkerID)
for _, workerInfo := range workerJobs {
addrToID[workerInfo.Address] = workerInfo.ID
jset.Merge(dax.NewSet(workerInfo.Jobs...))
@ -227,11 +207,7 @@ func (b *Balancer) addDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, db
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)
return b.store.WorkersJobs(tx, roleType, dbid)
}
func (b *Balancer) WorkerState(tx dax.Transaction, roleType dax.RoleType, addr dax.Address) (dax.WorkerInfo, error) {
@ -246,25 +222,95 @@ func (b *Balancer) RemoveJobs(tx dax.Transaction, roleType dax.RoleType, qtid da
return diffs.Output(), nil
}
func (b *Balancer) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) {
dbid := qdbid.DatabaseID
ws, err := b.store.WorkerService(tx, dbid)
// TODO can probably simplify this with better SQL
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.store.WorkersJobs(tx, roleType, qdbid.DatabaseID)
if err != nil {
return nil, errors.Wrap(err, "getting worker service")
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
}
}
diffs, err := b.balanceDatabase(tx, qdbid.DatabaseID, ws.ID)
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
}
// TODO can probably simplify this with better SQL
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])
dbid := qtid.QualifiedDatabaseID.DatabaseID
prefix := string(qtid.Key())
workerJobs, err := b.store.WorkersJobs(tx, roleType, dbid)
if err != nil {
return nil, errors.Wrapf(err, "getting worker jobs: (%s) %s", roleType, dbid)
}
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) BalanceDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerDiff, error) {
diffs, err := b.balanceDatabase(tx, qdbid.DatabaseID)
if err != nil {
return nil, errors.Wrap(err, "balancing")
}
return diffs.Output(), nil
}
func (b *Balancer) balanceDatabase(tx dax.Transaction, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) {
func (b *Balancer) balanceDatabase(tx dax.Transaction, dbid dax.DatabaseID) (controller.InternalDiffs, error) {
diffs := controller.NewInternalDiffs()
for _, role := range dax.AllRoleTypes {
diff, err := b.balanceDatabaseForRole(tx, role, dbid, svcID)
diff, err := b.balanceDatabaseForRole(tx, role, dbid)
if err != nil {
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", role, dbid)
}
@ -274,26 +320,26 @@ func (b *Balancer) balanceDatabase(tx dax.Transaction, dbid dax.DatabaseID, svcI
return diffs, nil
}
func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) {
func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID) (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.store.WorkerCount(tx, roleType, svcID); err != nil {
if cnt, err := b.store.WorkerCount(tx, roleType, dbid); err != nil {
return nil, errors.Wrapf(err, "getting worker count: (%s) %s", roleType, dbid)
} else if cnt == 0 {
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, dbid, svcID); err != nil {
if diff, err := b.processFreeJobs(tx, roleType, dbid); 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, dbid, svcID, diffs)
diff, err := b.balanceDatabaseJobs(tx, roleType, dbid, diffs)
if err != nil {
return nil, errors.Wrap(err, "balancing jobs")
}
@ -308,8 +354,8 @@ func (b *Balancer) balanceDatabaseForRole(tx dax.Transaction, roleType dax.RoleT
// 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, dbid dax.DatabaseID, svcID dax.WorkerServiceID, diffs controller.InternalDiffs) (controller.InternalDiffs, error) {
workerInfos, err := b.store.WorkersJobs(tx, roleType, svcID)
func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID, diffs controller.InternalDiffs) (controller.InternalDiffs, error) {
workerInfos, err := b.store.WorkersJobs(tx, roleType, dbid)
if err != nil {
return nil, errors.Wrapf(err, "getting current state: (%s) %s", roleType, dbid)
}
@ -327,7 +373,7 @@ func (b *Balancer) balanceDatabaseJobs(tx dax.Transaction, roleType dax.RoleType
numWorkersAboveMin := numJobs % numWorkers
removedJobs := make(dax.Jobs, 0)
addedJobs := make(map[string]dax.Jobs)
addedJobs := make(map[dax.WorkerID]dax.Jobs)
// remove jobs loop
//
@ -392,14 +438,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, dbid dax.DatabaseID, svcID dax.WorkerServiceID) (controller.InternalDiffs, error) {
func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, dbid dax.DatabaseID) (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, dbid, svcID, jobs...); err != nil {
if aj, err := b.addDatabaseJobs(tx, roleType, dbid, jobs...); err != nil {
return nil, errors.Wrapf(err, "adding jobs: %s", jobs)
} else {
diffs.Merge(aj)
@ -408,6 +454,13 @@ func (b *Balancer) processFreeJobs(tx dax.Transaction, roleType dax.RoleType, db
return diffs, nil
}
func (b *Balancer) ReadNode(tx dax.Transaction, addr dax.Address) (*dax.Node, error) {
return b.store.Worker(tx, addr)
}
func (b *Balancer) Nodes(tx dax.Transaction) ([]*dax.Node, error) {
return b.store.Workers(tx)
}
func (b *Balancer) CreateWorkerServiceProvider(tx dax.Transaction, sp dax.WorkerServiceProvider) error {
return b.store.CreateWorkerServiceProvider(tx, sp)
}

View file

@ -1,86 +0,0 @@
package balancer_test
import (
"context"
"testing"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller/sqldb"
"github.com/stretchr/testify/require"
)
func TestFreeJobService(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)
}
}()
// must have a database to do job stuff
schemar := sqldb.NewSchemar(nil)
err = schemar.CreateDatabase(tx,
&dax.QualifiedDatabase{
OrganizationID: orgID,
Database: dax.Database{ID: dbID, Name: dbName}})
require.NoError(t, err)
fjSvc := sqldb.NewFreeJobService(nil)
qdbid := dax.QualifiedDatabaseID{OrganizationID: orgID, DatabaseID: dbID}
qtid := dax.QualifiedTableID{
QualifiedDatabaseID: qdbid,
Name: tableName,
ID: tableID,
}
job1 := dax.Job(qtid.Key() + "job1")
job2 := dax.Job(qtid.Key() + "job2")
job3 := dax.Job(qtid.Key() + "job3")
node := &dax.Node{
Address: nodeAddr,
RoleTypes: []dax.RoleType{role},
}
err = fjSvc.CreateJobs(tx, role, qdbid, job1, job2, job3)
require.NoError(t, err)
err = fjSvc.DeleteJob(tx, role, qdbid, job2)
require.NoError(t, err)
jobs, err := fjSvc.ListJobs(tx, role, qdbid)
require.NoError(t, err)
require.ElementsMatch(t, dax.Jobs{job1, job3}, jobs)
workerReg := sqldb.NewWorkerRegistry(nil)
err = workerReg.AddWorker(tx, node)
require.NoError(t, err)
wjSvc := sqldb.NewWorkerJobService(nil)
err = wjSvc.CreateWorker(tx, role, qdbid, nodeAddr)
require.NoError(t, err)
err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job1)
require.NoError(t, err)
jobs, err = fjSvc.ListJobs(tx, role, qdbid)
require.NoError(t, err)
require.ElementsMatch(t, dax.Jobs{job3}, jobs)
err = fjSvc.MarkJobsAsFree(tx, role, qdbid, dax.Jobs{job1})
require.NoError(t, err)
jobs, err = fjSvc.ListJobs(tx, role, qdbid)
require.NoError(t, err)
require.ElementsMatch(t, dax.Jobs{job1, job3}, jobs)
err = fjSvc.DeleteJobsForTable(tx, role, qtid)
require.NoError(t, err)
jobs, err = fjSvc.ListJobs(tx, role, qdbid)
require.NoError(t, err)
require.Empty(t, jobs)
}

View file

@ -1,51 +0,0 @@
package balancer_test
import (
"context"
"testing"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller/sqldb"
"github.com/stretchr/testify/require"
)
func TestFreeWorkerService(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)
}
}()
node1 := &dax.Node{Address: nodeAddr, RoleTypes: dax.AllRoleTypes}
node2 := &dax.Node{Address: nodeAddr2, RoleTypes: dax.AllRoleTypes}
node3 := &dax.Node{Address: nodeAddr3, RoleTypes: dax.AllRoleTypes}
node4 := &dax.Node{Address: nodeAddr4, RoleTypes: dax.AllRoleTypes}
node5 := &dax.Node{Address: nodeAddr5, RoleTypes: dax.AllRoleTypes}
workerReg := sqldb.NewWorkerRegistry(nil)
// Add some workers.
require.NoError(t, workerReg.AddWorker(tx, node1))
require.NoError(t, workerReg.AddWorker(tx, node2))
require.NoError(t, workerReg.AddWorker(tx, node3))
require.NoError(t, workerReg.AddWorker(tx, node4))
require.NoError(t, workerReg.AddWorker(tx, node5))
// Remove one of the workers.
require.NoError(t, workerReg.RemoveWorker(tx, node2.Address))
fwSvc := sqldb.NewFreeWorkerService(nil)
addrs, err := fwSvc.ListWorkers(tx, role)
require.NoError(t, err)
require.ElementsMatch(t, dax.Addresses{nodeAddr, nodeAddr3, nodeAddr4, nodeAddr5}, addrs)
addrs, err = fwSvc.PopWorkers(tx, role, 2)
require.NoError(t, err)
require.Equal(t, 2, len(addrs))
require.NotEqual(t, addrs[0], addrs[1])
}

View file

@ -1,65 +0,0 @@
package balancer_test
import (
"context"
"testing"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller/sqldb"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
const (
nodeAddr = "myaddress"
nodeAddr2 = "myaddress2"
nodeAddr3 = "myaddress3"
nodeAddr4 = "myaddress4"
nodeAddr5 = "myaddress5"
)
func TestWorkerRegistry(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)
}
}()
workerReg := sqldb.NewWorkerRegistry(nil)
err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr, RoleTypes: []dax.RoleType{dax.RoleTypeCompute}})
require.NoError(t, err)
node, err := workerReg.Worker(tx, nodeAddr)
require.NoError(t, err)
require.EqualValues(t, nodeAddr, node.Address)
require.EqualValues(t, 1, len(node.RoleTypes))
require.EqualValues(t, "compute", node.RoleTypes[0])
err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr2, RoleTypes: []dax.RoleType{dax.RoleTypeTranslate, dax.RoleTypeCompute}})
require.NoError(t, err, "create node 2")
err = workerReg.AddWorker(tx, &dax.Node{Address: nodeAddr3, RoleTypes: []dax.RoleType{dax.RoleTypeCompute}})
require.NoError(t, err, "create node 3")
nodes, err := workerReg.Workers(tx)
require.NoError(t, err)
assert.EqualValues(t, 3, len(nodes))
for _, node := range nodes {
assert.Contains(t, node.RoleTypes, dax.RoleType("compute"), "node should have compute role but is: %+v", node)
}
err = workerReg.RemoveWorker(tx, nodeAddr2)
require.NoError(t, err, "deleting node")
nodes, err = workerReg.Workers(tx)
require.NoError(t, err)
require.EqualValues(t, 2, len(nodes))
for _, node := range nodes {
require.Contains(t, node.RoleTypes, dax.RoleType("compute"), "node should have compute role but is: %+v", node)
}
}

View file

@ -1,140 +0,0 @@
package balancer_test
import (
"context"
"testing"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller/sqldb"
"github.com/stretchr/testify/require"
)
const (
orgID = "orgid"
dbID = "blah"
dbName = "nameofdb"
role = "compute"
tableName = "tbl"
tableID = "tblid"
)
func TestWorkerJobService(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)
}
}()
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,
&dax.QualifiedDatabase{
OrganizationID: orgID,
Database: dax.Database{ID: dbID, Name: dbName}})
require.NoError(t, err)
wjSvc := sqldb.NewWorkerJobService(nil)
qdbid := dax.QualifiedDatabaseID{OrganizationID: orgID, DatabaseID: dbID}
node := &dax.Node{
Address: nodeAddr,
ServiceID: "mysvc1",
RoleTypes: []dax.RoleType{role},
}
// have to create a free worker before you can create a worker job worker
workerReg := sqldb.NewWorkerRegistry(nil)
err = workerReg.AddWorker(tx, node)
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
// the Delete method? Why yes, yes it is... thank you for asking. #TODO
qtid := dax.QualifiedTableID{
QualifiedDatabaseID: qdbid,
Name: tableName,
ID: tableID,
}
job1 := dax.Job(qtid.Key() + "job1")
job2 := dax.Job(qtid.Key() + "job2")
job3 := dax.Job(qtid.Key() + "job3")
fjSvc := sqldb.NewFreeJobService(nil)
err = fjSvc.CreateJobs(tx, role, qdbid, job1, job2, job3)
require.NoError(t, err)
err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job1, job2)
require.NoError(t, err)
jobs, err := wjSvc.ListJobs(tx, role, qdbid, nodeAddr)
require.NoError(t, err)
require.ElementsMatch(t, dax.Jobs{job1, job2}, jobs)
workerInfos, err := wjSvc.WorkersJobs(tx, role, qdbid)
require.NoError(t, err)
require.Equal(t, 1, len(workerInfos))
require.ElementsMatch(t, []dax.Job{job1, job2}, workerInfos[0].Jobs)
cnt, err := wjSvc.WorkerCount(tx, role, qdbid)
require.NoError(t, err)
require.Equal(t, 1, cnt)
addrs, err := wjSvc.ListWorkers(tx, role, qdbid)
require.NoError(t, err)
require.ElementsMatch(t, dax.Addresses{nodeAddr}, addrs)
err = wjSvc.AssignWorkerToJobs(tx, role, qdbid, nodeAddr, job3)
require.NoError(t, err)
jcs, err := wjSvc.JobCounts(tx, role, qdbid, nodeAddr)
require.NoError(t, err)
require.Equal(t, 3, jcs[nodeAddr])
err = wjSvc.DeleteJob(tx, role, qdbid, nodeAddr, job3)
require.NoError(t, err)
idiffs, err := wjSvc.DeleteJobsForTable(tx, role, qtid)
require.NoError(t, err)
workerDiffs := idiffs.Output()
require.Equal(t, 1, len(workerDiffs))
require.EqualValues(t, nodeAddr, workerDiffs[0].Address)
require.Empty(t, workerDiffs[0].AddedJobs)
require.ElementsMatch(t, []dax.Job{job1, job2}, workerDiffs[0].RemovedJobs)
jobs, err = wjSvc.ListJobs(tx, role, qdbid, nodeAddr)
require.NoError(t, err)
require.Empty(t, jobs)
dk := wjSvc.DatabaseForWorker(tx, nodeAddr)
require.EqualValues(t, "db__orgid__blah", dk)
err = wjSvc.ReleaseWorkers(tx, nodeAddr)
require.NoError(t, err)
addrs, err = wjSvc.ListWorkers(tx, role, qdbid)
require.NoError(t, err)
require.Empty(t, addrs)
}

View file

@ -7,12 +7,7 @@ import (
// NewBalancer returns a new instance of controller.Balancer.
func NewBalancer(log logger.Logger) *balancer.Balancer {
schemar := NewSchemar(log)
fjs := NewFreeJobService(log)
wjs := NewWorkerJobService(log)
fws := NewFreeWorkerService(log)
ns := NewWorkerRegistry(log)
store := NewStore(log)
return balancer.New(ns, fjs, wjs, fws, schemar, store, log)
return balancer.New(store, log)
}

View file

@ -1,115 +0,0 @@
package sqldb
import (
"fmt"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller/balancer"
"github.com/featurebasedb/featurebase/v3/dax/models"
"github.com/featurebasedb/featurebase/v3/logger"
"github.com/pkg/errors"
)
func NewFreeJobService(log logger.Logger) balancer.FreeJobService {
if log == nil {
log = logger.NopLogger
}
return &freeJobService{
log: log,
}
}
type freeJobService struct {
log logger.Logger
}
func (fj *freeJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job ...dax.Job) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
// jobNames is used as input to the "name in (...)" query.
jobNames := make([]interface{}, 0, len(job))
for i := range job {
jobNames = append(jobNames, job[i].Job())
}
// existing will contain the list of jobs which already exist.
existing := &models.Jobs{}
if err := dt.C.Where("name in (?)", jobNames...).All(existing); err != nil {
return errors.Wrap(err, "getting existing jobs")
}
jobs := make(models.Jobs, 0, len(job))
for _, j := range job {
// Check to be sure this job doesn't already exist.
if existing.Contains(j) {
continue
}
jobs = append(jobs, models.Job{
Name: j,
Role: roleType,
DatabaseID: qdbid.DatabaseID,
})
}
if len(jobs) == 0 {
return nil
}
err := dt.C.Create(jobs)
return errors.Wrap(err, "creating free jobs")
}
func (fj *freeJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, job dax.Job) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
err := dt.C.RawQuery("DELETE from jobs where role = ? and database_id = ? and name = ? and worker_id is NULL", roleType, qdbid.DatabaseID, job).Exec()
return errors.Wrap(err, "deleting")
}
func (fj *freeJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
err := dt.C.RawQuery("DELETE from jobs where role = ? and database_id = ? and name LIKE ? and worker_id is NULL",
roleType, qtid.DatabaseID, fmt.Sprintf("%s%%", qtid.Key())).Exec()
return errors.Wrap(err, "deleting")
}
func (fj *freeJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (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", roleType, qdbid.DatabaseID).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
}
// MarkJobsAsFree disassociates any worker that was previously assigned to this
// job.
func (fj *freeJobService) MarkJobsAsFree(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, jobs dax.Jobs) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
err := dt.C.RawQuery("UPDATE jobs SET worker_id = NULL WHERE role = ? and database_id = ?", roleType, qdbid.DatabaseID).Exec()
return errors.Wrap(err, "marking jobs free")
}

View file

@ -1,72 +0,0 @@
package sqldb
import (
"fmt"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller/balancer"
"github.com/featurebasedb/featurebase/v3/dax/models"
"github.com/featurebasedb/featurebase/v3/logger"
"github.com/pkg/errors"
)
func NewFreeWorkerService(log logger.Logger) balancer.FreeWorkerService {
if log == nil {
log = logger.NopLogger
}
return &freeWorkerService{
log: log,
}
}
type freeWorkerService struct {
log logger.Logger
}
func (fw *freeWorkerService) PopWorkers(tx dax.Transaction, roleType dax.RoleType, num int) ([]dax.Address, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
results := make([]struct {
Address dax.Address `db:"address"`
}, 0, num)
sel := fmt.Sprintf("select address from workers where role_%s = true and database_id is NULL limit ?", roleType)
err := dt.C.RawQuery(sel, num).All(&results)
if err != nil {
return nil, errors.Wrap(err, "querying")
}
if len(results) < num {
return nil, errors.Errorf("not enough free workers to get: wanted %d, have: %d", num, len(results))
}
ret := make([]dax.Address, num)
for i, res := range results {
ret[i] = res.Address
}
return ret, nil
}
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")
}
workers := make(models.Workers, 0)
where := fmt.Sprintf("role_%s = true and database_id is NULL", roleType)
err := dt.C.Select("address").Where(where).Order("address asc").All(&workers)
if err != nil {
return nil, errors.Wrap(err, "querying for free workers")
}
ret := make(dax.Addresses, len(workers))
for i, w := range workers {
ret[i] = w.Address
}
return ret, nil
}

View file

@ -5,7 +5,6 @@ import (
"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"
@ -62,24 +61,8 @@ func (s *store) AddWorker(tx dax.Transaction, node *dax.Node) (*dax.DatabaseID,
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) {
// WorkerCount returns the number of workers associated with the given database.
func (s *store) WorkerCount(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) (int, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
@ -158,7 +141,7 @@ func (s *store) WorkerJobs(tx dax.Transaction, role dax.RoleType, addr dax.Addre
return wi, nil
}
func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.WorkerServiceID) ([]dax.WorkerInfo, error) {
func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) ([]dax.WorkerInfo, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
@ -166,7 +149,7 @@ func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.Wor
results := []struct {
Address dax.Address `db:"address"`
WorkerID string `db:"worker_id"`
WorkerID dax.WorkerID `db:"worker_id"`
JobName nulls.String `db:"job_name"`
}{}
@ -176,14 +159,15 @@ func (s *store) WorkersJobs(tx dax.Transaction, role dax.RoleType, svcID dax.Wor
jobs.name as job_name
FROM
workers LEFT JOIN jobs on jobs.worker_id = workers.id
INNER JOIN worker_service ws ON ws.id = workers.service_id
WHERE
jobs.role = ?
workers.service_id = ?
ws.database_id = ?
ORDER BY
address ASC
job_name ASC`
err := dt.C.RawQuery(sql, role, svcID).All(&results)
err := dt.C.RawQuery(sql, role, dbid).All(&results)
if err != nil {
return nil, errors.Wrap(err, "querying")
}
@ -390,7 +374,7 @@ func (s *store) DeleteJobsForTable(tx dax.Transaction, role dax.RoleType, qtid d
return nil, errors.Wrap(err, "querying for jobs")
}
idiffs := make(balancer.InternalDiffs)
idiffs := make(controller.InternalDiffs)
ids := make([]uuid.UUID, 0, len(results))
for _, job := range results {
idiffs.Removed(job.Address, job.Name)
@ -403,3 +387,56 @@ func (s *store) DeleteJobsForTable(tx dax.Transaction, role dax.RoleType, qtid d
return idiffs, errors.Wrap(err, "deleting jobs")
}
func (s *store) Workers(tx dax.Transaction) ([]*dax.Node, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
workers := []*models.Worker{}
dt.C.Eager().Order("address asc").All(&workers)
ret := make([]*dax.Node, len(workers))
for i, worker := range workers {
ret[i] = &dax.Node{
Address: worker.Address,
RoleTypes: workerRoleTypes(worker),
}
}
return ret, nil
}
func (w *store) Worker(tx dax.Transaction, addr dax.Address) (*dax.Node, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
err := dt.C.Eager().Where("address = ?", addr).First(worker)
if err != nil {
return nil, errors.Wrapf(err, "getting worker: %s", addr)
}
return &dax.Node{
Address: worker.Address,
RoleTypes: workerRoleTypes(worker),
}, nil
}
func workerRoleTypes(worker *models.Worker) []dax.RoleType {
roleTypes := make([]dax.RoleType, 0)
if worker.RoleCompute {
roleTypes = append(roleTypes, dax.RoleTypeCompute)
}
if worker.RoleTranslate {
roleTypes = append(roleTypes, dax.RoleTypeTranslate)
}
if worker.RoleQuery {
roleTypes = append(roleTypes, dax.RoleTypeQuery)
}
return roleTypes
}

View file

@ -106,12 +106,6 @@ func TestStore(t *testing.T) {
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)
@ -126,5 +120,4 @@ func TestStore(t *testing.T) {
require.Equal(t, ws.ID, node.ServiceID)
require.Equal(t, expectedDBID, node.DatabaseID)
})
}

View file

@ -1,114 +0,0 @@
package sqldb
import (
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller"
"github.com/featurebasedb/featurebase/v3/dax/models"
"github.com/featurebasedb/featurebase/v3/logger"
"github.com/pkg/errors"
)
var _ controller.WorkerRegistry = (*workerRegistry)(nil)
func NewWorkerRegistry(log logger.Logger) *workerRegistry {
if log == nil {
log = logger.NopLogger
}
return &workerRegistry{
log: log,
}
}
type workerRegistry struct {
log logger.Logger
}
func (w *workerRegistry) AddWorker(tx dax.Transaction, node *dax.Node) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{
Address: node.Address,
ServiceID: node.ServiceID,
}
for _, roleType := range node.RoleTypes {
if err := worker.SetRole(roleType); err != nil {
return errors.Wrapf(err, "setting role: %s", roleType)
}
}
return dt.C.Create(worker)
}
func (w *workerRegistry) Worker(tx dax.Transaction, addr dax.Address) (*dax.Node, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
err := dt.C.Eager().Where("address = ?", addr).First(worker)
if err != nil {
return nil, errors.Wrapf(err, "getting worker: %s", addr)
}
return &dax.Node{
Address: worker.Address,
RoleTypes: workerRoleTypes(worker),
}, nil
}
func (w *workerRegistry) RemoveWorker(tx dax.Transaction, addr dax.Address) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
err := dt.C.Eager().Where("address = ?", addr).First(worker)
if isNoRowsError(err) {
return nil
} else if err != nil {
return errors.Wrapf(err, "finding worker: %s", addr)
}
err = dt.C.Destroy(worker)
return errors.Wrap(err, "destroying worker")
}
func (w *workerRegistry) Workers(tx dax.Transaction) ([]*dax.Node, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
workers := []*models.Worker{}
dt.C.Eager().Order("address asc").All(&workers)
ret := make([]*dax.Node, len(workers))
for i, worker := range workers {
ret[i] = &dax.Node{
Address: worker.Address,
RoleTypes: workerRoleTypes(worker),
}
}
return ret, nil
}
func workerRoleTypes(worker *models.Worker) []dax.RoleType {
roleTypes := make([]dax.RoleType, 0)
if worker.RoleCompute {
roleTypes = append(roleTypes, dax.RoleTypeCompute)
}
if worker.RoleTranslate {
roleTypes = append(roleTypes, dax.RoleTypeTranslate)
}
if worker.RoleQuery {
roleTypes = append(roleTypes, dax.RoleTypeQuery)
}
return roleTypes
}

View file

@ -1,285 +0,0 @@
package sqldb
import (
"fmt"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/controller/balancer"
"github.com/featurebasedb/featurebase/v3/dax/models"
"github.com/featurebasedb/featurebase/v3/logger"
"github.com/gofrs/uuid"
"github.com/pkg/errors"
)
func NewWorkerJobService(log logger.Logger) balancer.WorkerJobService {
if log == nil {
log = logger.NopLogger
}
return &workerJobService{
log: log,
}
}
type workerJobService struct {
log logger.Logger
}
// WorkersJobs returns all the workers for the database along with the jobs
// associated to each worker, even if the number of jobs is 0.
func (w *workerJobService) WorkersJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) ([]dax.WorkerInfo, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
// First, get all workers for the database.
workers := models.Workers{}
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")
}
// Then, get the jobs for each worker. Ideally, we would do this in a single
// sql query, but it wasn't clear how to do an Eager() LeftJoin() where
// there is a where clause condition on the right side of the join (in this
// case, `jobs.role = ?`).
ret := make([]dax.WorkerInfo, len(workers))
for i, worker := range workers {
ret[i].Address = worker.Address
jobs, err := jobsForWorker(dt, &worker, roleType)
if err != nil {
return nil, errors.Wrap(err, "getting jobs for worker")
}
ret[i].Jobs = jobs
}
return ret, nil
}
func jobsForWorker(dt *DaxTransaction, worker *models.Worker, roleType dax.RoleType) ([]dax.Job, error) {
jobs := models.Jobs{}
if err := dt.C.Where("worker_id = ? and role = ?", worker.ID, roleType).Order("name asc").All(&jobs); err != nil {
return nil, errors.Wrapf(err, "getting jobs for worker: %s", worker.ID)
}
ret := make([]dax.Job, len(jobs))
for i := range jobs {
ret[i] = jobs[i].Name
}
return ret, nil
}
func (w *workerJobService) WorkerCount(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (int, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
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) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
workers := models.Workers{}
sql := fmt.Sprintf("role_%s = true and database_id = ?", roleType)
err := dt.C.Select("address").Where(sql, qdbid.DatabaseID).Order("address asc").All(&workers)
if err != nil {
return nil, errors.Wrap(err, "getting workers")
}
ret := make(dax.Addresses, len(workers))
for i, wrkr := range workers {
ret[i] = wrkr.Address
}
return ret, nil
}
func (w *workerJobService) AssignWorkerToJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job ...dax.Job) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
sql := fmt.Sprintf("address = ? and role_%s = true", roleType)
err := dt.C.Where(sql, addr).First(worker)
if err != nil {
return errors.Wrapf(err, "getting worker: (%s) %s", roleType, addr)
}
jobs := models.Jobs{}
err = dt.C.RawQuery("UPDATE jobs SET worker_id = ? WHERE role = ? and name in (?) RETURNING jobs.ID, jobs.Name", worker.ID, roleType, job).All(&jobs)
if err != nil {
return errors.Wrap(err, "updating jobs")
}
// Assign jobs not in "jobs", and therefore didn't get updated by the
// previous sql statement.
toBeAssigned := jobsNotAssigned(job, jobs, roleType, worker, qdbid.DatabaseID)
if err := dt.C.Create(toBeAssigned); err != nil {
return errors.Wrap(err, "creating jobs")
}
return nil
}
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 {
if assignedJob.Name == incJob {
continue outer
}
}
toBeAssigned = append(toBeAssigned,
models.Job{
Name: incJob,
Role: roleType,
DatabaseID: dbid,
Worker: worker,
},
)
}
return toBeAssigned
}
func (w *workerJobService) DeleteJob(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job dax.Job) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
sql := fmt.Sprintf("role_%s = true and database_id = ? and address = ?", roleType)
err := dt.C.Select("id").Where(sql, qdbid.DatabaseID, addr).First(worker)
if err != nil {
return errors.Wrap(err, "getting worker")
}
jerb := &models.Job{}
dt.C.Select("id").Where("role = ? and worker_id = ? and name = ?", roleType, worker.ID, job).First(jerb)
if err != nil {
return errors.Wrap(err, "getting job")
}
err = dt.C.Destroy(jerb)
if err != nil {
return errors.Wrap(err, "destroying job")
}
return nil
}
func (w *workerJobService) DeleteJobsForTable(tx dax.Transaction, roleType dax.RoleType, qtid dax.QualifiedTableID) (balancer.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 ?", roleType, 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")
}
func (w *workerJobService) JobCounts(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addrs ...dax.Address) (map[dax.Address]int, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
results := []struct {
Address dax.Address `db:"address"`
Count int `db:"count"`
}{}
var err error
if len(addrs) == 0 {
qstring := `select address, count(*) as count
from workers w inner join jobs j on j.worker_id = w.id
where w.database_id = ? and w.role_%s = true
and j.role = ?
group by w.address`
sql := fmt.Sprintf(qstring, roleType)
err = dt.C.RawQuery(sql, qdbid.DatabaseID, roleType).All(&results)
} else {
qstring := `select address, count(*) as count
from workers w inner join jobs j on j.worker_id = w.id
where w.database_id = ? and w.role_%s = true
and j.role = ?
and w.address in (?)
group by w.address`
sql := fmt.Sprintf(qstring, roleType)
err = dt.C.RawQuery(sql, qdbid.DatabaseID, roleType, addrs).All(&results)
}
if err != nil {
return nil, errors.Wrap(err, "querying for jobs")
}
ret := make(map[dax.Address]int)
for _, res := range results {
ret[res.Address] = res.Count
}
return ret, nil
}
func (w *workerJobService) ListJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address) (dax.Jobs, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return nil, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
sql := fmt.Sprintf("role_%s = true and database_id = ? and address = ?", roleType)
err := dt.C.Where(sql, qdbid.DatabaseID, addr).First(worker)
if isNoRowsError(err) {
return nil, nil
} else if err != nil {
return nil, errors.Wrap(err, "getting worker")
}
return jobsForWorker(dt, worker, roleType)
}
func (w *workerJobService) DatabaseForWorker(tx dax.Transaction, addr dax.Address) dax.DatabaseKey {
dt, ok := tx.(*DaxTransaction)
if !ok {
panic("wrong transaction type passed to sqldb DatabaseForWorker")
}
db := &models.Database{}
err := dt.C.RawQuery("select d.ID, d.organization_id from databases d inner join workers w on d.id = w.database_id where w.address = ?", addr).First(db)
if isNoRowsError(err) {
return ""
} else if err != nil {
panic(err)
}
return dax.QualifiedDatabase{OrganizationID: dax.OrganizationID(db.OrganizationID), Database: dax.Database{ID: dax.DatabaseID(db.ID)}}.Key()
}

View file

@ -1,37 +0,0 @@
package sqldb
import (
"testing"
"github.com/featurebasedb/featurebase/v3/dax"
"github.com/featurebasedb/featurebase/v3/dax/models"
"github.com/gofrs/uuid"
"github.com/stretchr/testify/require"
)
func TestJobsNotUpdated(t *testing.T) {
u, _ := uuid.NewV4()
u2, _ := uuid.NewV4()
incJobs := []dax.Job{"job1", "job2", "job3", "job4"}
created := models.Jobs{
models.Job{
ID: u,
Name: "job2",
},
}
toCreate := jobsNotAssigned(incJobs, created, dax.RoleTypeCompute, &models.Worker{
ID: u2,
RoleCompute: true,
}, "dbid")
require.Equal(t, 3, len(toCreate))
// Test when 0 jobs are updated
toCreate = jobsNotAssigned(incJobs, models.Jobs{}, dax.RoleTypeCompute, &models.Worker{
ID: u2,
RoleCompute: true,
}, "dbid")
require.Equal(t, 4, len(toCreate))
}

View file

@ -36,11 +36,10 @@ type Store interface {
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)
WorkerCount(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
WorkersJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID) ([]dax.WorkerInfo, error)
AssignWorkerToJobs(tx dax.Transaction, role dax.RoleType, dbid dax.DatabaseID, workerID dax.WorkerID, 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)
@ -51,4 +50,6 @@ type Store interface {
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)
Workers(tx dax.Transaction) ([]*dax.Node, error)
Worker(tx dax.Transaction, addr dax.Address) (*dax.Node, error)
}

View file

@ -224,7 +224,6 @@ func TestDAXIntegration(t *testing.T) {
t.Run("Poller", func(t *testing.T) {
cfg := test.DefaultConfig()
cfg.Computer.N = 2
opt := server.OptCommandConfig(cfg)
mc := test.MustRunManagedCommand(t, opt)
@ -690,7 +689,6 @@ func TestDAXIntegration(t *testing.T) {
t.Run("DatabaseOptions", func(t *testing.T) {
cfg := test.DefaultConfig()
cfg.Computer.N = 4
opt := server.OptCommandConfig(cfg)
mc := test.MustRunManagedCommand(t, opt)
defer mc.Close()