add clustertests testing backup's retry

This commit is contained in:
Matthew Jaffee 2021-12-13 16:43:49 -06:00
parent 8486efaa79
commit cdf4bc4c88
5 changed files with 118 additions and 63 deletions

View file

@ -151,7 +151,7 @@ jobs:
- checkout-plus
- skip-if-root-unchanged
- setup_remote_docker
- run: make clustertests-build
- run: make clustertests
release:
executor:
name: golang

View file

@ -149,12 +149,10 @@ DOCKER_COMPOSE=internal/clustertests/docker-compose.yml
clustertests: vendor
docker-compose -f $(DOCKER_COMPOSE) down
docker-compose -f $(DOCKER_COMPOSE) build
docker-compose -f $(DOCKER_COMPOSE) up --exit-code-from=client1
docker-compose -f $(DOCKER_COMPOSE) up -d pilosa1 pilosa2 pilosa3
docker-compose -f $(DOCKER_COMPOSE) run client1
docker-compose -f $(DOCKER_COMPOSE) down
# Like clustertests, but rebuilds all images.
clustertests-build: vendor
docker-compose -f $(DOCKER_COMPOSE) down -v
docker-compose -f $(DOCKER_COMPOSE) up --exit-code-from=client1 --build
# Install Pilosa
install:

View file

@ -1751,26 +1751,22 @@ func (c *InternalClient) doWithRetry(req *http.Request) (*http.Response, error)
// start timer after first request, so if retryPeriod > 0 we
// pretty much always do at least one retry
start := time.Now()
for ; ; resp, err = c.httpClient.Do(req) {
if err != nil || resp.StatusCode < 200 || resp.StatusCode >= 300 {
if time.Since(start) > c.retryPeriod {
break
}
if err != nil {
c.log.Printf("retrying request due to error: '%v'", err)
} else {
if bod, readErr := ioutil.ReadAll(resp.Body); readErr != nil {
c.log.Printf("retrying request due to status: %d, error reading body: '%v', body: '%s'", resp.StatusCode, readErr, bod)
} else {
c.log.Printf("retrying request due to status: %d, body: '%s'", resp.StatusCode, bod)
}
}
time.Sleep(sleepDuration)
sleepDuration *= 2
} else {
for ; err != nil || resp.StatusCode < 200 || resp.StatusCode >= 300; resp, err = c.httpClient.Do(req) {
if time.Since(start) > c.retryPeriod {
break
}
if err != nil {
c.log.Printf("retrying request due to error: '%v'", err)
} else {
if bod, readErr := ioutil.ReadAll(resp.Body); readErr != nil {
c.log.Printf("retrying request due to status: %d, error reading body: '%v', body: '%s'", resp.StatusCode, readErr, bod)
} else {
c.log.Printf("retrying request due to status: %d, body: '%s'", resp.StatusCode, bod)
}
}
time.Sleep(sleepDuration)
sleepDuration *= 2
}
return resp, err
}

View file

@ -3,6 +3,8 @@ package clustertest
import (
"context"
"fmt"
"io/ioutil"
"os"
"os/exec"
"testing"
@ -30,56 +32,53 @@ func TestClusterStuff(t *testing.T) {
t.Fatalf("getting client: %v", err)
}
t.Run("long pause", func(t *testing.T) {
err := cli1.CreateIndex(context.Background(), "testidx", pilosa.IndexOptions{})
if err != nil {
t.Fatalf("creating index: %v", err)
}
err = cli1.CreateFieldWithOptions(context.Background(), "testidx", "testf", pilosa.FieldOptions{CacheType: pilosa.CacheTypeRanked, CacheSize: 100})
if err != nil {
t.Fatalf("creating field: %v", err)
}
if err := cli1.CreateIndex(context.Background(), "testidx", pilosa.IndexOptions{}); err != nil {
t.Fatalf("creating index: %v", err)
}
if err := cli1.CreateFieldWithOptions(context.Background(), "testidx", "testf", pilosa.FieldOptions{CacheType: pilosa.CacheTypeRanked, CacheSize: 100}); err != nil {
t.Fatalf("creating field: %v", err)
}
req := &pilosa.ImportRequest{
Index: "testidx",
Field: "testf",
}
req.ColumnIDs = make([]uint64, 10)
req.RowIDs = make([]uint64, 10)
req := &pilosa.ImportRequest{
Index: "testidx",
Field: "testf",
}
req.ColumnIDs = make([]uint64, 10)
req.RowIDs = make([]uint64, 10)
for i := 0; i < 1000; i++ {
req.RowIDs[i%10] = 0
req.ColumnIDs[i%10] = uint64((i/10)*pilosa.ShardWidth + i%10)
req.Shard = uint64(i / 10)
if i%10 == 9 {
err = cli1.Import(context.Background(), nil, req, &pilosa.ImportOptions{})
if err != nil {
t.Fatalf("importing: %v", err)
}
}
}
// Check query results from each node.
for i, cli := range []*picli.InternalClient{cli1, cli2, cli3} {
r, err := cli.Query(context.Background(), "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"})
for i := 0; i < 1000; i++ {
req.RowIDs[i%10] = 0
req.ColumnIDs[i%10] = uint64((i/10)*pilosa.ShardWidth + i%10)
req.Shard = uint64(i / 10)
if i%10 == 9 {
err = cli1.Import(context.Background(), nil, req, &pilosa.ImportOptions{})
if err != nil {
t.Fatalf("count querying pilosa%d: %v", i, err)
}
if r.Results[0].(uint64) != 1000 {
t.Fatalf("count on pilosa%d after import is %d", i, r.Results[0].(uint64))
t.Fatalf("importing: %v", err)
}
}
}
// Check query results from each node.
for i, cli := range []*picli.InternalClient{cli1, cli2, cli3} {
r, err := cli.Query(context.Background(), "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"})
if err != nil {
t.Fatalf("count querying pilosa%d: %v", i, err)
}
if r.Results[0].(uint64) != 1000 {
t.Fatalf("count on pilosa%d after import is %d", i, r.Results[0].(uint64))
}
}
t.Run("long pause", func(t *testing.T) {
pcmd := exec.Command("/pumba", "pause", "clustertests_pilosa3_1", "--duration", "10s")
pcmd.Stdout = os.Stdout
pcmd.Stderr = os.Stderr
t.Log("pausing pilosa3 for 10s")
err = pcmd.Start()
if err != nil {
if err := pcmd.Start(); err != nil {
t.Fatalf("starting pumba command: %v", err)
}
err = pcmd.Wait()
if err != nil {
if err := pcmd.Wait(); err != nil {
t.Fatalf("waiting on pumba pause cmd: %v", err)
}
@ -98,6 +97,63 @@ func TestClusterStuff(t *testing.T) {
}
}
})
t.Run("backup", func(t *testing.T) {
// do backup with node 1 down, but restart it after a few seconds
if err := sendCmd("docker", "stop", "clustertests_pilosa1_1"); err != nil {
t.Fatalf("sending stop command: %v", err)
}
var backupCmd *exec.Cmd
tmpdir, err := ioutil.TempDir("", "")
if err != nil {
t.Fatalf("getting tmp dir: %v", err)
}
if backupCmd, err = startCmd(
"featurebase", "backup", "--host=pilosa1:10101", fmt.Sprintf("--output=%s", tmpdir+"/backuptest")); err != nil {
t.Fatalf("sending backup command: %v", err)
}
time.Sleep(time.Second * 5)
if err = sendCmd("docker", "start", "clustertests_pilosa1_1"); err != nil {
t.Fatalf("sending start command: %v", err)
}
if err = backupCmd.Wait(); err != nil {
t.Fatalf("waiting on backup to finish: %v", err)
}
// now do backup with all nodes down and too short a timeout
// so it fails. Has be to be all 3 because the cluster has
// replicas=3 and the backup command will retry on replicas.
if err = sendCmd("docker", "stop", "clustertests_pilosa2_1"); err != nil {
t.Fatalf("sending stop command: %v", err)
}
if backupCmd, err = startCmd(
"featurebase", "backup", "--host=pilosa1:10101", fmt.Sprintf("--output=%s", tmpdir+"/backuptest2"), "--retry-period=0.5s"); err != nil {
t.Fatalf("sending second backup command: %v", err)
}
time.Sleep(time.Millisecond * 5) // want the backup to get started, then fail
if err = sendCmd("docker", "stop", "clustertests_pilosa1_1"); err != nil {
t.Fatalf("sending stop command: %v", err)
}
if err = sendCmd("docker", "stop", "clustertests_pilosa3_1"); err != nil {
t.Fatalf("sending stop command: %v", err)
}
time.Sleep(time.Second * 5)
if err = sendCmd("docker", "start", "clustertests_pilosa1_1"); err != nil {
t.Fatalf("sending start command: %v", err)
}
if err = sendCmd("docker", "start", "clustertests_pilosa2_1"); err != nil {
t.Fatalf("sending start command: %v", err)
}
if err = sendCmd("docker", "start", "clustertests_pilosa3_1"); err != nil {
t.Fatalf("sending start command: %v", err)
}
if err = backupCmd.Wait(); err == nil {
t.Fatal("backup command should have errored but didn't")
}
})
}
func waitForStatus(t *testing.T, stator func(context.Context) (string, error), status string, n int, sleep time.Duration) {

View file

@ -23,11 +23,16 @@ import (
"github.com/pkg/errors"
)
func sendCmd(cmd string, args ...string) error {
func startCmd(cmd string, args ...string) (*exec.Cmd, error) {
pcmd := exec.Command(cmd, args...)
pcmd.Stdout = os.Stdout
pcmd.Stderr = os.Stderr
err := pcmd.Start()
return pcmd, err
}
func sendCmd(cmd string, args ...string) error {
pcmd, err := startCmd(cmd, args...)
if err != nil {
return errors.Wrap(err, "starting cmd")
}