mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 10:54:59 +00:00
Merge pull request #991 from tgruben/single-http-client
refactored httpclient handling
This commit is contained in:
commit
29aca37d83
13 changed files with 118 additions and 97 deletions
30
client.go
30
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,
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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
|
||||
}
|
||||
|
|
|
|||
24
executor.go
24
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]
|
||||
}
|
||||
|
|
|
|||
11
fragment.go
11
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
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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())
|
||||
|
|
|
|||
21
holder.go
21
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
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
34
server.go
34
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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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(nil)
|
||||
executor := pilosa.NewExecutor(remoteClient)
|
||||
e := &Executor{Executor: executor}
|
||||
e.Holder = holder
|
||||
e.Cluster = cluster
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue