mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-06 00:25:55 +00:00
* 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
303 lines
9.3 KiB
Go
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()
|
|
}
|