diff --git a/Makefile b/Makefile index d5f1022d6..4ce8ac6ee 100644 --- a/Makefile +++ b/Makefile @@ -7,7 +7,7 @@ VERSION := $(shell git describe --tags 2> /dev/null || echo unknown) IDENTIFIER := $(VERSION)-$(GOOS)-$(GOARCH) CLONE_URL=github.com/pilosa/pilosa BUILD_TIME=`date -u +%FT%T%z` -LDFLAGS=-ldflags "-X github.com/pilosa/pilosa/cmd.Version=$(VERSION) -X github.com/pilosa/pilosa/cmd.BuildTime=$(BUILD_TIME)" +LDFLAGS=-ldflags "-X github.com/pilosa/pilosa.Version=$(VERSION) -X github.com/pilosa/pilosa.BuildTime=$(BUILD_TIME)" default: test pilosa diff --git a/client.go b/client.go index 96e725876..5df0385d6 100644 --- a/client.go +++ b/client.go @@ -87,6 +87,8 @@ func (c *Client) maxSliceByIndex(ctx context.Context, inverse bool) (map[string] return nil, err } + req.Header.Set("User-Agent", "pilosa/"+Version) + // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -119,6 +121,8 @@ func (c *Client) Schema(ctx context.Context) ([]*IndexInfo, error) { return nil, err } + req.Header.Set("User-Agent", "pilosa/"+Version) + // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -154,6 +158,7 @@ func (c *Client) CreateIndex(ctx context.Context, index string, opt IndexOptions req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/json") req.Header.Set("Accept", "application/json") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request against the host. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -195,6 +200,8 @@ func (c *Client) FragmentNodes(ctx context.Context, index string, slice uint64) return nil, err } + req.Header.Set("User-Agent", "pilosa/"+Version) + // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -242,6 +249,7 @@ func (c *Client) ExecuteQuery(ctx context.Context, index, query string, allowRed req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/x-protobuf") req.Header.Set("Accept", "application/x-protobuf") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request against the host. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -283,6 +291,8 @@ func (c *Client) ExecutePQL(ctx context.Context, index, query string) (interface if err != nil { return nil, err } + req.Header.Set("User-Agent", "pilosa/"+Version) + resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -362,6 +372,7 @@ func (c *Client) importNode(ctx context.Context, node *Node, buf []byte) error { req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/x-protobuf") req.Header.Set("Accept", "application/x-protobuf") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request against the host. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -438,6 +449,7 @@ func (c *Client) exportNodeCSV(ctx context.Context, node *Node, index, frame str return err } req.Header.Set("Accept", "text/csv") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request against the host. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -572,6 +584,8 @@ func (c *Client) backupSliceNode(ctx context.Context, index, frame, view string, return nil, err } + req.Header.Set("User-Agent", "pilosa/"+Version) + // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -657,6 +671,7 @@ func (c *Client) restoreSliceFrom(ctx context.Context, buf []byte, index, frame, return err } req.Header.Set("Content-Type", "application/octet-stream") + req.Header.Set("User-Agent", "pilosa/"+Version) resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -696,6 +711,7 @@ func (c *Client) CreateFrame(ctx context.Context, index, frame string, opt Frame req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Content-Type", "application/json") req.Header.Set("Accept", "application/json") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request against the host. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -738,6 +754,7 @@ func (c *Client) RestoreFrame(ctx context.Context, host, index, frame string) er return err } req.Header.Set("Content-Type", "application/octet-stream") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -767,6 +784,7 @@ func (c *Client) FrameViews(ctx context.Context, index, frame string) ([]string, return nil, err } req.Header.Set("Accept", "application/json") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request against the host. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -814,6 +832,8 @@ func (c *Client) FragmentBlocks(ctx context.Context, index, frame, view string, return nil, err } + req.Header.Set("User-Agent", "pilosa/"+Version) + // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -859,6 +879,7 @@ func (c *Client) BlockData(ctx context.Context, index, frame, view string, slice req.Header.Set("Content-Type", "application/protobuf") req.Header.Set("Content-Length", strconv.Itoa(len(buf))) req.Header.Set("Accept", "application/protobuf") + req.Header.Set("User-Agent", "pilosa/"+Version) resp, err := c.HTTPClient.Do(req.WithContext(ctx)) if err != nil { @@ -905,6 +926,7 @@ func (c *Client) ColumnAttrDiff(ctx context.Context, index string, blks []AttrBl return nil, err } req.Header.Set("Content-Type", "application/json") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) @@ -948,6 +970,7 @@ func (c *Client) RowAttrDiff(ctx context.Context, index, frame string, blks []At return nil, err } req.Header.Set("Content-Type", "application/json") + req.Header.Set("User-Agent", "pilosa/"+Version) // Execute request. resp, err := c.HTTPClient.Do(req.WithContext(ctx)) diff --git a/cluster.go b/cluster.go index 8c59059a1..70208f0c7 100644 --- a/cluster.go +++ b/cluster.go @@ -17,6 +17,7 @@ package pilosa import ( "encoding/binary" "hash/fnv" + "time" "github.com/pilosa/pilosa/internal" ) @@ -130,6 +131,9 @@ type Cluster struct { // The number of replicas a partition has. ReplicaN int + + // Threshold for logging long-running queries + LongQueryTime time.Duration } // NewCluster returns a new instance of Cluster with defaults. diff --git a/cmd/root.go b/cmd/root.go index d50dd05b4..228a623e4 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -19,22 +19,17 @@ import ( "io" "strings" + "github.com/pilosa/pilosa" "github.com/spf13/cobra" "github.com/spf13/pflag" "github.com/spf13/viper" ) -var ( - Version string - BuildTime string -) - // TODO maybe give this an Add method which will ensure two command // with same name aren't added var subcommandFns = map[string]func(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command{} func NewRootCommand(stdin io.Reader, stdout, stderr io.Writer) *cobra.Command { - setupVersionBuild() rc := &cobra.Command{ Use: "pilosa", Short: "Pilosa - A Distributed In-memory Binary Bitmap Index.", @@ -46,8 +41,8 @@ tools for administering pilosa, importing/exporting data, backing up, and more. Complete documentation is available at https://www.pilosa.com/docs/ -Version: ` + Version + ` -Build Time: ` + BuildTime + "\n", +Version: ` + pilosa.Version + ` +Build Time: ` + pilosa.BuildTime + "\n", PersistentPreRunE: func(cmd *cobra.Command, args []string) error { v := viper.New() err := setAllConfig(v, cmd.Flags(), "PILOSA") @@ -77,15 +72,6 @@ Build Time: ` + BuildTime + "\n", return rc } -func setupVersionBuild() { - if Version == "" { - Version = "v0.0.0" - } - if BuildTime == "" { - BuildTime = "not recorded" - } -} - // setAllConfig takes a FlagSet to be the definition of all configuration // options, as well as their defaults. It then reads from the command line, the // environment, and a config file (if specified), and applies the configuration diff --git a/cmd/server.go b/cmd/server.go index 833e567bb..0c85bc7a8 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -24,6 +24,7 @@ import ( "github.com/spf13/cobra" + "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/server" ) @@ -41,8 +42,7 @@ It will load existing data from the configured directory, and start listening client connections on the configured port.`, RunE: func(cmd *cobra.Command, args []string) error { - Server.Server.Handler.Version = Version - fmt.Fprintf(Server.Stderr, "Pilosa %s, build time %s\n", Version, BuildTime) + fmt.Fprintf(Server.Stderr, "Pilosa %s, build time %s\n", pilosa.Version, pilosa.BuildTime) // Start CPU profiling. if Server.CPUProfile != "" { @@ -93,6 +93,7 @@ on the configured port.`, flags.StringSliceVarP(&Server.Config.Cluster.Hosts, "cluster.hosts", "", []string{}, "Comma separated list of hosts in cluster.") flags.StringSliceVarP(&Server.Config.Cluster.InternalHosts, "cluster.internal-hosts", "", []string{}, "Comma separated list of hosts in cluster used for internal communication.") flags.DurationVarP((*time.Duration)(&Server.Config.Cluster.PollingInterval), "cluster.poll-interval", "", time.Minute, "Polling interval for cluster.") // TODO what actually is this? + flags.DurationVarP((*time.Duration)(&Server.Config.Cluster.LongQueryTime), "long-query-time", "", 10*time.Second, "Threshold for logging long-running queries (0 to disable)") flags.StringVarP(&Server.Config.Plugins.Path, "plugins.path", "", "", "Path to plugin directory.") flags.StringVar(&Server.Config.LogPath, "log-path", "", "Log path") flags.DurationVarP((*time.Duration)(&Server.Config.AntiEntropy.Interval), "anti-entropy.interval", "", time.Minute*10, "Interval at which to run anti-entropy routine.") diff --git a/config.go b/config.go index 3c2f39f6a..107d1ff10 100644 --- a/config.go +++ b/config.go @@ -46,6 +46,7 @@ type Config struct { PollingInterval Duration `toml:"polling-interval"` InternalPort string `toml:"internal-port"` GossipSeed string `toml:"gossip-seed"` + LongQueryTime Duration `toml:"long-query-time"` } `toml:"cluster"` Plugins struct { diff --git a/executor.go b/executor.go index a2c5ab8f1..d9a58750c 100644 --- a/executor.go +++ b/executor.go @@ -1029,6 +1029,7 @@ func (e *Executor) exec(ctx context.Context, node *Node, index string, q *pql.Qu // Require protobuf encoding. req.Header.Set("Accept", "application/x-protobuf") req.Header.Set("Content-Type", "application/x-protobuf") + req.Header.Set("User-Agent", "pilosa/"+Version) // Send request to remote node. resp, err := e.HTTPClient.Do(req) diff --git a/executor_test.go b/executor_test.go index 084511fb0..aa1574400 100644 --- a/executor_test.go +++ b/executor_test.go @@ -501,6 +501,7 @@ func TestExecutor_Execute_Remote_Bitmap(t *testing.T) { // The local node owns slice 1. hldr := MustOpenHolder() defer hldr.Close() + s.Handler.Holder = hldr.Holder hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).MustSetBits(10, (1*SliceWidth)+1) e := NewExecutor(hldr.Holder, c) @@ -528,6 +529,7 @@ func TestExecutor_Execute_Remote_Count(t *testing.T) { // Create local executor data. The local node owns slice 1. hldr := MustOpenHolder() defer hldr.Close() + s.Handler.Holder = hldr.Holder hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).MustSetBits(10, (2*SliceWidth)+1) hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).MustSetBits(10, (2*SliceWidth)+2) @@ -564,6 +566,7 @@ func TestExecutor_Execute_Remote_SetBit(t *testing.T) { // Create local executor data. hldr := MustOpenHolder() defer hldr.Close() + s.Handler.Holder = hldr.Holder // Create frame. if _, err := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{}).CreateFrame("f", pilosa.FrameOptions{}); err != nil { @@ -609,6 +612,7 @@ func TestExecutor_Execute_Remote_SetBit_With_Timestamp(t *testing.T) { // Create local executor data. hldr := MustOpenHolder() defer hldr.Close() + s.Handler.Holder = hldr.Holder // Create frame. if f, err := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{}).CreateFrame("f", pilosa.FrameOptions{}); err != nil { @@ -676,6 +680,7 @@ func TestExecutor_Execute_Remote_TopN(t *testing.T) { // Create local executor data on slice 2 & 4. hldr := MustOpenHolder() defer hldr.Close() + s.Handler.Holder = hldr.Holder hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).MustSetBits(30, (2*SliceWidth)+1) hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 4).MustSetBits(30, (4*SliceWidth)+2) diff --git a/handler.go b/handler.go index 3f7bb8845..6cf175e05 100644 --- a/handler.go +++ b/handler.go @@ -61,13 +61,23 @@ type Handler struct { Execute(context context.Context, index string, query *pql.Query, slices []uint64, opt *ExecOptions) ([]interface{}, error) } - // The version to report on the /version endpoint. - Version string - // The writer for any logging. LogOutput io.Writer } +// externalPrefixFlag denotes endpoints that are intended to be exposed to clients. +// This is used for stats tagging. +var externalPrefixFlag = map[string]bool{ + "schema": true, + "query": true, + "import": true, + "export": true, + "index": true, + "frame": true, + "nodes": true, + "version": true, +} + // NewHandler returns a new instance of Handler with a default logger. func NewHandler() *Handler { handler := &Handler{ @@ -126,7 +136,30 @@ func (h *Handler) methodNotAllowedHandler(w http.ResponseWriter, r *http.Request // ServeHTTP handles an HTTP request. func (h *Handler) ServeHTTP(w http.ResponseWriter, r *http.Request) { + t := time.Now() h.Router.ServeHTTP(w, r) + dif := time.Since(t) + + // Handle some stats tagging + statsTags := make([]string, 0, 3) + + if h.Cluster.LongQueryTime > 0 && dif > h.Cluster.LongQueryTime { + h.logger().Printf("%s %s %.03fs", r.Method, r.URL.String(), float64(dif)) + statsTags = append(statsTags, "slow_query") + } + + pathParts := strings.Split(r.URL.Path, "/") + endpointName := strings.Join(pathParts, "_") + + if externalPrefixFlag[pathParts[1]] { + statsTags = append(statsTags, "external") + } + + // useragent tag identifies internal/external endpoints + statsTags = append(statsTags, "useragent:"+r.UserAgent()) + + stats := h.Holder.Stats.WithTags(statsTags...) + stats.Histogram("http_"+endpointName, float64(dif)) } func (h *Handler) handleWebUI(w http.ResponseWriter, r *http.Request) { @@ -1286,7 +1319,7 @@ func (h *Handler) handleGetVersion(w http.ResponseWriter, r *http.Request) { if err := json.NewEncoder(w).Encode(struct { Version string `json:"version"` }{ - Version: h.Version, + Version: Version, }); err != nil { h.logger().Printf("write version response error: %s", err) } diff --git a/handler_test.go b/handler_test.go index c2aa50927..9e1895bed 100644 --- a/handler_test.go +++ b/handler_test.go @@ -36,8 +36,15 @@ import ( // Ensure the handler returns "not found" for invalid paths. func TestHandler_NotFound(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder + w := httptest.NewRecorder() - NewHandler().ServeHTTP(w, MustNewHTTPRequest("GET", "/no_such_path", nil)) + h.ServeHTTP(w, MustNewHTTPRequest("GET", "/no_such_path", nil)) if w.Code != http.StatusNotFound { t.Fatalf("invalid status: %d", w.Code) } @@ -69,6 +76,7 @@ func TestHandler_Schema(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, MustNewHTTPRequest("GET", "/schema", nil)) if w.Code != http.StatusOK { @@ -93,6 +101,7 @@ func TestHandler_MaxSlices(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, MustNewHTTPRequest("GET", "/slices/max", nil)) if w.Code != http.StatusOK { @@ -133,6 +142,7 @@ func TestHandler_MaxSlices_Inverse(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, MustNewHTTPRequest("GET", "/slices/max?inverse=true", nil)) if w.Code != http.StatusOK { @@ -144,7 +154,12 @@ func TestHandler_MaxSlices_Inverse(t *testing.T) { // Ensure the handler can accept URL arguments. func TestHandler_Query_Args_URL(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { if index != "idx0" { t.Fatalf("unexpected index: %s", index) @@ -167,7 +182,12 @@ func TestHandler_Query_Args_URL(t *testing.T) { // Ensure the handler can accept arguments via protobufs. func TestHandler_Query_Args_Protobuf(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { if index != "idx0" { t.Fatalf("unexpected index: %s", index) @@ -202,7 +222,14 @@ func TestHandler_Query_Args_Protobuf(t *testing.T) { // Ensure the handler returns an error when parsing bad arguments. func TestHandler_Query_Args_Err(t *testing.T) { w := httptest.NewRecorder() - NewHandler().ServeHTTP(w, MustNewHTTPRequest("POST", "/index/idx0/query?slices=a,b", strings.NewReader("Bitmap(id=100)"))) + hldr := MustOpenHolder() + defer hldr.Close() + + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder + + h.ServeHTTP(w, MustNewHTTPRequest("POST", "/index/idx0/query?slices=a,b", strings.NewReader("Bitmap(id=100)"))) if w.Code != http.StatusBadRequest { t.Fatalf("unexpected status code: %d", w.Code) } else if body := w.Body.String(); body != `{"error":"invalid slice argument"}`+"\n" { @@ -212,7 +239,12 @@ func TestHandler_Query_Args_Err(t *testing.T) { // Ensure the handler can execute a query with a uint64 response as JSON. func TestHandler_Query_Uint64_JSON(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{uint64(100)}, nil } @@ -228,7 +260,12 @@ func TestHandler_Query_Uint64_JSON(t *testing.T) { // Ensure the handler can execute a query with a uint64 response as protobufs. func TestHandler_Query_Uint64_Protobuf(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{uint64(100)}, nil } @@ -251,7 +288,12 @@ func TestHandler_Query_Uint64_Protobuf(t *testing.T) { // Ensure the handler can execute a query that returns a bitmap as JSON. func TestHandler_Query_Bitmap_JSON(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, 3, 66, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": 1, "d": true} @@ -284,6 +326,7 @@ func TestHandler_Query_Bitmap_ColumnAttrs_JSON(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, 3, 66, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": 1, "d": true} @@ -301,7 +344,12 @@ func TestHandler_Query_Bitmap_ColumnAttrs_JSON(t *testing.T) { // Ensure the handler can execute a query that returns a bitmap as protobuf. func TestHandler_Query_Bitmap_Protobuf(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": int64(1), "d": true} @@ -347,6 +395,7 @@ func TestHandler_Query_Bitmap_ColumnAttrs_Protobuf(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { bm := pilosa.NewBitmap(1, pilosa.SliceWidth+1) bm.Attrs = map[string]interface{}{"a": "b", "c": int64(1), "d": true} @@ -400,7 +449,12 @@ func TestHandler_Query_Bitmap_ColumnAttrs_Protobuf(t *testing.T) { // Ensure the handler can execute a query that returns pairs as JSON. func TestHandler_Query_Pairs_JSON(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{[]pilosa.Pair{ {ID: 1, Count: 2}, @@ -419,7 +473,12 @@ func TestHandler_Query_Pairs_JSON(t *testing.T) { // Ensure the handler can execute a query that returns pairs as protobuf. func TestHandler_Query_Pairs_Protobuf(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return []interface{}{[]pilosa.Pair{ {ID: 1, Count: 2}, @@ -445,7 +504,12 @@ func TestHandler_Query_Pairs_Protobuf(t *testing.T) { // Ensure the handler can return an error as JSON. func TestHandler_Query_Err_JSON(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return nil, errors.New("marker") } @@ -461,7 +525,12 @@ func TestHandler_Query_Err_JSON(t *testing.T) { // Ensure the handler can return an error as protobuf. func TestHandler_Query_Err_Protobuf(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder h.Executor.ExecuteFn = func(ctx context.Context, index string, query *pql.Query, slices []uint64, opt *pilosa.ExecOptions) ([]interface{}, error) { return nil, errors.New("marker") } @@ -484,8 +553,14 @@ func TestHandler_Query_Err_Protobuf(t *testing.T) { // Ensure the handler returns "method not allowed" for non-POST queries. func TestHandler_Query_MethodNotAllowed(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder w := httptest.NewRecorder() - NewHandler().ServeHTTP(w, MustNewHTTPRequest("GET", "/index/i/query", nil)) + h.ServeHTTP(w, MustNewHTTPRequest("GET", "/index/i/query", nil)) if w.Code != http.StatusMethodNotAllowed { t.Fatalf("invalid status: %d", w.Code) } @@ -493,7 +568,12 @@ func TestHandler_Query_MethodNotAllowed(t *testing.T) { // Ensure the handler returns an error if there is a parsing error.. func TestHandler_Query_ErrParse(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder w := httptest.NewRecorder() h.ServeHTTP(w, MustNewHTTPRequest("POST", "/index/idx0/query?slices=0,1", strings.NewReader("bad_fn("))) if w.Code != http.StatusBadRequest { @@ -550,6 +630,7 @@ func TestHandler_DeleteFrame(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, MustNewHTTPRequest("DELETE", "/index/i0/frame/f1", strings.NewReader(""))) if w.Code != http.StatusOK { @@ -569,6 +650,7 @@ func TestHandler_SetIndexTimeQuantum(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, MustNewHTTPRequest("PATCH", "/index/i0/time-quantum", strings.NewReader(`{"timeQuantum":"ymdh"}`))) if w.Code != http.StatusOK { @@ -592,6 +674,7 @@ func TestHandler_SetFrameTimeQuantum(t *testing.T) { h := NewHandler() h.Holder = hldr.Holder + h.Cluster = NewCluster(1) w := httptest.NewRecorder() h.ServeHTTP(w, MustNewHTTPRequest("PATCH", "/index/i0/frame/f1/time-quantum", strings.NewReader(`{"timeQuantum":"ymdh"}`))) if w.Code != http.StatusOK { @@ -753,22 +836,30 @@ func TestHandler_Fragment_BackupRestore(t *testing.T) { // Ensure the handler can retrieve the version. func TestHandler_Version(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() - h.Version = "1.0.0" + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder w := httptest.NewRecorder() r := MustNewHTTPRequest("GET", "/version", nil) h.ServeHTTP(w, r) if w.Code != http.StatusOK { t.Fatalf("unexpected status code: %d", w.Code) - } else if w.Body.String() != `{"version":"1.0.0"}`+"\n" { + } else if w.Body.String() != `{"version":"`+pilosa.Version+`"}`+"\n" { t.Fatalf("unexpected body: %q", w.Body.String()) } } // Ensure the handler can return a list of nodes for a fragment. func TestHandler_Fragment_Nodes(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Holder = hldr.Holder h.Cluster = NewCluster(3) h.Cluster.ReplicaN = 2 @@ -784,7 +875,12 @@ func TestHandler_Fragment_Nodes(t *testing.T) { // Ensure the handler can return expvars without panicking. func TestHandler_Expvars(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + h := NewHandler() + h.Cluster = NewCluster(1) + h.Holder = hldr.Holder w := httptest.NewRecorder() r := MustNewHTTPRequest("GET", "/debug/vars", nil) h.ServeHTTP(w, r) diff --git a/httpbroadcast/messenger.go b/httpbroadcast/messenger.go index 387f28585..09c3b7f4b 100644 --- a/httpbroadcast/messenger.go +++ b/httpbroadcast/messenger.go @@ -97,6 +97,7 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *pilosa.Node, msg []byte) error { // Require protobuf encoding. req.Header.Set("Content-Type", "application/x-protobuf") + req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) // Send request to remote node. resp, err := client.Do(req) diff --git a/server.go b/server.go index 60570bc0f..03ea0cbf8 100644 --- a/server.go +++ b/server.go @@ -64,6 +64,9 @@ type Server struct { PollingInterval time.Duration MetricInterval time.Duration + // Threshold for logging long queries + LongQueryTime time.Duration + LogOutput io.Writer } @@ -410,6 +413,7 @@ func checkMaxSlices(hostport string) (map[string]uint64, error) { // Require protobuf encoding. req.Header.Set("Accept", "application/x-protobuf") req.Header.Set("Content-Type", "application/x-protobuf") + req.Header.Set("User-Agent", "pilosa/"+Version) // Send request to remote node. resp, err := http.DefaultClient.Do(req) diff --git a/server/server.go b/server/server.go index e2360522d..d99a05463 100644 --- a/server/server.go +++ b/server/server.go @@ -193,6 +193,7 @@ func (m *Command) SetupServer() error { // Set configuration options. m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval) + m.Server.Cluster.LongQueryTime = time.Duration(m.Config.Cluster.LongQueryTime) return nil } diff --git a/version.go b/version.go new file mode 100644 index 000000000..3c7d19ca3 --- /dev/null +++ b/version.go @@ -0,0 +1,4 @@ +package pilosa + +var Version = "v0.0.0" +var BuildTime = "not recorded"