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") +}