diff --git a/Dockerfile b/Dockerfile index f8db1c4b1..e32437ef7 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,29 +1,27 @@ -FROM golang:1.8.0 - -MAINTAINER Pilosa Corp. +FROM golang:1.8.1 as builder ARG ldflags='' +ARG GLIDE="https://github.com/Masterminds/glide/releases/download/v0.12.3/glide-v0.12.3-linux-amd64.tar.gz" +ARG GLIDE_HASH="d6d3816c70fba716466e7381a9c06cb31565a3b87acb5bad9dd3beb0a9f9b0f8" + +COPY . /go/src/github.com/pilosa/pilosa + +RUN wget ${GLIDE} -O /go/glide.tar.gz -q \ + && tar xf /go/glide.tar.gz \ + && mv /go/linux-amd64/glide /go/bin \ + && [ "$(sha256sum /go/bin/glide | cut -d' ' -f1)" = "$GLIDE_HASH" ] \ + && cd /go/src/github.com/pilosa/pilosa \ + && make vendor \ + && CGO_ENABLED=0 go install -a -ldflags "$ldflags" github.com/pilosa/pilosa/cmd/pilosa + +FROM scratch + +LABEL maintainer "dev@pilosa.com" + +COPY --from=builder /go/bin/pilosa /pilosa EXPOSE 10101 VOLUME /data -RUN echo 'data-dir = "/data"' > /config - -RUN git clone --depth 1 https://github.com/Masterminds/glide.git /go/src/github.com/Masterminds/glide \ - && cd /go/src/github.com/Masterminds/glide \ - && git fetch --tags --depth 1 \ - && git checkout tags/v0.12.3 -b build \ - && make build \ - && mv ./glide /go/bin \ - && cd / \ - && rm -r /go/src/github.com/Masterminds/glide - -COPY . /go/src/github.com/pilosa/pilosa - -RUN cd /go/src/github.com/pilosa/pilosa \ - && make vendor \ - && CGO_ENABLED=0 go install -a -ldflags "$ldflags" github.com/pilosa/pilosa/cmd/pilosa \ - && rm -rf /go/src /go/pkg - -ENTRYPOINT ["/go/bin/pilosa"] -CMD ["server", "--config", "/config"] +ENTRYPOINT ["/pilosa"] +CMD ["server", "--data-dir", "/data"] diff --git a/Makefile b/Makefile index 4ce8ac6ee..6b9ff1dc7 100644 --- a/Makefile +++ b/Makefile @@ -62,5 +62,6 @@ endif docker: docker build -t "pilosa:$(VERSION)" \ - --build-arg ldflags="-X main.Version=$(VERSION) -X main.BuildTime=$(BUILD_TIME)" . + --build-arg ldflags="-X github.com/pilosa/pilosa/cmd.Version=$(VERSION) \ + -X github.com/pilosa/pilosa/cmd.BuildTime=$(BUILD_TIME)" . @echo "Created image: pilosa:$(VERSION)" diff --git a/README.md b/README.md index 377ea62aa..40eab2edf 100644 --- a/README.md +++ b/README.md @@ -1,8 +1,15 @@ - - - +

+ + + +

-[![Build Status](https://travis-ci.com/pilosa/pilosa.svg?token=Peb4jvQ3kLbjUEhpU5aR&branch=master)](https://travis-ci.com/pilosa/pilosa) +[![Build Status](https://travis-ci.org/pilosa/pilosa.svg?branch=master)](https://travis-ci.org/pilosa/pilosa) +[![GoDoc](https://godoc.org/github.com/pilosa/pilosa?status.svg)](https://godoc.org/github.com/pilosa/pilosa) +[![Go Report Card](https://goreportcard.com/badge/github.com/pilosa/pilosa)](https://goreportcard.com/report/github.com/pilosa/pilosa) +[![license](https://img.shields.io/github/license/pilosa/pilosa.svg)](https://github.com/pilosa/pilosa/blob/master/LICENSE) +[![CLA Assistant](https://cla-assistant.io/readme/badge/pilosa/pilosa)](https://cla-assistant.io/pilosa/pilosa) +[![GitHub release](https://img.shields.io/github/release/pilosa/pilosa.svg)](https://github.com/pilosa/pilosa/releases) ## An open source, distributed bitmap index. - [Docs](#docs) @@ -57,8 +64,8 @@ There are supported libraries for the following languages: ## Get Support -There are [several channels](https://www.pilosa.com/community/#support) availble for you to reach out to us for support. +There are [several channels](https://www.pilosa.com/community/#support) available for you to reach out to us for support. ## Contributing -Pilosa is an open source project. Please see our [Contributing Guide](https://www.pilosa.com/docs/contributing/) for information about how to get involved. +Pilosa is an open source project. Please see our [Contributing Guide](CONTRIBUTING.md) for information about how to get involved. diff --git a/client_test.go b/client_test.go index a89a41cce..6b5f43195 100644 --- a/client_test.go +++ b/client_test.go @@ -75,39 +75,65 @@ func TestClient_MultiNode(t *testing.T) { } // Create a dispersed set of bitmaps across 3 nodes such that each individual node and slice width increment would reveal a different TopN. - hldr[0].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 9).MustSetBits(100, (SliceWidth*9)+10) - hldr[0].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 9).MustSetBits(4, (SliceWidth*9)+10, (SliceWidth*9)+11, (SliceWidth*9)+12) - hldr[0].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 9).MustSetBits(4, (SliceWidth*9)+10, (SliceWidth*9)+11, (SliceWidth*9)+12, (SliceWidth*9)+13, (SliceWidth*9)+14, (SliceWidth*9)+15) - hldr[0].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 9).MustSetBits(2, (SliceWidth*9)+1, (SliceWidth*9)+2, (SliceWidth*9)+3, (SliceWidth*9)+4) - hldr[0].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 9).MustSetBits(3, (SliceWidth*9)+1, (SliceWidth*9)+2, (SliceWidth*9)+3, (SliceWidth*9)+4, (SliceWidth*9)+5) - hldr[0].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 9).MustSetBits(22, (SliceWidth*9)+1, (SliceWidth*9)+2, (SliceWidth*9)+10) + sliceNums := []uint64{1, 2, 6} + for i, num := range sliceNums { + owns := s[i].Handler.Handler.Cluster.OwnsSlices("i", 20, s[i].Host()) + ownsNum := false + for _, ownNum := range owns { + if ownNum == num { + ownsNum = true + break + } + } + if !ownsNum { + t.Fatalf("Trying to use slice %d on host %s, but it doesn't own that slice. It owns %s", num, s[i].Host(), owns) + } + } - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).MustSetBits(24, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12, (SliceWidth*6)+13, (SliceWidth*6)+14) - hldr[1].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).MustSetBits(99, 1, 2, 3, 4) - hldr[1].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).MustSetBits(100, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10) - hldr[1].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).MustSetBits(98, 1, 2, 3, 4, 5, 6) - hldr[1].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).MustSetBits(1, 4) - hldr[1].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).MustSetBits(22, 1, 2, 3, 4, 5) + baseBit0 := SliceWidth * sliceNums[0] + baseBit1 := SliceWidth * sliceNums[1] + baseBit2 := SliceWidth * sliceNums[2] - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).MustSetBits(20, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12, (SliceWidth*6)+13) - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).MustSetBits(21, (SliceWidth*6)+10) - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).MustSetBits(100, (SliceWidth*6)+10) - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).MustSetBits(99, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12) - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).MustSetBits(98, (SliceWidth*6)+10, (SliceWidth*6)+11) - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).MustSetBits(22, (SliceWidth*6)+10, (SliceWidth*6)+11, (SliceWidth*6)+12) + maxSlice := uint64(0) + for _, x := range sliceNums { + if x > maxSlice { + maxSlice = x + } + } + + hldr[0].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[0]).MustSetBits(100, baseBit0+10) + hldr[0].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[0]).MustSetBits(4, baseBit0+10, baseBit0+11, baseBit0+12) + hldr[0].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[0]).MustSetBits(4, baseBit0+10, baseBit0+11, baseBit0+12, baseBit0+13, baseBit0+14, baseBit0+15) + hldr[0].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[0]).MustSetBits(2, baseBit0+1, baseBit0+2, baseBit0+3, baseBit0+4) + hldr[0].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[0]).MustSetBits(3, baseBit0+1, baseBit0+2, baseBit0+3, baseBit0+4, baseBit0+5) + hldr[0].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[0]).MustSetBits(22, baseBit0+1, baseBit0+2, baseBit0+10) + + hldr[1].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[1]).MustSetBits(99, baseBit1+1, baseBit1+2, baseBit1+3, baseBit1+4) + hldr[1].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[1]).MustSetBits(100, baseBit1+1, baseBit1+2, baseBit1+3, baseBit1+4, baseBit1+5, baseBit1+6, baseBit1+7, baseBit1+8, baseBit1+9, baseBit1+10) + hldr[1].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[1]).MustSetBits(98, baseBit1+1, baseBit1+2, baseBit1+3, baseBit1+4, baseBit1+5, baseBit1+6) + hldr[1].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[1]).MustSetBits(1, baseBit1+4) + hldr[1].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[1]).MustSetBits(22, baseBit1+1, baseBit1+2, baseBit1+3, baseBit1+4, baseBit1+5) + + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).MustSetBits(24, baseBit2+10, baseBit2+11, baseBit2+12, baseBit2+13, baseBit2+14) + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).MustSetBits(20, baseBit2+10, baseBit2+11, baseBit2+12, baseBit2+13) + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).MustSetBits(21, baseBit2+10) + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).MustSetBits(100, baseBit2+10) + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).MustSetBits(99, baseBit2+10, baseBit2+11, baseBit2+12) + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).MustSetBits(98, baseBit2+10, baseBit2+11) + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).MustSetBits(22, baseBit2+10, baseBit2+11, baseBit2+12) // Rebuild the RankCache. // We have to do this to avoid the 10-second cache invalidation delay // built into cache.Invalidate() - hldr[0].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).RecalculateCache() - hldr[1].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 10).RecalculateCache() - hldr[2].MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 6).RecalculateCache() + hldr[0].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[0]).RecalculateCache() + hldr[1].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[1]).RecalculateCache() + hldr[2].MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, sliceNums[2]).RecalculateCache() // Connect to each node to compare results. client := make([]*Client, 3) client[0] = MustNewClient(s[0].Host()) - client[1] = MustNewClient(s[0].Host()) - client[2] = MustNewClient(s[0].Host()) + client[1] = MustNewClient(s[1].Host()) + client[2] = MustNewClient(s[2].Host()) topN := 4 q := fmt.Sprintf(`TopN(frame="%s", n=%d)`, "f", topN) @@ -120,15 +146,15 @@ func TestClient_MultiNode(t *testing.T) { // Check the results before every node has the correct max slice value. pairs := result.(internal.QueryResponse).Results[0].Pairs for _, pair := range pairs { - if pair.Key == 22 && pair.Count != 11 { + if pair.Key == 22 && pair.Count != 3 { t.Fatalf("Invalid Cluster wide MaxSlice prevents accurate calculation of %s", pair) } } // Set max slice to correct value. - hldr[0].Index("i").SetRemoteMaxSlice(10) - hldr[1].Index("i").SetRemoteMaxSlice(10) - hldr[2].Index("i").SetRemoteMaxSlice(10) + hldr[0].Index("i").SetRemoteMaxSlice(maxSlice) + hldr[1].Index("i").SetRemoteMaxSlice(maxSlice) + hldr[2].Index("i").SetRemoteMaxSlice(maxSlice) result, err = client[0].ExecuteQuery(context.Background(), "i", q, true) if err != nil { diff --git a/cluster.go b/cluster.go index 70208f0c7..981e8bdc2 100644 --- a/cluster.go +++ b/cluster.go @@ -24,7 +24,7 @@ import ( const ( // DefaultPartitionN is the default number of partitions in a cluster. - DefaultPartitionN = 16 + DefaultPartitionN = 256 // DefaultReplicaN is the default number of replicas per partition. DefaultReplicaN = 1 diff --git a/cluster_test.go b/cluster_test.go index 3caea3fe2..5e4668fe2 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -137,7 +137,7 @@ func TestCluster_OwnsSlices(t *testing.T) { c := NewCluster(5) slices := c.OwnsSlices("test", 10, "host2") - if !reflect.DeepEqual(slices, []uint64{1, 3, 10}) { + if !reflect.DeepEqual(slices, []uint64{0, 3, 6, 10}) { t.Fatalf("unexpected slices for node's index: %v", slices) } } diff --git a/cmd/server.go b/cmd/server.go index 0c85bc7a8..c077fbf85 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -89,6 +89,7 @@ on the configured port.`, flags.StringVarP(&Server.Config.DataDir, "data-dir", "d", "~/.pilosa", "Directory to store pilosa data files.") flags.StringVarP(&Server.Config.Host, "bind", "b", ":10101", "Default URI on which pilosa should listen.") + flags.IntVarP(&Server.Config.MaxWritesPerRequest, "max-writes-per-request", "", Server.Config.MaxWritesPerRequest, "Number of write commands per request.") flags.IntVarP(&Server.Config.Cluster.ReplicaN, "cluster.replicas", "", 1, "Number of hosts each piece of data should be stored on.") 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.") diff --git a/config.go b/config.go index 107d1ff10..8fa3f8152 100644 --- a/config.go +++ b/config.go @@ -31,6 +31,9 @@ const ( // DefaultMetrics sets the internal metrics to no op DefaultMetrics = "nop" + + // DefaultMaxWritesPerRequest is the default number of writes per request. + DefaultMaxWritesPerRequest = 5000 ) // Config represents the configuration for the command. @@ -57,6 +60,10 @@ type Config struct { Interval Duration `toml:"interval"` } `toml:"anti-entropy"` + // Limits the number of mutating commands that can be in a single request to + // the server. This includes SetBit, ClearBit, SetRowAttrs & SetColumnAttrs. + MaxWritesPerRequest int `toml:"max-writes-per-request"` + LogPath string `toml:"log-path"` Metric struct { @@ -69,7 +76,8 @@ type Config struct { // NewConfig returns an instance of Config with default options. func NewConfig() *Config { c := &Config{ - Host: DefaultHost + ":" + DefaultPort, + Host: DefaultHost + ":" + DefaultPort, + MaxWritesPerRequest: DefaultMaxWritesPerRequest, } c.Cluster.ReplicaN = DefaultReplicaN c.Cluster.Type = DefaultClusterType diff --git a/ctl/config.go b/ctl/config.go index 6de94bf8c..4062250e6 100644 --- a/ctl/config.go +++ b/ctl/config.go @@ -40,6 +40,7 @@ func (cmd *ConfigCommand) Run(ctx context.Context) error { fmt.Fprintln(cmd.Stdout, strings.TrimSpace(` data-dir = "~/.pilosa" bind = "localhost:10101" +max-writes-per-request = 5000 [cluster] poll-interval = "2m0s" diff --git a/executor.go b/executor.go index d9a58750c..755bfda53 100644 --- a/executor.go +++ b/executor.go @@ -49,6 +49,9 @@ type Executor struct { // Client used for remote HTTP requests. HTTPClient *http.Client + + // Maximum number of SetBit() or ClearBit() commands per request. + MaxWritesPerRequest int } // NewExecutor returns a new instance of Executor. @@ -65,6 +68,11 @@ func (e *Executor) Execute(ctx context.Context, index string, q *pql.Query, slic return nil, ErrIndexRequired } + // Verify that the number of writes do not exceed the maximum. + if e.MaxWritesPerRequest > 0 && q.WriteCallN() > e.MaxWritesPerRequest { + return nil, ErrTooManyWrites + } + // Default options. if opt == nil { opt = &ExecOptions{} @@ -497,19 +505,43 @@ func (e *Executor) executeRangeSlice(ctx context.Context, index string, c *pql.C frame = DefaultFrame } + // Retrieve column label. + idx := e.Holder.Index(index) + if idx == nil { + return nil, ErrIndexNotFound + } + columnLabel := idx.ColumnLabel() + // Retrieve base frame. - f := e.Holder.Frame(index, frame) + f := idx.Frame(frame) if f == nil { return nil, ErrFrameNotFound } rowLabel := f.RowLabel() - // Read row id. - rowID, _, err := c.UintArg(rowLabel) // TODO: why are we ignoring missing rowID? + // Read row & column id. + columnID, columnOK, err := c.UintArg(columnLabel) + if err != nil { + return nil, fmt.Errorf("executeRangeSlice - reading column: %v", err) + } + rowID, rowOK, err := c.UintArg(rowLabel) if err != nil { return nil, fmt.Errorf("executeRangeSlice - reading row: %v", err) } + // Determine view. + var id uint64 + var viewName string + if columnOK && rowOK { + return nil, fmt.Errorf("Range() cannot contain both %q and %q", columnLabel, rowLabel) + } else if !columnOK && !rowOK { + return nil, fmt.Errorf("Range() must specify either %q or %q", columnLabel, rowLabel) + } else if columnOK { + viewName, id = ViewInverse, columnID + } else { + viewName, id = ViewStandard, rowID + } + // Parse start time. startTimeStr, ok := c.Args["start"].(string) if !ok { @@ -538,12 +570,12 @@ func (e *Executor) executeRangeSlice(ctx context.Context, index string, c *pql.C // Union bitmaps across all time-based subframes. bm := &Bitmap{} - for _, view := range ViewsByTimeRange(ViewStandard, startTime, endTime, q) { + for _, view := range ViewsByTimeRange(viewName, startTime, endTime, q) { f := e.Holder.Fragment(index, frame, view, slice) if f == nil { continue } - bm = bm.Union(f.Row(rowID)) + bm = bm.Union(f.Row(id)) } f.Stats.Count("range", 1) return bm, nil @@ -1046,7 +1078,7 @@ func (e *Executor) exec(ctx context.Context, node *Node, index string, q *pql.Qu // Check status code. if resp.StatusCode != http.StatusOK { - return nil, fmt.Errorf("invalid status: code=%d, err=%s", resp.StatusCode, body) + return nil, fmt.Errorf("invalid status Executor.exec: code=%d, err=%s, req: %v", resp.StatusCode, body, req) } // Decode response object. diff --git a/executor_test.go b/executor_test.go index aa1574400..b69aaac37 100644 --- a/executor_test.go +++ b/executor_test.go @@ -276,15 +276,19 @@ func TestExecutor_Execute_TopN(t *testing.T) { defer hldr.Close() // Set bits for rows 0, 10, & 20 across two slices. - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth+2) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 5).SetBit(0, (5*SliceWidth)+100) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(10, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "other", pilosa.ViewStandard, 0).SetBit(0, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth+2) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 5).SetBit(0, (5*SliceWidth)+100) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(10, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "other", pilosa.ViewStandard, 0).SetBit(0, 0) + + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).RecalculateCache() + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).RecalculateCache() + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 5).RecalculateCache() // Execute query. e := NewExecutor(hldr.Holder, NewCluster(1)) @@ -302,12 +306,12 @@ func TestExecutor_Execute_TopN_fill(t *testing.T) { defer hldr.Close() // Set bits for rows 0, 10, & 20 across two slices. - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 2) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(1, SliceWidth+2) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(1, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 2) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(1, SliceWidth+2) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(1, SliceWidth) // Execute query. e := NewExecutor(hldr.Holder, NewCluster(1)) @@ -325,23 +329,23 @@ func TestExecutor_Execute_TopN_fill_small(t *testing.T) { hldr := MustOpenHolder() defer hldr.Close() - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).SetBit(0, 2*SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 3).SetBit(0, 3*SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 4).SetBit(0, 4*SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).SetBit(0, 2*SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 3).SetBit(0, 3*SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 4).SetBit(0, 4*SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(1, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(1, 1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(1, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(1, 1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(2, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(2, SliceWidth+1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(2, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(2, SliceWidth+1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).SetBit(3, 2*SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).SetBit(3, 2*SliceWidth+1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).SetBit(3, 2*SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).SetBit(3, 2*SliceWidth+1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 3).SetBit(4, 3*SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 3).SetBit(4, 3*SliceWidth+1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 3).SetBit(4, 3*SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 3).SetBit(4, 3*SliceWidth+1) // Execute query. e := NewExecutor(hldr.Holder, NewCluster(1)) @@ -360,19 +364,23 @@ func TestExecutor_Execute_TopN_Src(t *testing.T) { defer hldr.Close() // Set bits for rows 0, 10, & 20 across two slices. - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth+1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth+1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth+2) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(0, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth+1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth+1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(20, SliceWidth+2) // Create an intersecting row. - hldr.MustCreateFragmentIfNotExists("i", "other", pilosa.ViewStandard, 1).SetBit(100, SliceWidth) - hldr.MustCreateFragmentIfNotExists("i", "other", pilosa.ViewStandard, 1).SetBit(100, SliceWidth+1) - hldr.MustCreateFragmentIfNotExists("i", "other", pilosa.ViewStandard, 1).SetBit(100, SliceWidth+2) + hldr.MustCreateRankedFragmentIfNotExists("i", "other", pilosa.ViewStandard, 1).SetBit(100, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "other", pilosa.ViewStandard, 1).SetBit(100, SliceWidth+1) + hldr.MustCreateRankedFragmentIfNotExists("i", "other", pilosa.ViewStandard, 1).SetBit(100, SliceWidth+2) + + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).RecalculateCache() + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).RecalculateCache() + hldr.MustCreateRankedFragmentIfNotExists("i", "other", pilosa.ViewStandard, 1).RecalculateCache() // Execute query. e := NewExecutor(hldr.Holder, NewCluster(1)) @@ -392,9 +400,9 @@ func TestExecutor_Execute_TopN_Attr(t *testing.T) { // hldr := MustOpenHolder() defer hldr.Close() - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) if err := hldr.Frame("i", "f").RowAttrStore().SetAttrs(10, map[string]interface{}{"category": int64(123)}); err != nil { t.Fatal(err) @@ -415,9 +423,9 @@ func TestExecutor_Execute_TopN_Attr_Src(t *testing.T) { // hldr := MustOpenHolder() defer hldr.Close() - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) - hldr.MustCreateFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 0) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 0).SetBit(0, 1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 1).SetBit(10, SliceWidth) if err := hldr.Frame("i", "f").RowAttrStore().SetAttrs(10, map[string]interface{}{"category": uint64(123)}); err != nil { t.Fatal(err) @@ -437,36 +445,51 @@ func TestExecutor_Execute_TopN_Attr_Src(t *testing.T) { func TestExecutor_Execute_Range(t *testing.T) { hldr := MustOpenHolder() defer hldr.Close() + e := NewExecutor(hldr.Holder, NewCluster(1)) // Create index. index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{}) // Create frame. - f, err := index.CreateFrameIfNotExists("f", pilosa.FrameOptions{}) - if err != nil { - t.Fatal(err) - } else if err := f.SetTimeQuantum(pilosa.TimeQuantum("YMDH")); err != nil { + if _, err := index.CreateFrameIfNotExists("f", pilosa.FrameOptions{ + InverseEnabled: true, + TimeQuantum: pilosa.TimeQuantum("YMDH"), + }); err != nil { t.Fatal(err) } // Set bits. - f.MustSetBit(pilosa.ViewStandard, 1, 2, MustParseTimePtr("1999-12-31 00:00")) - f.MustSetBit(pilosa.ViewStandard, 1, 3, MustParseTimePtr("2000-01-01 00:00")) - f.MustSetBit(pilosa.ViewStandard, 1, 4, MustParseTimePtr("2000-01-02 00:00")) - f.MustSetBit(pilosa.ViewStandard, 1, 5, MustParseTimePtr("2000-02-01 00:00")) - f.MustSetBit(pilosa.ViewStandard, 1, 6, MustParseTimePtr("2001-01-01 00:00")) - f.MustSetBit(pilosa.ViewStandard, 1, 7, MustParseTimePtr("2002-01-01 02:00")) + if _, err := e.Execute(context.Background(), "i", MustParse(` + SetBit(frame=f, rowID=1, columnID=2, timestamp="1999-12-31T00:00") + SetBit(frame=f, rowID=1, columnID=3, timestamp="2000-01-01T00:00") + SetBit(frame=f, rowID=1, columnID=4, timestamp="2000-01-02T00:00") + SetBit(frame=f, rowID=1, columnID=5, timestamp="2000-02-01T00:00") + SetBit(frame=f, rowID=1, columnID=6, timestamp="2001-01-01T00:00") + SetBit(frame=f, rowID=1, columnID=7, timestamp="2002-01-01T02:00") - f.MustSetBit(pilosa.ViewStandard, 1, 2, MustParseTimePtr("1999-12-30 00:00")) // too early - f.MustSetBit(pilosa.ViewStandard, 1, 2, MustParseTimePtr("2002-02-01 00:00")) // too late - f.MustSetBit(pilosa.ViewStandard, 10, 2, MustParseTimePtr("2001-01-01 00:00")) // different row - - e := NewExecutor(hldr.Holder, NewCluster(1)) - if res, err := e.Execute(context.Background(), "i", MustParse(`Range(rowID=1, frame=f, start="1999-12-31T00:00", end="2002-01-01T03:00")`), nil, nil); err != nil { + SetBit(frame=f, rowID=1, columnID=2, timestamp="1999-12-30T00:00") + SetBit(frame=f, rowID=1, columnID=2, timestamp="2002-02-01T00:00") + SetBit(frame=f, rowID=10, columnID=2, timestamp="2001-01-01T00:00") + `), nil, nil); err != nil { t.Fatal(err) - } else if bits := res[0].(*pilosa.Bitmap).Bits(); !reflect.DeepEqual(bits, []uint64{2, 3, 4, 5, 6, 7}) { - t.Fatalf("unexpected bits: %+v", bits) } + + t.Run("Standard", func(t *testing.T) { + if res, err := e.Execute(context.Background(), "i", MustParse(`Range(rowID=1, frame=f, start="1999-12-31T00:00", end="2002-01-01T03:00")`), nil, nil); err != nil { + t.Fatal(err) + } else if bits := res[0].(*pilosa.Bitmap).Bits(); !reflect.DeepEqual(bits, []uint64{2, 3, 4, 5, 6, 7}) { + t.Fatalf("unexpected bits: %+v", bits) + } + }) + + t.Run("Inverse", func(t *testing.T) { + e := NewExecutor(hldr.Holder, NewCluster(1)) + if res, err := e.Execute(context.Background(), "i", MustParse(`Range(columnID=2, frame=f, start="1999-01-01T00:00", end="2003-01-01T00:00")`), nil, nil); err != nil { + t.Fatal(err) + } else if bits := res[0].(*pilosa.Bitmap).Bits(); !reflect.DeepEqual(bits, []uint64{1, 10}) { + t.Fatalf("unexpected bits: %+v", bits) + } + }) } // Ensure a remote query can return a bitmap. @@ -681,8 +704,8 @@ func TestExecutor_Execute_Remote_TopN(t *testing.T) { 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) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 2).MustSetBits(30, (2*SliceWidth)+1) + hldr.MustCreateRankedFragmentIfNotExists("i", "f", pilosa.ViewStandard, 4).MustSetBits(30, (4*SliceWidth)+2) e := NewExecutor(hldr.Holder, c) if res, err := e.Execute(context.Background(), "i", MustParse(`TopN(frame=f, n=3)`), nil, nil); err != nil { @@ -696,6 +719,17 @@ func TestExecutor_Execute_Remote_TopN(t *testing.T) { } } +// Ensure executor returns an error if too many writes are in a single request. +func TestExecutor_Execute_ErrMaxWritesPerRequest(t *testing.T) { + hldr := MustOpenHolder() + defer hldr.Close() + e := NewExecutor(hldr.Holder, NewCluster(1)) + e.MaxWritesPerRequest = 3 + if _, err := e.Execute(context.Background(), "i", MustParse(`SetBit() ClearBit() SetBit() SetBit()`), nil, nil); err != pilosa.ErrTooManyWrites { + t.Fatalf("unexpected error: %s", err) + } +} + // Executor represents a test wrapper for pilosa.Executor. type Executor struct { *pilosa.Executor diff --git a/fragment.go b/fragment.go index 1d8df780f..e0515978d 100644 --- a/fragment.go +++ b/fragment.go @@ -82,9 +82,9 @@ type Fragment struct { opN int // number of ops since snapshot // Cache for row counts. - cacheType string // passed in by frame + CacheType string // passed in by frame cache Cache - cacheSize uint32 + CacheSize uint32 // Cache containing full rows (not just counts). rowCache BitmapCache @@ -115,8 +115,8 @@ func NewFragment(path, index, frame, view string, slice uint64) *Fragment { frame: frame, view: view, slice: slice, - cacheType: DefaultCacheType, - cacheSize: DefaultCacheSize, + CacheType: DefaultCacheType, + CacheSize: DefaultCacheSize, LogOutput: ioutil.Discard, MaxOpN: DefaultFragmentMaxOpN, @@ -236,11 +236,11 @@ func (f *Fragment) openStorage() error { // openCache initializes the cache from row ids persisted to disk. func (f *Fragment) openCache() error { // Determine cache type from frame name. - switch f.cacheType { + switch f.CacheType { case CacheTypeRanked: - f.cache = NewRankCache(f.cacheSize) + f.cache = NewRankCache(f.CacheSize) case CacheTypeLRU: - f.cache = NewLRUCache(f.cacheSize) + f.cache = NewLRUCache(f.CacheSize) default: return ErrInvalidCacheType } diff --git a/fragment_test.go b/fragment_test.go index 5f3da53e2..f6b36f11e 100644 --- a/fragment_test.go +++ b/fragment_test.go @@ -37,7 +37,7 @@ const SliceWidth = pilosa.SliceWidth // Ensure a fragment can set a bit and retrieve it. func TestFragment_SetBit(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() // Set bits on the fragment. @@ -68,7 +68,7 @@ func TestFragment_SetBit(t *testing.T) { // Ensure a fragment can clear a set bit. func TestFragment_ClearBit(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() // Set and then clear bits on the fragment. @@ -95,7 +95,7 @@ func TestFragment_ClearBit(t *testing.T) { // Ensure a fragment can snapshot correctly. func TestFragment_Snapshot(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() // Set and then clear bits on the fragment. @@ -124,7 +124,7 @@ func TestFragment_Snapshot(t *testing.T) { // Ensure a fragment can iterate over all bits in order. func TestFragment_ForEachBit(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() // Set bits on the fragment. @@ -153,13 +153,13 @@ func TestFragment_ForEachBit(t *testing.T) { // Ensure a fragment can return the top n results. func TestFragment_Top(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeRanked) defer f.Close() - // Set bits on the rows 100, 101, & 102. f.MustSetBits(100, 1, 3, 200) f.MustSetBits(101, 1) f.MustSetBits(102, 1, 2) + f.RecalculateCache() // Retrieve top rows. if pairs, err := f.Top(pilosa.TopOptions{N: 2}); err != nil { @@ -175,14 +175,14 @@ func TestFragment_Top(t *testing.T) { // Ensure a fragment can filter rows when retrieving the top n rows. func TestFragment_Top_Filter(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeRanked) defer f.Close() // Set bits on the rows 100, 101, & 102. f.MustSetBits(100, 1, 3, 200) f.MustSetBits(101, 1) f.MustSetBits(102, 1, 2) - + f.RecalculateCache() // Assign attributes. f.RowAttrStore.SetAttrs(101, map[string]interface{}{"x": uint64(10)}) f.RowAttrStore.SetAttrs(102, map[string]interface{}{"x": uint64(20)}) @@ -205,7 +205,7 @@ func TestFragment_Top_Filter(t *testing.T) { // Ensure a fragment can return top rows that intersect with an input row. func TestFragment_TopN_Intersect(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeRanked) defer f.Close() // Create an intersecting input row. @@ -216,6 +216,7 @@ func TestFragment_TopN_Intersect(t *testing.T) { f.MustSetBits(101, 1, 2, 3, 4) // three intersections f.MustSetBits(102, 1, 2, 4, 5, 6) // two intersections f.MustSetBits(103, 1000, 1001, 1002) // no intersection + f.RecalculateCache() // Retrieve top rows. if pairs, err := f.Top(pilosa.TopOptions{N: 3, Src: src}); err != nil { @@ -235,7 +236,7 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) { t.Skip("short mode") } - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeRanked) defer f.Close() // Create an intersecting input row. @@ -250,6 +251,7 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) { f.MustSetBits(i, j) } } + f.RecalculateCache() // Retrieve top rows. if pairs, err := f.Top(pilosa.TopOptions{N: 10, Src: src}); err != nil { @@ -272,7 +274,7 @@ func TestFragment_TopN_Intersect_Large(t *testing.T) { // Ensure a fragment can return top rows when specified by ID. func TestFragment_TopN_IDs(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeRanked) defer f.Close() // Set bits on various rows. @@ -360,7 +362,7 @@ func TestFragment_TopN_CacheSize(t *testing.T) { // Ensure fragment can return a checksum for its blocks. func TestFragment_Checksum(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() // Retrieve checksum and set bits. @@ -379,7 +381,7 @@ func TestFragment_Checksum(t *testing.T) { // Ensure fragment can return a checksum for a given block. func TestFragment_Blocks(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() // Retrieve initial checksum. @@ -417,7 +419,7 @@ func TestFragment_Blocks(t *testing.T) { // Ensure fragment returns an empty checksum if no data exists for a block. func TestFragment_Blocks_Empty(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() // Set bits on a different block. @@ -435,7 +437,7 @@ func TestFragment_Blocks_Empty(t *testing.T) { // Ensure a fragment's cache can be persisted between restarts. func TestFragment_LRUCache_Persistence(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeLRU) defer f.Close() // Set bits on the fragment. @@ -520,7 +522,7 @@ func TestFragment_RankCache_Persistence(t *testing.T) { // Ensure a fragment can be copied to another fragment. func TestFragment_WriteTo_ReadFrom(t *testing.T) { - f0 := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f0 := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f0.Close() // Set and then clear bits on the fragment. @@ -545,7 +547,7 @@ func TestFragment_WriteTo_ReadFrom(t *testing.T) { } // Read into another fragment. - f1 := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f1 := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") if rn, err := f1.ReadFrom(&buf); err != nil { t.Fatal(err) } else if wn != rn { @@ -594,7 +596,7 @@ func BenchmarkFragment_Blocks(b *testing.B) { } func BenchmarkFragment_IntersectionCount(b *testing.B) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, "") defer f.Close() f.MaxOpN = math.MaxInt32 @@ -631,7 +633,7 @@ type Fragment struct { } // NewFragment returns a new instance of Fragment with a temporary path. -func NewFragment(index, frame, view string, slice uint64) *Fragment { +func NewFragment(index, frame, view string, slice uint64, cacheType string) *Fragment { file, err := ioutil.TempFile("", "pilosa-fragment-") if err != nil { panic(err) @@ -642,13 +644,18 @@ func NewFragment(index, frame, view string, slice uint64) *Fragment { Fragment: pilosa.NewFragment(file.Name(), index, frame, view, slice), RowAttrStore: MustOpenAttrStore(), } + f.Fragment.CacheType = cacheType f.Fragment.RowAttrStore = f.RowAttrStore.AttrStore return f } // MustOpenFragment creates and opens an fragment at a temporary path. Panic on error. -func MustOpenFragment(index, frame, view string, slice uint64) *Fragment { - f := NewFragment(index, frame, view, slice) +func MustOpenFragment(index, frame, view string, slice uint64, cacheType string) *Fragment { + if cacheType == "" { + cacheType = pilosa.DefaultCacheType + } + f := NewFragment(index, frame, view, slice, cacheType) + if err := f.Open(); err != nil { panic(err) } @@ -665,12 +672,14 @@ func (f *Fragment) Close() error { // Reopen closes the fragment and reopens it as a new instance. func (f *Fragment) Reopen() error { + cacheType := f.Fragment.CacheType path := f.Path() if err := f.Fragment.Close(); err != nil { return err } f.Fragment = pilosa.NewFragment(path, f.Index(), f.Frame(), f.View(), f.Slice()) + f.Fragment.CacheType = cacheType f.Fragment.RowAttrStore = f.RowAttrStore.AttrStore if err := f.Open(); err != nil { return err @@ -734,7 +743,7 @@ func GenerateImportFill(rowN int, pct float64) (rowIDs, columnIDs []uint64) { } func TestFragment_Tanimoto(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeRanked) defer f.Close() src := pilosa.NewBitmap(1, 2, 3) @@ -743,6 +752,7 @@ func TestFragment_Tanimoto(t *testing.T) { f.MustSetBits(100, 1, 3, 2, 200) f.MustSetBits(101, 1, 3) f.MustSetBits(102, 1, 2, 10, 12) + f.RecalculateCache() if pairs, err := f.Top(pilosa.TopOptions{TanimotoThreshold: 50, Src: src}); err != nil { t.Fatal(err) @@ -756,7 +766,7 @@ func TestFragment_Tanimoto(t *testing.T) { } func TestFragment_Zero_Tanimoto(t *testing.T) { - f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0) + f := MustOpenFragment("i", "f", pilosa.ViewStandard, 0, pilosa.CacheTypeRanked) defer f.Close() src := pilosa.NewBitmap(1, 2, 3) @@ -765,6 +775,7 @@ func TestFragment_Zero_Tanimoto(t *testing.T) { f.MustSetBits(100, 1, 3, 2, 200) f.MustSetBits(101, 1, 3) f.MustSetBits(102, 1, 2, 10, 12) + f.RecalculateCache() if pairs, err := f.Top(pilosa.TopOptions{TanimotoThreshold: 0, Src: src}); err != nil { t.Fatal(err) diff --git a/frame.go b/frame.go index 6afff11de..487caea1f 100644 --- a/frame.go +++ b/frame.go @@ -32,7 +32,7 @@ import ( // Default frame settings. const ( DefaultRowLabel = "rowID" - DefaultCacheType = CacheTypeLRU + DefaultCacheType = CacheTypeRanked DefaultInverseEnabled = false // Default ranked frame cache diff --git a/handler.go b/handler.go index 6cf175e05..69117a80e 100644 --- a/handler.go +++ b/handler.go @@ -171,7 +171,7 @@ func (h *Handler) handleWebUI(w http.ResponseWriter, r *http.Request) { statikFS, err := fs.New() if err != nil { h.writeQueryResponse(w, r, &QueryResponse{Err: err}) - fmt.Println("Pilosa WebUI is not available. Please run `make generate-statik` before building Pilosa with `make install`.") + h.logger().Println("Pilosa WebUI is not available. Please run `make generate-statik` before building Pilosa with `make install`.") return } http.FileServer(statikFS).ServeHTTP(w, r) @@ -261,7 +261,12 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) { // Set appropriate status code, if there is an error. if resp.Err != nil { - w.WriteHeader(http.StatusInternalServerError) + switch resp.Err { + case ErrTooManyWrites: + w.WriteHeader(http.StatusRequestEntityTooLarge) + default: + w.WriteHeader(http.StatusInternalServerError) + } } // Write response back to client. diff --git a/handler_test.go b/handler_test.go index 9e1895bed..740b03390 100644 --- a/handler_test.go +++ b/handler_test.go @@ -868,7 +868,7 @@ func TestHandler_Fragment_Nodes(t *testing.T) { h.ServeHTTP(w, r) if w.Code != http.StatusOK { t.Fatalf("unexpected status code: %d", w.Code) - } else if w.Body.String() != `[{"host":"host1","internalHost":""},{"host":"host2","internalHost":""}]`+"\n" { + } else if w.Body.String() != `[{"host":"host2","internalHost":""},{"host":"host0","internalHost":""}]`+"\n" { t.Fatalf("unexpected body: %q", w.Body.String()) } } diff --git a/holder_test.go b/holder_test.go index d7c90fccd..52288670a 100644 --- a/holder_test.go +++ b/holder_test.go @@ -246,3 +246,21 @@ func (h *Holder) MustCreateFragmentIfNotExists(index, frame, view string, slice } return &Fragment{Fragment: frag} } + +// MustCreateRankedFragmentIfNotExists returns a given fragment with a ranked cache. Panic on error. +func (h *Holder) MustCreateRankedFragmentIfNotExists(index, frame, view string, slice uint64) *Fragment { + idx := h.MustCreateIndexIfNotExists(index, pilosa.IndexOptions{}) + f, err := idx.CreateFrameIfNotExists(frame, pilosa.FrameOptions{CacheType: pilosa.CacheTypeRanked}) + if err != nil { + panic(err) + } + v, err := f.CreateViewIfNotExists(view) + if err != nil { + panic(err) + } + frag, err := v.CreateFragmentIfNotExists(slice) + if err != nil { + panic(err) + } + return &Fragment{Fragment: frag} +} diff --git a/httpbroadcast/messenger.go b/httpbroadcast/messenger.go index 09c3b7f4b..43a7e3ec7 100644 --- a/httpbroadcast/messenger.go +++ b/httpbroadcast/messenger.go @@ -115,7 +115,7 @@ func (h *HTTPBroadcaster) sendNodeMessage(node *pilosa.Node, msg []byte) error { // Check status code. if resp.StatusCode != http.StatusOK { - return fmt.Errorf("invalid status: code=%d, err=%s", resp.StatusCode, body) + return fmt.Errorf("invalid status sendNodeMessage: code=%d, err=%s, req=%v", resp.StatusCode, body, req) } return nil diff --git a/index.go b/index.go index e400b19ce..c596fbf81 100644 --- a/index.go +++ b/index.go @@ -380,6 +380,11 @@ func (i *Index) createFrame(name string, opt FrameOptions) (*Frame, error) { return nil, ErrInvalidCacheType } + // Validate that row label does not match column label. + if i.columnLabel == opt.RowLabel || (opt.RowLabel == "" && i.columnLabel == DefaultRowLabel) { + return nil, ErrColumnRowLabelEqual + } + // Initialize frame. f, err := i.newFrame(i.FramePath(name), name) if err != nil { diff --git a/index_test.go b/index_test.go index 7986855e6..bf59236c6 100644 --- a/index_test.go +++ b/index_test.go @@ -48,42 +48,70 @@ func TestIndex_CreateFrameIfNotExists(t *testing.T) { } } -// Ensure index is assigned the correct time quantum on creation. -func TestIndex_CreateFrame_TimeQuantum(t *testing.T) { - t.Run("Explicit", func(t *testing.T) { - index := MustOpenIndex() - defer index.Close() +func TestIndex_CreateFrame(t *testing.T) { + // Ensure time quantum can be set appropriately on a new frame. + t.Run("TimeQuantum", func(t *testing.T) { + t.Run("Explicit", func(t *testing.T) { + index := MustOpenIndex() + defer index.Close() - // Set index time quantum. - if err := index.SetTimeQuantum(pilosa.TimeQuantum("YM")); err != nil { - t.Fatal(err) - } + // Set index time quantum. + if err := index.SetTimeQuantum(pilosa.TimeQuantum("YM")); err != nil { + t.Fatal(err) + } - // Create frame with explicit quantum. - f, err := index.CreateFrame("f", pilosa.FrameOptions{TimeQuantum: pilosa.TimeQuantum("YMDH")}) - if err != nil { - t.Fatal(err) - } else if q := f.TimeQuantum(); q != pilosa.TimeQuantum("YMDH") { - t.Fatalf("unexpected frame time quantum: %s", q) - } + // Create frame with explicit quantum. + f, err := index.CreateFrame("f", pilosa.FrameOptions{TimeQuantum: pilosa.TimeQuantum("YMDH")}) + if err != nil { + t.Fatal(err) + } else if q := f.TimeQuantum(); q != pilosa.TimeQuantum("YMDH") { + t.Fatalf("unexpected frame time quantum: %s", q) + } + }) + + t.Run("Inherited", func(t *testing.T) { + index := MustOpenIndex() + defer index.Close() + + // Set index time quantum. + if err := index.SetTimeQuantum(pilosa.TimeQuantum("YM")); err != nil { + t.Fatal(err) + } + + // Create frame. + f, err := index.CreateFrame("f", pilosa.FrameOptions{}) + if err != nil { + t.Fatal(err) + } else if q := f.TimeQuantum(); q != pilosa.TimeQuantum("YM") { + t.Fatalf("unexpected frame time quantum: %s", q) + } + }) }) - t.Run("Inherited", func(t *testing.T) { - index := MustOpenIndex() - defer index.Close() + // Ensure frame cannot be created with a matching row label. + t.Run("ErrColumnRowLabelEqual", func(t *testing.T) { + t.Run("Explicit", func(t *testing.T) { + index := MustOpenIndex() + defer index.Close() - // Set index time quantum. - if err := index.SetTimeQuantum(pilosa.TimeQuantum("YM")); err != nil { - t.Fatal(err) - } + _, err := index.CreateFrame("f", pilosa.FrameOptions{RowLabel: pilosa.DefaultColumnLabel}) + if err != pilosa.ErrColumnRowLabelEqual { + t.Fatalf("unexpected error: %s", err) + } + }) - // Create frame. - f, err := index.CreateFrame("f", pilosa.FrameOptions{}) - if err != nil { - t.Fatal(err) - } else if q := f.TimeQuantum(); q != pilosa.TimeQuantum("YM") { - t.Fatalf("unexpected frame time quantum: %s", q) - } + t.Run("Default", func(t *testing.T) { + index := MustOpenIndex() + defer index.Close() + if err := index.SetColumnLabel(pilosa.DefaultRowLabel); err != nil { + t.Fatal(err) + } + + _, err := index.CreateFrame("f", pilosa.FrameOptions{}) + if err != pilosa.ErrColumnRowLabelEqual { + t.Fatalf("unexpected error: %s", err) + } + }) }) } diff --git a/internal/LICENSE b/internal/LICENSE new file mode 100644 index 000000000..d64569567 --- /dev/null +++ b/internal/LICENSE @@ -0,0 +1,202 @@ + + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + + END OF TERMS AND CONDITIONS + + APPENDIX: How to apply the Apache License to your work. + + To apply the Apache License to your work, attach the following + boilerplate notice, with the fields enclosed by brackets "[]" + replaced with your own identifying information. (Don't include + the brackets!) The text should be enclosed in the appropriate + comment syntax for the file format. We also recommend that a + file or class name and description of purpose be included on the + same "printed page" as the copyright notice for easier + identification within third-party archives. + + Copyright [yyyy] [name of copyright owner] + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. diff --git a/internal/private.pb.go b/internal/private.pb.go index c5d7be01e..eba920486 100644 --- a/internal/private.pb.go +++ b/internal/private.pb.go @@ -1,17 +1,3 @@ -// Copyright 2017 Pilosa Corp. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - // Code generated by protoc-gen-gogo. // source: private.proto // DO NOT EDIT! diff --git a/internal/public.pb.go b/internal/public.pb.go index 8499042fd..9eb005c17 100644 --- a/internal/public.pb.go +++ b/internal/public.pb.go @@ -1,17 +1,3 @@ -// Copyright 2017 Pilosa Corp. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. - // Code generated by protoc-gen-gogo. // source: public.proto // DO NOT EDIT! @@ -2590,7 +2576,7 @@ func init() { proto.RegisterFile("public.proto", fileDescriptorPublic) } var fileDescriptorPublic = []byte{ // 576 bytes of a gzipped FileDescriptorProto - 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x8c, 0x54, 0x4b, 0x8e, 0xd3, 0x40, + 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x09, 0x6e, 0x88, 0x02, 0xff, 0x8c, 0x54, 0x4b, 0x8e, 0xd3, 0x40, 0x10, 0xa5, 0x63, 0xe7, 0x57, 0xf9, 0x28, 0x6a, 0xf1, 0xb1, 0x10, 0x8a, 0x2c, 0x8b, 0x85, 0x57, 0x19, 0x69, 0x38, 0x00, 0xc2, 0x49, 0x46, 0xb2, 0x10, 0x23, 0xa6, 0x33, 0xb0, 0xf7, 0xcc, 0xb4, 0x06, 0x4b, 0xfe, 0xd1, 0xdd, 0x16, 0xe4, 0x00, 0xec, 0x91, 0xd8, 0x70, 0x03, 0x38, 0x0a, 0x4b, diff --git a/pilosa.go b/pilosa.go index dd7402677..a192364ed 100644 --- a/pilosa.go +++ b/pilosa.go @@ -34,6 +34,7 @@ var ( ErrFrameExists = errors.New("frame already exists") ErrFrameNotFound = errors.New("frame not found") ErrFrameInverseDisabled = errors.New("frame inverse disabled") + ErrColumnRowLabelEqual = errors.New("column and row labels cannot be equal") ErrInvalidView = errors.New("invalid view") ErrInvalidCacheType = errors.New("invalid cache type") @@ -44,6 +45,7 @@ var ( // ErrFragmentNotFound is returned when a fragment does not exist. ErrFragmentNotFound = errors.New("fragment not found") ErrQueryRequired = errors.New("query required") + ErrTooManyWrites = errors.New("too many write commands") ) // Regular expression to validate index and frame names. diff --git a/pql/ast.go b/pql/ast.go index b83b35731..19b5381ed 100644 --- a/pql/ast.go +++ b/pql/ast.go @@ -28,6 +28,18 @@ type Query struct { Calls []*Call } +// WriteCallN returns the number of mutating calls. +func (q *Query) WriteCallN() int { + var n int + for _, call := range q.Calls { + switch call.Name { + case "SetBit", "ClearBit", "SetRowAttrs", "SetColumnAttrs": + n++ + } + } + return n +} + // String returns a string representation of the query. func (q *Query) String() string { a := make([]string, len(q.Calls)) diff --git a/roaring/assembly_amd64.s b/roaring/assembly_amd64.s index b41091ad8..5b65e0635 100644 --- a/roaring/assembly_amd64.s +++ b/roaring/assembly_amd64.s @@ -1,6 +1,6 @@ #include "textflag.h" -TEXT ·hasAsm(SB),4,$0 +TEXT ·hasAsm(SB),4,$0-1 MOVQ $1, AX CPUID SHRQ $23, CX @@ -8,14 +8,14 @@ TEXT ·hasAsm(SB),4,$0 MOVB CX, ret+0(FP) RET -TEXT ·POPCNTQ(SB),NOSPLIT,$0-8 - MOVQ x+0(FP), BP +TEXT ·POPCNTQ(SB),NOSPLIT,$0-16 + MOVQ memory+0(FP), BP POPCNTQ BP, BX MOVQ BX, ret+8(FP) RET -TEXT ·BSFQ(SB),NOSPLIT,$0-8 - MOVQ x+0(FP), BP +TEXT ·BSFQ(SB),NOSPLIT,$0-16 + MOVQ memory+0(FP), BP BSFQ BP, BX MOVQ BX, ret+8(FP) RET @@ -24,8 +24,8 @@ TEXT ·BSFQ(SB),NOSPLIT,$0-8 TEXT ·popcntSliceAsm(SB),4,$0-32 XORQ AX, AX -MOVQ s+0(FP), SI -MOVQ s+8(FP), CX +MOVQ s_base+0(FP), SI +MOVQ s_len+8(FP), CX TESTQ CX, CX JZ popcntSliceEnd popcntSliceLoop: @@ -39,8 +39,8 @@ RET TEXT ·popcntMaskSliceAsm(SB),4,$0-56 XORQ AX, AX -MOVQ s+0(FP), SI -MOVQ s+8(FP), CX +MOVQ s_base+0(FP), SI +MOVQ s_len+8(FP), CX TESTQ CX, CX JZ popcntMaskSliceEnd MOVQ m+24(FP), DI @@ -59,8 +59,8 @@ RET TEXT ·popcntAndSliceAsm(SB),4,$0-56 XORQ AX, AX -MOVQ s+0(FP), SI -MOVQ s+8(FP), CX +MOVQ s_base+0(FP), SI +MOVQ s_len+8(FP), CX TESTQ CX, CX JZ popcntAndSliceEnd MOVQ m+24(FP), DI @@ -78,8 +78,8 @@ RET TEXT ·popcntOrSliceAsm(SB),4,$0-56 XORQ AX, AX -MOVQ s+0(FP), SI -MOVQ s+8(FP), CX +MOVQ s_base+0(FP), SI +MOVQ s_len+8(FP), CX TESTQ CX, CX JZ popcntOrSliceEnd MOVQ m+24(FP), DI @@ -97,8 +97,8 @@ RET TEXT ·popcntXorSliceAsm(SB),4,$0-56 XORQ AX, AX -MOVQ s+0(FP), SI -MOVQ s+8(FP), CX +MOVQ s_base+0(FP), SI +MOVQ s_len+8(FP), CX TESTQ CX, CX JZ popcntXorSliceEnd MOVQ m+24(FP), DI diff --git a/server.go b/server.go index 03ea0cbf8..8f28f849b 100644 --- a/server.go +++ b/server.go @@ -67,6 +67,9 @@ type Server struct { // Threshold for logging long queries LongQueryTime time.Duration + // Misc options. + MaxWritesPerRequest int + LogOutput io.Writer } @@ -136,6 +139,7 @@ func (s *Server) Open() error { e.Holder = s.Holder e.Host = s.Host e.Cluster = s.Cluster + e.MaxWritesPerRequest = s.MaxWritesPerRequest // Initialize HTTP handler. s.Handler.Broadcaster = s.Broadcaster @@ -430,7 +434,7 @@ func checkMaxSlices(hostport string) (map[string]uint64, error) { // Check status code. if resp.StatusCode != http.StatusOK { - return nil, fmt.Errorf("invalid status: code=%d, err=%s", resp.StatusCode, body) + return nil, fmt.Errorf("invalid status checkMaxSlices: code=%d, err=%s, req=%v", resp.StatusCode, body, req) } // Decode response object. diff --git a/server/server.go b/server/server.go index d99a05463..a3a6d060e 100644 --- a/server/server.go +++ b/server/server.go @@ -141,6 +141,9 @@ func (m *Command) SetupServer() error { return err } + // Copy configuration flags. + m.Server.MaxWritesPerRequest = m.Config.MaxWritesPerRequest + m.Server.Host, err = normalizeHost(m.Config.Host) if err != nil { return err diff --git a/server/server_test.go b/server/server_test.go index d981a4cf3..d677569e0 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -291,14 +291,14 @@ func TestMain_FrameRestore(t *testing.T) { // Create frames. client := m0.Client() - if err := client.CreateIndex(context.Background(), "x", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { + if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { t.Fatal(err) - } else if err := client.CreateFrame(context.Background(), "x", "f", pilosa.FrameOptions{}); err != nil { + } else if err := client.CreateFrame(context.Background(), "i", "f", pilosa.FrameOptions{}); err != nil { t.Fatal(err) } // Write data on first cluster. - if _, err := m0.Query("x", "", ` + if _, err := m0.Query("i", "", ` SetBit(rowID=1, frame="f", columnID=100) SetBit(rowID=1, frame="f", columnID=1000) SetBit(rowID=1, frame="f", columnID=100000) @@ -311,7 +311,7 @@ func TestMain_FrameRestore(t *testing.T) { } // Query row on first cluster. - if res, err := m0.Query("x", "", `Bitmap(rowID=1, frame="f")`); err != nil { + if res, err := m0.Query("i", "", `Bitmap(rowID=1, frame="f")`); err != nil { t.Fatal(err) } else if res != `{"results":[{"attrs":{},"bits":[100,1000,100000,200000,400000,600000,800000]}]}`+"\n" { t.Fatalf("unexpected result: %s", res) @@ -325,16 +325,16 @@ func TestMain_FrameRestore(t *testing.T) { client, err := pilosa.NewClient(m2.Server.Host) if err != nil { t.Fatal(err) - } else if err := m2.Client().CreateIndex(context.Background(), "x", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { + } else if err := m2.Client().CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists { t.Fatal(err) - } else if err := m2.Client().CreateFrame(context.Background(), "x", "f", pilosa.FrameOptions{}); err != nil { + } else if err := m2.Client().CreateFrame(context.Background(), "i", "f", pilosa.FrameOptions{}); err != nil { t.Fatal(err) - } else if err := client.RestoreFrame(context.Background(), m0.Server.Host, "x", "f"); err != nil { + } else if err := client.RestoreFrame(context.Background(), m0.Server.Host, "i", "f"); err != nil { t.Fatal(err) } // Query row on second cluster. - if res, err := m2.Query("x", "", `Bitmap(rowID=1, frame="f")`); err != nil { + if res, err := m2.Query("i", "", `Bitmap(rowID=1, frame="f")`); err != nil { t.Fatal(err) } else if res != `{"results":[{"attrs":{},"bits":[100,1000,100000,200000,400000,600000,800000]}]}`+"\n" { t.Fatalf("unexpected result: %s", res) diff --git a/view.go b/view.go index 3f26185e4..2e2ca32ab 100644 --- a/view.go +++ b/view.go @@ -255,8 +255,8 @@ func (v *View) createFragmentIfNotExists(slice uint64) (*Fragment, error) { func (v *View) newFragment(path string, slice uint64) *Fragment { frag := NewFragment(path, v.index, v.frame, v.name, slice) - frag.cacheType = v.cacheType - frag.cacheSize = v.cacheSize + frag.CacheType = v.cacheType + frag.CacheSize = v.cacheSize frag.LogOutput = v.LogOutput frag.stats = v.stats.WithTags(fmt.Sprintf("slice:%d", slice)) return frag diff --git a/webui/assets/main.js b/webui/assets/main.js index 4e17db19a..e18c562fb 100644 --- a/webui/assets/main.js +++ b/webui/assets/main.js @@ -116,42 +116,90 @@ class REPL { process_query(query) { var xhr = new XMLHttpRequest(); + var url, data, request, command_name; var e = document.getElementById("index-dropdown"); var indexname = e.options[e.selectedIndex].text; - xhr.open('POST', '/index/' + indexname + '/query'); + var repl = this; + if (query.startsWith(":")) { + var parsed_query = parse_query(query, indexname); + if (Object.keys(parsed_query).length === 0) { + repl.create_single_output({ + "input": query, + "output": "invalid query", + "status": 400, + "indexname": indexname, + }); + return; + } else { + // set selectedIndex from dropdown list + if (parsed_query.command === "use") { + for (var i = 0; i < e.options.length; i++) { + if (e.options[i].text === parsed_query.command_name) { + e.selectedIndex = i; + break; + } + } + return; + } + url = parsed_query.url; + data = parsed_query.data; + request = parsed_query.request; + command_name = parsed_query.command_name; + } + } + else { + url = '/index/' + indexname + '/query' + request = "POST" + data = query + } + xhr.open(request, url); xhr.setRequestHeader('Content-Type', 'application/text'); - var repl = this var start_time = new Date().getTime(); - xhr.send(query) - xhr.onload = function() { - var end_time = new Date().getTime() + xhr.onload = function () { + var end_time = new Date().getTime(); repl.result_number++ - repl.createSingleOutput({ - "input": query, - "output": xhr.responseText, - "indexname": indexname, - "querytime_ms": end_time - start_time, - }) - } + repl.create_single_output({ + "input": query, + "output": xhr.responseText, + "status": xhr.status, + "indexname": indexname, + "querytime_ms": end_time - start_time, + }); + }; + xhr.send(data); + // Remove index from dropdown with delete index command + if (request === 'DELETE' && url === '/index/' + command_name){ + for (var i = 0; i < e.options.length; i++) { + if (e.options[i].text === command_name) { + e.remove(i); + break; + } + } + } } - createSingleOutput(res) { - var node = document.createElement("div"); - node.classList.add('output'); - var output_string = res['output'] - var output_json = JSON.parse(output_string) - var result_class = "result-output" - var getting_started_errors = [ - 'index not found', - 'frame not found', - ] - - if("error" in output_json) { - result_class = "result-error" - if(getting_started_errors.indexOf(output_json['error']) >= 0) { - output_string += `
+ create_single_output(res) { + var node = document.createElement("div"); + node.classList.add('output'); + var output_string = res['output'] + var result_class = "result-output" + var getting_started_errors = [ + 'index not found', + 'frame not found', + ] + var output_json; + if (isJSON(output_string)) { + output_json = JSON.parse(output_string) + } + // handle output formatting + if (res["status"] != 200) { + result_class = "result-error"; + if (output_json) { + if ("error" in output_json) { + if (getting_started_errors.indexOf(output_json['error']) >= 0) { + output_string += `

Just getting started? Try this:
$ curl -XPOST "http://127.0.0.1:10101/index/test" -d '{"options": {"columnLabel": "col"}}' # create index "test"
@@ -159,8 +207,10 @@ class REPL { # Select "test" in the index dropdown above
SetBit(row=0, col=0, frame=foo) # Use PQL to set a bit ` + } + } + } } - } var markup =` @@ -184,7 +234,9 @@ class REPL {
${output_string} -
+ +
Expand
+
@@ -195,8 +247,22 @@ class REPL {
` - node.innerHTML = markup; - this.output.insertBefore(node, this.output.firstChild) + node.innerHTML = markup; + this.output.insertBefore(node, this.output.firstChild); + + // Expand when overflow + var element = this.output.firstChild.getElementsByClassName(result_class)[0]; + var expand = this.output.firstChild.getElementsByClassName("expand")[0]; + if (element.clientHeight < element.scrollHeight) { + expand.style.display = 'block'; + } else { + expand.style.display = 'none'; + } + expand.onclick = function () { + element.style.height = element.scrollHeight + "px"; + expand.style.display = 'none'; + return false; + }; } populate_index_dropdown() { @@ -367,7 +433,7 @@ function check_anchor_uri() { } } -Date.prototype.today = function () { +Date.prototype.today = function () { return this.getFullYear() +"/"+ (((this.getMonth()+1) < 10)?"0":"") + (this.getMonth()+1) +"/"+ ((this.getDate() < 10)?"0":"") + this.getDate(); } @@ -388,3 +454,88 @@ repl.bind_events() input.focus() check_anchor_uri() + +function isJSON(str) { + try { + JSON.parse(str) + } catch (e) { + return false + } + return true +} + +function parse_query(query, indexname) { + var keys = query.replace(/\s+/g, " ").split(" "); + var command = keys[0]; + var command_type = keys[1]; + var command_name = keys[2]; + var option_str = keys.slice(3, keys.length) + var options = parse_options(option_str); + if (command !== ":use") { + if (!command_name){ + return {} + } + } + + var parsed_query = {}; + parsed_query["command"] = command.substr(1, command.length); + parsed_query["command_name"] = command_name; + switch (command) { + case ":create": + parsed_query["request"] = "POST"; + if(Object.keys(options).length === 0) { + parsed_query["data"] = ""; + } else { + var opts = {"options":{}}; + for (var o in options) { + opts.options[o] = options[o] + } + parsed_query["data"] = JSON.stringify(opts); + } + switch (command_type){ + case "index": + parsed_query["url"] = '/index/' + command_name; + break; + case "frame": + parsed_query["url"] = '/index/' + indexname + '/frame/' + command_name; + break + } + break; + case ":delete": + parsed_query["request"] = "DELETE"; + switch (command_type){ + case "index": + parsed_query["url"] = '/index/' + command_name; + parsed_query["data"] = ""; + break; + case "frame": + parsed_query["url"] = '/index/' + indexname + '/frame/' + command_name; + parsed_query["data"] = ""; + break; + } + break; + case ":use": + parsed_query["command_name"] = keys[1]; + break; + default: + return {} + } + return parsed_query; +} + +function parse_options(option_str) { + var int_keys = ["cacheSize"]; + var bool_keys = ["inverseEnabled"]; + var options = {}; + for (var i = 0; i < option_str.length; i++) { + var parts = option_str[i].split('='); + if (int_keys.indexOf(parts[0]) !== -1 ){ + options[parts[0]] = Number(parts[1]) + } else if (bool_keys.indexOf(parts[0]) !== -1){ + options[parts[0]] = (parts[1] == "true") + } else { + options[parts[0]] = parts[1] + } + } + return options; +} \ No newline at end of file diff --git a/webui/assets/style.css b/webui/assets/style.css index a3a493e16..515335ef9 100644 --- a/webui/assets/style.css +++ b/webui/assets/style.css @@ -203,8 +203,6 @@ em{ display: block; } - - .result-io-header{ display: flex; align-items: center; @@ -214,6 +212,7 @@ em{ .result-input, .result-output, .result-error{ + height: 60px; border-radius: 2px; background-color: #fafafa; border: solid 1.5px #e4eff4; @@ -224,6 +223,8 @@ em{ color: #102445; padding: 15px; margin-bottom: 15px; + word-break: break-all; + overflow:hidden; } @@ -290,4 +291,8 @@ td{ .number { color: darkorange; } .boolean { color: blue; } .null { color: magenta; } -.key { color: red; } \ No newline at end of file +.key { color: red; } + +.expand { + text-align: center; +} \ No newline at end of file