diff --git a/dax/controller/balancer/boltdb/balancer.go b/dax/controller/balancer/boltdb/balancer.go index bb2e85711..e151aa4db 100644 --- a/dax/controller/balancer/boltdb/balancer.go +++ b/dax/controller/balancer/boltdb/balancer.go @@ -112,7 +112,13 @@ func (w *workerJobService) getWorkerInfos(tx dax.Transaction, roleType dax.RoleT // Deserialize rows into WorkerInfo objects. workerInfos := make(dax.WorkerInfos, 0) - prefix := []byte(fmt.Sprintf(prefixFmtWorkersDB, roleType, qdbid.Key())) + 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 { diff --git a/dax/controller/controller.go b/dax/controller/controller.go index 1046b2f93..c5244a8a2 100644 --- a/dax/controller/controller.go +++ b/dax/controller/controller.go @@ -201,24 +201,31 @@ func (c *Controller) RegisterNodes(ctx context.Context, nodes ...*dax.Node) erro // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) // For the addresses which are being added, set their method to "reset". - for i := range addressMethods { + for i := range addrMethods { for j := range nodes { - if addressMethods[i].address == nodes[j].Address { - addressMethods[i].method = dax.DirectiveMethodReset + if addrMethods[i].address == nodes[j].Address { + addrMethods[i].method = dax.DirectiveMethodReset } } } - // Get the current job assignments for this worker and send that to the node - // as a Directive. - if err := c.sendDirectives(tx, addressMethods...); err != nil { + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") + } + + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { return NewErrDirectiveSendFailure(err.Error()) } - return tx.Commit() + return nil } // RegisterNode adds a node to the controller's list of registered @@ -343,27 +350,34 @@ func (c *Controller) DeregisterNodes(ctx context.Context, addresses ...dax.Addre // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) // For the addresses which are being removed, set their method to "reset". - for i := range addressMethods { + for i := range addrMethods { for j := range addresses { - if addressMethods[i].address == addresses[j] { - addressMethods[i].method = dax.DirectiveMethodReset + if addrMethods[i].address == addresses[j] { + addrMethods[i].method = dax.DirectiveMethodReset } } } - // Get the current job assignments for these workers and send them to the - // nodes as Directives. - if err := c.sendDirectives(tx, addressMethods...); err != nil { + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") + } + + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { return NewErrDirectiveSendFailure(err.Error()) } - return tx.Commit() + return nil } -func (c *Controller) nodesTranslateReadOrWrite(tx dax.Transaction, role *dax.TranslateRole, qdbid dax.QualifiedDatabaseID, createMissing bool, asWrite bool) ([]dax.AssignedNode, bool, error) { +func (c *Controller) nodesTranslateReadOrWrite(tx dax.Transaction, role *dax.TranslateRole, qdbid dax.QualifiedDatabaseID, createMissing bool, asWrite bool) ([]dax.AssignedNode, bool, []*dax.Directive, error) { qtid := role.TableKey.QualifiedTableID() roleType := dax.RoleTypeTranslate @@ -375,7 +389,7 @@ func (c *Controller) nodesTranslateReadOrWrite(tx dax.Transaction, role *dax.Tra workers, err := c.Balancer.WorkersForJobs(tx, roleType, qdbid, inJobs.Sorted()...) if err != nil { - return nil, false, errors.Wrap(err, "getting workers for jobs") + return nil, false, nil, errors.Wrap(err, "getting workers for jobs") } outJobs := dax.NewSet[dax.Job]() @@ -387,9 +401,11 @@ func (c *Controller) nodesTranslateReadOrWrite(tx dax.Transaction, role *dax.Tra missed := inJobs.Minus(outJobs).Sorted() if !createMissing && len(missed) > 0 { - return nil, false, NewErrUnassignedJobs(missed) + return nil, false, nil, NewErrUnassignedJobs(missed) } + var directives []*dax.Directive + // If any provided jobs were not returned in the WorkersForJobs request, // then create those. if createMissing && len(missed) > 0 { @@ -398,7 +414,7 @@ func (c *Controller) nodesTranslateReadOrWrite(tx dax.Transaction, role *dax.Tra // directives sent) to workers. In that case, we need to abort this // method run and notify the caller to rety as a write. if !asWrite { - return nil, true, nil + return nil, true, nil, nil } sort.Slice(missed, func(i, j int) bool { return missed[i] < missed[j] }) @@ -407,11 +423,11 @@ func (c *Controller) nodesTranslateReadOrWrite(tx dax.Transaction, role *dax.Tra for _, job := range missed { j, err := decodePartition(job) if err != nil { - return nil, false, NewErrInternal(err.Error()) + return nil, false, nil, NewErrInternal(err.Error()) } diffs, err := c.Balancer.AddJobs(tx, roleType, qtid, j.Job()) if err != nil { - return nil, false, errors.Wrap(err, "adding job") + return nil, false, nil, errors.Wrap(err, "adding job") } for _, diff := range diffs { workerSet.Add(dax.Address(diff.Address)) @@ -420,26 +436,27 @@ func (c *Controller) nodesTranslateReadOrWrite(tx dax.Transaction, role *dax.Tra // Convert the slice of addresses into a slice of addressMethod // containing the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return nil, false, NewErrDirectiveSendFailure(err.Error()) + directives, err = c.buildDirectives(tx, addrMethods) + if err != nil { + return nil, false, nil, errors.Wrap(err, "building directives") } // Re-run WorkersForJobs. workers, err = c.Balancer.WorkersForJobs(tx, roleType, qdbid, inJobs.Sorted()...) if err != nil { - return nil, false, errors.Wrap(err, "getting workers for jobs") + return nil, false, nil, errors.Wrap(err, "getting workers for jobs") } } nodes, err := c.translateWorkersToAssignedNodes(tx, workers) - return nodes, false, errors.Wrap(err, "converting to assigned nodes") + return nodes, false, directives, errors.Wrap(err, "converting to assigned nodes") } // nodesComputeReadOrWrite contains the logic for the c.nodesCompute() method, // but it supports being called with either a read or write lock. -func (c *Controller) nodesComputeReadOrWrite(tx dax.Transaction, role *dax.ComputeRole, qdbid dax.QualifiedDatabaseID, createMissing bool, asWrite bool) ([]dax.AssignedNode, bool, error) { +func (c *Controller) nodesComputeReadOrWrite(tx dax.Transaction, role *dax.ComputeRole, qdbid dax.QualifiedDatabaseID, createMissing bool, asWrite bool) ([]dax.AssignedNode, bool, []*dax.Directive, error) { qtid := role.TableKey.QualifiedTableID() roleType := dax.RoleTypeCompute @@ -451,7 +468,7 @@ func (c *Controller) nodesComputeReadOrWrite(tx dax.Transaction, role *dax.Compu workers, err := c.Balancer.WorkersForJobs(tx, roleType, qdbid, inJobs.Sorted()...) if err != nil { - return nil, false, errors.Wrap(err, "getting workers for jobs") + return nil, false, nil, errors.Wrap(err, "getting workers for jobs") } // figure out if any jobs in the role have no workers assigned @@ -464,9 +481,11 @@ func (c *Controller) nodesComputeReadOrWrite(tx dax.Transaction, role *dax.Compu missed := inJobs.Minus(outJobs).Sorted() if !createMissing && len(missed) > 0 { - return nil, false, NewErrUnassignedJobs(missed) + return nil, false, nil, NewErrUnassignedJobs(missed) } + var directives []*dax.Directive + // If any provided jobs were not returned in the WorkersForJobs request, // then create those. if createMissing && len(missed) > 0 { @@ -475,7 +494,7 @@ func (c *Controller) nodesComputeReadOrWrite(tx dax.Transaction, role *dax.Compu // sent) to workers. In that case, we need to abort this method run and // notify the caller to rety as a write. if !asWrite { - return nil, true, nil + return nil, true, nil, nil } sort.Slice(missed, func(i, j int) bool { return missed[i] < missed[j] }) @@ -484,11 +503,11 @@ func (c *Controller) nodesComputeReadOrWrite(tx dax.Transaction, role *dax.Compu for _, job := range missed { j, err := decodeShard(job) if err != nil { - return nil, false, NewErrInternal(err.Error()) + return nil, false, nil, NewErrInternal(err.Error()) } diffs, err := c.Balancer.AddJobs(tx, roleType, qtid, j.Job()) if err != nil { - return nil, false, errors.Wrap(err, "adding job") + return nil, false, nil, errors.Wrap(err, "adding job") } for _, diff := range diffs { workerSet.Add(dax.Address(diff.Address)) @@ -497,21 +516,22 @@ func (c *Controller) nodesComputeReadOrWrite(tx dax.Transaction, role *dax.Compu // Convert the slice of addresses into a slice of addressMethod // containing the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return nil, false, NewErrDirectiveSendFailure(err.Error()) + directives, err = c.buildDirectives(tx, addrMethods) + if err != nil { + return nil, false, nil, errors.Wrap(err, "building directives") } // Re-run WorkersForJobs. workers, err = c.Balancer.WorkersForJobs(tx, roleType, qdbid, inJobs.Sorted()...) if err != nil { - return nil, false, errors.Wrap(err, "getting workers for jobs") + return nil, false, nil, errors.Wrap(err, "getting workers for jobs") } } nodes, err := c.computeWorkersToAssignedNodes(tx, workers) - return nodes, false, errors.Wrap(err, "converting to assigned nodes") + return nodes, false, directives, errors.Wrap(err, "converting to assigned nodes") } func (c *Controller) computeWorkersToAssignedNodes(tx dax.Transaction, workers []dax.WorkerInfo) ([]dax.AssignedNode, error) { @@ -627,14 +647,21 @@ func (c *Controller) DropDatabase(ctx context.Context, qdbid dax.QualifiedDataba // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - // Finally, send the directives based on the dropTable calls. - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") } - return tx.Commit() + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return NewErrDirectiveSendFailure(err.Error()) + } + return nil } // DatabaseByName returns the database for the given name. @@ -693,12 +720,20 @@ func (c *Controller) SetDatabaseOption(ctx context.Context, qdbid dax.QualifiedD // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") } - return tx.Commit() + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return NewErrDirectiveSendFailure(err.Error()) + } + return nil } func (c *Controller) Databases(ctx context.Context, orgID dax.OrganizationID, ids ...dax.DatabaseID) ([]*dax.QualifiedDatabase, error) { @@ -715,6 +750,7 @@ func (c *Controller) Databases(ctx context.Context, orgID dax.OrganizationID, id // CreateTable adds a table to the schemar, and then sends directives to all // affected nodes based on the change. func (c *Controller) CreateTable(ctx context.Context, qtbl *dax.QualifiedTable) error { + c.logger.Debugf("CreateTable %+v", qtbl) tx, err := c.BoltDB.BeginTx(ctx, true) if err != nil { return errors.Wrap(err, "beginning tx") @@ -731,6 +767,8 @@ func (c *Controller) CreateTable(ctx context.Context, qtbl *dax.QualifiedTable) return errors.Wrapf(err, "creating table: %s", qtbl) } + var addrMethods []addressMethod + // If the table is keyed, add partitions to the balancer. if qtbl.StringKeys() { // workerSet maintains the set of workers which have a job assignment change @@ -758,11 +796,7 @@ func (c *Controller) CreateTable(ctx context.Context, qtbl *dax.QualifiedTable) // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) - } + addrMethods = applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) } // This is more FieldVersion hackery. Even if the table is not keyed, we @@ -790,14 +824,22 @@ func (c *Controller) CreateTable(ctx context.Context, qtbl *dax.QualifiedTable) // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) - } + addrMethods = applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) } - return tx.Commit() + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") + } + + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return NewErrDirectiveSendFailure(err.Error()) + } + return nil } // DropTable removes a table from the schema and sends directives to all affected @@ -816,13 +858,21 @@ func (c *Controller) DropTable(ctx context.Context, qtid dax.QualifiedTableID) e // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(addrs.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(addrs.SortedSlice(), dax.DirectiveMethodDiff) - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") } - return tx.Commit() + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return NewErrDirectiveSendFailure(err.Error()) + } + return nil } // dropTable removes a table from the schema and sends directives to all affected @@ -921,35 +971,32 @@ func (c *Controller) RemoveShards(ctx context.Context, qtid dax.QualifiedTableID // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) - } - - return tx.Commit() -} - -// sendDirectives sends a directive (based on the current balancer state) to -// each of the nodes provided. -func (c *Controller) sendDirectives(tx dax.Transaction, addrs ...addressMethod) error { - // If nodes is empty, return early. - if len(addrs) == 0 { - return nil - } - - directives, err := c.buildDirectives(tx, addrs) + directives, err := c.buildDirectives(tx, addrMethods) if err != nil { return errors.Wrap(err, "building directives") } + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return NewErrDirectiveSendFailure(err.Error()) + } + + return nil +} + +func (c *Controller) sendDirectives2(ctx context.Context, directives []*dax.Directive) error { errs := make([]error, len(directives)) var eg errgroup.Group for i, dir := range directives { i := i dir := dir eg.Go(func() error { - errs[i] = c.Director.SendDirective(tx.Context(), dir) + errs[i] = c.Director.SendDirective(ctx, dir) if dir.IsEmpty() { errs[i] = nil } @@ -957,8 +1004,7 @@ func (c *Controller) sendDirectives(tx dax.Transaction, addrs ...addressMethod) }) } - err = eg.Wait() - if err != nil { + if err := eg.Wait(); err != nil { errCount := 0 for _, err := range errs { if err != nil { @@ -1009,6 +1055,11 @@ func applyAddressMethod(addrs []dax.Address, method dax.DirectiveMethod) []addre // buildDirectives builds a list of directives for the given addrs (i.e. nodes) // using information (i.e. current state) from the balancers. func (c *Controller) buildDirectives(tx dax.Transaction, addrs []addressMethod) ([]*dax.Directive, error) { + // If nodes is empty, return early. + if len(addrs) == 0 { + return nil, nil + } + directives := make([]*dax.Directive, len(addrs)) for i, addressMethod := range addrs { @@ -1345,7 +1396,7 @@ func (c *Controller) ComputeNodes(ctx context.Context, qtid dax.QualifiedTableID return computeNodes, nil } - assignedNodes, _, err := c.nodesComputeReadOrWrite(tx, role, qdbid, false, false) + assignedNodes, _, _, err := c.nodesComputeReadOrWrite(tx, role, qdbid, false, false) if err != nil { return nil, errors.Wrap(err, "getting compute nodes read or write") } @@ -1407,7 +1458,7 @@ func (c *Controller) TranslateNodes(ctx context.Context, qtid dax.QualifiedTable return translateNodes, nil } - assignedNodes, _, err := c.nodesTranslateReadOrWrite(tx, role, qdbid, false, false) + assignedNodes, _, _, err := c.nodesTranslateReadOrWrite(tx, role, qdbid, false, false) if err != nil { return nil, errors.Wrap(err, "getting translate nodes read or write") } @@ -1474,14 +1525,15 @@ func (c *Controller) IngestPartition(ctx context.Context, qtid dax.QualifiedTabl // Verify that the table exists. if _, err := c.Schemar.Table(tx, qtid); err != nil { - return "", err + return "", errors.Wrap(err, "getting table") } - nodes, retryAsWrite, err := c.nodesTranslateReadOrWrite(tx, role, qdbid, true, false) + nodes, retryAsWrite, _, err := c.nodesTranslateReadOrWrite(tx, role, qdbid, true, false) if err != nil { return "", errors.Wrap(err, "getting translate nodes read or write") } + var directives []*dax.Directive // If it's writable, and we couldn't find all the partitions with just a // read, try again with a write transaction. if retryAsWrite { @@ -1493,7 +1545,7 @@ func (c *Controller) IngestPartition(ctx context.Context, qtid dax.QualifiedTabl } defer tx.Rollback() - nodes, _, err = c.nodesTranslateReadOrWrite(tx, role, qdbid, true, true) + nodes, _, directives, err = c.nodesTranslateReadOrWrite(tx, role, qdbid, true, true) if err != nil { return "", errors.Wrap(err, "getting translate nodes read or write retry") } @@ -1531,7 +1583,15 @@ func (c *Controller) IngestPartition(ctx context.Context, qtid dax.QualifiedTabl // Only commit if the transaction is writable. if retryAsWrite { - return node.Address, tx.Commit() + if err := tx.Commit(); err != nil { + return node.Address, errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return node.Address, NewErrDirectiveSendFailure(err.Error()) + } + + return node.Address, nil } return node.Address, nil @@ -1561,11 +1621,13 @@ func (c *Controller) IngestShard(ctx context.Context, qtid dax.QualifiedTableID, return "", err } - nodes, retryAsWrite, err := c.nodesComputeReadOrWrite(tx, role, qdbid, true, false) + nodes, retryAsWrite, _, err := c.nodesComputeReadOrWrite(tx, role, qdbid, true, false) if err != nil { return "", errors.Wrap(err, "getting compute nodes read or write") } + var directives []*dax.Directive + // If it's writable, and we couldn't find all the partitions with just a // read, try again with a write transaction. if retryAsWrite { @@ -1577,7 +1639,7 @@ func (c *Controller) IngestShard(ctx context.Context, qtid dax.QualifiedTableID, } defer tx.Rollback() - nodes, _, err = c.nodesComputeReadOrWrite(tx, role, qdbid, true, true) + nodes, _, directives, err = c.nodesComputeReadOrWrite(tx, role, qdbid, true, true) if err != nil { return "", errors.Wrap(err, "getting compute nodes read or write retry") } @@ -1612,7 +1674,13 @@ func (c *Controller) IngestShard(ctx context.Context, qtid dax.QualifiedTableID, // Only commit if the transaction is writable. if retryAsWrite { - return node.Address, tx.Commit() + if err := tx.Commit(); err != nil { + return node.Address, errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return node.Address, NewErrDirectiveSendFailure(err.Error()) + } } return node.Address, nil @@ -1671,14 +1739,21 @@ func (c *Controller) CreateField(ctx context.Context, qtid dax.QualifiedTableID, // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - // Send a directive to any compute node responsible for this field. - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") } - return tx.Commit() + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return NewErrDirectiveSendFailure(err.Error()) + } + return nil } func (c *Controller) DropField(ctx context.Context, qtid dax.QualifiedTableID, fldName dax.FieldName) error { @@ -1732,14 +1807,21 @@ func (c *Controller) DropField(ctx context.Context, qtid dax.QualifiedTableID, f // Convert the slice of addresses into a slice of addressMethod containing // the appropriate method. - addressMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) + addrMethods := applyAddressMethod(workerSet.SortedSlice(), dax.DirectiveMethodDiff) - // Send a directive to any compute node responsible for this field. - if err := c.sendDirectives(tx, addressMethods...); err != nil { - return NewErrDirectiveSendFailure(err.Error()) + directives, err := c.buildDirectives(tx, addrMethods) + if err != nil { + return errors.Wrap(err, "building directives") } - return tx.Commit() + if err := tx.Commit(); err != nil { + return errors.Wrap(err, "committing") + } + + if err := c.sendDirectives2(ctx, directives); err != nil { + return NewErrDirectiveSendFailure(err.Error()) + } + return nil } ////////////////////////////////// @@ -1763,6 +1845,16 @@ func (c *Controller) DebugNodes(ctx context.Context) ([]*dax.Node, error) { return c.Balancer.Nodes(tx) } +func (c *Controller) CurrentState(ctx context.Context) ([]dax.WorkerInfo, error) { + tx, err := c.BoltDB.BeginTx(ctx, false) + if err != nil { + return nil, errors.Wrap(err, "beginning tx") + } + defer tx.Rollback() + + return c.Balancer.CurrentState(tx, "", dax.QualifiedDatabaseID{}) +} + // sanitizeQTID populates Table.ID (by looking up the table, by name, in // schemar) for a given table having only a Name value, but no ID. func (c *Controller) sanitizeQTID(tx dax.Transaction, qtid *dax.QualifiedTableID) error { diff --git a/dax/controller/http/director.go b/dax/controller/http/director.go index 429b80d72..9f213e905 100644 --- a/dax/controller/http/director.go +++ b/dax/controller/http/director.go @@ -71,6 +71,7 @@ type DirectorConfig struct { func (d *Director) SendDirective(ctx context.Context, dir *dax.Directive) error { url := fmt.Sprintf("%s/%s", dir.Address.WithScheme("http"), d.directivePath) d.logger.Printf("SEND HTTP directive to: %s\n", url) + d.logger.Debugf("Directive: %+v", dir) // Encode the request. postBody, err := json.Marshal(dir) diff --git a/dax/controller/http/handler.go b/dax/controller/http/handler.go index edc863a87..79c9794b7 100644 --- a/dax/controller/http/handler.go +++ b/dax/controller/http/handler.go @@ -51,6 +51,7 @@ func Handler(c *controller.Controller) http.Handler { // debug endpoints router.HandleFunc("/debug/nodes", server.getDebugNodes).Methods("GET").Name("GetDebugNodes") + router.HandleFunc("/debug/balancer", server.getDebugBalancer).Methods("GET").Name("getDebugBalancer") return router } @@ -758,6 +759,19 @@ func (s *server) getDebugNodes(w http.ResponseWriter, r *http.Request) { } +func (s *server) getDebugBalancer(w http.ResponseWriter, r *http.Request) { + nodes, err := s.controller.CurrentState(r.Context()) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + if err := json.NewEncoder(w).Encode(nodes); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + +} + // ComputeNodesRequest is used to specify the table/shards to consider in the // ComputeNodes method call. If IsWrite is true, shards which are not currently // being managed by the underlying Controller will be added to (registered with) diff --git a/disco/disco.go b/disco/disco.go index 551d03f6e..285c3ba32 100644 --- a/disco/disco.go +++ b/disco/disco.go @@ -271,12 +271,6 @@ func (*nopSchemator) CreateView(ctx context.Context, index, field, view string) // DeleteView is a no-op implementation of the Schemator DeleteView method. func (*nopSchemator) DeleteView(ctx context.Context, index, field, view string) error { return nil } -// InMemSchemator represents a Schemator that manages the schema in memory. The -// intention is that this would be used for testing. -var InMemSchemator Schemator = &inMemSchemator{ - schema: make(Schema), -} - type inMemSchemator struct { mu sync.RWMutex schema Schema @@ -451,8 +445,10 @@ func (s *inMemSchemator) DeleteView(ctx context.Context, index, field, view stri return nil } -var InMemSharder Sharder = &inMemSharder{ - shards: make(map[string][]byte), +func NewInMemSharder() *inMemSharder { + return &inMemSharder{ + shards: make(map[string][]byte), + } } type inMemSharder struct { diff --git a/holder.go b/holder.go index 7a0751bc9..3d8459583 100644 --- a/holder.go +++ b/holder.go @@ -275,7 +275,7 @@ func DefaultHolderConfig() *HolderConfig { TranslationSyncer: NopTranslationSyncer, Serializer: GobSerializer, Schemator: disco.NewInMemSchemator(), - Sharder: disco.InMemSharder, + Sharder: disco.NewInMemSharder(), CacheFlushInterval: defaultCacheFlushInterval, Logger: logger.NopLogger, StorageConfig: storage.NewDefaultConfig(), diff --git a/server/server.go b/server/server.go index 0ebdb513d..fb50cb2c8 100644 --- a/server/server.go +++ b/server/server.go @@ -606,8 +606,8 @@ func (m *Command) setupServer() error { disco.NewLocalNoder([]*disco.Node{ {ID: nodeID, URI: *advertiseURI, IsPrimary: true, State: disco.NodeStateStarted}, }), - disco.InMemSharder, - disco.InMemSchemator, + disco.NewInMemSharder(), + disco.NewInMemSchemator(), ), pilosa.OptServerNodeID(nodeID), )