diff --git a/cluster.go b/cluster.go index 9695a4dd1..429c1c970 100644 --- a/cluster.go +++ b/cluster.go @@ -1431,11 +1431,11 @@ func (c *cluster) unprotectedGenerateResizeJobByAction(nodeAction nodeAction) (* } for _, node := range toCluster.nodes { - // If a host doesn't need to request data, mark it as complete. - if len(fragmentSourcesByNode[node.ID]) == 0 && len(translationSourcesByNode[node.ID]) == 0 { - j.IDs[node.ID] = true - continue - } + // We may send a resize instruction that has no sources that + // the node needs to read from -- for instance, if there's no + // data in any fragments it would process. But it still needs + // to get the NodeStatus to pick up the schema so it knows + // about existing indexes. instr := &ResizeInstruction{ JobID: j.ID, Node: toCluster.unprotectedNodeByID(node.ID), @@ -2410,7 +2410,7 @@ func (c *cluster) translateIndexIDs(ctx context.Context, indexName string, ids [ } func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idSet map[uint64]struct{}) (map[uint64]string, error) { - idMap := make(map[uint64]string) + idMap := make(map[uint64]string, len(idSet)) index := c.holder.Index(indexName) if index == nil { diff --git a/server/cluster_test.go b/server/cluster_test.go index 0899b8e97..d788de803 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -199,22 +199,20 @@ func TestClusterResize_AddNode(t *testing.T) { t.Fatal(err) } + col := pilosa.ShardWidth + 20 + // Write data on first node. - if _, err := m0.Query(t, "i", "", ` + if _, err := m0.Queryf(t, "i", "", ` Set(1, f=1) - Set(1300000, f=1) - `); err != nil { + Set(%d, f=1) + `, col); err != nil { t.Fatal(err) } // exp is the expected result for the Row queries that follow. - exp := `{"results":[{"attrs":{},"columns":[1,1300000]}]}` + "\n" + exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col) // Verify the data exists on the single node. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 m1 := test.NewCommandNode(false) @@ -233,16 +231,57 @@ func TestClusterResize_AddNode(t *testing.T) { } // Verify the data exists on both nodes. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) + m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) + }) + t.Run("OneShard", func(t *testing.T) { + // Configure node0 + m0 := test.MustRunCluster(t, 1)[0] + defer m0.Close() + + seed := m0.GossipAddress() + + // Create a client for each node. + client0 := m0.Client() + + // Create indexes and fields on one node. + if err := client0.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } - if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil { + } else if err := client0.CreateField(context.Background(), "i", "f"); err != nil { t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) } + + // Write data on first node. + if _, err := m0.Query(t, "i", "", ` + Set(1, f=1) + `); err != nil { + t.Fatal(err) + } + // exp is the expected result for the Row queries that follow. + exp := `{"results":[{"attrs":{},"columns":[1]}]}` + + // Verify the data exists on the single node. + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) + + // Configure node1 + m1 := test.NewCommandNode(false) + m1.Config.Gossip.Port = "0" + m1.Config.Gossip.Seeds = []string{seed} + err := m1.Start() + if err != nil { + t.Fatalf("starting second main: %v", err) + } + defer m1.Close() + + if !checkClusterState(m0, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node0 cluster state: %s", m0.API.State()) + } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { + t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) + } + + // Verify the data exists on both nodes. + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) + m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) }) t.Run("SkippedShard", func(t *testing.T) { // Configure node0 @@ -261,23 +300,21 @@ func TestClusterResize_AddNode(t *testing.T) { t.Fatal(err) } + col := pilosa.ShardWidth*2 + 20 + // Write data on first node. Note that no data is placed on shard 1. - if _, err := m0.Query(t, "i", "", ` + if _, err := m0.Queryf(t, "i", "", ` Set(1, f=1) - Set(2400000, f=1) - `); err != nil { + Set(%d, f=1) + `, col); err != nil { t.Fatal(err) } // exp is the expected result for the Row queries that follow. - exp := `{"results":[{"attrs":{},"columns":[1,2400000]}]}` + "\n" + exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col) // Verify the data exists on the single node. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 m1 := test.NewCommandNode(false) @@ -296,16 +333,8 @@ func TestClusterResize_AddNode(t *testing.T) { } // Verify the data exists on both nodes. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } - if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) + m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) }) } @@ -371,23 +400,21 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { t.Fatal(err) } + col := pilosa.ShardWidth + 20 + // Write data on first node. - if _, err := m0.Query(t, "i", "", ` + if _, err := m0.Queryf(t, "i", "", ` Set(1, f=1) - Set(1300000, f=1) - `); err != nil { + Set(%d, f=1) + `, col); err != nil { t.Fatal(err) } // exp is the expected result for the Row queries that follow. - exp := `{"results":[{"attrs":{},"columns":[1,1300000]}]}` + "\n" + exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col) // Verify the data exists on the single node. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 m1 := test.NewCommandNode(false) @@ -411,16 +438,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } // Verify the data exists on both nodes. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } - if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) + m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) }) t.Run("SkippedShard", func(t *testing.T) { // Configure node0 @@ -439,23 +458,21 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { t.Fatal(err) } + col := pilosa.ShardWidth*2 + 20 + // Write data on first node. Note that no data is placed on shard 1. - if _, err := m0.Query(t, "i", "", ` + if _, err := m0.Queryf(t, "i", "", ` Set(1, f=1) - Set(2400000, f=1) - `); err != nil { + Set(%d, f=1) + `, col); err != nil { t.Fatal(err) } // exp is the expected result for the Row queries that follow. - exp := `{"results":[{"attrs":{},"columns":[1,2400000]}]}` + "\n" + exp := fmt.Sprintf(`{"results":[{"attrs":{},"columns":[1,%d]}]}`, col) // Verify the data exists on the single node. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 m1 := test.NewCommandNode(false) @@ -479,16 +496,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } // Verify the data exists on both nodes. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } - if res, err := m1.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) + m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) }) t.Run("WithIndexKeys", func(t *testing.T) { // Configure node0 @@ -516,14 +525,10 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } // exp is the expected result for the Row queries that follow. - exp := `{"results":[{"attrs":{},"columns":[],"keys":["col2","col1"]}]}` + "\n" + exp := `{"results":[{"attrs":{},"columns":[],"keys":["col2","col1"]}]}` // Verify the data exists on the single node. - if res, err := m0.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("unexpected result: %s", res) - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 m1 := test.NewCommandNode(false) @@ -545,15 +550,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { } else if !checkClusterState(m1, pilosa.ClusterStateNormal, 1000) { t.Fatalf("unexpected node1 cluster state: %s", m1.API.State()) } - - // Verify the data exists on both nodes. - for i, node := range []*test.Command{m0, m1} { - if res, err := node.Query(t, "i", "", `Row(f=1)`); err != nil { - t.Fatal(err) - } else if res != exp { - t.Fatalf("node%d expected: %s, but got: %s", i, exp, res) - } - } + m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) + m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) }) } diff --git a/server/handler_test.go b/server/handler_test.go index e893aa371..c318056a1 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -1290,6 +1290,7 @@ func TestCluster_TranslateStore(t *testing.T) { if err != nil { t.Fatalf("starting cluster 0: %v", err) } + defer cluster[0].Close() test.Do(t, "POST", cluster[0].URL()+"/index/i0", "{\"options\": {\"keys\": true}}") } @@ -1306,6 +1307,7 @@ func TestClusterTranslator(t *testing.T) { if err != nil { t.Fatalf("starting cluster 0: %v", err) } + defer cluster[0].Close() cluster[1] = test.NewCommandNode(false, server.OptCommandServerOptions( pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore), @@ -1318,6 +1320,7 @@ func TestClusterTranslator(t *testing.T) { if err != nil { t.Fatalf("starting cluster 1: %v", err) } + defer cluster[1].Close() test.Do(t, "POST", cluster[0].URL()+"/index/i0", "{\"options\": {\"keys\": true}}") test.Do(t, "POST", cluster[0].URL()+"/index/i0/field/f0", "{\"options\": {\"keys\": true}}") diff --git a/server/server.go b/server/server.go index 902abf472..df78a33bf 100644 --- a/server/server.go +++ b/server/server.go @@ -406,6 +406,17 @@ func (m *Command) setupNetworking() error { // get the host portion of addr to use for binding gossipHost := m.listenURI.Host m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) + if err != nil && gossipPort >= 32768 { + // In testing, we sometimes try to reuse an ephemeral port. + // Which probably works. If it doesn't, this test will take + // about a minute longer because we'll come back in from a + // new port. See also the gossip config in gossip/gossip.go. + // TODO: Maybe make that more configurable here. + m.logger.Printf("ephemeral port %d already occupied, switching to :0 (%v)", gossipPort, err) + m.Config.Gossip.Port = "0" + gossipPort = 0 + m.gossipTransport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) + } if err != nil { return errors.Wrap(err, "getting transport") } diff --git a/server/server_test.go b/server/server_test.go index 96d0fa527..3534fe0d0 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -630,15 +630,9 @@ func TestClusteringNodesReplica1(t *testing.T) { cluster := test.MustRunCluster(t, 3) defer cluster.Close() - var wait = true - for wait { - wait = false - for _, node := range cluster { - if node.API.State() != pilosa.ClusterStateNormal { - wait = true - } - } - time.Sleep(time.Millisecond * 1) + err := cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond) + if err != nil { + t.Fatalf("starting cluster: %v", err) } if err := cluster[2].Command.Close(); err != nil { @@ -665,14 +659,9 @@ func TestClusteringNodesReplica1(t *testing.T) { t.Fatalf("restarting node 2: %v", err) } - for wait { - wait = false - for _, node := range cluster { - if node.API.State() != pilosa.ClusterStateNormal { - wait = true - } - } - time.Sleep(time.Millisecond) + err = cluster.AwaitState(pilosa.ClusterStateNormal, 200*time.Millisecond) + if err != nil { + t.Fatalf("resuming normal operations: %v", err) } } @@ -685,24 +674,20 @@ func TestClusteringNodesReplica2(t *testing.T) { if err != nil { t.Fatalf("starting cluster: %v", err) } + defer cluster.Close() - var wait = true - for wait { - wait = false - for _, node := range cluster { - if node.API.State() != pilosa.ClusterStateNormal { - wait = true - } - } - time.Sleep(time.Millisecond * 1) + err = cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond) + if err != nil { + t.Fatalf("starting cluster: %v", err) } if err := cluster[2].Command.Close(); err != nil { t.Fatalf("closing third node: %v", err) } - if cluster[0].API.State() != pilosa.ClusterStateDegraded { - t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond) + if err != nil { + t.Fatalf("after closing first server: %v", err) } // confirm that cluster keeps accepting queries if replication > 1 @@ -715,8 +700,9 @@ func TestClusteringNodesReplica2(t *testing.T) { t.Fatalf("closing 2nd node: %v", err) } - if cluster[0].API.State() != pilosa.ClusterStateStarting { - t.Fatalf("expected state to be Starting, but got %s", cluster[0].API.State()) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateStarting, 100*time.Millisecond) + if err != nil { + t.Fatalf("after closing second server: %v", err) } if _, err := cluster[0].API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { @@ -739,8 +725,9 @@ func TestClusteringNodesReplica2(t *testing.T) { t.Fatalf("restarting node 2: %v", err) } - if cluster[0].API.State() != pilosa.ClusterStateDegraded { - t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond) + if err != nil { + t.Fatalf("after restarting first server: %v", err) } // Create new main with the same config. @@ -756,19 +743,12 @@ func TestClusteringNodesReplica2(t *testing.T) { // Run new program. if err := cluster[1].Start(); err != nil { - t.Fatalf("restarting node 2: %v", err) + t.Fatalf("restarting node 1: %v", err) } - defer cluster.Close() - - for wait { - wait = false - for _, node := range cluster { - if node.API.State() != pilosa.ClusterStateNormal { - wait = true - } - } - time.Sleep(time.Millisecond) + err = cluster.AwaitState(pilosa.ClusterStateNormal, 200*time.Microsecond) + if err != nil { + t.Fatalf("resuming normal operations: %v", err) } } @@ -781,32 +761,38 @@ func TestRemoveNodeAfterItDies(t *testing.T) { if err != nil { t.Fatalf("starting cluster: %v", err) } + // The anonymous function is necessary so that the slice + // passed to Close() as a receiver is the modified value + // of cluster, because we're removing the last entry from it + // below. + defer func() { + cluster.Close() + }() - var wait = true - for wait { - wait = false - for _, node := range cluster { - if node.API.State() != pilosa.ClusterStateNormal { - wait = true - } - } - time.Sleep(time.Millisecond * 1) + err = cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond) + if err != nil { + t.Fatalf("starting cluster: %v", err) } - if err := cluster[2].Command.Close(); err != nil { + // prevent double-closing cluster[2] from the deferred Close above + disabled, cluster := cluster[2], cluster[:2] + + if err := disabled.Command.Close(); err != nil { t.Fatalf("closing third node: %v", err) } - if cluster[0].API.State() != pilosa.ClusterStateDegraded { - t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond) + if err != nil { + t.Fatalf("starting cluster: %v", err) } - if _, err := cluster[0].API.RemoveNode(cluster[2].API.Node().ID); err != nil { + if _, err := cluster[0].API.RemoveNode(disabled.API.Node().ID); err != nil { t.Fatalf("removing failed node: %v", err) } - if cluster[0].API.State() != pilosa.ClusterStateNormal { - t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateNormal, 100*time.Millisecond) + if err != nil { + t.Fatalf("removing disabled node: %v", err) } hosts := cluster[0].API.Hosts(context.Background()) @@ -824,16 +810,10 @@ func TestRemoveConcurrentIndexCreation(t *testing.T) { if err != nil { t.Fatalf("starting cluster: %v", err) } - - var wait = true - for wait { - wait = false - for _, node := range cluster { - if node.API.State() != pilosa.ClusterStateNormal { - wait = true - } - } - time.Sleep(time.Millisecond * 1) + defer cluster.Close() + err = cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond) + if err != nil { + t.Fatalf("starting cluster: %v", err) } errc := make(chan error) @@ -846,11 +826,9 @@ func TestRemoveConcurrentIndexCreation(t *testing.T) { t.Fatalf("removing node: %v", err) } - for i := 0; cluster[0].API.State() != pilosa.ClusterStateNormal; i++ { - time.Sleep(time.Millisecond) - if i > 10 { - t.Fatalf("expected state to be DEGRADED, but got %s", cluster[0].API.State()) - } + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateNormal, 100*time.Millisecond) + if err != nil { + t.Fatalf("starting cluster: %v", err) } hosts := cluster[0].API.Hosts(context.Background()) diff --git a/test/cluster.go b/test/cluster.go index 8cd2a6991..babb3890f 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -16,9 +16,9 @@ package test import ( "context" + "fmt" "io/ioutil" "path" - "runtime" "strconv" "strings" "testing" @@ -140,10 +140,53 @@ func (c Cluster) Close() error { return nil } -// MustNewCluster creates a new cluster +// AwaitState waits for the cluster coordinator (assumed to be the first +// node) to reach a specified state. +func (c Cluster) AwaitCoordinatorState(expectedState string, timeout time.Duration) error { + if len(c) < 1 { + return errors.New("can't await coordinator state on an empty cluster") + } + return c[:1].AwaitState(expectedState, timeout) +} + +// ExceptionalState returns an error if any node in the cluster is not +// in the expected state. +func (c Cluster) ExceptionalState(expectedState string) error { + for _, node := range c { + state := node.API.State() + if state != expectedState { + return fmt.Errorf("node %q: state %s", node.ID(), state) + } + } + return nil +} + +// AwaitState waits for the whole cluster to reach a specified state. +func (c Cluster) AwaitState(expectedState string, timeout time.Duration) (err error) { + if len(c) < 1 { + return errors.New("can't await state of an empty cluster") + } + startTime := time.Now() + var elapsed time.Duration + for elapsed = 0; elapsed <= timeout; elapsed = time.Since(startTime) { + // Counterintuitive: We're returning if the err *is* nil, + // meaning we've reached the expected state. + if err = c.ExceptionalState(expectedState); err == nil { + return err + } + time.Sleep(1 * time.Millisecond) + } + return fmt.Errorf("waited %v for cluster to reach state %q: %v", + elapsed, expectedState, err) +} + +// MustNewCluster creates a new cluster. If opts contains only one +// slice of command options, those options are used with every node. +// If it is empty, default options are used. Otherwise, it must contain size +// slices of command options, which are used with corresponding nodes. func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster { tb.Helper() - c, err := newCluster(size, opts...) + c, err := newCluster(tb, size, opts...) if err != nil { tb.Fatalf("new cluster: %v", err) } @@ -151,7 +194,7 @@ func MustNewCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu } // newCluster creates a new cluster -func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) { +func newCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (Cluster, error) { if size == 0 { return nil, errors.New("cluster must contain at least one node") } @@ -160,26 +203,7 @@ func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) { } cluster := make(Cluster, size) - // try to find a Test function to use as the "name" for our node. - name := "node" - callers := make([]uintptr, 10) - n := runtime.Callers(2, callers) - callers = callers[:n] - for _, pc := range callers { - fn := runtime.FuncForPC(pc) - if fn != nil { - fnName := fn.Name() - sections := strings.Split(fnName, ".") - if len(sections) > 1 { - fnName = sections[2] - } - if strings.HasPrefix(fnName, "Test") { - name = "test" + fnName[4:] - break - } - } - } - _ = name + name := tb.Name() for i := 0; i < size; i++ { var commandOpts []server.CommandOption if len(opts) > 0 { @@ -197,8 +221,8 @@ func newCluster(size int, opts ...[]server.CommandOption) (Cluster, error) { } // runCluster creates and starts a new cluster -func runCluster(size int, opts ...[]server.CommandOption) (Cluster, error) { - cluster, err := newCluster(size, opts...) +func runCluster(tb testing.TB, size int, opts ...[]server.CommandOption) (Cluster, error) { + cluster, err := newCluster(tb, size, opts...) if err != nil { return nil, errors.Wrap(err, "new cluster") } @@ -209,7 +233,8 @@ func runCluster(size int, opts ...[]server.CommandOption) (Cluster, error) { return cluster, nil } -// MustRunCluster creates and starts a new cluster +// MustRunCluster creates and starts a new cluster. The opts parameter +// is slightly magical; see MustNewCluster. func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster { // We want tests to default to using the in-memory translate store, so we // prepend opts with that functional option. If a different translate store @@ -217,7 +242,7 @@ func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu opts = prependOpts(opts) tb.Helper() - c, err := runCluster(size, opts...) + c, err := runCluster(tb, size, opts...) if err != nil { tb.Fatalf("run cluster: %v", err) } diff --git a/test/pilosa.go b/test/pilosa.go index 0beb5b970..e592f058e 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -195,6 +195,9 @@ func (m *Command) MustRecalculateCaches(tb testing.TB) { // URL returns the base URL string for accessing the running program. func (m *Command) URL() string { return m.API.Node().URI.String() } +// ID returns the node ID used by the running program. +func (m *Command) ID() string { return m.API.Node().ID } + // Client returns a client to connect to the program. func (m *Command) Client() *http.InternalClient { return m.Server.InternalClient().(*http.InternalClient) @@ -202,13 +205,41 @@ func (m *Command) Client() *http.InternalClient { // Query executes a query against the program through the HTTP API. func (m *Command) Query(t *testing.T, index, rawQuery, query string) (string, error) { - resp := Do(t, "POST", m.URL()+fmt.Sprintf("/index/%s/query?", index)+rawQuery, query) + resp := Do(t, "POST", fmt.Sprintf("%s/index/%s/query?%s", m.URL(), index, rawQuery), query) if resp.StatusCode != gohttp.StatusOK { return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body) } return resp.Body, nil } +// Queryf is like Query, but with a format string. +func (m *Command) Queryf(t *testing.T, index, rawQuery, query string, params ...interface{}) (string, error) { + query = fmt.Sprintf(query, params...) + resp := Do(t, "POST", fmt.Sprintf("%s/index/%s/query?%s", m.URL(), index, rawQuery), query) + if resp.StatusCode != gohttp.StatusOK { + return "", fmt.Errorf("invalid status: %d, body=%s", resp.StatusCode, resp.Body) + } + return resp.Body, nil +} + +// QueryExpect executes a query against the program through the HTTP API, and +// confirms that it got an expected response. +func (m *Command) QueryExpect(t *testing.T, index, rawQuery, query string, expected string) { + resp := Do(t, "POST", fmt.Sprintf("%s/index/%s/query?%s", m.URL(), index, rawQuery), query) + if resp.StatusCode != gohttp.StatusOK { + t.Fatalf("invalid status from %s: %d, body=%q", m.ID(), resp.StatusCode, resp.Body) + } + last := len(resp.Body) - 1 + // Trim trailing newline so we don't need it to be present in the expected data. + if last >= 0 && resp.Body[last] == '\n' { + resp.Body = resp.Body[:last] + } + + if resp.Body != expected { + t.Fatalf("node %s, query %q: expected response %s, got %s", m.ID(), query, expected, resp.Body) + } +} + func (m *Command) QueryProtobuf(indexName string, query string) (*pilosa.QueryResponse, error) { var ser proto.Serializer queryReq := &pilosa.QueryRequest{ @@ -261,8 +292,6 @@ func (m *Command) RecalculateCaches(t *testing.T) error { return nil } -//////////////////////////////////////////////////////////////////////////////////// - // Do executes http.Do() with an http.NewRequest(). func Do(t *testing.T, method, urlStr string, body string) *httpResponse { t.Helper() diff --git a/translate.go b/translate.go index 157c05e2a..045d9c03b 100644 --- a/translate.go +++ b/translate.go @@ -16,7 +16,9 @@ package pilosa import ( "context" + "encoding/json" "io" + "io/ioutil" "sync" "github.com/pkg/errors" @@ -414,20 +416,43 @@ func (s *InMemTranslateStore) EntryReader(ctx context.Context, offset uint64) (T return newInMemTranslateEntryReader(ctx, s, offset), nil } -// WriteTo ensures that the TranslateStore implements io.WriterTo. -// It's not important that this be implemented. It would really -// only be necessary if we wanted to test cluster resizing while using -// an in-memory translate store. +// WriteTo implements io.WriterTo. It's not efficient or careful, but we +// don't expect to use InMemTranslateStore much, it's mostly there to +// avoid disk load during testing. func (s *InMemTranslateStore) WriteTo(w io.Writer) (int64, error) { - return 0, nil // TODO: try to use ErrNotImplemented + bytes, err := json.Marshal(s.keysByID) + if err != nil { + return 0, err + } + n, err := w.Write(bytes) + return int64(n), err } -// ReadFrom ensures that the TranslateStore implements io.ReaderFrom. -// It's not important that this be implemented. It would really -// only be necessary if we wanted to test cluster resizing while using -// an in-memory translate store. -func (s *InMemTranslateStore) ReadFrom(r io.Reader) (int64, error) { - return 0, nil // TODO: try to use ErrNotImplemented +// ReadFrom implements io.ReaderFrom. It's not efficient or careful, but we +// don't expect to use InMemTranslateStore much, it's mostly there to +// avoid disk load during testing. +func (s *InMemTranslateStore) ReadFrom(r io.Reader) (count int64, err error) { + var bytes []byte + bytes, err = ioutil.ReadAll(r) + count = int64(len(bytes)) + if err != nil { + return count, err + } + var keysByID map[uint64]string + err = json.Unmarshal(bytes, &keysByID) + if err != nil { + return count, err + } + s.maxID = 0 + s.keysByID = keysByID + s.idsByKey = make(map[string]uint64, len(s.keysByID)) + for k, v := range s.keysByID { + s.idsByKey[v] = k + if k > s.maxID { + s.maxID = k + } + } + return count, nil } // MaxID returns the highest identifier in the store.