From ba5438e02b0a808e876958834b21f940fb4f5b65 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Wed, 23 Sep 2020 15:21:33 +0200 Subject: [PATCH] Trabnslate field IDs on coordinator --- api.go | 2 +- cluster.go | 40 +++++++++++++++++++++ executor.go | 28 +++------------ translator_test.go | 90 ++++++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 135 insertions(+), 25 deletions(-) diff --git a/api.go b/api.go index 28860e894..878744e12 100644 --- a/api.go +++ b/api.go @@ -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 } } diff --git a/cluster.go b/cluster.go index 3946d1024..6e5450cf4 100644 --- a/cluster.go +++ b/cluster.go @@ -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 { diff --git a/executor.go b/executor.go index b3360aae5..0a5b207ae 100644 --- a/executor.go +++ b/executor.go @@ -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") } diff --git a/translator_test.go b/translator_test.go index 9a7607973..3de856772 100644 --- a/translator_test.go +++ b/translator_test.go @@ -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) + } + } +}