diff --git a/gossip/gossip.go b/gossip/gossip.go index a41d05065..0f4427d3c 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -106,7 +106,8 @@ func (g *memberSet) Open() (err error) { // Close attempts to gracefully leave the cluster, and finally calls shutdown // after (at most) a timeout period. func (g *memberSet) Close() error { - g.eventReceiver.Close() + defer g.eventReceiver.Close() + leaveErr := g.memberlist.Leave(5 * time.Second) shutdownErr := g.memberlist.Shutdown() if leaveErr != nil || shutdownErr != nil { diff --git a/server.go b/server.go index 3f70b9fe9..b3af0a97a 100644 --- a/server.go +++ b/server.go @@ -678,51 +678,57 @@ func (s *Server) Open() error { // Close closes the server and waits for it to shutdown. func (s *Server) Close() error { - fmt.Println("--- disco: server close:", s.disCo.ID()) - errE := s.executor.Close() + select { + case <-s.closing: + return nil + default: - // Notify goroutines to stop. - close(s.closing) - s.wg.Wait() - var errh, errd error - var errhs error - var errc error + fmt.Println("--- disco: server close:", s.disCo.ID()) + errE := s.executor.Close() - if s.cluster != nil { - errc = s.cluster.close() - } - errhs = s.syncer.stopTranslationSync() - if s.disCo != nil { - fmt.Println("--- disco: try close:", s.disCo.ID()) - errd = s.disCo.Close() - fmt.Println("--- disco: closed", s.disCo.ID(), errd) - } - if s.holder != nil { - errh = s.holder.Close() - } - if s.snapshotQueue != nil { - s.holder.SnapshotQueue = nil - s.snapshotQueue.Stop() - s.snapshotQueue = nil - } + // Notify goroutines to stop. + close(s.closing) + s.wg.Wait() + var errh, errd error + var errhs error + var errc error - // prefer to return holder error over cluster - // error. This order is somewhat arbitrary. It would be better if we had - // some way to combine all the errors, but probably not important enough to - // warrant the extra complexity. - if errh != nil { - return errors.Wrap(errh, "closing holder") + if s.cluster != nil { + errc = s.cluster.close() + } + errhs = s.syncer.stopTranslationSync() + if s.disCo != nil { + fmt.Println("--- disco: try close:", s.disCo.ID()) + errd = s.disCo.Close() + fmt.Println("--- disco: closed", s.disCo.ID(), errd) + } + if s.holder != nil { + errh = s.holder.Close() + } + if s.snapshotQueue != nil { + s.holder.SnapshotQueue = nil + s.snapshotQueue.Stop() + s.snapshotQueue = nil + } + + // prefer to return holder error over cluster + // error. This order is somewhat arbitrary. It would be better if we had + // some way to combine all the errors, but probably not important enough to + // warrant the extra complexity. + if errh != nil { + return errors.Wrap(errh, "closing holder") + } + if errhs != nil { + return errors.Wrap(errhs, "terminating holder translation sync") + } + if errc != nil { + return errors.Wrap(errc, "closing cluster") + } + if errd != nil { + return errors.Wrap(errd, "closing disco") + } + return errors.Wrap(errE, "closing executor") } - if errhs != nil { - return errors.Wrap(errhs, "terminating holder translation sync") - } - if errc != nil { - return errors.Wrap(errc, "closing cluster") - } - if errd != nil { - return errors.Wrap(errd, "closing disco") - } - return errors.Wrap(errE, "closing executor") } // NodeID returns the server's node id. diff --git a/server/server.go b/server/server.go index 29c33f44c..883e1cc63 100644 --- a/server/server.go +++ b/server/server.go @@ -545,30 +545,36 @@ func (m *Command) GossipTransport() *gossip.Transport { // Close shuts down the server. func (m *Command) Close() error { - defer close(m.done) - eg := errgroup.Group{} - m.grpcServer.Stop() - eg.Go(m.Handler.Close) - eg.Go(m.Server.Close) - eg.Go(m.API.Close) - eg.Go(m.pgserver.Close) - if m.gossipMemberSet != nil { - eg.Go(m.gossipMemberSet.Close) - } - if closer, ok := m.logOutput.(io.Closer); ok { - // If closer is os.Stdout or os.Stderr, don't close it. - if closer != os.Stdout && closer != os.Stderr { - eg.Go(closer.Close) + select { + case <-m.done: + return nil + default: + + defer close(m.done) + eg := errgroup.Group{} + m.grpcServer.Stop() + eg.Go(m.Handler.Close) + eg.Go(m.Server.Close) + eg.Go(m.API.Close) + eg.Go(m.pgserver.Close) + if m.gossipMemberSet != nil { + eg.Go(m.gossipMemberSet.Close) } + if closer, ok := m.logOutput.(io.Closer); ok { + // If closer is os.Stdout or os.Stderr, don't close it. + if closer != os.Stdout && closer != os.Stderr { + eg.Go(closer.Close) + } + } + + // prevent the closed sockets from being re-injected into etcd. + m.Config.DisCo.LPeerSocket = nil + m.Config.DisCo.LClientSocket = nil + + err := eg.Wait() + _ = testhook.Closed(pilosa.NewAuditor(), m, nil) + return errors.Wrap(err, "closing everything") } - - // prevent the closed sockets from being re-injected into etcd. - m.Config.DisCo.LPeerSocket = nil - m.Config.DisCo.LClientSocket = nil - - err := eg.Wait() - _ = testhook.Closed(pilosa.NewAuditor(), m, nil) - return errors.Wrap(err, "closing everything") } // newStatsClient creates a stats client from the config diff --git a/server/server_test.go b/server/server_test.go index f1926729c..eb90915c1 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -630,39 +630,22 @@ func TestClusteringNodesReplica1(t *testing.T) { cluster := test.MustRunCluster(t, 3) defer cluster.Close() - err := cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond) - if err != nil { + if err := cluster.AwaitState(pilosa.ClusterStateNormal, 100*time.Millisecond); err != nil { t.Fatalf("starting cluster: %v", err) } - if err := cluster.GetNode(2).Command.Close(); err != nil { + if err := cluster.GetNonCoordinator().Command.Close(); err != nil { t.Fatalf("closing third node: %v", err) } + if err := cluster.AwaitCoordinatorState(pilosa.ClusterStateStarting, 30*time.Second); err != nil { + t.Fatalf("starting cluster: %v", err) + } + // confirm that cluster stops accepting queries after one node closes - if _, err := cluster.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { + if _, err := cluster.GetCoordinator().API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { t.Fatalf("got unexpected error querying an incomplete cluster: %v", err) } - - // Create new main with the same config. - config := cluster.GetNode(2).Command.Config - config.Translation.MapSize = 100000 - - // this isn't necessary, but makes the test run way faster - config.Gossip.Port = strconv.Itoa(int(cluster.GetNode(2).Command.GossipTransport().URI.Port)) - - cluster.GetNode(2).Command = server.NewCommand(cluster.GetNode(2).Stdin, cluster.GetNode(2).Stdout, cluster.GetNode(2).Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore))) - cluster.GetNode(2).Command.Config = config - - // Run new program. - if err := cluster.GetNode(2).Start(); err != nil { - t.Fatalf("restarting node 2: %v", err) - } - - err = cluster.AwaitState(pilosa.ClusterStateNormal, 200*time.Millisecond) - if err != nil { - t.Fatalf("resuming normal operations: %v", err) - } } func TestClusteringNodesReplica2(t *testing.T) { @@ -681,75 +664,35 @@ func TestClusteringNodesReplica2(t *testing.T) { t.Fatalf("starting cluster: %v", err) } - if err := cluster.GetNode(2).Command.Close(); err != nil { + coord, others := cluster.GetCoordinator(), cluster.GetNonCoordinators() + + if err := others[0].Close(); err != nil { t.Fatalf("closing third node: %v", err) } - err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 30*time.Second) if err != nil { t.Fatalf("after closing first server: %v", err) } // confirm that cluster keeps accepting queries if replication > 1 - if _, err := cluster.GetNode(0).API.CreateIndex(context.Background(), "anewindex", pilosa.IndexOptions{}); err != nil { + if _, err := coord.API.CreateIndex(context.Background(), "anewindex", pilosa.IndexOptions{}); err != nil { t.Fatalf("got unexpected error creating index: %v", err) } // confirm that cluster stops accepting queries if 2 nodes fail and replication == 2 - if err := cluster.GetNode(1).Command.Close(); err != nil { + if err := others[1].Close(); err != nil { t.Fatalf("closing 2nd node: %v", err) } - err = cluster.AwaitCoordinatorState(pilosa.ClusterStateStarting, 100*time.Millisecond) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateStarting, 30*time.Second) if err != nil { t.Fatalf("after closing second server: %v", err) } - if _, err := cluster.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { + if _, err := coord.API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { t.Fatalf("got unexpected error querying an incomplete cluster: %v", err) } - - // Create new main with the same config. - config := cluster.GetNode(2).Command.Config - config.Translation.MapSize = 100000 - // config.Bind = cluster.GetNode(2).API.Node().URI.HostPort() - - // this isn't necessary, but makes the test run way faster - config.Gossip.Port = strconv.Itoa(int(cluster.GetNode(2).Command.GossipTransport().URI.Port)) - - cluster.GetNode(2).Command = server.NewCommand(cluster.GetNode(2).Stdin, cluster.GetNode(2).Stdout, cluster.GetNode(2).Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore))) - cluster.GetNode(2).Command.Config = config - - // Run new program. - if err := cluster.GetNode(2).Start(); err != nil { - t.Fatalf("restarting node 2: %v", err) - } - - 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. - config = cluster.GetNode(1).Command.Config - // config.Bind = cluster.GetNode(1).API.Node().URI.HostPort() - config.Translation.MapSize = 100000 - - // this isn't necessary, but makes the test run way faster - config.Gossip.Port = strconv.Itoa(int(cluster.GetNode(1).Command.GossipTransport().URI.Port)) - - cluster.GetNode(1).Command = server.NewCommand(cluster.GetNode(1).Stdin, cluster.GetNode(1).Stdout, cluster.GetNode(1).Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore))) - cluster.GetNode(1).Command.Config = config - - // Run new program. - if err := cluster.GetNode(1).Start(); err != nil { - t.Fatalf("restarting node 1: %v", err) - } - - err = cluster.AwaitState(pilosa.ClusterStateNormal, 200*time.Microsecond) - if err != nil { - t.Fatalf("resuming normal operations: %v", err) - } } func TestRemoveNodeAfterItDies(t *testing.T) { @@ -774,27 +717,28 @@ func TestRemoveNodeAfterItDies(t *testing.T) { t.Fatalf("starting cluster: %v", err) } + coord, others := cluster.GetCoordinator(), cluster.GetNonCoordinators() // prevent double-closing cluster.GetNode(2) from the deferred Close above - disabled := cluster.GetNode(2) - if err := cluster.CloseAndRemove(2); err != nil { + disabled := others[0] + if err := disabled.Close(); err != nil { t.Fatalf("closing third node: %v", err) } - err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 100*time.Millisecond) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateDegraded, 30*time.Second) if err != nil { t.Fatalf("starting cluster: %v", err) } - if _, err := cluster.GetNode(0).API.RemoveNode(disabled.API.Node().ID); err != nil { + if _, err := coord.API.RemoveNode(disabled.API.Node().ID); err != nil { t.Fatalf("removing failed node: %v", err) } - err = cluster.AwaitCoordinatorState(pilosa.ClusterStateNormal, 100*time.Millisecond) + err = cluster.AwaitCoordinatorState(pilosa.ClusterStateNormal, 30*time.Second) if err != nil { t.Fatalf("removing disabled node: %v", err) } - hosts := cluster.GetNode(0).API.Hosts(context.Background()) + hosts := coord.API.Hosts(context.Background()) if len(hosts) != 2 { t.Fatalf("unexpected hosts: %v", hosts) } diff --git a/test/cluster.go b/test/cluster.go index cd22bfff7..9031e766f 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -478,7 +478,7 @@ func (c *Cluster) AwaitCoordinatorState(expectedState string, timeout time.Durat if len(c.Nodes) < 1 { return errors.New("can't await coordinator state on an empty cluster") } - onlyCoordinator := &Cluster{Nodes: c.Nodes[:1]} + onlyCoordinator := &Cluster{Nodes: []*Command{c.GetCoordinator()}} return onlyCoordinator.AwaitState(expectedState, timeout) }