Merge pull request #24 from kuba--/delete-views

Fix TestAPI_Import. Delete views when deleting a field.
This commit is contained in:
Kuba Podgórski 2021-02-09 16:39:58 +01:00 committed by GitHub
commit 8591599b4c
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
4 changed files with 78 additions and 101 deletions

4
api.go
View file

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

View file

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

View file

@ -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,

View file

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