From f051f7ccea6674b2ebe168f72b6e5923eebb1945 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Fri, 1 Dec 2017 16:00:12 -0600 Subject: [PATCH 1/2] refactored httpclient handling --- client.go | 30 ++++-------------------------- client_test.go | 34 +++++++++++++++++++++------------- ctl/common.go | 8 ++++---- executor.go | 24 +++++++++++------------- fragment.go | 11 ++++++----- handler.go | 8 ++++---- holder.go | 21 +++++++++++---------- holder_test.go | 11 ++++++----- server.go | 34 +++++++++++++++++++++++++--------- server/server.go | 13 +++++++++++-- server/server_test.go | 6 +++--- test/client.go | 6 ++++-- test/executor.go | 2 +- 13 files changed, 111 insertions(+), 97 deletions(-) diff --git a/client.go b/client.go index 666447586..e44d9c5de 100644 --- a/client.go +++ b/client.go @@ -25,7 +25,6 @@ import ( "io/ioutil" "log" "math/rand" - "net" "net/http" "net/url" "sort" @@ -46,14 +45,13 @@ type ClientOptions struct { // InternalHTTPClient represents a client to the Pilosa cluster. type InternalHTTPClient struct { defaultURI *URI - options *ClientOptions // The client to use for HTTP communication. HTTPClient *http.Client } // NewInternalHTTPClient returns a new instance of InternalHTTPClient to connect to host. -func NewInternalHTTPClient(host string, options *ClientOptions) (*InternalHTTPClient, error) { +func NewInternalHTTPClient(host string, remoteClient *http.Client) (*InternalHTTPClient, error) { if host == "" { return nil, ErrHostRequired } @@ -63,34 +61,14 @@ func NewInternalHTTPClient(host string, options *ClientOptions) (*InternalHTTPCl return nil, err } - client := NewInternalHTTPClientFromURI(uri, options) + client := NewInternalHTTPClientFromURI(uri, remoteClient) return client, nil } -func NewInternalHTTPClientFromURI(defaultURI *URI, options *ClientOptions) *InternalHTTPClient { - if options == nil { - options = &ClientOptions{} - } - transport := &http.Transport{ - Proxy: http.ProxyFromEnvironment, - DialContext: (&net.Dialer{ - Timeout: 30 * time.Second, - KeepAlive: 30 * time.Second, - DualStack: true, - }).DialContext, - MaxIdleConns: 1000, - MaxIdleConnsPerHost: 200, - IdleConnTimeout: 90 * time.Second, - TLSHandshakeTimeout: 10 * time.Second, - ExpectContinueTimeout: 1 * time.Second, - } - if options.TLS != nil { - transport.TLSClientConfig = options.TLS - } - client := &http.Client{Transport: transport} +func NewInternalHTTPClientFromURI(defaultURI *URI, remoteClient *http.Client) *InternalHTTPClient { return &InternalHTTPClient{ defaultURI: defaultURI, - HTTPClient: client, + HTTPClient: remoteClient, } } diff --git a/client_test.go b/client_test.go index 613d37659..8feec09ea 100644 --- a/client_test.go +++ b/client_test.go @@ -18,6 +18,7 @@ import ( "bytes" "context" "fmt" + "net/http" "reflect" "testing" @@ -43,6 +44,13 @@ func createCluster(c *pilosa.Cluster) ([]*test.Server, []*test.Holder) { return server, hldr } +var defaultClient *http.Client + +func init() { + defaultClient = pilosa.GetHTTPClient(nil) + +} + // Test distributed TopN Row count across 3 nodes. func TestClient_MultiNode(t *testing.T) { cluster := test.NewCluster(3) @@ -54,7 +62,7 @@ func TestClient_MultiNode(t *testing.T) { } s[0].Handler.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { - e := pilosa.NewExecutor(nil) + e := pilosa.NewExecutor(defaultClient) e.Holder = hldr[0].Holder e.Scheme = cluster.Nodes[0].Scheme e.Host = cluster.Nodes[0].Host @@ -62,7 +70,7 @@ func TestClient_MultiNode(t *testing.T) { return e.Execute(ctx, index, query, slices, opt) } s[1].Handler.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { - e := pilosa.NewExecutor(nil) + e := pilosa.NewExecutor(defaultClient) e.Holder = hldr[1].Holder e.Scheme = cluster.Nodes[1].Scheme e.Host = cluster.Nodes[1].Host @@ -70,7 +78,7 @@ func TestClient_MultiNode(t *testing.T) { return e.Execute(ctx, index, query, slices, opt) } s[2].Handler.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { - e := pilosa.NewExecutor(nil) + e := pilosa.NewExecutor(defaultClient) e.Holder = hldr[2].Holder e.Scheme = cluster.Nodes[2].Scheme e.Host = cluster.Nodes[2].Host @@ -135,9 +143,9 @@ func TestClient_MultiNode(t *testing.T) { // Connect to each node to compare results. client := make([]*test.Client, 3) - client[0] = test.MustNewClient(s[0].Host()) - client[1] = test.MustNewClient(s[1].Host()) - client[2] = test.MustNewClient(s[2].Host()) + client[0] = test.MustNewClient(s[0].Host(), defaultClient) + client[1] = test.MustNewClient(s[1].Host(), defaultClient) + client[2] = test.MustNewClient(s[2].Host(), defaultClient) topN := 4 queryRequest := &internal.QueryRequest{ @@ -218,7 +226,7 @@ func TestClient_Import(t *testing.T) { s.Handler.Holder = hldr.Holder // Send import request. - c := test.MustNewClient(s.Host()) + c := test.MustNewClient(s.Host(), defaultClient) if err := c.Import(context.Background(), "i", "f", 0, []pilosa.Bit{ {RowID: 0, ColumnID: 1}, {RowID: 0, ColumnID: 5}, @@ -269,7 +277,7 @@ func TestClient_ImportInverseEnabled(t *testing.T) { s.Handler.Holder = hldr.Holder // Send import request. - c := test.MustNewClient(s.Host()) + c := test.MustNewClient(s.Host(), defaultClient) if err := c.Import(context.Background(), "i", "f", 0, []pilosa.Bit{ {RowID: 0, ColumnID: 1}, {RowID: 0, ColumnID: 5}, @@ -318,7 +326,7 @@ func TestClient_ImportValue(t *testing.T) { s.Handler.Holder = hldr.Holder // Send import request. - c := test.MustNewClient(s.Host()) + c := test.MustNewClient(s.Host(), defaultClient) if err := c.ImportValue(context.Background(), "i", "f", fld.Name, 0, []pilosa.FieldValue{ {ColumnID: 1, Value: -10}, {ColumnID: 2, Value: 20}, @@ -355,7 +363,7 @@ func TestClient_BackupRestore(t *testing.T) { s.Handler.Cluster.Nodes[0].Host = s.Host() s.Handler.Holder = hldr.Holder - c := test.MustNewClient(s.Host()) + c := test.MustNewClient(s.Host(), defaultClient) // Backup from frame. var buf bytes.Buffer @@ -420,7 +428,7 @@ func TestClient_BackupInverseView(t *testing.T) { s.Handler.Cluster.Nodes[0].Host = s.Host() s.Handler.Holder = hldr.Holder - c := test.MustNewClient(s.Host()) + c := test.MustNewClient(s.Host(), defaultClient) // Backup from frame. var buf bytes.Buffer @@ -457,7 +465,7 @@ func TestClient_BackupInvalidView(t *testing.T) { s.Handler.Cluster.Nodes[0].Host = s.Host() s.Handler.Holder = hldr.Holder - c := test.MustNewClient(s.Host()) + c := test.MustNewClient(s.Host(), defaultClient) // Backup from frame. var buf bytes.Buffer @@ -487,7 +495,7 @@ func TestClient_FragmentBlocks(t *testing.T) { s.Handler.Holder = hldr.Holder // Retrieve blocks. - c := test.MustNewClient(s.Host()) + c := test.MustNewClient(s.Host(), defaultClient) blocks, err := c.FragmentBlocks(context.Background(), "i", "f", pilosa.ViewStandard, 0) if err != nil { t.Fatal(err) diff --git a/ctl/common.go b/ctl/common.go index dc64c0bee..826462b8c 100644 --- a/ctl/common.go +++ b/ctl/common.go @@ -2,6 +2,7 @@ package ctl import ( "crypto/tls" + "github.com/pilosa/pilosa" "github.com/spf13/pflag" ) @@ -22,19 +23,18 @@ func SetTLSConfig(flags *pflag.FlagSet, certificatePath *string, certificateKeyP // CommandClient returns a pilosa.InternalHTTPClient for the command func CommandClient(cmd CommandWithTLSSupport) (*pilosa.InternalHTTPClient, error) { tlsConfig := cmd.TLSConfiguration() - var clientOptions *pilosa.ClientOptions + var TLSConfig *tls.Config if tlsConfig.CertificatePath != "" && tlsConfig.CertificateKeyPath != "" { cert, err := tls.LoadX509KeyPair(tlsConfig.CertificatePath, tlsConfig.CertificateKeyPath) if err != nil { return nil, err } - TLSConfig := &tls.Config{ + TLSConfig = &tls.Config{ Certificates: []tls.Certificate{cert}, InsecureSkipVerify: tlsConfig.SkipVerify, } - clientOptions = &pilosa.ClientOptions{TLS: TLSConfig} } - client, err := pilosa.NewInternalHTTPClient(cmd.TLSHost(), clientOptions) + client, err := pilosa.NewInternalHTTPClient(cmd.TLSHost(), pilosa.GetHTTPClient(TLSConfig)) if err != nil { return nil, err } diff --git a/executor.go b/executor.go index 2f5acf3c0..a617eb184 100644 --- a/executor.go +++ b/executor.go @@ -18,6 +18,7 @@ import ( "context" "errors" "fmt" + "net/http" "sort" "time" @@ -51,12 +52,9 @@ type Executor struct { } // NewExecutor returns a new instance of Executor. -func NewExecutor(clientOptions *ClientOptions) *Executor { - if clientOptions == nil { - clientOptions = &ClientOptions{} - } +func NewExecutor(remoteClient *http.Client) *Executor { return &Executor{ - client: NewInternalHTTPClientFromURI(nil, clientOptions), + client: NewInternalHTTPClientFromURI(nil, remoteClient), } } @@ -968,7 +966,7 @@ func (e *Executor) executeClearBitView(ctx context.Context, index string, c *pql } // Forward call to remote node otherwise. - if res, err := e.exec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt); err != nil { + if res, err := e.remoteExec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt); err != nil { return false, err } else { ret = res[0].(bool) @@ -1074,7 +1072,7 @@ func (e *Executor) executeSetBitView(ctx context.Context, index string, c *pql.C } // Forward call to remote node otherwise. - if res, err := e.exec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt); err != nil { + if res, err := e.remoteExec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt); err != nil { return false, err } else { ret = res[0].(bool) @@ -1141,7 +1139,7 @@ func (e *Executor) executeSetFieldValue(ctx context.Context, index string, c *pq resp := make(chan error, len(nodes)) for _, node := range nodes { go func(node *Node) { - _, err := e.exec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt) + _, err := e.remoteExec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt) resp <- err }(node) } @@ -1199,7 +1197,7 @@ func (e *Executor) executeSetRowAttrs(ctx context.Context, index string, c *pql. resp := make(chan error, len(nodes)) for _, node := range nodes { go func(node *Node) { - _, err := e.exec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt) + _, err := e.remoteExec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt) resp <- err }(node) } @@ -1286,7 +1284,7 @@ func (e *Executor) executeBulkSetRowAttrs(ctx context.Context, index string, cal resp := make(chan error, len(nodes)) for _, node := range nodes { go func(node *Node) { - _, err := e.exec(ctx, node, index, &pql.Query{Calls: calls}, nil, opt) + _, err := e.remoteExec(ctx, node, index, &pql.Query{Calls: calls}, nil, opt) resp <- err }(node) } @@ -1345,7 +1343,7 @@ func (e *Executor) executeSetColumnAttrs(ctx context.Context, index string, c *p resp := make(chan error, len(nodes)) for _, node := range nodes { go func(node *Node) { - _, err := e.exec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt) + _, err := e.remoteExec(ctx, node, index, &pql.Query{Calls: []*pql.Call{c}}, nil, opt) resp <- err }(node) } @@ -1361,7 +1359,7 @@ func (e *Executor) executeSetColumnAttrs(ctx context.Context, index string, c *p } // exec executes a PQL query remotely for a set of slices on a node. -func (e *Executor) exec(ctx context.Context, node *Node, index string, q *pql.Query, slices []uint64, opt *ExecOptions) (results []interface{}, err error) { +func (e *Executor) remoteExec(ctx context.Context, node *Node, index string, q *pql.Query, slices []uint64, opt *ExecOptions) (results []interface{}, err error) { // Encode request object. pbreq := &internal.QueryRequest{ Query: q.String(), @@ -1511,7 +1509,7 @@ func (e *Executor) mapper(ctx context.Context, ch chan mapResponse, nodes []*Nod if n.Host == e.Host { resp.result, resp.err = e.mapperLocal(ctx, nodeSlices, mapFn, reduceFn) } else if !opt.Remote { - results, err := e.exec(ctx, n, index, &pql.Query{Calls: []*pql.Call{c}}, nodeSlices, opt) + results, err := e.remoteExec(ctx, n, index, &pql.Query{Calls: []*pql.Call{c}}, nodeSlices, opt) if len(results) > 0 { resp.result = results[0] } diff --git a/fragment.go b/fragment.go index e5fbb7988..073bdfa84 100644 --- a/fragment.go +++ b/fragment.go @@ -28,6 +28,7 @@ import ( "io" "io/ioutil" "log" + "net/http" "os" "sort" "sync" @@ -1677,9 +1678,9 @@ func (h *blockHasher) WriteValue(v uint64) { type FragmentSyncer struct { Fragment *Fragment - Host string - Cluster *Cluster - ClientOptions *ClientOptions + Host string + Cluster *Cluster + RemoteClient *http.Client Closing <-chan struct{} } @@ -1714,7 +1715,7 @@ func (s *FragmentSyncer) SyncFragment() error { } // Retrieve remote blocks. - client, err := NewInternalHTTPClient(node.Host, s.ClientOptions) + client, err := NewInternalHTTPClient(node.Host, s.RemoteClient) if err != nil { return err } @@ -1793,7 +1794,7 @@ func (s *FragmentSyncer) syncBlock(id int) error { return nil } - client, err := NewInternalHTTPClient(node.Host, s.ClientOptions) + client, err := NewInternalHTTPClient(node.Host, s.RemoteClient) if err != nil { return err } diff --git a/handler.go b/handler.go index f3d9fe9d3..ad9e02e2e 100644 --- a/handler.go +++ b/handler.go @@ -56,9 +56,9 @@ type Handler struct { StatusHandler StatusHandler // Local hostname & cluster configuration. - URI *URI - Cluster *Cluster - ClientOptions *ClientOptions + URI *URI + Cluster *Cluster + RemoteClient *http.Client Router *mux.Router @@ -1506,7 +1506,7 @@ func (h *Handler) handlePostFrameRestore(w http.ResponseWriter, r *http.Request) } // Create a client for the remote cluster. - client := NewInternalHTTPClientFromURI(host, h.ClientOptions) + client := NewInternalHTTPClientFromURI(host, h.RemoteClient) // Determine the maximum number of slices. maxSlices, err := client.MaxSliceByIndex(r.Context()) diff --git a/holder.go b/holder.go index 4241dcd85..7cbb8ca42 100644 --- a/holder.go +++ b/holder.go @@ -20,6 +20,7 @@ import ( "fmt" "io" "log" + "net/http" "os" "path/filepath" "sort" @@ -430,9 +431,9 @@ func (h *Holder) logger() *log.Logger { return log.New(h.LogOutput, "", log.Lstd type HolderSyncer struct { Holder *Holder - URI *URI - Cluster *Cluster - ClientOptions *ClientOptions + URI *URI + Cluster *Cluster + RemoteClient *http.Client // Signals that the sync should stop. Closing <-chan struct{} @@ -518,7 +519,7 @@ func (s *HolderSyncer) syncIndex(index string) error { // Sync with every other host. for _, node := range Nodes(s.Cluster.Nodes).FilterHost(s.URI.HostPort()) { - client, err := NewInternalHTTPClient(node.Host, s.ClientOptions) + client, err := NewInternalHTTPClient(node.Host, s.RemoteClient) if err != nil { return err } @@ -563,7 +564,7 @@ func (s *HolderSyncer) syncFrame(index, name string) error { // Sync with every other host. for _, node := range Nodes(s.Cluster.Nodes).FilterHost(s.URI.HostPort()) { - client, err := NewInternalHTTPClient(node.Host, s.ClientOptions) + client, err := NewInternalHTTPClient(node.Host, s.RemoteClient) if err != nil { return err } @@ -616,11 +617,11 @@ func (s *HolderSyncer) syncFragment(index, frame, view string, slice uint64) err // Sync fragments together. fs := FragmentSyncer{ - Fragment: frag, - Host: s.URI.HostPort(), - Cluster: s.Cluster, - Closing: s.Closing, - ClientOptions: s.ClientOptions, + Fragment: frag, + Host: s.URI.HostPort(), + Cluster: s.Cluster, + Closing: s.Closing, + RemoteClient: s.RemoteClient, } if err := fs.SyncFragment(); err != nil { return err diff --git a/holder_test.go b/holder_test.go index 7020a1b8c..6cbfdf150 100644 --- a/holder_test.go +++ b/holder_test.go @@ -320,7 +320,7 @@ func TestHolder_DeleteIndex(t *testing.T) { // Ensure holder can sync with a remote holder. func TestHolderSyncer_SyncHolder(t *testing.T) { cluster := test.NewCluster(2) - + client := pilosa.GetHTTPClient(nil) // Create a local holder. hldr0 := test.MustOpenHolder() defer hldr0.Close() @@ -332,7 +332,7 @@ func TestHolderSyncer_SyncHolder(t *testing.T) { defer s.Close() s.Handler.Holder = hldr1.Holder s.Handler.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { - e := pilosa.NewExecutor(nil) + e := pilosa.NewExecutor(client) e.Holder = hldr1.Holder e.Scheme = cluster.Nodes[1].Scheme e.Host = cluster.Nodes[1].Host @@ -400,9 +400,10 @@ func TestHolderSyncer_SyncHolder(t *testing.T) { t.Fatal(err) } syncer := pilosa.HolderSyncer{ - Holder: hldr0.Holder, - URI: uri, - Cluster: cluster, + Holder: hldr0.Holder, + URI: uri, + Cluster: cluster, + RemoteClient: pilosa.GetHTTPClient(nil), } if err := syncer.SyncHolder(); err != nil { diff --git a/server.go b/server.go index a21b408c3..ffaad0b71 100644 --- a/server.go +++ b/server.go @@ -58,6 +58,7 @@ type Server struct { Handler *Handler Broadcaster Broadcaster BroadcastReceiver BroadcastReceiver + RemoteClient *http.Client // Cluster configuration. // Host is replaced with actual host after opening if port is ":0". @@ -167,10 +168,10 @@ func (s *Server) Open() error { } // Create default HTTP client - s.createDefaultClient() + s.createDefaultClient(s.RemoteClient) // Create executor for executing queries. - e := NewExecutor(&ClientOptions{TLS: s.TLS}) + e := NewExecutor(s.RemoteClient) e.Holder = s.Holder e.Scheme = s.URI.Scheme() e.Host = s.URI.HostPort() @@ -229,6 +230,25 @@ func (s *Server) Addr() net.Addr { } return s.ln.Addr() } +func GetHTTPClient(t *tls.Config) *http.Client { + transport := &http.Transport{ + Proxy: http.ProxyFromEnvironment, + DialContext: (&net.Dialer{ + Timeout: 30 * time.Second, + KeepAlive: 30 * time.Second, + DualStack: true, + }).DialContext, + MaxIdleConns: 1000, + MaxIdleConnsPerHost: 200, + IdleConnTimeout: 90 * time.Second, + TLSHandshakeTimeout: 10 * time.Second, + ExpectContinueTimeout: 1 * time.Second, + } + if t != nil { + transport.TLSClientConfig = t + } + return &http.Client{Transport: transport} +} // Logger returns a logger that writes to LogOutput func (s *Server) Logger() *log.Logger { return log.New(s.LogOutput, "", log.LstdFlags) } @@ -256,7 +276,7 @@ func (s *Server) monitorAntiEntropy() { syncer.URI = s.URI syncer.Cluster = s.Cluster syncer.Closing = s.closing - syncer.ClientOptions = &ClientOptions{TLS: s.TLS} + syncer.RemoteClient = s.RemoteClient // Sync holders. if err := syncer.SyncHolder(); err != nil { @@ -603,12 +623,8 @@ func (s *Server) monitorRuntime() { } } -func (s *Server) createDefaultClient() { - transport := &http.Transport{} - if s.TLS != nil { - transport.TLSClientConfig = s.TLS - } - s.defaultClient = NewInternalHTTPClientFromURI(nil, &ClientOptions{TLS: s.TLS}) +func (s *Server) createDefaultClient(remoteClient *http.Client) { + s.defaultClient = NewInternalHTTPClientFromURI(nil, remoteClient) } // CountOpenFiles on operating systems that support lsof. diff --git a/server/server.go b/server/server.go index 7cc8c5dec..a5f17f6d7 100644 --- a/server/server.go +++ b/server/server.go @@ -31,10 +31,11 @@ import ( "crypto/tls" + "io/ioutil" + "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/gossip" "github.com/pilosa/pilosa/statsd" - "io/ioutil" ) func init() { @@ -161,6 +162,7 @@ func (m *Command) SetupServer() error { m.Server.MaxWritesPerRequest = m.Config.MaxWritesPerRequest // Setup TLS + var TLSConfig *tls.Config if uri.Scheme() == "https" { if m.Config.TLS.CertificatePath == "" { return errors.New("certificate path is required for TLS sockets") @@ -176,8 +178,15 @@ func (m *Command) SetupServer() error { Certificates: []tls.Certificate{cert}, InsecureSkipVerify: m.Config.TLS.SkipVerify, } - m.Server.Handler.ClientOptions = &pilosa.ClientOptions{TLS: m.Server.TLS} + + // TODO Review this location + + TLSConfig = m.Server.TLS + } + c := pilosa.GetHTTPClient(TLSConfig) + m.Server.RemoteClient = c + m.Server.Handler.RemoteClient = c // Set internal port (string). gossipPortStr := pilosa.DefaultGossipPort diff --git a/server/server_test.go b/server/server_test.go index d22abbd1b..28886bc67 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -50,7 +50,7 @@ func TestMain_Set_Quick(t *testing.T) { defer m.Close() // Create client. - client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), nil) + client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) if err != nil { t.Fatal(err) } @@ -323,7 +323,7 @@ func TestMain_FrameRestore(t *testing.T) { defer m2.Close() // Import from first cluster. - client, err := pilosa.NewInternalHTTPClient(m2.Server.URI.HostPort(), nil) + client, err := pilosa.NewInternalHTTPClient(m2.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) if err != nil { t.Fatal(err) } else if err := m2.Client().CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { @@ -672,7 +672,7 @@ func (m *Main) URL() string { return "http://" + m.Server.Addr().String() } // Client returns a client to connect to the program. func (m *Main) Client() *pilosa.InternalHTTPClient { - client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), nil) + client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) if err != nil { panic(err) } diff --git a/test/client.go b/test/client.go index 4ea2fbbad..5eac8cb79 100644 --- a/test/client.go +++ b/test/client.go @@ -1,6 +1,8 @@ package test import ( + "net/http" + "github.com/pilosa/pilosa" ) @@ -10,8 +12,8 @@ type Client struct { } // MustNewClient returns a new instance of Client. Panic on error. -func MustNewClient(host string) *Client { - c, err := pilosa.NewInternalHTTPClient(host, nil) +func MustNewClient(host string, h *http.Client) *Client { + c, err := pilosa.NewInternalHTTPClient(host, h) if err != nil { panic(err) } diff --git a/test/executor.go b/test/executor.go index be370908d..1cddfd521 100644 --- a/test/executor.go +++ b/test/executor.go @@ -15,7 +15,7 @@ type Executor struct { // NewExecutor returns a new instance of Executor. // The executor always matches the hostname of the first cluster node. func NewExecutor(holder *pilosa.Holder, cluster *pilosa.Cluster) *Executor { - executor := pilosa.NewExecutor(nil) + executor := pilosa.NewExecutor(pilosa.GetHTTPClient(nil)) e := &Executor{Executor: executor} e.Holder = holder e.Cluster = cluster From e7d64c4f48e39f15ca55b0cd3b5c9fb3b5fca0f5 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 4 Dec 2017 12:17:41 -0600 Subject: [PATCH 2/2] limit httpclient instances on executor tests --- test/executor.go | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/test/executor.go b/test/executor.go index 1cddfd521..f00692361 100644 --- a/test/executor.go +++ b/test/executor.go @@ -1,6 +1,7 @@ package test import ( + "net/http" "strings" "github.com/pilosa/pilosa" @@ -12,10 +13,16 @@ type Executor struct { *pilosa.Executor } +var remoteClient *http.Client + +func init() { + remoteClient = pilosa.GetHTTPClient(nil) +} + // NewExecutor returns a new instance of Executor. // The executor always matches the hostname of the first cluster node. func NewExecutor(holder *pilosa.Holder, cluster *pilosa.Cluster) *Executor { - executor := pilosa.NewExecutor(pilosa.GetHTTPClient(nil)) + executor := pilosa.NewExecutor(remoteClient) e := &Executor{Executor: executor} e.Holder = holder e.Cluster = cluster