featurebase/dax/controller/sqldb/workerjob.go
Matthew Jaffee 8fe73146c8
Sqldb rip boltdb (#2341)
* serverless sqldb use same env for test config as normal

* rip boltdb implementation of controller backend out

it was replaced by postgres and no longer works properly.

This involved migrating a number of tests which only worked with
boltdb, which exposed several ways in which the postgres
implementation had slightly different behavior from the bolt
one:
1. ordering of results in some cases, and
2. (more importantly) erroring when a record to delete was not
found. The bolt implementation silently ignored it when things to
delete weren't found, so we make some changes to match that behavior.

Also stopped propagating CreatedAt and UpdatedAt from DB tables into
dax types. These were breaking existing tests. Perhaps it would be
better to actually use them, but for now they will only exist at the
DB level.

This change set also moves the insertion of the directive_versions
record out of migrations and into the startup/connection code. Having
this in the migrations was a bit ugly because you couldn't just
truncate all the tables and have everything work from
scratch. Inserting it during startup is fairly innocuous, and will
just continue on if it already exists.

* update directive_version test

I changed the initial value to 0 so that the first version that gets
sent out is 1
2023-03-22 08:54:13 -05:00

303 lines
9.3 KiB
Go

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
}
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")
}
workers := models.Workers{}
err := dt.C.Eager().Where("role = ? and database_id = ?", roleType, qdbid.DatabaseID).Order("address asc").All(&workers)
if err != nil {
return nil, errors.Wrap(err, "getting workers")
}
ret := make([]dax.WorkerInfo, len(workers))
for i, worker := range workers {
ret[i].Address = worker.Address
ret[i].Jobs = make([]dax.Job, len(worker.Jobs))
for j, job := range worker.Jobs {
ret[i].Jobs[j] = job.Name
}
}
return ret, nil
}
func (w *workerJobService) WorkerCount(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID) (int, error) {
dt, ok := tx.(*DaxTransaction)
if !ok {
return 0, dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
cnt, err := dt.C.Where("role = ? and database_id = ?", roleType, qdbid.DatabaseID).Count(worker)
return cnt, 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{}
err := dt.C.Select("address").Where("role = ? and database_id = ?", roleType, 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) 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{}
err := dt.C.RawQuery("UPDATE workers SET database_id = ? WHERE role = ? and address = ? RETURNING workers.ID", qdbid.DatabaseID, roleType, addr).First(worker)
return errors.Wrap(err, "associating worker to database")
}
func (w *workerJobService) FreeWorkers(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) DeleteWorker(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{}
err := dt.C.Where("address = ? and role = ? and database_id = ?", addr, roleType, qdbid.DatabaseID).First(worker)
if isNoRowsError(err) {
return nil
} else if err != nil {
return errors.Wrap(err, "getting worker")
}
err = dt.C.Destroy(worker)
return errors.Wrap(err, "deleting worker")
}
func (w *workerJobService) CreateJobs(tx dax.Transaction, roleType dax.RoleType, qdbid dax.QualifiedDatabaseID, addr dax.Address, job ...dax.Job) error {
dt, ok := tx.(*DaxTransaction)
if !ok {
return dax.NewErrInvalidTransaction("*sqldb.DaxTransaction")
}
worker := &models.Worker{}
err := dt.C.Where("address = ? and role = ?", addr, roleType).First(worker)
if err != nil {
return errors.Wrap(err, "getting worker")
}
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")
}
// create jobs not in "jobs"
toCreate := jobsNotUpdated(job, jobs, worker)
err = dt.C.Create(toCreate)
if err != nil {
return errors.Wrap(err, "creating jobs")
}
return errors.Wrap(err, "creating jobs")
}
func jobsNotUpdated(incomingJobs []dax.Job, created models.Jobs, worker *models.Worker) (toCreate models.Jobs) {
outer:
for _, incJob := range incomingJobs {
for _, createdJob := range created {
if createdJob.Name == incJob {
continue outer
}
}
toCreate = append(toCreate,
models.Job{
Name: incJob,
Role: worker.Role,
DatabaseID: dax.DatabaseID(worker.DatabaseID.String),
Worker: worker,
},
)
}
return toCreate
}
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{}
err := dt.C.Select("id").Where("role = ? and database_id = ? and address = ?", roleType, 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, addr ...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(addr) == 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 = ?
group by w.address`
err = dt.C.RawQuery(qstring, qdbid.DatabaseID, roleType).All(&results)
} else {
qstring := `select address, count(*) as count
from workers w inner join jobs j on j.worker_id = w.id
where w.address in (?) and w.database_id = ? and w.role = ?
group by w.address`
err = dt.C.RawQuery(qstring, addr, qdbid.DatabaseID, roleType).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{}
err := dt.C.Eager().Where("role = ? and database_id = ? and address = ?", roleType, qdbid.DatabaseID, addr).First(worker)
if isNoRowsError(err) {
return nil, nil
} else if err != nil {
return nil, errors.Wrap(err, "getting worker")
}
ret := make(dax.Jobs, len(worker.Jobs))
// jobs are ordered by "name asc" defined on the worker model.
for i, job := range worker.Jobs {
ret[i] = job.Name
}
return ret, nil
}
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()
}