From 0f9706691475266b6b7ee8e703772a3b2f8331a0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Tue, 9 Feb 2021 13:36:53 +0100 Subject: [PATCH 1/3] Fix TestAPI_Import. Delete views when deleting a field. --- etcd/embed.go | 79 ++++++++++++++++++++++++++++++++------------------- holder.go | 4 ++- 2 files changed, 52 insertions(+), 31 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index e671442cd..2870bf051 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -482,13 +482,7 @@ func (e *Etcd) DeleteNode(ctx context.Context, nodeID string) error { } func (e *Etcd) Schema(ctx context.Context) (map[string]*disco.Index, error) { - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "Schema: creating client") - } - defer cli.Close() - - keys, vals, err := e.getKey(ctx, cli, schemaPrefix) + keys, vals, err := e.getKey(ctx, schemaPrefix) if err != nil { return nil, err } @@ -580,33 +574,31 @@ func (e *Etcd) CreateIndex(ctx context.Context, name string, val []byte) error { } func (e *Etcd) Index(ctx context.Context, name string) ([]byte, error) { - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "Index: creating client") - } - defer cli.Close() - - return e.getKeyBytes(ctx, cli, schemaPrefix+name) + return e.getKeyBytes(ctx, schemaPrefix+name) } func (e *Etcd) DeleteIndex(ctx context.Context, name string) error { - // Delete any fields below the index path. - if err := e.delKey(ctx, schemaPrefix+name+"/", true); err != nil { - return errors.Wrap(err, "deleting index fields") - } - // Delete the index. - return e.delKey(ctx, schemaPrefix+name, false) -} - -func (e *Etcd) Field(ctx context.Context, indexName string, name string) ([]byte, error) { cli, err := e.client() if err != nil { - return nil, errors.Wrap(err, "GetField: creating client") + return errors.Wrap(err, "DeleteIndex: creating client") } defer cli.Close() + key := schemaPrefix + name + // Deleting index and fields in one transaction. + _, err = cli.KV.Txn(ctx). + If(clientv3.Compare(clientv3.Version(key), ">", -1)). + Then( + clientv3.OpDelete(key+"/", clientv3.WithPrefix()), // deleting index fields + clientv3.OpDelete(key), // deleting index + ).Commit() + + return errors.Wrap(err, "DeleteIndex") +} + +func (e *Etcd) Field(ctx context.Context, indexName string, name string) ([]byte, error) { key := schemaPrefix + indexName + "/" + name - return e.getKeyBytes(ctx, cli, key) + return e.getKeyBytes(ctx, key) } func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, val []byte) error { @@ -639,7 +631,22 @@ func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, v } func (e *Etcd) DeleteField(ctx context.Context, indexname string, name string) error { - return e.delKey(ctx, schemaPrefix+indexname+"/"+name, false) + cli, err := e.client() + if err != nil { + return errors.Wrap(err, "DeleteField: creating client") + } + defer cli.Close() + + key := schemaPrefix + indexname + "/" + name + // Deleting field and views in one transaction. + _, err = cli.KV.Txn(ctx). + If(clientv3.Compare(clientv3.Version(key), ">", -1)). + Then( + clientv3.OpDelete(key+"/", clientv3.WithPrefix()), // deleting field views + clientv3.OpDelete(key), // deleting field + ).Commit() + + return errors.Wrap(err, "DeleteField") } func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpOption) error { @@ -656,7 +663,13 @@ func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpO return nil } -func (e *Etcd) getKeyBytes(ctx context.Context, cli *clientv3.Client, key string) ([]byte, error) { +func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { + cli, err := e.client() + if err != nil { + return nil, errors.Wrap(err, "getKeyBytes: creates a new client") + } + defer cli.Close() + // Get the current value for the key. resp, err := cli.Get(ctx, key) if err != nil { @@ -671,7 +684,13 @@ func (e *Etcd) getKeyBytes(ctx context.Context, cli *clientv3.Client, key string return resp.Kvs[0].Value, nil } -func (e *Etcd) getKey(ctx context.Context, cli *clientv3.Client, key string) ([]string, [][]byte, error) { +func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, error) { + cli, err := e.client() + if err != nil { + return nil, nil, errors.Wrap(err, "getKey: creates a new client") + } + defer cli.Close() + resp, err := cli.KV.Txn(ctx). If(clientv3.Compare(clientv3.Version(key), ">", -1)). Then(clientv3.OpGet(key, clientv3.WithPrefix())). @@ -700,9 +719,9 @@ func (e *Etcd) getKey(ctx context.Context, cli *clientv3.Client, key string) ([] } func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) error { - cli, err := clientv3.NewFromURLs(e.e.Server.Cluster().ClientURLs()) + cli, err := e.client() if err != nil { - return errors.Wrap(err, "delKey") + return errors.Wrap(err, "delKey: creates a new client") } defer cli.Close() diff --git a/holder.go b/holder.go index 835be2f6d..ae0b6575a 100644 --- a/holder.go +++ b/holder.go @@ -896,7 +896,9 @@ func (h *Holder) schema(ctx context.Context, includeViews bool) ([]*IndexInfo, e if err != nil { return nil, errors.Wrap(err, "decoding CreateFieldMessage") } - + if fieldName == existenceFieldName { + continue + } fi := &FieldInfo{ Name: fieldName, CreatedAt: createFieldMessage.CreatedAt, From 5d9fd0906cd863b4a2dfcaa3662bcf2f3e52b316 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Tue, 9 Feb 2021 15:51:27 +0100 Subject: [PATCH 2/3] Fix TestClusterResize_AddNodeConcurrentIndex --- server/cluster_test.go | 92 +++++++++++------------------------------- 1 file changed, 23 insertions(+), 69 deletions(-) diff --git a/server/cluster_test.go b/server/cluster_test.go index b5047ad72..eac20c76e 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -394,8 +394,11 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { skipTestUnderBlueGreenWithRoaring(t) t.Run("WithIndex", func(t *testing.T) { + c := test.MustRunCluster(t, 2) + defer c.Close() + // Configure node0 - m0 := test.MustRunCluster(t, 1).GetNode(0) + m0 := c.GetNode(0) defer m0.Close() // Create a client for each node. @@ -415,18 +418,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { }() // Configure node1 - m1 := test.NewCommandNode(t) - if err := port.GetListeners(func(lsns []*net.TCPListener) error { - portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) - - m1.Config.Etcd = portsCfg[0].Etcd - m1.Config.Name = portsCfg[0].Name - m1.Config.Cluster.Name = portsCfg[0].Cluster.Name - m1.Config.BindGRPC = portsCfg[0].BindGRPC - return m1.Start() - }, 3, 10); err != nil { - t.Fatalf("starting second main: %v", err) - } + m1 := c.GetNode(1) defer m1.Close() state0, err0 := m0.API.State() @@ -441,9 +433,13 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { t.Fatalf("error from index creation: %v", err) } }) + t.Run("ContinuousShards", func(t *testing.T) { + c := test.MustRunCluster(t, 2) + defer c.Close() + // Configure node0 - m0 := test.MustRunCluster(t, 1).GetNode(0) + m0 := c.GetNode(0) defer m0.Close() // Create a client for each node. @@ -473,23 +469,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 - m1 := test.NewCommandNode(t) - if err := port.GetListeners(func(lsns []*net.TCPListener) error { - portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) - - m1.Config.Etcd = portsCfg[0].Etcd - m1.Config.Name = portsCfg[0].Name - m1.Config.Cluster.Name = portsCfg[0].Cluster.Name - m1.Config.BindGRPC = portsCfg[0].BindGRPC - return m1.Start() - }, 3, 10); err != nil { - t.Fatalf("starting second main: %v", err) - } - errc := make(chan error, 1) - go func() { - _, err := m0.API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{}) - errc <- err - }() + m1 := c.GetNode(1) defer m1.Close() state0, err0 := m0.API.State() @@ -504,10 +484,13 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) }) + t.Run("SkippedShard", func(t *testing.T) { + c := test.MustRunCluster(t, 2) + defer c.Close() // Configure node0 - m0 := test.MustRunCluster(t, 1).GetNode(0) + m0 := c.GetNode(0) defer m0.Close() // Create a client for each node. @@ -537,24 +520,7 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 - m1 := test.NewCommandNode(t) - if err := port.GetListeners(func(lsns []*net.TCPListener) error { - portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) - - m1.Config.Etcd = portsCfg[0].Etcd - m1.Config.Name = portsCfg[0].Name - m1.Config.Cluster.Name = portsCfg[0].Cluster.Name - m1.Config.BindGRPC = portsCfg[0].BindGRPC - - errc := make(chan error, 1) - go func() { - _, err := m0.API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{}) - errc <- err - }() - return m1.Start() - }, 3, 10); err != nil { - t.Fatalf("starting second main: %v", err) - } + m1 := c.GetNode(1) defer m1.Close() state0, err0 := m0.API.State() @@ -569,9 +535,13 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) m1.QueryExpect(t, "i", "", `Row(f=1)`, exp) }) + t.Run("WithIndexKeys", func(t *testing.T) { + c := test.MustRunCluster(t, 2) + defer c.Close() + // Configure node0 - m0 := test.MustRunCluster(t, 1).GetNode(0) + m0 := c.GetNode(0) defer m0.Close() // Create a client for each node. @@ -599,24 +569,8 @@ func TestClusterResize_AddNodeConcurrentIndex(t *testing.T) { m0.QueryExpect(t, "i", "", `Row(f=1)`, exp) // Configure node1 - m1 := test.NewCommandNode(t) - if err := port.GetListeners(func(lsns []*net.TCPListener) error { - portsCfg := test.GenPortsConfig(test.NewPorts(lsns)) - - m1.Config.Etcd = portsCfg[0].Etcd - m1.Config.Name = portsCfg[0].Name - m1.Config.Cluster.Name = portsCfg[0].Cluster.Name - m1.Config.BindGRPC = portsCfg[0].BindGRPC - - errc := make(chan error, 1) - go func() { - _, err := m0.API.CreateIndex(context.Background(), "blah", pilosa.IndexOptions{}) - errc <- err - }() - return m1.Start() - }, 3, 10); err != nil { - t.Fatalf("starting second main: %v", err) - } + m1 := c.GetNode(1) + defer m1.Close() state0, err0 := m0.API.State() state1, err1 := m1.API.State() From d827f6314809dc6c98bfb0c4b8e8f666ffa83b0d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Tue, 9 Feb 2021 16:21:25 +0100 Subject: [PATCH 3/3] Fix ApplySchema API --- api.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/api.go b/api.go index 0fd7bcb0d..4906b80b8 100644 --- a/api.go +++ b/api.go @@ -1022,10 +1022,12 @@ func (api *API) ApplySchema(ctx context.Context, s *Schema, remote bool) error { } } } - if !remote { nodes := api.cluster.Nodes() for i, node := range nodes { + if node.ID == api.Node().ID { + continue + } err := api.server.defaultClient.PostSchema(ctx, &node.URI, s, true) if err != nil { return errors.Wrapf(err, "forwarding post schema to node %d of %d", i+1, len(nodes))