Merge pull request #892 from kuba--/fix-545

Translate field IDs on coordinator
This commit is contained in:
Kuba Podgórski 2020-09-23 17:49:15 +02:00 committed by GitHub
commit 8b855b97d0
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
4 changed files with 135 additions and 25 deletions

2
api.go
View file

@ -1827,7 +1827,7 @@ func (api *API) TranslateIDs(ctx context.Context, r io.Reader) (_ []byte, err er
if err != nil {
return nil, err
}
} else if keys, err = field.TranslateStore().TranslateIDs(req.IDs); err != nil {
} else if keys, err = api.cluster.translateFieldListIDs(field, req.IDs); err != nil {
return nil, err
}
}

View file

@ -2377,6 +2377,46 @@ func (c *cluster) translateFieldKeys(ctx context.Context, field *Field, keys []s
return ids, nil
}
func (c *cluster) translateFieldIDs(field *Field, ids map[uint64]struct{}) (map[uint64]string, error) {
idList := make([]uint64, len(ids))
{
i := 0
for id := range ids {
idList[i] = id
i++
}
}
keyList, err := c.translateFieldListIDs(field, idList)
if err != nil {
return nil, err
}
mapped := make(map[uint64]string, len(idList))
for i, key := range keyList {
mapped[idList[i]] = key
}
return mapped, nil
}
func (c *cluster) translateFieldListIDs(field *Field, ids []uint64) (keys []string, err error) {
coordinator := c.coordinatorNode()
if coordinator == nil {
return nil, errors.Errorf("translating field(%s/%s) ids(%v) - cannot find coordinator node", field.Index(), field.Name(), ids)
}
if c.Node.ID == coordinator.ID {
keys, err = field.TranslateStore().TranslateIDs(ids)
} else {
keys, err = c.InternalClient.TranslateIDsNode(context.Background(), &coordinator.URI, field.Index(), field.Name(), ids)
}
if err != nil {
return nil, errors.Wrapf(err, "translating field(%s/%s) ids(%v)", field.Index(), field.Name(), ids)
}
return keys, err
}
func (c *cluster) translateIndexKey(ctx context.Context, indexName string, key string, writable bool) (uint64, error) {
keyMap, err := c.translateIndexKeySet(ctx, indexName, map[string]struct{}{key: struct{}{}}, writable)
if err != nil {

View file

@ -5130,26 +5130,6 @@ func (e *executor) collectResultIDs(index string, idx *Index, call *pql.Call, re
return nil
}
func (e *executor) translateFieldIDs(field *Field, ids map[uint64]struct{}) (map[uint64]string, error) {
idList := make([]uint64, len(ids))
{
i := 0
for id := range ids {
idList[i] = id
i++
}
}
keyList, err := field.TranslateStore().TranslateIDs(idList)
if err != nil {
return nil, err
}
mapped := make(map[uint64]string, len(idList))
for i, key := range keyList {
mapped[idList[i]] = key
}
return mapped, nil
}
// preTranslateMatrixSet translates the IDs of a set field in an extracted matrix.
func (e *executor) preTranslateMatrixSet(mat ExtractedIDMatrix, fieldIdx uint, field *Field) (map[uint64]string, error) {
ids := make(map[uint64]struct{}, len(mat.Columns))
@ -5159,7 +5139,7 @@ func (e *executor) preTranslateMatrixSet(mat ExtractedIDMatrix, fieldIdx uint, f
}
}
return e.translateFieldIDs(field, ids)
return e.Cluster.translateFieldIDs(field, ids)
}
func (e *executor) translateResult(ctx context.Context, index string, idx *Index, call *pql.Call, result interface{}, idSet map[uint64]string) (_ interface{}, err error) {
@ -5245,7 +5225,7 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index
for i := range result.Pairs {
ids[i] = result.Pairs[i].ID
}
keys, err := field.TranslateStore().TranslateIDs(ids)
keys, err := e.Cluster.translateFieldListIDs(field, ids)
if err != nil {
return nil, err
}
@ -5296,7 +5276,7 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index
fieldTranslations := make(map[string]map[uint64]string)
for field, ids := range fieldIDs {
trans, err := e.translateFieldIDs(field, ids)
trans, err := e.Cluster.translateFieldIDs(field, ids)
if err != nil {
return nil, errors.Wrapf(err, "translating IDs in field %q", field.Name())
}
@ -5348,7 +5328,7 @@ func (e *executor) translateResult(ctx context.Context, index string, idx *Index
if field := idx.Field(fieldName); field == nil {
return nil, newNotFoundError(ErrFieldNotFound, fieldName)
} else if field.Keys() {
keys, err := field.TranslateStore().TranslateIDs(result)
keys, err := e.Cluster.translateFieldListIDs(field, result)
if err != nil {
return nil, errors.Wrap(err, "translating row IDs")
}

View file

@ -605,3 +605,93 @@ func TestTranslation_Coordinator(t *testing.T) {
}
})
}
func TestTranslation_TranslateIDsOnCluster(t *testing.T) {
c := test.MustRunCluster(t, 4,
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(true),
pilosa.OptServerNodeID("node0"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node1"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node2"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerIsCoordinator(false),
pilosa.OptServerNodeID("node3"),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
)},
)
defer c.Close()
node0 := c.GetNode(0)
node3 := c.GetNode(3)
ctx := context.Background()
idx, fld := "i", "f"
// Create an index with keys.
if _, err := node0.API.CreateIndex(ctx, idx, pilosa.IndexOptions{Keys: true}); err != nil {
t.Fatal(err)
}
// Create an index with keys.
if _, err := node0.API.CreateField(ctx, idx, fld, pilosa.OptFieldKeys()); err != nil {
t.Fatal(err)
}
keys := []string{"k0", "k1", "k2", "k3", "k4", "k5", "k6", "k7", "k8", "k9"}
// write a new key and get id
req, err := node0.API.Serializer.Marshal(&pilosa.TranslateKeysRequest{
Index: idx,
Field: fld,
Keys: keys,
NotWritable: false,
})
if err != nil {
t.Fatal(err)
}
if buf, err := node0.API.TranslateKeys(ctx, bytes.NewReader(req)); err != nil {
t.Fatal(err)
} else {
var (
respKeys pilosa.TranslateKeysResponse
respIDs pilosa.TranslateIDsResponse
)
if err = node0.API.Serializer.Unmarshal(buf, &respKeys); err != nil {
t.Fatal(err)
}
ids := respKeys.IDs
// translate ids
req, err = node3.API.Serializer.Marshal(&pilosa.TranslateIDsRequest{
Index: idx,
Field: fld,
IDs: ids,
})
if err != nil {
t.Fatal(err)
}
if buf, err = node3.API.TranslateIDs(ctx, bytes.NewReader(req)); err != nil {
t.Fatal(err)
}
if err = node3.API.Serializer.Unmarshal(buf, &respIDs); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(respIDs.Keys, keys) {
t.Fatalf("TranslateIDs(%+v): expected: %+v, got: %+v", ids, keys, respIDs.Keys)
}
}
}