mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-08 03:47:51 +00:00
Merge pull request #1437 from kuba--/check-state
Check state in shardsByNode once stator is implemented
This commit is contained in:
commit
e724ad53f2
7 changed files with 79 additions and 23 deletions
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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) {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue