mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
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
This commit is contained in:
parent
4648aa9477
commit
c24a5e77ba
7 changed files with 424 additions and 457 deletions
|
|
@ -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
|
||||
|
|
|
|||
8
Makefile
8
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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
}
|
||||
}
|
||||
412
internal/clustertests/pause_node_test.go
Normal file
412
internal/clustertests/pause_node_test.go
Normal file
|
|
@ -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")
|
||||
}
|
||||
Loading…
Add table
Reference in a new issue