mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
clean up TODOs. adds a control channel for on-demand snapshotting
(cherry picked from commit 4cc1667399)
This commit is contained in:
parent
e09a9dca47
commit
36f2dcce9e
3 changed files with 9 additions and 64 deletions
|
|
@ -45,6 +45,7 @@ type Controller struct {
|
|||
registrationBatchTimeout time.Duration
|
||||
nodeChan chan *dax.Node
|
||||
snappingTurtleTimeout time.Duration
|
||||
snapControl chan struct{}
|
||||
stopping chan struct{}
|
||||
|
||||
logger logger.Logger
|
||||
|
|
@ -70,6 +71,7 @@ func New(cfg Config) *Controller {
|
|||
registrationBatchTimeout: cfg.RegistrationBatchTimeout,
|
||||
nodeChan: make(chan *dax.Node, 10),
|
||||
snappingTurtleTimeout: cfg.SnappingTurtleTimeout,
|
||||
snapControl: make(chan struct{}),
|
||||
|
||||
stopping: make(chan struct{}),
|
||||
logger: logger.NopLogger,
|
||||
|
|
@ -107,7 +109,7 @@ func New(cfg Config) *Controller {
|
|||
// Run starts long running subroutines.
|
||||
func (c *Controller) Run() error {
|
||||
go c.nodeRegistrationRoutine(c.nodeChan, c.registrationBatchTimeout)
|
||||
go c.snappingTurtleRoutine(c.snappingTurtleTimeout)
|
||||
go c.snappingTurtleRoutine(c.snappingTurtleTimeout, c.snapControl)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
|
@ -1206,51 +1208,11 @@ func (c *Controller) InitializePoller(ctx context.Context) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// SnapshotTable snapshots a table.
|
||||
// TODO(jaffee): do we need to re-implement this w/o version store? Or can we just remove entirely?
|
||||
// SnapshotTable snapshots a table. It might also snapshot everything
|
||||
// else... no guarantees here, only used in tests as of this writing.
|
||||
func (c *Controller) SnapshotTable(ctx context.Context, qtid dax.QualifiedTableID) error {
|
||||
c.snapControl <- struct{}{}
|
||||
return nil
|
||||
// shards, ok, err := c.versionStore.Shards(ctx, qtid)
|
||||
// if err != nil {
|
||||
// return errors.Wrap(err, "getting shards from version store")
|
||||
// } else if !ok {
|
||||
// return errors.New(errors.ErrUncoded, "got false back from versionStore.Shards")
|
||||
// }
|
||||
|
||||
// for _, shard := range shards {
|
||||
// if err := c.SnapshotShardData(ctx, qtid, shard.Num); err != nil {
|
||||
// return errors.Wrapf(err, "snapshotting shard data: qtid: %s, shard: %d", qtid, shard.Num)
|
||||
// }
|
||||
// }
|
||||
|
||||
// partitions, ok, err := c.versionStore.Partitions(ctx, qtid)
|
||||
// if err != nil {
|
||||
// return errors.Wrap(err, "getting partitions from version store")
|
||||
// } else if !ok {
|
||||
// return errors.New(errors.ErrUncoded, "got false back from versionStore.Partitions")
|
||||
// }
|
||||
|
||||
// for _, part := range partitions {
|
||||
// if err := c.SnapshotTableKeys(ctx, qtid, part.Num); err != nil {
|
||||
// return errors.Wrapf(err, "snapshotting table keys: qtid: %s, partition: %d", qtid, part.Num)
|
||||
// }
|
||||
// }
|
||||
|
||||
// fields, ok, err := c.versionStore.Fields(ctx, qtid)
|
||||
// if err != nil {
|
||||
// return errors.Wrap(err, "getting fields from version store")
|
||||
// } else if !ok {
|
||||
// return errors.New(errors.ErrUncoded, "got false back from versionStore.Fields")
|
||||
// }
|
||||
|
||||
// for _, fld := range fields {
|
||||
// if fld.Name != "_id" {
|
||||
// if err := c.SnapshotFieldKeys(ctx, qtid, fld.Name); err != nil {
|
||||
// return errors.Wrapf(err, "snapshotting field keys: qtid: %s, field: %s", qtid, fld.Name)
|
||||
// }
|
||||
// }
|
||||
// }
|
||||
// return nil
|
||||
}
|
||||
|
||||
// SnapshotShardData forces the compute node responsible for the given shard to
|
||||
|
|
@ -1477,8 +1439,6 @@ func (c *Controller) CreateField(ctx context.Context, qtid dax.QualifiedTableID,
|
|||
workerSet.Add(dax.Address(w.ID))
|
||||
}
|
||||
|
||||
// TODO if there isn't already a worker for partition 0, do we need to add it?
|
||||
|
||||
// Get the list of workers responsible for shard data for this table.
|
||||
if state, err := c.ComputeBalancer.CurrentState(ctx); err != nil {
|
||||
return errors.Wrap(err, "getting current compute state")
|
||||
|
|
|
|||
|
|
@ -108,16 +108,6 @@ func TestController(t *testing.T) {
|
|||
exp = []*dax.Directive{}
|
||||
assert.Equal(t, exp, director.flush())
|
||||
|
||||
// Add the same non-keyed table again.
|
||||
// TODO(jaffee) figure out what this should do
|
||||
// err := con.CreateTable(ctx, tbl0)
|
||||
// if assert.Error(t, err) {
|
||||
// assert.True(t, errors.Is(err, dax.ErrTableIDExists))
|
||||
// }
|
||||
|
||||
exp = []*dax.Directive{}
|
||||
assert.Equal(t, exp, director.flush())
|
||||
|
||||
// Add a shard.
|
||||
assert.NoError(t, con.AddShards(ctx, tbl0.QualifiedID(), 0))
|
||||
|
||||
|
|
@ -817,13 +807,6 @@ func TestController(t *testing.T) {
|
|||
assert.True(t, errors.Is(err, dax.ErrTableIDDoesNotExist))
|
||||
}
|
||||
|
||||
// Add shards to a table which doesn't exist.
|
||||
// TODO(jaffee) figure out what this should do
|
||||
// err = con.AddShards(ctx, invalidQtid, dax.NewShardNums(1, 2)...)
|
||||
// if assert.Error(t, err) {
|
||||
// assert.True(t, errors.Is(err, dax.ErrTableIDDoesNotExist))
|
||||
// }
|
||||
|
||||
// Register an invalid node.
|
||||
nodeX := &dax.Node{
|
||||
Address: "",
|
||||
|
|
|
|||
|
|
@ -7,7 +7,7 @@ import (
|
|||
"github.com/molecula/featurebase/v3/dax"
|
||||
)
|
||||
|
||||
func (c *Controller) snappingTurtleRoutine(period time.Duration) {
|
||||
func (c *Controller) snappingTurtleRoutine(period time.Duration, control chan struct{}) {
|
||||
if period == 0 {
|
||||
return // disable automatic snapshotting
|
||||
}
|
||||
|
|
@ -20,6 +20,8 @@ func (c *Controller) snappingTurtleRoutine(period time.Duration) {
|
|||
return
|
||||
case <-ticker.C:
|
||||
c.snapAll()
|
||||
case <-control:
|
||||
c.snapAll()
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue