From bfc24a1745b896df5d7c366460f4a05c13cee8e8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Wed, 17 Feb 2021 16:01:22 +0100 Subject: [PATCH 1/2] Check state in shardsByNode once stator is implemented --- cluster.go | 4 ++-- encoding/proto/proto.go | 5 +++-- executor.go | 5 ++--- server.go | 2 +- topology/node.go | 11 ++++++----- 5 files changed, 14 insertions(+), 13 deletions(-) diff --git a/cluster.go b/cluster.go index 37ba01b5f..96b03be7c 100644 --- a/cluster.go +++ b/cluster.go @@ -726,10 +726,10 @@ func (c *cluster) Nodes() []*topology.Node { s, err := c.stator.NodeState(context.Background(), node.ID) if err != nil { // TODO should we delete this? - copiedNodes[i].State = string(disco.NodeStateUnknown) + copiedNodes[i].State = disco.NodeStateUnknown continue } - copiedNodes[i].State = string(s) + copiedNodes[i].State = s } return result } diff --git a/encoding/proto/proto.go b/encoding/proto/proto.go index f091c77ae..465a44b63 100644 --- a/encoding/proto/proto.go +++ b/encoding/proto/proto.go @@ -21,6 +21,7 @@ import ( "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/v2" + "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/internal" pnet "github.com/pilosa/pilosa/v2/net" "github.com/pilosa/pilosa/v2/pql" @@ -708,7 +709,7 @@ func (s Serializer) encodeNode(m *topology.Node) *internal.Node { return &internal.Node{ ID: n.ID, URI: s.encodeURI(n.URI), - State: n.State, + State: string(n.State), GRPCURI: s.encodeURI(n.GRPCURI), } } @@ -1077,7 +1078,7 @@ func (s Serializer) decodeNode(node *internal.Node, m *topology.Node) { m.ID = node.ID s.decodeURI(node.URI, &m.URI) s.decodeURI(node.GRPCURI, &m.GRPCURI) - m.State = node.State + m.State = disco.NodeState(node.State) } func (s Serializer) decodeURI(i *internal.URI, m *pnet.URI) { diff --git a/executor.go b/executor.go index f52a5083d..4fccfa9b3 100644 --- a/executor.go +++ b/executor.go @@ -26,6 +26,7 @@ import ( "time" "unsafe" + "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/pql" pb "github.com/pilosa/pilosa/v2/proto" "github.com/pilosa/pilosa/v2/roaring" @@ -5527,9 +5528,7 @@ loop: // If the node being considered is in any state other than STARTED, // then exclude it from the map. This way, one of that node's // healthy replicas will be included instead. - // TODO: check state once stator is implemented - //if topology.Nodes(nodes).ContainsID(node.ID) && node.State == disco.NodeStateStarted { - if topology.Nodes(nodes).ContainsID(node.ID) { + if topology.Nodes(nodes).ContainsID(node.ID) && node.State == disco.NodeStateStarted { m[node] = append(m[node], shard) continue loop } diff --git a/server.go b/server.go index 2803a2c5d..ee783d2ec 100644 --- a/server.go +++ b/server.go @@ -565,7 +565,7 @@ func (s *Server) Open() error { ID: s.nodeID, URI: s.uri, GRPCURI: s.grpcURI, - State: string(disco.NodeStateUnknown), + State: disco.NodeStateUnknown, IsPrimary: s.IsPrimary(), } diff --git a/topology/node.go b/topology/node.go index c691c18a3..6fde62c7c 100644 --- a/topology/node.go +++ b/topology/node.go @@ -17,16 +17,17 @@ package topology import ( "fmt" + "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/net" ) // Node represents a node in the cluster. type Node struct { - ID string `json:"id"` - URI net.URI `json:"uri"` - GRPCURI net.URI `json:"grpc-uri"` - IsPrimary bool `json:"isPrimary"` - State string `json:"state"` + ID string `json:"id"` + URI net.URI `json:"uri"` + GRPCURI net.URI `json:"grpc-uri"` + IsPrimary bool `json:"isPrimary"` + State disco.NodeState `json:"state"` } func (n *Node) Clone() *Node { From 484f709621ae97b02412f98111c2f9ddb9032d2d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 18 Feb 2021 12:33:01 +0100 Subject: [PATCH 2/2] Add regression test --- server/server_test.go | 45 ++++++++++++++++++++++++++++++++++++++++++- test/cluster.go | 30 ++++++++++++++++++++--------- 2 files changed, 65 insertions(+), 10 deletions(-) diff --git a/server/server_test.go b/server/server_test.go index 4fcc7f253..126e08ba4 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -1219,7 +1219,7 @@ func TestClusterCreatedAtRace(t *testing.T) { for _, com := range cluster.Nodes { nodes := com.API.Hosts(context.Background()) for _, n := range nodes { - if n.State != string(disco.NodeStateStarted) { + if n.State != disco.NodeStateStarted { t.Fatalf("unexpected node state (%s) after upping cluster: %v", n.State, nodes) } } @@ -1268,3 +1268,46 @@ func TestClusterCreatedAtRace(t *testing.T) { }) } } + +func TestClusterQueryCountInDegraded(t *testing.T) { + cluster := test.MustNewCluster(t, 3) + for _, c := range cluster.Nodes { + c.Config.Cluster.ReplicaN = 2 + } + err := cluster.Start() + if err != nil { + t.Fatalf("starting cluster: %v", err) + } + defer cluster.Close() + + p := cluster.GetPrimary() + if err := p.Client().CreateIndex(context.Background(), "i", pilosa.IndexOptions{TrackExistence: true}); err != nil { + t.Fatal(err) + } else if err := p.Client().CreateField(context.Background(), "i", "f"); err != nil { + t.Fatal(err) + } + + np := cluster.GetNonPrimary() + // Write some data + for i := 0; i < 10; i++ { + if _, err := np.Query(t, "i", "", fmt.Sprintf(`Set(%d, f=1)`, i*pilosa.ShardWidth+1)); err != nil { + t.Fatal(err) + } + } + + if err := p.Close(); err != nil { + t.Fatal(err) + } + + if err := np.AwaitState(disco.ClusterStateDegraded, 30*time.Second); err != nil { + t.Fatal(err) + } + if resp, err := np.Client().Query(context.Background(), "i", &pilosa.QueryRequest{ + Index: "i", + Query: "Count(All())", + }); err != nil { + t.Fatal(err) + } else { + t.Logf("%+v", resp) + } +} diff --git a/test/cluster.go b/test/cluster.go index c3d0d116b..988c53927 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -132,11 +132,7 @@ func (c *Cluster) GetNode(n int) *Command { return c.Nodes[ids[n].idx] } -// GetCoordinator gets the node which has been determined to be the coordinator. -// This used to be node0 in tests, but since implementing etcd, the coordinator -// can be any node in the cluster, so we have to use this method in tests which -// need to act on the coordinator. -func (c *Cluster) GetCoordinator() *Command { +func (c *Cluster) GetPrimary() *Command { for _, n := range c.Nodes { if n.IsPrimary() { return n @@ -145,8 +141,7 @@ func (c *Cluster) GetCoordinator() *Command { return nil } -// GetNonCoordinator gets first first non-coordinator node in the list of nodes. -func (c *Cluster) GetNonCoordinator() *Command { +func (c *Cluster) GetNonPrimary() *Command { for _, n := range c.Nodes { if !n.IsPrimary() { return n @@ -155,8 +150,7 @@ func (c *Cluster) GetNonCoordinator() *Command { return nil } -// GetNonCoordinators gets all nodes except the coordinator. -func (c *Cluster) GetNonCoordinators() []*Command { +func (c *Cluster) GetNonPrimaries() []*Command { rtn := make([]*Command, 0) for _, n := range c.Nodes { if !n.IsPrimary() { @@ -166,6 +160,24 @@ func (c *Cluster) GetNonCoordinators() []*Command { return rtn } +// GetCoordinator gets the node which has been determined to be the coordinator. +// This used to be node0 in tests, but since implementing etcd, the coordinator +// can be any node in the cluster, so we have to use this method in tests which +// need to act on the coordinator. +func (c *Cluster) GetCoordinator() *Command { + return c.GetPrimary() +} + +// GetNonCoordinator gets first first non-coordinator node in the list of nodes. +func (c *Cluster) GetNonCoordinator() *Command { + return c.GetNonPrimary() +} + +// GetNonCoordinators gets all nodes except the coordinator. +func (c *Cluster) GetNonCoordinators() []*Command { + return c.GetNonPrimaries() +} + // nodePlace represents a node's ID and its index into the c.Nodes slice. type nodePlace struct { id string