From b3e86e839492b40a94b36ce68be50de5b6d54495 Mon Sep 17 00:00:00 2001 From: Ben Johnson Date: Sat, 14 Dec 2019 14:11:59 -0700 Subject: [PATCH] refactoring stores back into index/field --- api.go | 104 +++++++++++- cluster.go | 341 ++++++++++++++++---------------------- encoding/proto/proto.go | 20 +++ executor.go | 12 +- executor_internal_test.go | 4 - field.go | 17 ++ holder.go | 4 + http/client.go | 2 +- http/handler.go | 25 +++ index.go | 26 +++ server.go | 4 +- server/grpc.go | 4 +- 12 files changed, 342 insertions(+), 221 deletions(-) diff --git a/api.go b/api.go index 03c3213e8..381387dbb 100644 --- a/api.go +++ b/api.go @@ -545,7 +545,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin var err error if field.keys() { - if rowStr, err = api.cluster.translateFieldID(indexName, fieldName, rowID); err != nil { + if rowStr, err = field.TranslateStore().TranslateID(rowID); err != nil { return errors.Wrap(err, "translating row") } } else { @@ -553,7 +553,9 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin } if index.Keys() { - if colStr, err = api.cluster.translateIndexPartitionID(ctx, indexName, api.cluster.idPartition(indexName, columnID), columnID); err != nil { + if store := index.TranslateStore(api.cluster.idPartition(indexName, columnID)); store == nil { + return errors.Wrap(err, "partition does not exist") + } else if colStr, err = store.TranslateID(columnID); err != nil { return errors.Wrap(err, "translating column") } } else { @@ -963,7 +965,7 @@ func (api *API) Import(ctx context.Context, req *ImportRequest, opts ...ImportOp if len(req.RowIDs) != 0 { return errors.New("row ids cannot be used because field uses string keys") } - if req.RowIDs, err = api.cluster.translateFieldKeys(req.Index, req.Field, req.RowKeys); err != nil { + if req.RowIDs, err = field.TranslateStore().TranslateKeys(req.RowKeys); err != nil { return errors.Wrap(err, "translating rows") } } @@ -1368,10 +1370,67 @@ func (api *API) Info() serverInfo { } // GetTranslateEntryReader provides an entry reader for key translation logs starting at offset. -func (api *API) GetTranslateEntryReader(ctx context.Context, offsets TranslateOffsetMap) (TranslateEntryReader, error) { +func (api *API) GetTranslateEntryReader(ctx context.Context, offsets TranslateOffsetMap) (_ TranslateEntryReader, err error) { span, ctx := tracing.StartSpanFromContext(ctx, "API.GetTranslateEntryReader") defer span.Finish() - return api.cluster.translateEntryReader(ctx, offsets) + + // Ensure all readers are cleaned up if any error. + var a []TranslateEntryReader + defer func() { + if err != nil { + for i := range a { + a[i].Close() + } + } + }() + + // Fetch all index partition readers. + for indexName, indexMap := range offsets { + index := api.holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + for partitionID, offset := range indexMap.Partitions { + store := index.TranslateStore(partitionID) + if store == nil { + return nil, ErrTranslateStoreNotFound + } + + r, err := store.EntryReader(ctx, uint64(offset)) + if err != nil { + return nil, errors.Wrap(err, "index partition translate reader") + } + a = append(a, r) + } + } + + // Fetch all field readers. + for indexName, indexMap := range offsets { + index := api.holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + for fieldName, offset := range indexMap.Fields { + field := index.Field(fieldName) + if field == nil { + return nil, ErrIndexNotFound + } + + r, err := field.TranslateStore().EntryReader(ctx, uint64(offset)) + if err != nil { + return nil, errors.Wrap(err, "field translate reader") + } + a = append(a, r) + } + } + + return NewMultiTranslateEntryReader(ctx, a), nil +} + +func (api *API) TranslateIndexKey(ctx context.Context, indexName string, key string) (uint64, error) { + return api.cluster.translateIndexKey(ctx, indexName, key) } // TranslateKeys handles a TranslateKeyRequest. @@ -1390,7 +1449,9 @@ func (api *API) TranslateKeys(ctx context.Context, r io.Reader) (_ []byte, err e return nil, err } } else { - if ids, err = api.cluster.translateFieldKeys(req.Index, req.Field, req.Keys); err != nil { + if field := api.holder.Field(req.Index, req.Field); field == nil { + return nil, ErrFieldNotFound + } else if ids, err = field.TranslateStore().TranslateKeys(req.Keys); err != nil { return nil, err } } @@ -1403,6 +1464,37 @@ func (api *API) TranslateKeys(ctx context.Context, r io.Reader) (_ []byte, err e return buf, nil } +// TranslateIDs handles a TranslateIDRequest. +func (api *API) TranslateIDs(ctx context.Context, r io.Reader) (_ []byte, err error) { + var req TranslateIDsRequest + if buf, err := ioutil.ReadAll(r); err != nil { + return nil, NewBadRequestError(errors.Wrap(err, "read translate ids request error")) + } else if err := api.Serializer.Unmarshal(buf, &req); err != nil { + return nil, NewBadRequestError(errors.Wrap(err, "unmarshal translate ids request error")) + } + + // Lookup store for either index or field and translate ids. + var keys []string + if req.Field == "" { + if keys, err = api.cluster.translateIndexIDs(ctx, req.Index, req.IDs); err != nil { + return nil, err + } + } else { + if field := api.holder.Field(req.Index, req.Field); field == nil { + return nil, ErrFieldNotFound + } else if keys, err = field.TranslateStore().TranslateIDs(req.IDs); err != nil { + return nil, err + } + } + + // Encode response. + buf, err := api.Serializer.Marshal(&TranslateIDsResponse{Keys: keys}) + if err != nil { + return nil, errors.Wrap(err, "translate ids response encoding error") + } + return buf, nil +} + // PrimaryReplicaNodeURL returns the URL of the cluster's primary replica. func (api *API) PrimaryReplicaNodeURL() url.URL { node := api.cluster.PrimaryReplicaNode() diff --git a/cluster.go b/cluster.go index 92516fb7c..f14ff96a3 100644 --- a/cluster.go +++ b/cluster.go @@ -215,10 +215,6 @@ type cluster struct { // nolint: maligned holder *Holder broadcaster broadcaster - // translation stores - indexTranslateStoreMap map[string]map[int]TranslateStore // store by index+partition - fieldTranslateStoreMap map[string]map[string]TranslateStore // store by index+field - joiningLeavingNodes chan nodeAction // joining is held open until this node @@ -240,9 +236,7 @@ type cluster struct { // nolint: maligned InternalClient InternalClient - // Instantiates new translation stores - OpenTranslateStore OpenTranslateStoreFunc - OpenTranslateReader OpenTranslateReaderFunc + // OpenTranslateReader OpenTranslateReaderFunc } // newCluster returns a new instance of Cluster with defaults. @@ -257,13 +251,8 @@ func newCluster() *cluster { closing: make(chan struct{}), joining: make(chan struct{}), - indexTranslateStoreMap: make(map[string]map[int]TranslateStore), - fieldTranslateStoreMap: make(map[string]map[string]TranslateStore), - InternalClient: newNopInternalClient(), - OpenTranslateStore: OpenInMemTranslateStore, - logger: logger.NopLogger, } } @@ -1011,10 +1000,6 @@ func (c *cluster) setup() error { if err != nil { return errors.Wrap(err, "adding local node") } - - if err := c.updateTranslateStores(); err != nil { - return err - } return nil } @@ -2024,11 +2009,6 @@ func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error { } } - // Open appropriate translate stores. - if err := c.updateTranslateStores(); err != nil { - return err - } - c.unprotectedSetState(cs.State) c.markAsJoined() @@ -2085,6 +2065,145 @@ func (c *cluster) setStatic(hosts []string) error { return nil } +func (c *cluster) translateIndexKey(ctx context.Context, indexName string, key string) (uint64, error) { + keyMap, err := c.translateIndexKeySet(ctx, indexName, map[string]struct{}{key: struct{}{}}) + if err != nil { + return 0, err + } + return keyMap[key], nil +} + +func (c *cluster) translateIndexKeys(ctx context.Context, indexName string, keys []string) ([]uint64, error) { + keySet := make(map[string]struct{}) + for _, key := range keys { + keySet[key] = struct{}{} + } + + keyMap, err := c.translateIndexKeySet(ctx, indexName, keySet) + if err != nil { + return nil, err + } + + ids := make([]uint64, len(keys)) + for i := range keys { + ids[i] = keyMap[keys[i]] + } + return ids, nil +} + +func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, keySet map[string]struct{}) (map[string]uint64, error) { + keyMap := make(map[string]uint64) + + idx := c.holder.Index(indexName) + if idx == nil { + return nil, ErrIndexNotFound + } + + // Split keys by partition. + keysByPartition := make(map[int][]string, c.partitionN) + for key := range keySet { + partitionID := c.keyPartition(indexName, key) + keysByPartition[partitionID] = append(keysByPartition[partitionID], key) + } + + // Translate keys by partition. + var g errgroup.Group + var mu sync.Mutex + for partitionID := range keysByPartition { + keys := keysByPartition[partitionID] + + g.Go(func() (err error) { + var ids []uint64 + if c.ownsPartition(c.Node.ID, partitionID) { + if ids, err = idx.TranslateStore(partitionID).TranslateKeys(keys); err != nil { + return err + } + } else { + nodes := c.partitionNodes(partitionID) + if ids, err = c.InternalClient.TranslateKeysNode(ctx, &nodes[0].URI, indexName, "", keys); err != nil { + return err + } + } + + mu.Lock() + defer mu.Unlock() + for i := range keys { + keyMap[keys[i]] = ids[i] + } + return nil + }) + } + if err := g.Wait(); err != nil { + return nil, err + } + return keyMap, nil +} + +func (c *cluster) translateIndexIDs(ctx context.Context, indexName string, ids []uint64) ([]string, error) { + idSet := make(map[uint64]struct{}) + for _, id := range ids { + idSet[id] = struct{}{} + } + + idMap, err := c.translateIndexIDSet(ctx, indexName, idSet) + if err != nil { + return nil, err + } + + keys := make([]string, len(ids)) + for i := range ids { + keys[i] = idMap[ids[i]] + } + return keys, nil +} + +func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idSet map[uint64]struct{}) (map[uint64]string, error) { + idMap := make(map[uint64]string) + + index := c.holder.Index(indexName) + if index == nil { + return nil, ErrIndexNotFound + } + + // Split ids by partition. + idsByPartition := make(map[int][]uint64, c.partitionN) + for id := range idSet { + partitionID := c.idPartition(indexName, id) + idsByPartition[partitionID] = append(idsByPartition[partitionID], id) + } + + // Translate ids by partition. + var g errgroup.Group + var mu sync.Mutex + for partitionID, ids := range idsByPartition { + g.Go(func() (err error) { + var keys []string + if c.ownsPartition(c.Node.ID, partitionID) { + if keys, err = index.TranslateStore(partitionID).TranslateIDs(ids); err != nil { + return err + } + } else { + nodes := c.partitionNodes(partitionID) + if keys, err = c.InternalClient.TranslateIDsNode(ctx, &nodes[0].URI, indexName, "", ids); err != nil { + return err + } + } + + mu.Lock() + defer mu.Unlock() + for i := range ids { + idMap[ids[i]] = keys[i] + } + return nil + }) + } + if err := g.Wait(); err != nil { + return nil, err + } + return idMap, nil +} + +/* // updateTranslateStores starts stores for partitions & fields owned by this // node and stops stores for ones this node does not own. func (c *cluster) updateTranslateStores() error { @@ -2115,23 +2234,6 @@ func (c *cluster) closeUnownedTranslateStores() error { } } - // Close index partition stores. - for indexName, m := range c.fieldTranslateStoreMap { - idx := c.holder.Index(indexName) - - for fieldName, store := range m { - if idx != nil && idx.Field(fieldName) != nil { - continue - } - if err := store.Close(); err != nil { - return err - } - delete(m, fieldName) - } - if len(c.fieldTranslateStoreMap) == 0 { - delete(c.fieldTranslateStoreMap, indexName) - } - } return nil } @@ -2157,26 +2259,6 @@ func (c *cluster) openOwnedTranslateStores() error { } } - // Open field stores. - for _, index := range c.holder.Indexes() { - m := c.fieldTranslateStoreMap[index.Name()] - if m == nil { - m = make(map[string]TranslateStore) - c.fieldTranslateStoreMap[index.Name()] = m - } - - for _, field := range index.Fields() { - if m[field.Name()] != nil { - continue - } - store, err := c.OpenTranslateStore(field.TranslateStorePath(), index.Name(), field.Name(), 0) - if err != nil { - return err - } - m[field.Name()] = store - } - } - return nil } @@ -2196,101 +2278,7 @@ func (c *cluster) fieldTranslateStore(index, field string) TranslateStore { return m[field] } -func (c *cluster) translateIndexKeys(ctx context.Context, index string, keys []string) ([]uint64, error) { - keySet := make(map[string]struct{}) - for _, key := range keys { - keySet[key] = struct{}{} - } - keyMap, err := c.translateIndexKeySet(ctx, index, keySet) - if err != nil { - return nil, err - } - - ids := make([]uint64, len(keys)) - for i := range keys { - ids[i] = keyMap[keys[i]] - } - return ids, nil -} - -func (c *cluster) translateIndexKeySet(ctx context.Context, index string, keySet map[string]struct{}) (map[string]uint64, error) { - keyMap := make(map[string]uint64) - - // Split keys by partition. - keysByPartition := make(map[int][]string, c.partitionN) - for key := range keySet { - partitionID := c.keyPartition(index, key) - keysByPartition[partitionID] = append(keysByPartition[partitionID], key) - } - - // Translate keys by partition. - var g errgroup.Group - var mu sync.Mutex - for partitionID := range keysByPartition { - keys := keysByPartition[partitionID] - - g.Go(func() (err error) { - var ids []uint64 - if c.ownsPartition(c.Node.ID, partitionID) { - if ids, err = c.translateIndexPartitionKeys(ctx, index, partitionID, keys); err != nil { - return err - } - } else { - nodes := c.partitionNodes(partitionID) - if ids, err = c.InternalClient.TranslateKeysNode(ctx, &nodes[0].URI, index, "", keys); err != nil { - return err - } - } - - mu.Lock() - defer mu.Unlock() - for i := range keys { - keyMap[keys[i]] = ids[i] - } - return nil - }) - } - if err := g.Wait(); err != nil { - return nil, err - } - return keyMap, nil -} - -func (c *cluster) translateIndexIDSet(ctx context.Context, index string, idSet map[uint64]struct{}) (map[uint64]string, error) { - idMap := make(map[uint64]string) - - // Split ids by partition. - idsByPartition := make(map[int][]uint64, c.partitionN) - for id := range idSet { - partitionID := c.idPartition(index, id) - idsByPartition[partitionID] = append(idsByPartition[partitionID], id) - } - - // Translate ids by partition. - var g errgroup.Group - var mu sync.Mutex - for partitionID, ids := range idsByPartition { - g.Go(func() error { - nodes := c.partitionNodes(partitionID) - keys, err := c.InternalClient.TranslateIDsNode(ctx, &nodes[0].URI, index, "", ids) - if err != nil { - return err - } - - mu.Lock() - defer mu.Unlock() - for i := range ids { - idMap[ids[i]] = keys[i] - } - return nil - }) - } - if err := g.Wait(); err != nil { - return nil, err - } - return idMap, nil -} func (c *cluster) translateIndexPartitionKeys(ctx context.Context, index string, partitionID int, keys []string) ([]uint64, error) { s := c.indexPartitionTranslateStore(index, partitionID) @@ -2388,54 +2376,7 @@ func (c *cluster) setTranslateStoreReadOnly(v bool) { } } } - -// translateEntryReader returns a reader that merges all index & field reader -// that are specified in the offsets map. -func (c *cluster) translateEntryReader(ctx context.Context, offsets TranslateOffsetMap) (_ TranslateEntryReader, err error) { - // Ensure all readers are cleaned up if any error. - var a []TranslateEntryReader - defer func() { - if err != nil { - for i := range a { - a[i].Close() - } - } - }() - - // Fetch all index partition readers. - for indexName, indexMap := range offsets { - for partitionID, offset := range indexMap.Partitions { - store := c.indexPartitionTranslateStore(indexName, partitionID) - if store == nil { - return nil, ErrTranslateStoreNotFound - } - - r, err := store.EntryReader(ctx, uint64(offset)) - if err != nil { - return nil, errors.Wrap(err, "index partition translate reader") - } - a = append(a, r) - } - } - - // Fetch all field readers. - for indexName, indexMap := range offsets { - for fieldName, offset := range indexMap.Fields { - store := c.fieldTranslateStore(indexName, fieldName) - if store == nil { - return nil, ErrTranslateStoreNotFound - } - - r, err := store.EntryReader(ctx, uint64(offset)) - if err != nil { - return nil, errors.Wrap(err, "field translate reader") - } - a = append(a, r) - } - } - - return NewMultiTranslateEntryReader(ctx, a), nil -} +*/ /* // holderTranslateStoreReplicator manages the replication of translation store diff --git a/encoding/proto/proto.go b/encoding/proto/proto.go index d75b784e1..74313e4cf 100644 --- a/encoding/proto/proto.go +++ b/encoding/proto/proto.go @@ -273,6 +273,22 @@ func (Serializer) Unmarshal(buf []byte, m pilosa.Message) error { } decodeTranslateKeysResponse(msg, mt) return nil + case *pilosa.TranslateIDsRequest: + msg := &internal.TranslateIDsRequest{} + err := proto.Unmarshal(buf, msg) + if err != nil { + return errors.Wrap(err, "unmarshaling TranslateIDsRequest") + } + decodeTranslateIDsRequest(msg, mt) + return nil + case *pilosa.TranslateIDsResponse: + msg := &internal.TranslateIDsResponse{} + err := proto.Unmarshal(buf, msg) + if err != nil { + return errors.Wrap(err, "unmarshaling TranslateIDsResponse") + } + decodeTranslateIDsResponse(msg, mt) + return nil default: panic(fmt.Sprintf("unhandled pilosa.Message of type %T: %#v", mt, m)) } @@ -338,6 +354,10 @@ func encodeToProto(m pilosa.Message) proto.Message { return encodeTranslateKeysRequest(mt) case *pilosa.TranslateKeysResponse: return encodeTranslateKeysResponse(mt) + case *pilosa.TranslateIDsRequest: + return encodeTranslateIDsRequest(mt) + case *pilosa.TranslateIDsResponse: + return encodeTranslateIDsResponse(mt) } return nil } diff --git a/executor.go b/executor.go index ddec4e18b..a011dc5f8 100644 --- a/executor.go +++ b/executor.go @@ -3656,7 +3656,7 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call, keyMap m return errors.Errorf("row value must be a string or non-negative integer, but got: %v of %[1]T", c.Args[rowKey]) } } else if value := callArgString(c, rowKey); value != "" { - id, err := e.Cluster.translateFieldKey(index, fieldName, value) + id, err := field.TranslateStore().TranslateKey(value) if err != nil { return err } @@ -3752,7 +3752,7 @@ func (e *executor) translateGroupByCall(index string, idx *Index, c *pql.Call, k if !ok { return errors.New("prev value must be a string when field 'keys' option enabled") } - id, err := e.Cluster.translateFieldKey(index, field.Name(), prevStr) + id, err := field.TranslateStore().TranslateKey(prevStr) if err != nil { return errors.Wrapf(err, "translating row key '%s'", prevStr) } @@ -3857,7 +3857,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res return nil, fmt.Errorf("field %q not found", fieldName) } if field.keys() { - key, err := field.translateStore.TranslateID(result.Pair.ID) + key, err := field.TranslateStore().TranslateID(result.Pair.ID) if err != nil { return nil, err } @@ -3881,7 +3881,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res if field.keys() { other := make([]Pair, len(result.Pairs)) for i := range result.Pairs { - key, err := field.translateStore.TranslateID(result.Pairs[i].ID) + key, err := field.TranslateStore().TranslateID(result.Pairs[i].ID) if err != nil { return nil, err } @@ -3908,7 +3908,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res return nil, ErrFieldNotFound } if field.keys() { - key, err := e.Cluster.translateFieldID(index, g.Field, g.RowID) + key, err := field.TranslateStore().TranslateID(g.RowID) if err != nil { return nil, errors.Wrap(err, "translating row ID in Group") } @@ -3939,7 +3939,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res } else if field.keys() { other.Keys = make([]string, len(result)) for i, id := range result { - key, err := e.Cluster.translateFieldID(index, fieldName, id) + key, err := field.TranslateStore().TranslateID(id) if err != nil { return nil, errors.Wrap(err, "translating row ID") } diff --git a/executor_internal_test.go b/executor_internal_test.go index 16588c1e2..c6ea6066a 100644 --- a/executor_internal_test.go +++ b/executor_internal_test.go @@ -54,10 +54,6 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) { t.Fatalf("creating fields %v, %v, %v", erra, errb, errc) } - if err := cluster.updateTranslateStores(); err != nil { - t.Fatal(err) - } - query, err := pql.ParseString(`GroupBy(Rows(ak), Rows(b), Rows(ck), previous=["la", 0, "ha"], having=Condition(count > 10))`) if err != nil { t.Fatalf("parsing query: %v", err) diff --git a/field.go b/field.go index 842be406d..41b1ed3b7 100644 --- a/field.go +++ b/field.go @@ -89,6 +89,11 @@ type Field struct { logger logger.Logger snapshotQueue snapshotQueue + + translateStore TranslateStore + + // Instantiates new translation stores + OpenTranslateStore OpenTranslateStoreFunc } // FieldOption is a functional option type for pilosa.fieldOptions. @@ -260,6 +265,8 @@ func newField(path, index, name string, opts FieldOption) (*Field, error) { remoteAvailableShards: roaring.NewBitmap(), logger: logger.NopLogger, + + OpenTranslateStore: OpenInMemTranslateStore, } return f, nil } @@ -278,6 +285,11 @@ func (f *Field) TranslateStorePath() string { return filepath.Join(f.path, "keys") } +// TranslateStore returns the field's translation store. +func (f *Field) TranslateStore() TranslateStore { + return f.translateStore +} + // RowAttrStore returns the attribute storage. func (f *Field) RowAttrStore() AttrStore { return f.rowAttrStore } @@ -468,6 +480,11 @@ func (f *Field) Open() error { return errors.Wrap(err, "opening attrstore") } + f.logger.Debugf("open translate store for index/field: %s/%s", f.index, f.name) + if f.translateStore, err = f.OpenTranslateStore(f.TranslateStorePath(), f.index, f.name, 0); err != nil { + return errors.Wrap(err, "opening field translate store") + } + return nil }(); err != nil { f.Close() diff --git a/holder.go b/holder.go index 2bc300ec1..d4abe0064 100644 --- a/holder.go +++ b/holder.go @@ -79,6 +79,10 @@ type Holder struct { // Manages replication from the primary node. primaryTranslateNode *Node + + // Instantiates new translation stores + OpenTranslateStore OpenTranslateStoreFunc + OpenTranslateReader OpenTranslateReaderFunc } // lockedChan looks a little ridiculous admittedly, but exists for good reason. diff --git a/http/client.go b/http/client.go index b99341c15..769d18d46 100644 --- a/http/client.go +++ b/http/client.go @@ -1183,7 +1183,7 @@ func (c *InternalClient) TranslateIDsNode(ctx context.Context, uri *pilosa.URI, return nil, pilosa.ErrIndexRequired } - buf, err := c.serializer.Marshal(pilosa.TranslateIDsRequest{ + buf, err := c.serializer.Marshal(&pilosa.TranslateIDsRequest{ Index: index, Field: field, IDs: ids, diff --git a/http/handler.go b/http/handler.go index 30d296a40..375fa09bb 100644 --- a/http/handler.go +++ b/http/handler.go @@ -309,6 +309,7 @@ func newRouter(handler *Handler) *mux.Router { router.HandleFunc("/internal/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST").Name("PostIndexAttrDiff") router.HandleFunc("/internal/translate/data", handler.handlePostTranslateData).Methods("POST").Name("PostTranslateData") router.HandleFunc("/internal/translate/keys", handler.handlePostTranslateKeys).Methods("POST").Name("PostTranslateKeys") + router.HandleFunc("/internal/translate/ids", handler.handlePostTranslateIDs).Methods("POST").Name("PostTranslateIDs") router.HandleFunc("/internal/index/{index}/field/{field}/attr/diff", handler.handlePostFieldAttrDiff).Methods("POST").Name("PostFieldAttrDiff") router.HandleFunc("/internal/index/{index}/field/{field}/remote-available-shards/{shardID}", handler.handleDeleteRemoteAvailableShard).Methods("DELETE") router.HandleFunc("/internal/nodes", handler.handleGetNodes).Methods("GET").Name("GetNodes") @@ -1757,6 +1758,30 @@ func (h *Handler) handlePostTranslateKeys(w http.ResponseWriter, r *http.Request http.Error(w, fmt.Sprintf("translate keys: %v", err), http.StatusInternalServerError) } + // Write response. + _, err = w.Write(buf) + if err != nil { + h.logger.Printf("writing translate keys response: %v", err) + return + } +} + +func (h *Handler) handlePostTranslateIDs(w http.ResponseWriter, r *http.Request) { + // Verify that request is only communicating over protobufs. + if r.Header.Get("Content-Type") != "application/x-protobuf" { + http.Error(w, "Unsupported media type", http.StatusUnsupportedMediaType) + return + } else if r.Header.Get("Accept") != "application/x-protobuf" { + http.Error(w, "Not acceptable", http.StatusNotAcceptable) + return + } + + buf, err := h.api.TranslateIDs(r.Context(), r.Body) + if err != nil { + http.Error(w, fmt.Sprintf("translate ids: %v", err), http.StatusInternalServerError) + return + } + // Write response. _, err = w.Write(buf) if err != nil { diff --git a/index.go b/index.go index 57533b22a..29b5fc56b 100644 --- a/index.go +++ b/index.go @@ -61,6 +61,12 @@ type Index struct { // Used for notifying holder when a field is added. holder *Holder + + // Per-partition translation stores + translateStores map[int]TranslateStore + + // Instantiates new translation stores + OpenTranslateStore OpenTranslateStoreFunc } // NewIndex returns a new instance of Index. @@ -82,6 +88,10 @@ func NewIndex(path, name string) (*Index, error) { Stats: stats.NopStatsClient, logger: logger.NopLogger, trackExistence: true, + + translateStores: make(map[int]TranslateStore), + + OpenTranslateStore: OpenInMemTranslateStore, }, nil } @@ -96,6 +106,11 @@ func (i *Index) TranslateStorePath(partitionID int) string { return filepath.Join(i.path, "keys", strconv.Itoa(partitionID)) } +// TranslateStore returns the translation store for a given partition. +func (i *Index) TranslateStore(partitionID int) TranslateStore { + return i.translateStores[partitionID] +} + // Keys returns true if the index uses string keys. func (i *Index) Keys() bool { return i.keys } @@ -145,6 +160,16 @@ func (i *Index) Open() (err error) { return errors.Wrap(err, "opening attrstore") } + // TODO(BBJ): Support non-default partition counts. + i.logger.Debugf("open translate store for index: %s", i.name) + for partitionID := 0; partitionID < defaultPartitionN; partitionID++ { + store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID) + if err != nil { + return errors.Wrap(err, "opening index translate store") + } + i.translateStores[partitionID] = store + } + return nil } @@ -445,6 +470,7 @@ func (i *Index) newField(path, name string) (*Field, error) { if i.snapshotQueue != nil { f.snapshotQueue = i.snapshotQueue } + f.OpenTranslateStore = i.OpenTranslateStore return f, nil } diff --git a/server.go b/server.go index d89000bb9..96a4d5cde 100644 --- a/server.go +++ b/server.go @@ -284,7 +284,7 @@ func OptServerClusterHasher(h Hasher) ServerOption { // used to specify the translation data store type. func OptServerOpenTranslateStore(fn OpenTranslateStoreFunc) ServerOption { return func(s *Server) error { - s.cluster.OpenTranslateStore = fn + s.holder.OpenTranslateStore = fn return nil } } @@ -293,7 +293,7 @@ func OptServerOpenTranslateStore(fn OpenTranslateStoreFunc) ServerOption { // used to specify the remote translation data reader. func OptServerOpenTranslateReader(fn OpenTranslateReaderFunc) ServerOption { return func(s *Server) error { - s.cluster.OpenTranslateReader = fn + s.holder.OpenTranslateReader = fn return nil } } diff --git a/server/grpc.go b/server/grpc.go index a0030d6f5..7749fff1a 100644 --- a/server/grpc.go +++ b/server/grpc.go @@ -421,7 +421,7 @@ func (h grpcHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSer case "int": // Translate column key. - id, err := index.TranslateStore().TranslateKey(col) + id, err := h.api.TranslateIndexKey(context.Background(), index.Name(), col) if err != nil { return errors.Wrap(err, "translating column key") } @@ -438,7 +438,7 @@ func (h grpcHandler) Inspect(req *pb.InspectRequest, stream pb.Pilosa_InspectSer } case "decimal": // Translate column key. - id, err := index.TranslateStore().TranslateKey(col) + id, err := h.api.TranslateIndexKey(context.Background(), index.Name(), col) if err != nil { return errors.Wrap(err, "translating column key") }