From ff800131cb3b2cc8b8643c6722fac471323465a0 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 19 Dec 2018 11:57:47 -0600 Subject: [PATCH 1/2] ensure internal client closes all response bodies to avoid leaking connections/goroutines --- http/client.go | 22 +++++++++++++++++----- 1 file changed, 17 insertions(+), 5 deletions(-) diff --git a/http/client.go b/http/client.go index 46800fd70..78382b81b 100644 --- a/http/client.go +++ b/http/client.go @@ -164,7 +164,7 @@ func (c *InternalClient) CreateIndex(ctx context.Context, index string, opt pilo } return err } - return nil + return errors.Wrap(resp.Body.Close(), "closing response body") } // FragmentNodes returns a list of nodes that own a shard. @@ -705,6 +705,9 @@ func (c *InternalClient) exportNodeCSV(ctx context.Context, node *pilosa.Node, i return nil } +// RetrieveShardFromURI returns a ReadCloser which contains the data of the +// specified shard from the specified node. Caller *must* close the returned +// ReadCloser or risk leaking goroutines/tcp connections. func (c *InternalClient) RetrieveShardFromURI(ctx context.Context, index, field, view string, shard uint64, uri pilosa.URI) (io.ReadCloser, error) { span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.RetrieveShardFromURI") defer span.Finish() @@ -800,7 +803,7 @@ func (c *InternalClient) CreateFieldWithOptions(ctx context.Context, index, fiel return err } - return nil + return errors.Wrap(resp.Body.Close(), "closing response body") } // FragmentBlocks returns a list of block checksums for a fragment on a host. @@ -994,15 +997,24 @@ func (c *InternalClient) SendMessage(ctx context.Context, uri *pilosa.URI, msg [ req.Header.Set("Accept", "application/json") // Execute request. - _, err = c.executeRequest(req.WithContext(ctx)) - return err + resp, err := c.executeRequest(req.WithContext(ctx)) + if err != nil { + return errors.Wrap(err, "executing request") + } + return errors.Wrap(resp.Body.Close(), "closing response body") } -// executeRequest executes the given request and checks the Response +// executeRequest executes the given request and checks the Response. For +// responses with non-2XX status, the body is read and closed, and an error is +// returned. If the error is nil, the caller must ensure that the response body +// is closed. func (c *InternalClient) executeRequest(req *http.Request) (*http.Response, error) { tracing.GlobalTracer.InjectHTTPHeaders(req) resp, err := c.httpClient.Do(req) if err != nil { + if resp != nil { + resp.Body.Close() + } return nil, errors.Wrap(err, "executing request") } if resp.StatusCode < 200 || resp.StatusCode >= 300 { From 31ab3d37b8e17ae6418493b3d9811c03c3992d15 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Wed, 19 Dec 2018 17:21:57 -0600 Subject: [PATCH 2/2] add stress tests these produce an issue on master, but it is fixed on this branch --- pilosa.go | 4 +- server/server_test.go | 105 ++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 108 insertions(+), 1 deletion(-) diff --git a/pilosa.go b/pilosa.go index fe61d8969..550696a31 100644 --- a/pilosa.go +++ b/pilosa.go @@ -56,7 +56,9 @@ var ( ErrQueryTimeout = errors.New("query timeout") ErrTooManyWrites = errors.New("too many write commands") - ErrClusterDoesNotOwnShard = errors.New("cluster does not own shard") + // TODO(2.0) poorly named - used when a *node* doesn't own a shard. Probably + // we won't need this error at all by 2.0 though. + ErrClusterDoesNotOwnShard = errors.New("node does not own shard") ErrNodeIDNotExists = errors.New("node with provided ID does not exist") ErrNodeNotCoordinator = errors.New("node is not the coordinator") diff --git a/server/server_test.go b/server/server_test.go index 627e83862..044ae4720 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -15,8 +15,10 @@ package server_test import ( + "bytes" "context" "encoding/json" + "flag" "fmt" "io/ioutil" "math/rand" @@ -28,13 +30,22 @@ import ( "testing/quick" "time" + "golang.org/x/sync/errgroup" + "github.com/pelletier/go-toml" "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/http" + "github.com/pilosa/pilosa/roaring" "github.com/pilosa/pilosa/server" "github.com/pilosa/pilosa/test" ) +var runStress bool + +func init() { // nolint: gochecknoinits + flag.BoolVar(&runStress, "stress", false, "Enable stress tests (time consuming)") +} + // Ensure program can process queries and maintain consistency. func TestMain_Set_Quick(t *testing.T) { if testing.Short() { @@ -754,3 +765,97 @@ func TestClusterQueriesAfterRestart(t *testing.T) { } // TODO: confirm that things keep working if a node is hard-closed (no nodeLeave event) and immediately restarted with a different address. + +func TestClusterExhaustingConnections(t *testing.T) { + if !runStress { + t.Skip("stress") + } + cluster := test.MustRunCluster(t, 5) + defer cluster.Close() + cmd1 := cluster[1] + + for _, com := range cluster { + nodes := com.API.Hosts(context.Background()) + for _, n := range nodes { + if n.State != "READY" { + t.Fatalf("unexpected node state after upping cluster: %v", nodes) + } + } + } + + cmd1.MustCreateIndex(t, "testidx", pilosa.IndexOptions{}) + cmd1.MustCreateField(t, "testidx", "testfield", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 10)) + + eg := errgroup.Group{} + for i := 0; i < 20; i++ { + i := i + eg.Go(func() error { + for j := i; j < 10000; j += 20 { + _, err := cluster[i%5].API.Query(context.Background(), &pilosa.QueryRequest{ + Index: "testidx", + Query: fmt.Sprintf("Set(%d, testfield=0)", j*pilosa.ShardWidth), + }) + if err != nil { + return err + } + } + return nil + }) + } + err := eg.Wait() + if err != nil { + t.Fatalf("setting lots of shards: %v", err) + } +} + +func TestClusterExhaustingConnectionsImport(t *testing.T) { + if !runStress { + t.Skip("stress") + } + cluster := test.MustRunCluster(t, 5) + defer cluster.Close() + cmd1 := cluster[1] + + for _, com := range cluster { + nodes := com.API.Hosts(context.Background()) + for _, n := range nodes { + if n.State != "READY" { + t.Fatalf("unexpected node state after upping cluster: %v", nodes) + } + } + } + + cmd1.MustCreateIndex(t, "testidx", pilosa.IndexOptions{}) + cmd1.MustCreateField(t, "testidx", "testfield", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 10)) + + bm := roaring.NewBitmap() + bm.DirectAdd(0) + buf := &bytes.Buffer{} + bm.WriteTo(buf) + data := buf.Bytes() + + eg := errgroup.Group{} + for i := uint64(0); i < 20; i++ { + i := i + eg.Go(func() error { + for j := i; j < 10000; j += 20 { + if (j-i)%1000 == 0 { + fmt.Printf("%d is %.2f%% done.\n", i, float64(j-i)*100/100000) + } + err := cluster[i%5].API.ImportRoaring(context.Background(), "testidx", "testfield", j, false, &pilosa.ImportRoaringRequest{ + Views: map[string][]byte{ + "": data, + }, + }) + if err != nil { + return err + } + } + return nil + }) + } + err := eg.Wait() + if err != nil { + t.Fatalf("setting lots of shards: %v", err) + } +}