From 36f2dcce9e35891c49b602a2458977e5a050af3d Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Tue, 3 Jan 2023 13:57:01 -0600 Subject: [PATCH] clean up TODOs. adds a control channel for on-demand snapshotting (cherry picked from commit 4cc1667399340a5c0de5f17f1eb0a0ec20fae891) --- dax/mds/controller/controller.go | 52 ++++----------------------- dax/mds/controller/controller_test.go | 17 --------- dax/mds/controller/snapping_turtle.go | 4 ++- 3 files changed, 9 insertions(+), 64 deletions(-) diff --git a/dax/mds/controller/controller.go b/dax/mds/controller/controller.go index 529537db7..df71b5bda 100644 --- a/dax/mds/controller/controller.go +++ b/dax/mds/controller/controller.go @@ -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") diff --git a/dax/mds/controller/controller_test.go b/dax/mds/controller/controller_test.go index 4f36b6d1e..e57148b16 100644 --- a/dax/mds/controller/controller_test.go +++ b/dax/mds/controller/controller_test.go @@ -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: "", diff --git a/dax/mds/controller/snapping_turtle.go b/dax/mds/controller/snapping_turtle.go index 4c592a37c..153fba832 100644 --- a/dax/mds/controller/snapping_turtle.go +++ b/dax/mds/controller/snapping_turtle.go @@ -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() } }