From b7126a58598653c48a6aacb1a3890940637e3ee0 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Tue, 28 Sep 2021 13:16:24 -0600 Subject: [PATCH 1/5] Allow StmtRows.Scan() for more types --- planner.go | 68 ++++++++++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 66 insertions(+), 2 deletions(-) diff --git a/planner.go b/planner.go index 196bfa817..2a418b6c3 100644 --- a/planner.go +++ b/planner.go @@ -645,11 +645,12 @@ func (rs *StmtRows) Columns() []*StmtColumn { return rs.node.Columns() } - /* +/* func (rs *StmtRows) Row() int64 { return rs.node.Row()[0].(int64) } - */ +*/ + func (rs *StmtRows) Next() bool { if rs.err != nil { return false @@ -696,6 +697,15 @@ func (rs *StmtRows) Scan(dst ...interface{}) error { // Copy row value to scan destination. switch v := row[i].(type) { + case bool: + switch p := dst[i].(type) { + case *bool: + *p = v + case *interface{}: + *p = v + default: + return fmt.Errorf("cannot scan %T value into %T destination at index %d", v, p, i) + } case int64: switch p := dst[i].(type) { case *int: @@ -711,6 +721,48 @@ func (rs *StmtRows) Scan(dst ...interface{}) error { default: return fmt.Errorf("cannot scan %T value into %T destination at index %d", v, p, i) } + case uint64: + switch p := dst[i].(type) { + case *int: + *p = int(v) + case *int64: + *p = int64(v) + case *uint: + *p = uint(v) + case *uint64: + *p = uint64(v) + case *interface{}: + *p = v + default: + return fmt.Errorf("cannot scan %T value into %T destination at index %d", v, p, i) + } + case []uint64: + switch p := dst[i].(type) { + case *[]uint64: + *p = []uint64(v) + case *interface{}: + *p = joinUint64Slice(v) + default: + return fmt.Errorf("cannot scan %T value into %T destination at index %d", v, p, i) + } + case string: + switch p := dst[i].(type) { + case *string: + *p = v + case *interface{}: + *p = v + default: + return fmt.Errorf("cannot scan %T value into %T destination at index %d", v, p, i) + } + case []string: + switch p := dst[i].(type) { + case *[]string: + *p = []string(v) + case *interface{}: + *p = strings.Join(v, ",") + default: + return fmt.Errorf("cannot scan %T value into %T destination at index %d", v, p, i) + } default: return fmt.Errorf("unexpected %T value at index %d", v, i) } @@ -1113,3 +1165,15 @@ func stringSliceIndex(a []string, v string) int { } return -1 } + +func joinUint64Slice(a []uint64) string { + b := []byte("[") + for i, v := range a { + b = strconv.AppendUint(b, v, 10) + if i < len(a)-1 { + b = append(b, ',') + } + } + b = append(b, ']') + return string(b) +} From 3fe28ca7b9b36cd1c55370397979ebe6c049d6a6 Mon Sep 17 00:00:00 2001 From: Souhaila Noor Date: Tue, 28 Sep 2021 16:34:52 -0500 Subject: [PATCH 2/5] renamed webUI from Pilosa to FeatureBase and updated version --- Makefile | 6 +++++- http/handler.go | 2 +- 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/Makefile b/Makefile index 4a89aa66c..5f7de2e36 100644 --- a/Makefile +++ b/Makefile @@ -3,6 +3,10 @@ CLONE_URL=github.com/pilosa/pilosa MOD_VERSION=v2 VERSION := $(shell git describe --tags 2> /dev/null || echo unknown) +MAJOR_VERSION=$(shell cut -d '.' -f1 <<<$(VERSION) | cut -d 'v' -f2) +MAJOR_VERSION_INCREMENT=$(shell echo `expr $(MAJOR_VERSION) + 1`) +MINOR_VERSION=$(shell cut -d'.' -f2-3 <<<$(VERSION)) +RELEASE_VERSION=$(shell echo v$(MAJOR_VERSION_INCREMENT).$(MINOR_VERSION)) VARIANT = Molecula GO=go GOOS=$(shell $(GO) env GOOS) @@ -13,7 +17,7 @@ BRANCH_ID := $(BRANCH)-$(GOOS)-$(GOARCH) BUILD_TIME := $(shell date -u +%FT%T%z) SHARD_WIDTH = 20 COMMIT := $(shell git describe --exact-match >/dev/null 2>&1 || git rev-parse --short HEAD) -LDFLAGS="-X github.com/molecula/featurebase/v2.Version=$(VERSION) -X github.com/molecula/featurebase/v2.BuildTime=$(BUILD_TIME) -X github.com/molecula/featurebase/v2.Variant=$(VARIANT) -X github.com/molecula/featurebase/v2.Commit=$(COMMIT) -X github.com/molecula/featurebase/v2.TrialDeadline=$(TRIAL_DEADLINE)" +LDFLAGS="-X github.com/molecula/featurebase/v2.Version=$(RELEASE_VERSION) -X github.com/molecula/featurebase/v2.BuildTime=$(BUILD_TIME) -X github.com/molecula/featurebase/v2.Variant=$(VARIANT) -X github.com/molecula/featurebase/v2.Commit=$(COMMIT) -X github.com/molecula/featurebase/v2.TrialDeadline=$(TRIAL_DEADLINE)" GO_VERSION=1.16.3 DOCKER_BUILD= # set to 1 to use `docker-build` instead of `build` when creating a release BUILD_TAGS += shardwidth$(SHARD_WIDTH) diff --git a/http/handler.go b/http/handler.go index 05e97f583..59a3c3162 100644 --- a/http/handler.go +++ b/http/handler.go @@ -527,7 +527,7 @@ func newStatikHandler(h *Handler) statikHandler { func (s statikHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { if strings.HasPrefix(r.UserAgent(), "curl") { - msg := "Welcome. Pilosa v" + s.handler.api.Version() + " is running. Visit https://www.pilosa.com/docs/ for more information." + msg := "Welcome. FeatureBase v" + s.handler.api.Version() + " is running. Visit https://docs.molecula.cloud for more information." if s.statikFS != nil { msg += " Try the Web UI by visiting this URL in your browser." } From 49e8c1c69fbadec52561094f71996ed285768f1f Mon Sep 17 00:00:00 2001 From: Souhaila Noor Date: Tue, 28 Sep 2021 16:59:54 -0500 Subject: [PATCH 3/5] undid the changes pushed earlier for incrementing the release version --- Makefile | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/Makefile b/Makefile index 5f7de2e36..4a89aa66c 100644 --- a/Makefile +++ b/Makefile @@ -3,10 +3,6 @@ CLONE_URL=github.com/pilosa/pilosa MOD_VERSION=v2 VERSION := $(shell git describe --tags 2> /dev/null || echo unknown) -MAJOR_VERSION=$(shell cut -d '.' -f1 <<<$(VERSION) | cut -d 'v' -f2) -MAJOR_VERSION_INCREMENT=$(shell echo `expr $(MAJOR_VERSION) + 1`) -MINOR_VERSION=$(shell cut -d'.' -f2-3 <<<$(VERSION)) -RELEASE_VERSION=$(shell echo v$(MAJOR_VERSION_INCREMENT).$(MINOR_VERSION)) VARIANT = Molecula GO=go GOOS=$(shell $(GO) env GOOS) @@ -17,7 +13,7 @@ BRANCH_ID := $(BRANCH)-$(GOOS)-$(GOARCH) BUILD_TIME := $(shell date -u +%FT%T%z) SHARD_WIDTH = 20 COMMIT := $(shell git describe --exact-match >/dev/null 2>&1 || git rev-parse --short HEAD) -LDFLAGS="-X github.com/molecula/featurebase/v2.Version=$(RELEASE_VERSION) -X github.com/molecula/featurebase/v2.BuildTime=$(BUILD_TIME) -X github.com/molecula/featurebase/v2.Variant=$(VARIANT) -X github.com/molecula/featurebase/v2.Commit=$(COMMIT) -X github.com/molecula/featurebase/v2.TrialDeadline=$(TRIAL_DEADLINE)" +LDFLAGS="-X github.com/molecula/featurebase/v2.Version=$(VERSION) -X github.com/molecula/featurebase/v2.BuildTime=$(BUILD_TIME) -X github.com/molecula/featurebase/v2.Variant=$(VARIANT) -X github.com/molecula/featurebase/v2.Commit=$(COMMIT) -X github.com/molecula/featurebase/v2.TrialDeadline=$(TRIAL_DEADLINE)" GO_VERSION=1.16.3 DOCKER_BUILD= # set to 1 to use `docker-build` instead of `build` when creating a release BUILD_TAGS += shardwidth$(SHARD_WIDTH) From 3ef25e4a161579ae3ae1efa1628e2245412480ad Mon Sep 17 00:00:00 2001 From: Seebs Date: Fri, 24 Sep 2021 15:31:33 -0500 Subject: [PATCH 4/5] rework executor's per-shard union to use UnionInPlace The actual code here is mostly jaffee's, but I've reworked it some. This doesn't directly seem to be using UnionInPlace, but really it is. The actual logic inside (*Row).Union is a mess and probably silly in a few ways, but hardly matters. The important part is that, instead of calling it once per child as we get them, we gather all of them at once and then call it on all of them. That gets us a call to (*Row).Union that does a very elaborate dance to compute a call to (*rowSegment).Union on the only segment present in each of those rows, which then does a simpler thing to call (*Bitmap).Union() with the first response as a receiver and the rest as parameters, and THAT then ends up calling either unionIntoTargetSingle() if there's only one other bitmap, or using UnionInPlace on a Freeze() of the first bitmap, which gets us (we hope) the benefits of the fancy UnionInPlace logic. Every part of this is a reminder that we really need to replace roaring and also the Row/rowSegment stuff some day. --- executor.go | 22 +++++++++++----------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/executor.go b/executor.go index a46ca6d2a..ac9f4455e 100644 --- a/executor.go +++ b/executor.go @@ -4750,25 +4750,25 @@ func (e *executor) executeIntersectShard(ctx context.Context, qcx *Qcx, index st } // executeUnionShard executes a union() call for a local shard. -func (e *executor) executeUnionShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (_ *Row, err error) { +func (e *executor) executeUnionShard(ctx context.Context, qcx *Qcx, index string, c *pql.Call, shard uint64) (out *Row, err error) { span, ctx := tracing.StartSpanFromContext(ctx, "Executor.executeUnionShard") defer span.Finish() - other := NewRow() + if len(c.Children) == 0 { + return NewRow(), nil + } + if len(c.Children) == 1 { + return e.executeBitmapCallShard(ctx, qcx, index, c.Children[0], shard) + } + // we have at least two, so... + rows := make([]*Row, len(c.Children)) for i, input := range c.Children { - row, err := e.executeBitmapCallShard(ctx, qcx, index, input, shard) + rows[i], err = e.executeBitmapCallShard(ctx, qcx, index, input, shard) if err != nil { return nil, err } - - if i == 0 { - other = row - } else { - other = other.Union(row) - } } - other.invalidateCount() - return other, nil + return rows[0].Union(rows[1:]...), nil } // executeXorShard executes a xor() call for a local shard. From c24a5e77ba36c8129738b85ed91e6cb6ebd6d9a8 Mon Sep 17 00:00:00 2001 From: nagamocha3000 Date: Thu, 23 Sep 2021 17:43:16 +0300 Subject: [PATCH 5/5] Test paused node picks up once cluster state is back to normal This adds the following test: 1. cluster comes up (node 1,2,3), status normal 2. Pause node 3 3. Insert keys making sure to filter out the keys that will go to the paused node 4. Wait for status to become degraded 5. Unpause node 3 6. Wait for status to get back to normal 7. Check that keys were replicated to all 3 nodes --- Dockerfile-clustertests | 4 + Makefile | 8 - internal/clustertests/cluster_test.go | 2 +- .../docker-compose-index-key-replication.yml | 73 ---- internal/clustertests/docker-compose.yml | 7 + .../index_key_replication_test.go | 375 ---------------- internal/clustertests/pause_node_test.go | 412 ++++++++++++++++++ 7 files changed, 424 insertions(+), 457 deletions(-) delete mode 100644 internal/clustertests/docker-compose-index-key-replication.yml delete mode 100644 internal/clustertests/index_key_replication_test.go create mode 100644 internal/clustertests/pause_node_test.go diff --git a/Dockerfile-clustertests b/Dockerfile-clustertests index 91fde88fb..8c60f11bf 100644 --- a/Dockerfile-clustertests +++ b/Dockerfile-clustertests @@ -14,6 +14,10 @@ RUN cd /go/src/github.com/molecula/featurebase \ ADD https://github.com/alexei-led/pumba/releases/download/0.6.0/pumba_linux_amd64 /pumba RUN chmod +x /pumba +# add docker client to pause/unpause nodes +RUN apt update +RUN apt install -y docker.io + RUN cp /go/bin/featurebase /featurebase COPY LICENSE /LICENSE diff --git a/Makefile b/Makefile index 4a89aa66c..130cae3ce 100644 --- a/Makefile +++ b/Makefile @@ -138,14 +138,6 @@ clustertests: vendor docker-compose -f $(DOCKER_COMPOSE) build docker-compose -f $(DOCKER_COMPOSE) up --exit-code-from=client1 - -DOCKER_COMPOSE_INDEX_KEY_REPLICATION=internal/clustertests/docker-compose-index-key-replication.yml -# Check clustertests target for more info -clustertests-index-key-replication: vendor - docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) down - docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) build - docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) up --exit-code-from=client1 - # Like clustertests, but rebuilds all images. clustertests-build: vendor docker-compose -f $(DOCKER_COMPOSE) down -v diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index faad15337..3726bcc3b 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -115,7 +115,7 @@ func waitForStatus(t *testing.T, stator func(context.Context) (string, error), s if err != nil { t.Logf("Status (try %d/%d): %v (retrying in %s)", i, n, err, sleep.String()) } else { - t.Logf("Status (try %d/%d): %s (retrying in %s)", i, n, s, sleep.String()) + t.Logf("Status (try %d/%d): curr: %s, expect: %s (retrying in %s)", i, n, s, status, sleep.String()) } if s == status { return diff --git a/internal/clustertests/docker-compose-index-key-replication.yml b/internal/clustertests/docker-compose-index-key-replication.yml deleted file mode 100644 index 6f84d5947..000000000 --- a/internal/clustertests/docker-compose-index-key-replication.yml +++ /dev/null @@ -1,73 +0,0 @@ -version: '2' -services: - pilosa1: - build: - context: ../.. - dockerfile: Dockerfile-clustertests - image: ptest - ports: - - "33455:10101" - environment: - - PILOSA_NAME=pilosa1 - - PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201 - - PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa1:10201 - - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 - - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa1:10301 - - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 - - PILOSA_CLUSTER_REPLICAS=3 - networks: - - pilosanet - command: - - "/featurebase server --bind pilosa1:10101" - pilosa2: - build: - context: ../.. - dockerfile: Dockerfile-clustertests - image: ptest - ports: - - "33456:10101" - environment: - - PILOSA_NAME=pilosa2 - - PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201 - - PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa2:10201 - - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 - - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa2:10301 - - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 - - PILOSA_CLUSTER_REPLICAS=3 - networks: - - pilosanet - command: - - "/featurebase server --bind pilosa2:10101" - pilosa3: - build: - context: ../.. - dockerfile: Dockerfile-clustertests - image: ptest - ports: - - "33457:10101" - environment: - - PILOSA_NAME=pilosa3 - - PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201 - - PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa3:10201 - - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 - - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa3:10301 - - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 - - PILOSA_CLUSTER_REPLICAS=3 - networks: - - pilosanet - command: - - "/featurebase server --bind pilosa3:10101" - client1: - build: - context: . - environment: - - ENABLE_PILOSA_CLUSTER_TESTS_FOR_INDEX_KEY_REPLICATION=1 - - GO111MODULE=on - networks: - - pilosanet - volumes: - - /var/run/docker.sock:/var/run/docker.sock - command: - - "cd /go/src/github.com/molecula/featurebase/ && go test -mod=vendor -v -run=IndexKey -count=1 github.com/molecula/featurebase/v2/internal/clustertests" -networks: - pilosanet: diff --git a/internal/clustertests/docker-compose.yml b/internal/clustertests/docker-compose.yml index 0c1bf33c1..4192be6f8 100644 --- a/internal/clustertests/docker-compose.yml +++ b/internal/clustertests/docker-compose.yml @@ -14,6 +14,7 @@ services: - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa1:10301 - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 + - PILOSA_CLUSTER_REPLICAS=3 networks: - pilosanet command: @@ -32,6 +33,7 @@ services: - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa2:10301 - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 + - PILOSA_CLUSTER_REPLICAS=3 networks: - pilosanet command: @@ -50,6 +52,7 @@ services: - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa3:10301 - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 + - PILOSA_CLUSTER_REPLICAS=3 networks: - pilosanet command: @@ -57,6 +60,10 @@ services: client1: build: context: . + depends_on: + - "pilosa1" + - "pilosa2" + - "pilosa3" environment: - ENABLE_PILOSA_CLUSTER_TESTS=1 - GO111MODULE=on diff --git a/internal/clustertests/index_key_replication_test.go b/internal/clustertests/index_key_replication_test.go deleted file mode 100644 index 54c3dfa63..000000000 --- a/internal/clustertests/index_key_replication_test.go +++ /dev/null @@ -1,375 +0,0 @@ -// 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. - -package clustertest - -import ( - "context" - "crypto/tls" - "fmt" - "math/rand" - "os" - "os/exec" - "sync" - "testing" - "time" - - pilosa "github.com/molecula/featurebase/v2" - "github.com/molecula/featurebase/v2/disco" - "github.com/molecula/featurebase/v2/http" - picli "github.com/molecula/featurebase/v2/http" - "github.com/molecula/featurebase/v2/net" - "github.com/molecula/featurebase/v2/topology" -) - -// index -> key -> ids from all replicas -type translationRes map[string]map[string][]uint64 - -func defaultTranslationResults(indexes []string) translationRes { - res := make(translationRes) - for _, index := range indexes { - res[index] = make(map[string][]uint64) - } - return res -} - -func verify(allRes translationRes, indexes []string, replicasN, count int) error { - received := 0 - allIDsSame := func(ids []uint64) bool { - if len(ids) <= 1 { - // trivially true - return true - } - id := ids[0] - for _, other := range ids { - if id != other { - return false - } - } - return true - } - for _, index := range indexes { - res, ok := allRes[index] - if !ok { - return fmt.Errorf("Expected index '%s' but not present", index) - } - for key, ids := range res { - // first id is after successful translation, rest are from - // translate readers - replicationCount := len(ids) - 1 - received += replicationCount - if replicationCount != replicasN { - return fmt.Errorf("Count for ids for key '%s'(%d), index '%s' not equal than replicasN(%d)", key, replicationCount, index, replicasN) - } - if !allIDsSame(ids) { - return fmt.Errorf("Expected all ids for key '%s', index '%s' to be the same across the cluster %+v", key, index, ids) - } - } - } - - if received != count { - return fmt.Errorf("Expected %d count of keys, received %d", count, received) - } - - return nil -} -func getURIsFromAddresses(addrs []string) ([]*net.URI, error) { - uris := make([]*net.URI, 0, len(addrs)) - for _, addr := range addrs { - uri, err := net.NewURIFromAddress(addr) - if err != nil { - return nil, err - } - uris = append(uris, uri) - } - return uris, nil -} - -func getClients(addrs []string) ([]*http.InternalClient, error) { - clients := make([]*http.InternalClient, 0, len(addrs)) - for _, addr := range addrs { - c, err := picli.NewInternalClient(addr, picli.GetHTTPClient(nil)) - if err != nil { - return nil, err - } - clients = append(clients, c) - } - return clients, nil -} - -func genIndexNames(indexCount int) []string { - indexNames := make([]string, 0, indexCount) - for i := 1; i <= indexCount; i++ { - indexNames = append(indexNames, fmt.Sprintf("idx-%d", i)) - } - return indexNames -} - -func parseDuration(t *testing.T, s string) time.Duration { - parsed, err := time.ParseDuration(s) - if err != nil { - t.Fatal(err) - } - return parsed -} - -func durationToSeconds(d time.Duration) string { - s := int(d.Seconds()) - return fmt.Sprintf("%ds", s) -} - -func TestIndexKeyReplication(t *testing.T) { - if os.Getenv("ENABLE_PILOSA_CLUSTER_TESTS_FOR_INDEX_KEY_REPLICATION") != "1" { - t.Skip("pilosa cluster tests for index key replication are not enabled") - } - // configurations for test - replicasN := 3 - indexCount := 4 - intervalDurationArg := "100ms" - totalInsertionDurationArg := "10s" - coolOffDurationArg := "5s" - numKeysToInsertPerDuration := 100 - addresses := []string{"pilosa1:10101", "pilosa2:10101", "pilosa3:10101"} - - intervalDuration := parseDuration(t, intervalDurationArg) - totalInsertionDuration := parseDuration(t, totalInsertionDurationArg) - coolOffDuration := parseDuration(t, coolOffDurationArg) - - indexes := genIndexNames(indexCount) - clients, err := getClients(addresses) - cli := clients[0] - if err != nil { - t.Fatalf("on init clients from addresses: %v, %v", addresses, err) - } - uris, err := getURIsFromAddresses(addresses) - if err != nil { - t.Fatalf("on init clients from addresses: %v, %v", addresses, err) - } - ctx := context.Background() - - // create index keyed - for _, index := range indexes { - err = cli.EnsureIndex(ctx, index, pilosa.IndexOptions{ - Keys: true, - }) - if err != nil { - t.Fatalf("creating/asserting index: %v", err) - } - } - - // set up index keys to insert - keysInserted := 0 - allRes := defaultTranslationResults(indexes) - { - ctxForIndexKeyCreation, cancelFurtherInsertions := context.WithCancel(ctx) - var wg sync.WaitGroup - var mu = &sync.Mutex{} - wg.Add(len(indexes)) - seed := time.Now().UnixNano() - t.Logf("start inserting index keys for %v, seed(%v)", totalInsertionDuration, seed) - rng := rand.New(rand.NewSource(seed)) - for _, index := range indexes { - translations := allRes[index] - go func(index string, translations map[string][]uint64) { - keys := make([]string, numKeysToInsertPerDuration) - offset := 1 - ticker := time.NewTicker(intervalDuration) - defer func() { - ticker.Stop() - wg.Done() - }() - for { - select { - case <-ticker.C: - for i := 0; i < numKeysToInsertPerDuration; i++ { - keys[i] = fmt.Sprintf("key-%d", i+offset) - } - // pick random node to send key creation to - r := rng.Intn(len(addresses)) - cli, uri := clients[r], uris[r] - - // insert index keys - transmap, err := cli.CreateIndexKeysNode(ctxForIndexKeyCreation, uri, index, keys...) - if err != nil { - if err == context.Canceled { - return - } else { - t.Logf("creating index keys for index(%s) send to node(%s): %v", index, uri.String(), err) - continue - } - } - for key, id := range transmap { - translations[key] = append(translations[key], id) - } - mu.Lock() - keysInserted += len(transmap) - mu.Unlock() - offset += numKeysToInsertPerDuration - case <-ctxForIndexKeyCreation.Done(): - return - } - } - }(index, translations) - } - // inject fault - pcmd := exec.Command("/pumba", "netem", "--duration", - durationToSeconds(totalInsertionDuration), - "loss", "--percent", "50", "--correlation", "60", - "clustertests_pilosa3_1") - pcmd.Stdout = os.Stdout - pcmd.Stderr = os.Stderr - t.Logf("sending pumba fault injection cmd: %v", pcmd.String()) - err = pcmd.Start() - if err != nil { - t.Fatalf("starting pumba command: %v", err) - } - err = pcmd.Wait() - if err != nil { - t.Fatalf("waiting on pumba pause cmd: %v", err) - } - - // wait for index keys to be created - t.Logf("start wait to complete index key creation") - time.Sleep(totalInsertionDuration) - cancelFurtherInsertions() - wg.Wait() - t.Logf("done with inserting index keys. Total keys inserted: %d", keysInserted) - } - - t.Logf("start cool off period: %v\n", coolOffDuration) - time.Sleep(coolOffDuration) - t.Log("done with cool off period, waiting for stability") - waitForStatus(t, clients[0].Status, string(disco.ClusterStateNormal), 30, 1*time.Second) - t.Log("done with waiting for stability, starting verifying persistence of index keys") - - // get all nodes - nodes, err := cli.Nodes(ctx) - if err != nil { - t.Fatal(err) - } - - // prepare translate offset maps - nodeMaps := make(map[string]pilosa.TranslateOffsetMap) - for _, n := range nodes { - nodeMaps[n.ID] = make(pilosa.TranslateOffsetMap) - } - schema, err := cli.Schema(ctx) - if err != nil { - t.Fatal(err) - } - for _, indexInfo := range schema { - index := indexInfo.Name - isKeyed := indexInfo.Options.Keys - if !isKeyed { - continue - } - - partitionN := topology.DefaultPartitionN - for partition := 0; partition < partitionN; partition++ { - nodes, err := cli.PartitionNodes(ctx, partition) - if err != nil { - t.Fatal(err) - } - for _, n := range nodes { - m := nodeMaps[n.ID] - m.SetIndexPartitionOffset(index, partition, 0) - } - } - } - - // open translate reader - readers := make([]pilosa.TranslateEntryReader, 0, len(nodes)) - closeAllTranslateReaders := func() []error { - var errs []error - for _, tr := range readers { - err := tr.Close() - if err != nil { - errs = append(errs, err) - } - } - return errs - } - var tlsConfig *tls.Config = nil - var wg sync.WaitGroup - entries := make(chan *pilosa.TranslateEntry) - for _, n := range nodes { - client := http.GetHTTPClient(tlsConfig) - nodeURL := n.URI.String() - offsets := nodeMaps[n.ID] - openTranslateReader := http.GetOpenTranslateReaderFunc(client) - tr, err := openTranslateReader(ctx, nodeURL, offsets) - if err != nil { - t.Fatal(err) - } - wg.Add(1) - readers = append(readers, tr) - go func(node *topology.Node, tr pilosa.TranslateEntryReader) { - defer func() { - wg.Done() - }() - for { - var entry pilosa.TranslateEntry - err := tr.ReadEntry(&entry) - if err != nil { - // TODO also ignore transient http: read on clsoed response body errors - // for now just print error - if err != context.Canceled { - fmt.Printf("node(%s) On read from translate entry reader: %v", node.URI.String(), err) - } - return - } - // ignore field keys - if entry.Field != "" { - continue - } - entries <- &entry - } - }(n, tr) - } - - // receive all translation entries made - count := replicasN * keysInserted - i := 0 - for { - entry := <-entries - res, ok := allRes[entry.Index] - if !ok { - // ignore indexes we did not create for this test - continue - } - ids, ok := res[entry.Key] - if !ok { - // ignore keys we did not insert for the test - continue - } - res[entry.Key] = append(ids, entry.ID) - i++ - if i == count { - break - } - } - errs := closeAllTranslateReaders() - wg.Wait() - for _, err := range errs { - if err != nil { - t.Errorf("on close translate readers: %v", err) - } - } - close(entries) - - err = verify(allRes, indexes, replicasN, count) - if err != nil { - t.Fatal(err) - } -} diff --git a/internal/clustertests/pause_node_test.go b/internal/clustertests/pause_node_test.go new file mode 100644 index 000000000..e9c4e822a --- /dev/null +++ b/internal/clustertests/pause_node_test.go @@ -0,0 +1,412 @@ +// 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. + +package clustertest + +import ( + "context" + "fmt" + "hash/fnv" + "io/ioutil" + "math/rand" + "os" + "os/exec" + "path/filepath" + "strconv" + "testing" + "time" + + pilosa "github.com/molecula/featurebase/v2" + boltdb "github.com/molecula/featurebase/v2/boltdb" + "github.com/molecula/featurebase/v2/disco" + "github.com/molecula/featurebase/v2/http" + picli "github.com/molecula/featurebase/v2/http" + "github.com/molecula/featurebase/v2/net" + "github.com/molecula/featurebase/v2/topology" + "github.com/pkg/errors" +) + +func sendCmd(cmd string, args ...string) error { + pcmd := exec.Command(cmd, args...) + pcmd.Stdout = os.Stdout + pcmd.Stderr = os.Stderr + err := pcmd.Start() + if err != nil { + return errors.Wrap(err, "starting cmd") + } + err = pcmd.Wait() + if err != nil { + return errors.Wrap(err, "waiting on cmd") + } + return nil +} + +func unpauseNode(node string) error { + unpauseArgs := []string{"container", "unpause", "clustertests_" + node + "_1"} + return sendCmd("docker", unpauseArgs...) +} + +func pauseNode(node string) error { + pauseArgs := []string{"container", "pause", "clustertests_" + node + "_1"} + return sendCmd("docker", pauseArgs...) +} + +type keyInserter struct { + client *http.InternalClient + uri *net.URI + index string + keys []string +} + +func (ki keyInserter) insertKeys(ctx context.Context) (map[string]uint64, error) { + ts, err := ki.client.CreateIndexKeysNode(ctx, ki.uri, ki.index, ki.keys...) + return ts, err +} + +func getAddress(node string) string { + return node + ":10101" +} + +func getClients(addrs []string) ([]*http.InternalClient, error) { + clients := make([]*http.InternalClient, 0, len(addrs)) + for _, addr := range addrs { + c, err := picli.NewInternalClient(addr, picli.GetHTTPClient(nil)) + if err != nil { + return nil, err + } + clients = append(clients, c) + } + return clients, nil +} + +func getURIsFromAddresses(addrs []string) ([]*net.URI, error) { + uris := make([]*net.URI, 0, len(addrs)) + for _, addr := range addrs { + uri, err := net.NewURIFromAddress(addr) + if err != nil { + return nil, err + } + uris = append(uris, uri) + } + return uris, nil +} + +func readIndexTranslateData(ctx context.Context, client *picli.InternalClient, dirPath, index string, partition int) error { + // read translateStore contents from endpoint + r, err := client.IndexTranslateDataReader(ctx, index, partition) + if err != nil { + return err + } + buf, err := ioutil.ReadAll(r) + if err != nil { + return err + } + r.Close() + + // create file and write contents to it + filename := strconv.FormatInt(int64(partition), 10) + filePath := filepath.Join(dirPath, filename) + file, err := os.Create(filePath) + if err != nil { + return err + } + _, err = file.Write(buf) + if err != nil { + file.Close() + return err + } + err = file.Sync() + if err != nil { + file.Close() + return err + } + file.Close() + return nil +} + +func openTranslateStores(dirPath, index string) (map[int]pilosa.TranslateStore, error) { + dirEntries, err := ioutil.ReadDir(dirPath) + if err != nil { + return nil, err + } + + // in case of error, close any translateStore that has been opened + rollback := make([]pilosa.TranslateStore, 0, topology.DefaultPartitionN) + defer func() { + for _, ts := range rollback { + _ = ts.Close() + } + }() + + // filter out non-file entries + filePaths := make([]string, 0, len(dirEntries)) + for _, entry := range dirEntries { + if entry.Mode().IsDir() { + continue + } + filePath := filepath.Join(dirPath, entry.Name()) + filePaths = append(filePaths, filePath) + } + + translateStores := make(map[int]pilosa.TranslateStore) + for _, filePath := range filePaths { + // extract partition number from filename + file := filepath.Base(filePath) + partition, err := strconv.Atoi(file) + if err != nil { + return nil, err + } + // open bolt db + ts, err := boltdb.OpenTranslateStore(filePath, index, "", partition, topology.DefaultPartitionN) + ts.SetReadOnly(true) + if err != nil { + return nil, err + } + translateStores[partition] = ts + rollback = append(rollback, ts) + } + rollback = nil + + return translateStores, nil +} + +var errOpRetriable = errors.New("If operation failed on this error, it can be retried") + +func verifyNodeHasGivenKeys(ctx context.Context, node, index, dirPath string, keys []string) error { + // get client that's connected to node + address := getAddress(node) + client, err := picli.NewInternalClient(address, picli.GetHTTPClient(nil)) + if err != nil { + return err + } + + // create dir to store boltdbs for this node + nodeDirPath := filepath.Join(dirPath, node) + err = os.Mkdir(nodeDirPath, 0755) + if err != nil { + return err + } + + // read in all the translate stores for each partition + for partition := 0; partition < topology.DefaultPartitionN; partition++ { + err := readIndexTranslateData(ctx, client, nodeDirPath, index, partition) + if err != nil { + return err + } + } + + // open all the translate stores + translateStores, err := openTranslateStores(nodeDirPath, index) + if err != nil { + return err + } + + // close all the translate stores on complete + defer func() { + for _, ts := range translateStores { + ts.Close() + } + }() + + // merge all translations + merged := make(map[string]uint64) + for _, ts := range translateStores { + entries, err := ts.FindKeys(keys...) + if err != nil { + return err + } + for key, id := range entries { + merged[key] = id + } + } + + // check that all expected keys present in node + if len(merged) != len(keys) { + msg := fmt.Sprintf("entries in node %s: %d. keys inserted: %d", + node, len(merged), len(keys)) + return errors.Wrap(errOpRetriable, msg) + } + for _, k := range keys { + if _, ok := merged[k]; !ok { + msg := fmt.Sprintf("Key '%s' not present in node %s", k, node) + return errors.Wrap(errOpRetriable, msg) + } + } + + return nil +} + +func genKeys(count, maxTries int, keyToNode func(string) string, filterOut []string) []string { + exclusionSet := make(map[string]struct{}) + for _, n := range filterOut { + exclusionSet[n] = struct{}{} + } + keys := make([]string, 0, count) + i := 0 + for { + if i >= maxTries { + break + } + key := fmt.Sprintf("key-%d", i) + i++ + // get primary node for this key + node := keyToNode(key) + // check if we should exclude this key + if _, ok := exclusionSet[node]; ok { + continue + } + // add key + keys = append(keys, key) + if len(keys) >= count { + return keys + } + } + return keys +} + +func TestPauseReplica(t *testing.T) { + if os.Getenv("ENABLE_PILOSA_CLUSTER_TESTS") != "1" { + t.Skip("pilosa cluster tests for replication when a replica is paused are not enabled") + } + // configurations for test + nodeNames := []string{"pilosa1", "pilosa2", "pilosa3"} + nodeToPause := "pilosa3" + addresses := make([]string, len(nodeNames)) + for i, node := range nodeNames { + addresses[i] = getAddress(node) + } + clients, err := getClients(addresses) + if err != nil { + t.Fatalf("on init clients from addresses: %v, %v", addresses, err) + } + cli := clients[0] + uris, err := getURIsFromAddresses(addresses) + if err != nil { + t.Fatalf("on init clients from addresses: %v, %v", addresses, err) + } + uri := uris[0] + + ctx := context.Background() + ctx, cancel := context.WithCancel(ctx) + + t.Log("start Client") + + // first achieve normal cluster status + waitForStatus(t, cli.Status, string(disco.ClusterStateNormal), 30, 1*time.Second) + + // create keyed index + rng := rand.New(rand.NewSource(time.Now().UnixNano())) + index := fmt.Sprintf("keyed-index-%d", rng.Int63()) + err = cli.EnsureIndex(ctx, index, pilosa.IndexOptions{ + Keys: true, + }) + if err != nil { + t.Fatalf("creating/asserting index: %v", err) + } + + // generate mapping from partition to primary node + partitionToNode := make([]string, topology.DefaultPartitionN) + for partition := 0; partition < topology.DefaultPartitionN; partition++ { + nodes, err := cli.PartitionNodes(ctx, partition) + if err != nil { + t.Fatal(err) + } + partitionToNode[partition] = nodes[0].URI.Host + } + + // generate random keys to insert + // no guarantee that all keys expected to be generated don't fall + // in nodes to be filtered out + keyCount := 100 + maxKeyGenTries := 1000 + filterOutKeysFromTheseNodes := []string{nodeToPause} + keyToNode := func(key string) string { + // get partition for this key + h := fnv.New64a() + _, _ = h.Write([]byte(index)) + _, _ = h.Write([]byte(key)) + partition := int(h.Sum64() % uint64(topology.DefaultPartitionN)) + // get node for this partition + return partitionToNode[partition] + } + keys := genKeys(keyCount, maxKeyGenTries, keyToNode, filterOutKeysFromTheseNodes) + + // pause node + t.Logf("pause %s", nodeToPause) + err = pauseNode(nodeToPause) + if err != nil { + t.Fatalf("error on pause node %s: %v", nodeToPause, err) + } + + // insert keys + t.Log("start insert") + ts, err := keyInserter{ + client: cli, + uri: uri, + index: index, + keys: keys, + }.insertKeys(ctx) + if err != nil { + t.Fatalf("Error: inserting index keys for index(%s) send to node(%s): %v", index, uri.String(), err) + } + t.Logf("successfully end insert: %v", len(ts)) + + // wait for cluster status to be non-normal + waitForStatus(t, cli.Status, string(disco.ClusterStateDegraded), 30, 1*time.Second) + + // wait for cluster status to get back to normal + t.Logf("unpause %s", nodeToPause) + err = unpauseNode(nodeToPause) + if err != nil { + t.Fatalf("error on unpause node %s: %v", nodeToPause, err) + } + waitForStatus(t, cli.Status, string(disco.ClusterStateNormal), 30, 1*time.Second) + + // set up directory to store keys + basePath := "." + keysDirName := "index_keys" + dirPath, err := filepath.Abs(basePath) + if err != nil { + t.Fatal(err) + } + dirPath = filepath.Join(dirPath, keysDirName) + err = os.Mkdir(dirPath, 0755) + if err != nil { + t.Fatal(err) + } + + maxRetries := 10 + durationInBetweenRetries := 5 * time.Second + for _, node := range nodeNames { + try := 1 + retries: + for { + err = verifyNodeHasGivenKeys(ctx, node, index, dirPath, keys) + if err == nil { + break retries + } + if try <= maxRetries && errors.Is(err, errOpRetriable) { + try++ + t.Logf("node %s, retry verify key replication: (%d/%d) after %v\n", + node, try, maxRetries, durationInBetweenRetries) + time.Sleep(durationInBetweenRetries) + continue + } + t.Fatal(errors.Wrap(err, fmt.Sprintf("try: (%d/%d)", try, maxRetries))) + } + } + + cancel() + t.Log("Done") +}