From 1a87e7c9ddd12eeca855579a03cb386e8d33efe5 Mon Sep 17 00:00:00 2001 From: rachithrr Date: Fri, 11 Feb 2022 15:12:52 -0600 Subject: [PATCH 1/5] FB-1159: Try to test etcd retry code -This code should cover the retry logic. -The ingest is set to 100,000 records, because it's difficult to cause leader change with lower number of records to ingest. -Additional node: By applying stress on other two nodes while ingesting simultaniously causes "context erorr" logic which can fail many tests. --- internal/clustertests/cluster_test.go | 82 +++++++++++++++++++++++++++ 1 file changed, 82 insertions(+) diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index 490121d0d..f61124e69 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -22,6 +22,7 @@ import ( "github.com/molecula/featurebase/v3/encoding/proto" "github.com/molecula/featurebase/v3/logger" "github.com/pkg/errors" + "golang.org/x/sync/errgroup" ) // container turns a docker-compose service name into a container ID @@ -288,6 +289,87 @@ func TestClusterStuff(t *testing.T) { }) } +func ingestRandomData(ctx context.Context, cli *pilosa.InternalClient, index, field string) error { + if err := cli.CreateIndex(ctx, index, pilosa.IndexOptions{}); err != nil { + return fmt.Errorf("creating index: %v", err) + } + if err := cli.CreateFieldWithOptions(ctx, index, field, pilosa.FieldOptions{CacheType: pilosa.CacheTypeRanked, CacheSize: 100}); err != nil { + return fmt.Errorf("creating field: %v", err) + } + + req := &pilosa.ImportRequest{ + Index: index, + Field: field, + } + req.ColumnIDs = make([]uint64, 10) + req.RowIDs = make([]uint64, 10) + + for i := 0; i < 100000; 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 := cli.Import(ctx, nil, req, &pilosa.ImportOptions{}) + if err != nil { + return fmt.Errorf("import error: %v", err) + } + } + } + return nil +} + +func TestRetryLogic(t *testing.T) { + if os.Getenv("ENABLE_PILOSA_CLUSTER_TESTS") != "1" { + t.Skip("pilosa cluster tests are not enabled") + } + ctx := context.Background() + auth := false + if os.Getenv("ENABLE_AUTH") == "1" { + auth = true + } + if auth { + token := GetAuthToken(t) + ctx = context.WithValue(ctx, "token", "Bearer "+token) + } + + cli1, err := pilosa.NewInternalClient("pilosa1:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) + if err != nil { + t.Fatalf("getting client: %v", err) + } + + g := new(errgroup.Group) + runCmd("pumba", "stress", "-d 10s", container(t, "pilosa2")) + g.Go(func() error { + return ingestRandomData(ctx, cli1, "testidx1", "testfield1") + }) + if err = sendCmd("docker", "pause", container(t, "pilosa3")); err != nil { + t.Fatalf("sending docker pause") + } + runCmd("pumba", "stress", "-d 10s", container(t, "pilosa2")) + if err = sendCmd("docker", "pause", container(t, "pilosa1")); err != nil { + t.Fatalf("sending docker pause") + } + time.Sleep(10 * time.Second) + if err = sendCmd("docker", "unpause", container(t, "pilosa3")); err != nil { + t.Fatalf("sending docker unpause") + } + if err = sendCmd("docker", "pause", container(t, "pilosa2")); err != nil { + t.Fatalf("sending docker pause") + } + if err = sendCmd("docker", "unpause", container(t, "pilosa2")); err != nil { + t.Fatalf("sending docker unpause") + } + if err = sendCmd("docker", "unpause", container(t, "pilosa1")); err != nil { + t.Fatalf("sending docker unpause") + } + time.Sleep(10 * time.Second) + if err = g.Wait(); err != nil { + t.Fatal(err) + } + waitForStatus(t, cli1.Status, string(disco.ClusterStateNormal), 30, time.Second, ctx) + +} + func waitForStatus(t *testing.T, stator func(context.Context) (string, error), status string, n int, sleep time.Duration, ctx context.Context) { t.Helper() From ce776e97fc4364baf550bcbb1e2e09699e232d6d Mon Sep 17 00:00:00 2001 From: rachithrr Date: Thu, 10 Mar 2022 09:43:52 -0600 Subject: [PATCH 2/5] adding more chaos --- internal/clustertests/cluster_test.go | 74 +++++++++++++++++++++------ 1 file changed, 59 insertions(+), 15 deletions(-) diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index f61124e69..d50ad73d5 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -289,7 +289,7 @@ func TestClusterStuff(t *testing.T) { }) } -func ingestRandomData(ctx context.Context, cli *pilosa.InternalClient, index, field string) error { +func ingestRandomData(ctx context.Context, cli *pilosa.InternalClient, index, field string, size int) error { if err := cli.CreateIndex(ctx, index, pilosa.IndexOptions{}); err != nil { return fmt.Errorf("creating index: %v", err) } @@ -304,7 +304,7 @@ func ingestRandomData(ctx context.Context, cli *pilosa.InternalClient, index, fi req.ColumnIDs = make([]uint64, 10) req.RowIDs = make([]uint64, 10) - for i := 0; i < 100000; i++ { + for i := 0; i < size; i++ { req.RowIDs[i%10] = 0 req.ColumnIDs[i%10] = uint64((i/10)*pilosa.ShardWidth + i%10) req.Shard = uint64(i / 10) @@ -336,38 +336,82 @@ func TestRetryLogic(t *testing.T) { if err != nil { t.Fatalf("getting client: %v", err) } + cli2, err := pilosa.NewInternalClient("pilosa2:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) + if err != nil { + t.Fatalf("getting client: %v", err) + } + cli3, err := pilosa.NewInternalClient("pilosa3:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) + if err != nil { + t.Fatalf("getting client: %v", err) + } g := new(errgroup.Group) - runCmd("pumba", "stress", "-d 10s", container(t, "pilosa2")) g.Go(func() error { - return ingestRandomData(ctx, cli1, "testidx1", "testfield1") + return ingestRandomData(ctx, cli1, "testidx1", "testfield1", 100000) + }) + if err = sendCmd("docker", "pause", container(t, "pilosa2")); err != nil { + t.Fatalf("sending docker pause %v", err) + } + if err = sendCmd("docker", "pause", container(t, "pilosa3")); err != nil { + t.Fatalf("sending docker pause %v", err) + } + time.Sleep(6 * time.Second) + if err = sendCmd("docker", "unpause", container(t, "pilosa2")); err != nil { + t.Fatalf("sending docker pause %v", err) + } + if err = sendCmd("docker", "unpause", container(t, "pilosa3")); err != nil { + t.Fatalf("sending docker pause %v", err) + } + time.Sleep(10 * time.Second) + g.Go(func() error { + return ingestRandomData(ctx, cli2, "testidx2", "testfield2", 10000) }) if err = sendCmd("docker", "pause", container(t, "pilosa3")); err != nil { - t.Fatalf("sending docker pause") + t.Fatalf("sending docker pause %v", err) } - runCmd("pumba", "stress", "-d 10s", container(t, "pilosa2")) if err = sendCmd("docker", "pause", container(t, "pilosa1")); err != nil { - t.Fatalf("sending docker pause") + t.Fatalf("sending docker pause %v", err) } - time.Sleep(10 * time.Second) + time.Sleep(6 * time.Second) if err = sendCmd("docker", "unpause", container(t, "pilosa3")); err != nil { - t.Fatalf("sending docker unpause") + t.Fatalf("sending docker pause %v", err) } + time.Sleep(6 * time.Second) if err = sendCmd("docker", "pause", container(t, "pilosa2")); err != nil { - t.Fatalf("sending docker pause") + t.Fatalf("sending docker pause %v", err) + } + time.Sleep(6 * time.Second) + if err = sendCmd("docker", "unpause", container(t, "pilosa1")); err != nil { + t.Fatalf("sending docker pause %v", err) } if err = sendCmd("docker", "unpause", container(t, "pilosa2")); err != nil { - t.Fatalf("sending docker unpause") + t.Fatalf("sending docker pause %v", err) } - if err = sendCmd("docker", "unpause", container(t, "pilosa1")); err != nil { - t.Fatalf("sending docker unpause") - } - time.Sleep(10 * time.Second) if err = g.Wait(); err != nil { t.Fatal(err) } waitForStatus(t, cli1.Status, string(disco.ClusterStateNormal), 30, time.Second, ctx) + for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { + r, err := cli.Query(ctx, "testidx1", &pilosa.QueryRequest{Index: "testidx1", Query: "Count(Row(testfield1 = 0))"}) + if err != nil { + t.Fatalf("count querying pilosa1 %v", err) + } + fmt.Println("count = ", r.Results[0].(uint64)) + if r.Results[0].(uint64) != 100000 { + t.Fatalf("count on pilosa%d after import is %d", i, r.Results[0].(uint64)) + } + } + for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { + r, err := cli.Query(ctx, "testidx2", &pilosa.QueryRequest{Index: "testidx2", Query: "Count(Row(testfield2 = 0))"}) + if err != nil { + t.Fatalf("count querying pilosa1 %v", err) + } + fmt.Println("count = ", r.Results[0].(uint64)) + if r.Results[0].(uint64) != 10000 { + t.Fatalf("count on pilosa%d after import is %d", i, r.Results[0].(uint64)) + } + } } func waitForStatus(t *testing.T, stator func(context.Context) (string, error), status string, n int, sleep time.Duration, ctx context.Context) { From bed631212961475a2227a3c5f77460885ae99e6e Mon Sep 17 00:00:00 2001 From: rachithrr Date: Thu, 10 Mar 2022 11:40:38 -0600 Subject: [PATCH 3/5] clean up --- internal/clustertests/cluster_test.go | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index d50ad73d5..6b836bad7 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -392,12 +392,12 @@ func TestRetryLogic(t *testing.T) { } waitForStatus(t, cli1.Status, string(disco.ClusterStateNormal), 30, time.Second, ctx) + // check data in all three nodes. for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { r, err := cli.Query(ctx, "testidx1", &pilosa.QueryRequest{Index: "testidx1", Query: "Count(Row(testfield1 = 0))"}) if err != nil { - t.Fatalf("count querying pilosa1 %v", err) + t.Fatalf("count querying pilosa%d, %v", i, err) } - fmt.Println("count = ", r.Results[0].(uint64)) if r.Results[0].(uint64) != 100000 { t.Fatalf("count on pilosa%d after import is %d", i, r.Results[0].(uint64)) } @@ -405,9 +405,8 @@ func TestRetryLogic(t *testing.T) { for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { r, err := cli.Query(ctx, "testidx2", &pilosa.QueryRequest{Index: "testidx2", Query: "Count(Row(testfield2 = 0))"}) if err != nil { - t.Fatalf("count querying pilosa1 %v", err) + t.Fatalf("count querying pilosa%d, %v", i, err) } - fmt.Println("count = ", r.Results[0].(uint64)) if r.Results[0].(uint64) != 10000 { t.Fatalf("count on pilosa%d after import is %d", i, r.Results[0].(uint64)) } From e2ce9afce95efa7e84a99e976c0dc6aa6e8a7939 Mon Sep 17 00:00:00 2001 From: rachithrr Date: Thu, 10 Mar 2022 16:25:59 -0600 Subject: [PATCH 4/5] sonarcloud fix --- internal/clustertests/cluster_test.go | 94 +++++++++++---------------- 1 file changed, 38 insertions(+), 56 deletions(-) diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index 6b836bad7..7d16e1c32 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -19,7 +19,6 @@ import ( "github.com/molecula/featurebase/v3/authn" "github.com/molecula/featurebase/v3/ctl" "github.com/molecula/featurebase/v3/disco" - "github.com/molecula/featurebase/v3/encoding/proto" "github.com/molecula/featurebase/v3/logger" "github.com/pkg/errors" "golang.org/x/sync/errgroup" @@ -92,15 +91,8 @@ func TestClusterStuff(t *testing.T) { auth = true } - cli1, err := pilosa.NewInternalClient("pilosa1:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) - if err != nil { - t.Fatalf("getting client: %v", err) - } - cli2, err := pilosa.NewInternalClient("pilosa2:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) - if err != nil { - t.Fatalf("getting client: %v", err) - } - cli3, err := pilosa.NewInternalClient("pilosa3:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) + var addrs = []string{"pilosa1:10101", "pilosa2:10101", "pilosa3:10101"} + cli, err := getClients(addrs) if err != nil { t.Fatalf("getting client: %v", err) } @@ -112,10 +104,10 @@ func TestClusterStuff(t *testing.T) { ctx = context.WithValue(ctx, "token", "Bearer "+token) } - if err := cli1.CreateIndex(ctx, "testidx", pilosa.IndexOptions{}); err != nil { + if err := cli[0].CreateIndex(ctx, "testidx", pilosa.IndexOptions{}); err != nil { t.Fatalf("creating index: %v", err) } - if err := cli1.CreateFieldWithOptions(ctx, "testidx", "testf", pilosa.FieldOptions{CacheType: pilosa.CacheTypeRanked, CacheSize: 100}); err != nil { + if err := cli[0].CreateFieldWithOptions(ctx, "testidx", "testf", pilosa.FieldOptions{CacheType: pilosa.CacheTypeRanked, CacheSize: 100}); err != nil { t.Fatalf("creating field: %v", err) } @@ -131,7 +123,7 @@ func TestClusterStuff(t *testing.T) { req.ColumnIDs[i%10] = uint64((i/10)*pilosa.ShardWidth + i%10) req.Shard = uint64(i / 10) if i%10 == 9 { - err = cli1.Import(ctx, nil, req, &pilosa.ImportOptions{}) + err = cli[0].Import(ctx, nil, req, &pilosa.ImportOptions{}) if err != nil { t.Fatalf("importing: %v", err) } @@ -139,8 +131,8 @@ func TestClusterStuff(t *testing.T) { } // Check query results from each node. - for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { - r, err := cli.Query(ctx, "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"}) + for i, c := range cli { + r, err := c.Query(ctx, "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"}) if err != nil { t.Fatalf("count querying pilosa%d: %v", i, err) } @@ -158,12 +150,12 @@ func TestClusterStuff(t *testing.T) { t.Fatalf("sending unpause: %v", err) } t.Log("done with pause, waiting for stability") - waitForStatus(t, cli1.Status, string(disco.ClusterStateNormal), 30, time.Second, ctx) + waitForStatus(t, cli[0].Status, string(disco.ClusterStateNormal), 30, time.Second, ctx) t.Log("done waiting for stability") // Check query results from each node. - for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { - r, err := cli.Query(ctx, "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"}) + for i, c := range cli { + r, err := c.Query(ctx, "testidx", &pilosa.QueryRequest{Index: "testidx", Query: "Count(Row(testf=0))"}) if err != nil { t.Fatalf("count querying pilosa%d: %v", i, err) } @@ -332,78 +324,68 @@ func TestRetryLogic(t *testing.T) { ctx = context.WithValue(ctx, "token", "Bearer "+token) } - cli1, err := pilosa.NewInternalClient("pilosa1:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) + var addrs = []string{"pilosa1:10101", "pilosa2:10101", "pilosa3:10101"} + cli, err := getClients(addrs) if err != nil { t.Fatalf("getting client: %v", err) } - cli2, err := pilosa.NewInternalClient("pilosa2:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) - if err != nil { - t.Fatalf("getting client: %v", err) - } - cli3, err := pilosa.NewInternalClient("pilosa3:10101", pilosa.GetHTTPClient(nil), pilosa.WithSerializer(proto.Serializer{})) - if err != nil { - t.Fatalf("getting client: %v", err) - } - g := new(errgroup.Group) g.Go(func() error { - return ingestRandomData(ctx, cli1, "testidx1", "testfield1", 100000) + return ingestRandomData(ctx, cli[0], "testidx1", "testfield1", 100000) }) - if err = sendCmd("docker", "pause", container(t, "pilosa2")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := pauseNode(t, "pilosa2"); err != nil { + t.Fatalf("sending pause command: %v", err) } - if err = sendCmd("docker", "pause", container(t, "pilosa3")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := pauseNode(t, "pilosa3"); err != nil { + t.Fatalf("sending pause command: %v", err) } time.Sleep(6 * time.Second) - if err = sendCmd("docker", "unpause", container(t, "pilosa2")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := unpauseNode(t, "pilosa2"); err != nil { + t.Fatalf("sending pause command: %v", err) } - if err = sendCmd("docker", "unpause", container(t, "pilosa3")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := unpauseNode(t, "pilosa3"); err != nil { + t.Fatalf("sending pause command: %v", err) } time.Sleep(10 * time.Second) g.Go(func() error { - return ingestRandomData(ctx, cli2, "testidx2", "testfield2", 10000) + return ingestRandomData(ctx, cli[1], "testidx2", "testfield2", 10000) }) - if err = sendCmd("docker", "pause", container(t, "pilosa3")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := pauseNode(t, "pilosa3"); err != nil { + t.Fatalf("sending pause command: %v", err) } - if err = sendCmd("docker", "pause", container(t, "pilosa1")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := pauseNode(t, "pilosa1"); err != nil { + t.Fatalf("sending pause command: %v", err) } time.Sleep(6 * time.Second) - if err = sendCmd("docker", "unpause", container(t, "pilosa3")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := unpauseNode(t, "pilosa3"); err != nil { + t.Fatalf("sending pause command: %v", err) } time.Sleep(6 * time.Second) - if err = sendCmd("docker", "pause", container(t, "pilosa2")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := pauseNode(t, "pilosa2"); err != nil { + t.Fatalf("sending pause command: %v", err) } time.Sleep(6 * time.Second) - if err = sendCmd("docker", "unpause", container(t, "pilosa1")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := unpauseNode(t, "pilosa1"); err != nil { + t.Fatalf("sending pause command: %v", err) } - if err = sendCmd("docker", "unpause", container(t, "pilosa2")); err != nil { - t.Fatalf("sending docker pause %v", err) + if err := unpauseNode(t, "pilosa2"); err != nil { + t.Fatalf("sending pause command: %v", err) } if err = g.Wait(); err != nil { t.Fatal(err) } - waitForStatus(t, cli1.Status, string(disco.ClusterStateNormal), 30, time.Second, ctx) + waitForStatus(t, cli[0].Status, string(disco.ClusterStateNormal), 30, time.Second, ctx) // check data in all three nodes. - for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { - r, err := cli.Query(ctx, "testidx1", &pilosa.QueryRequest{Index: "testidx1", Query: "Count(Row(testfield1 = 0))"}) + for i, c := range cli { + r, err := c.Query(ctx, "testidx1", &pilosa.QueryRequest{Index: "testidx1", Query: "Count(Row(testfield1 = 0))"}) if err != nil { t.Fatalf("count querying pilosa%d, %v", i, err) } if r.Results[0].(uint64) != 100000 { t.Fatalf("count on pilosa%d after import is %d", i, r.Results[0].(uint64)) } - } - for i, cli := range []*pilosa.InternalClient{cli1, cli2, cli3} { - r, err := cli.Query(ctx, "testidx2", &pilosa.QueryRequest{Index: "testidx2", Query: "Count(Row(testfield2 = 0))"}) + r, err = c.Query(ctx, "testidx2", &pilosa.QueryRequest{Index: "testidx2", Query: "Count(Row(testfield2 = 0))"}) if err != nil { t.Fatalf("count querying pilosa%d, %v", i, err) } From 6a2a0cf00852a8c125e406d7ec97950742c6fdef Mon Sep 17 00:00:00 2001 From: Hoang Pham Date: Fri, 11 Mar 2022 11:37:36 -0600 Subject: [PATCH 5/5] FB-1188 - TTL - run TtlRemoval on server start, better ttl log info message --- server.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/server.go b/server.go index d68b0c40d..12efb6375 100644 --- a/server.go +++ b/server.go @@ -844,6 +844,8 @@ func (s *Server) monitorResetTranslationSync() { func (s *Server) monitorTtl() { ctx := context.Background() + // Run TtlRemoval on server start + s.TtlRemoval(ctx) ticker := time.NewTicker(s.ttlRemovalInterval) for { select { @@ -879,7 +881,7 @@ func (s *Server) TtlRemoval(ctx context.Context) { if err != nil { s.logger.Errorf("ttl delete view: %s", err) } - s.logger.Infof("ttl deleted view: %s", view.name) + s.logger.Infof("ttl deleted - index: %s, field: %s, view: %s ", index.name, field.name, view.name) } } }