Merge branch 'master' into FB-1185

This commit is contained in:
tgruben 2022-03-14 13:00:24 -05:00 • committed by GitHub
commit a54fded24a
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
2 changed files with 128 additions and 19 deletions

View file

@ -19,9 +19,9 @@ 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"
)
// container turns a docker-compose service name into a container ID
@ -91,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)
}
@ -111,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)
}
@ -130,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)
}
@ -138,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)
}
@ -157,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)
}
@ -288,6 +281,120 @@ func TestClusterStuff(t *testing.T) {
})
}
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)
}
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 < size; 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)
}
var addrs = []string{"pilosa1:10101", "pilosa2:10101", "pilosa3:10101"}
cli, err := getClients(addrs)
if err != nil {
t.Fatalf("getting client: %v", err)
}
g := new(errgroup.Group)
g.Go(func() error {
return ingestRandomData(ctx, cli[0], "testidx1", "testfield1", 100000)
})
if err := pauseNode(t, "pilosa2"); err != nil {
t.Fatalf("sending pause command: %v", err)
}
if err := pauseNode(t, "pilosa3"); err != nil {
t.Fatalf("sending pause command: %v", err)
}
time.Sleep(6 * time.Second)
if err := unpauseNode(t, "pilosa2"); err != nil {
t.Fatalf("sending pause command: %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, cli[1], "testidx2", "testfield2", 10000)
})
if err := pauseNode(t, "pilosa3"); err != nil {
t.Fatalf("sending pause command: %v", err)
}
if err := pauseNode(t, "pilosa1"); err != nil {
t.Fatalf("sending pause command: %v", err)
}
time.Sleep(6 * time.Second)
if err := unpauseNode(t, "pilosa3"); err != nil {
t.Fatalf("sending pause command: %v", err)
}
time.Sleep(6 * time.Second)
if err := pauseNode(t, "pilosa2"); err != nil {
t.Fatalf("sending pause command: %v", err)
}
time.Sleep(6 * time.Second)
if err := unpauseNode(t, "pilosa1"); err != nil {
t.Fatalf("sending pause command: %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, cli[0].Status, string(disco.ClusterStateNormal), 30, time.Second, ctx)
// check data in all three nodes.
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))
}
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)
}
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) {
t.Helper()

View file

@ -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)
}
}
}