Merge pull request #1795 from jaffee/close-client-responses

Ensure internal client closes all response bodies
This commit is contained in:
Matthew Jaffee 2018-12-20 10:41:27 -06:00 • committed by GitHub
commit b2ca5f1f26
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
3 changed files with 125 additions and 6 deletions

View file

@ -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 {

View file

@ -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")

View file

@ -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)
}
}