Translation store refactor

This commit is contained in:
Ben Johnson 2019-10-05 10:53:54 -06:00
parent f9e478722c
commit e844e1ad75
No known key found for this signature in database
GPG key ID: 81741CD251883081
28 changed files with 1982 additions and 2906 deletions

71
api.go
View file

@ -23,6 +23,7 @@ import (
"fmt"
"io"
"io/ioutil"
"net/url"
"strconv"
"strings"
"sync"
@ -540,7 +541,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
var err error
if field.keys() {
if rowStr, err = api.holder.translateFile.TranslateRowToString(index.Name(), field.Name(), rowID); err != nil {
if rowStr, err = field.translateStore.TranslateID(rowID); err != nil {
return errors.Wrap(err, "translating row")
}
} else {
@ -548,7 +549,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin
}
if index.Keys() {
if colStr, err = api.holder.translateFile.TranslateColumnToString(index.Name(), columnID); err != nil {
if colStr, err = index.translateStore.TranslateID(columnID); err != nil {
return errors.Wrap(err, "translating column")
}
} else {
@ -944,7 +945,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.holder.translateFile.TranslateRowsToUint64(index.Name(), field.Name(), req.RowKeys); err != nil {
if req.RowIDs, err = field.translateStore.TranslateKeys(req.RowKeys); err != nil {
return errors.Wrap(err, "translating rows")
}
}
@ -954,7 +955,7 @@ func (api *API) Import(ctx context.Context, req *ImportRequest, opts ...ImportOp
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = api.holder.translateFile.TranslateColumnsToUint64(index.Name(), req.ColumnKeys); err != nil {
if req.ColumnIDs, err = index.translateStore.TranslateKeys(req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
}
}
@ -1055,7 +1056,7 @@ func (api *API) ImportValue(ctx context.Context, req *ImportValueRequest, opts .
if len(req.ColumnIDs) != 0 {
return errors.New("column ids cannot be used because index uses string keys")
}
if req.ColumnIDs, err = api.holder.translateFile.TranslateColumnsToUint64(index.Name(), req.ColumnKeys); err != nil {
if req.ColumnIDs, err = index.translateStore.TranslateKeys(req.ColumnKeys); err != nil {
return errors.Wrap(err, "translating columns")
}
@ -1255,22 +1256,6 @@ func (api *API) ResizeAbort() error {
return errors.Wrap(err, "complete current job")
}
// GetTranslateData provides a reader for key translation logs starting at offset.
func (api *API) GetTranslateData(ctx context.Context, offset int64) (io.ReadCloser, error) {
span, ctx := tracing.StartSpanFromContext(ctx, "API.GetTranslateData")
defer span.Finish()
rc, err := api.holder.translateFile.Reader(ctx, offset)
if err != nil {
return nil, errors.Wrap(err, "read from translate store")
}
// Ensure reader is closed when the client disconnects.
go func() { <-ctx.Done(); rc.Close() }()
return rc, nil
}
// State returns the cluster state which is usually "NORMAL", but could be
// "STARTING", "RESIZING", or potentially others. See cluster.go for more
// details.
@ -1300,37 +1285,49 @@ 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) {
span, ctx := tracing.StartSpanFromContext(ctx, "API.GetTranslateEntryReader")
defer span.Finish()
return api.holder.TranslateEntryReader(ctx, offsets)
}
// TranslateKeys handles a TranslateKeyRequest.
func (api *API) TranslateKeys(body io.Reader) ([]byte, error) {
reqBytes, err := ioutil.ReadAll(body)
if err != nil {
return nil, NewBadRequestError(errors.Wrap(err, "read body error"))
}
func (api *API) TranslateKeys(r io.Reader) ([]byte, error) {
var req TranslateKeysRequest
if err := api.Serializer.Unmarshal(reqBytes, &req); err != nil {
return nil, NewBadRequestError(errors.Wrap(err, "unmarshal body error"))
if buf, err := ioutil.ReadAll(r); err != nil {
return nil, NewBadRequestError(errors.Wrap(err, "read translate keys request error"))
} else if err := api.Serializer.Unmarshal(buf, &req); err != nil {
return nil, NewBadRequestError(errors.Wrap(err, "unmarshal translate keys request error"))
}
var ids []uint64
if req.Field == "" {
ids, err = api.holder.translateFile.TranslateColumnsToUint64(req.Index, req.Keys)
} else {
ids, err = api.holder.translateFile.TranslateRowsToUint64(req.Index, req.Field, req.Keys)
// Lookup store for either index or field and translate keys.
store, err := api.holder.TranslateStore(req.Index, req.Field)
if err != nil {
return nil, err
}
ids, err := store.TranslateKeys(req.Keys)
if err != nil {
return nil, err
}
resp := TranslateKeysResponse{
IDs: ids,
}
// Encode response.
buf, err := api.Serializer.Marshal(&resp)
buf, err := api.Serializer.Marshal(&TranslateKeysResponse{IDs: ids})
if err != nil {
return nil, errors.Wrap(err, "translate keys 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()
if node == nil {
return url.URL{}
}
return node.URI.URL()
}
type serverInfo struct {
ShardWidth uint64 `json:"shardWidth"`
Memory uint64 `json:"memory"`

View file

@ -21,8 +21,11 @@ import (
"reflect"
"strings"
"testing"
"time"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/boltdb"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
)
@ -33,11 +36,15 @@ func TestAPI_Import(t *testing.T) {
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node0"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.OpenTranslateReader),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node1"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.OpenTranslateReader),
)},
)
defer c.Close()
@ -92,16 +99,20 @@ func TestAPI_Import(t *testing.T) {
if res, err := m0.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
t.Fatalf("unexpected column keys: %+v", keys)
t.Fatalf("unexpected column keys: %#v", keys)
}
// Query node1.
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
if err := test.RetryUntil(5*time.Second, func() error {
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
return err
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
return fmt.Errorf("unexpected column keys: %#v", keys)
}
return nil
}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
t.Fatalf("unexpected column keys: %+v", keys)
}
})
// Relies on the previous test creating an index with TrackExistence and
@ -178,11 +189,13 @@ func TestAPI_ImportValue(t *testing.T) {
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node0"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
pilosa.OptServerOpenTranslateReader(http.OpenTranslateReader),
)},
[]server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerNodeID("node1"),
pilosa.OptServerClusterHasher(&offsetModHasher{}),
pilosa.OptServerOpenTranslateReader(http.OpenTranslateReader),
)},
)
defer c.Close()
@ -234,10 +247,15 @@ func TestAPI_ImportValue(t *testing.T) {
}
// Query node1.
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
if err := test.RetryUntil(5*time.Second, func() error {
if res, err := m1.API.Query(ctx, &pilosa.QueryRequest{Index: index, Query: pql}); err != nil {
return err
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
return fmt.Errorf("unexpected column keys: %+v", keys)
}
return nil
}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, colKeys) {
t.Fatalf("unexpected column keys: %+v", keys)
}
})
}

399
boltdb/translate.go Normal file
View file

@ -0,0 +1,399 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package boltdb
import (
"context"
"os"
"path/filepath"
"sync"
"time"
"github.com/boltdb/bolt"
"github.com/pilosa/pilosa/v2"
"github.com/pkg/errors"
)
var (
// ErrTranslateStoreClosed is returned when reading from an TranslateEntryReader
// and the underlying store is closed.
ErrTranslateStoreClosed = errors.New("boltdb: translate store closing")
)
// OpenTranslateStore opens and initializes a boltdb translation store.
func OpenTranslateStore(path, index, field string) (pilosa.TranslateStore, error) {
s := NewTranslateStore(index, field)
s.Path = path
if err := s.Open(); err != nil {
return nil, err
}
return s, nil
}
// Ensure type implements interface.
var _ pilosa.TranslateStore = &TranslateStore{}
// TranslateStore is an on-disk storage engine for translating string-to-uint64 values.
type TranslateStore struct {
mu sync.RWMutex
db *bolt.DB
index string
field string
once sync.Once
closing chan struct{}
readOnly bool
writeNotify chan struct{}
// File path to database file.
Path string
}
// NewTranslateStore returns a new instance of TranslateStore.
func NewTranslateStore(index, field string) *TranslateStore {
return &TranslateStore{
index: index,
field: field,
closing: make(chan struct{}),
writeNotify: make(chan struct{}),
}
}
// Open opens the translate file.
func (s *TranslateStore) Open() (err error) {
if err := os.MkdirAll(filepath.Dir(s.Path), 0777); err != nil {
return errors.Wrapf(err, "mkdir %s", filepath.Dir(s.Path))
} else if s.db, err = bolt.Open(s.Path, 0666, &bolt.Options{Timeout: 1 * time.Second}); err != nil {
return errors.Wrapf(err, "open file: %s", err)
}
// Initialize buckets.
if err := s.db.Update(func(tx *bolt.Tx) error {
if _, err := tx.CreateBucketIfNotExists([]byte("keys")); err != nil {
return err
} else if _, err := tx.CreateBucketIfNotExists([]byte("ids")); err != nil {
return err
}
return nil
}); err != nil {
s.db.Close()
return err
}
return nil
}
// Close closes the underlying database.
func (s *TranslateStore) Close() (err error) {
s.once.Do(func() { close(s.closing) })
if s.db != nil {
if err := s.db.Close(); err != nil {
return err
}
}
return nil
}
// ReadOnly returns true if the store is in read-only mode.
func (s *TranslateStore) ReadOnly() bool {
s.mu.RLock()
defer s.mu.RUnlock()
return s.readOnly
}
// SetReadOnly toggles whether store is in read-only mode.
func (s *TranslateStore) SetReadOnly(v bool) {
s.mu.Lock()
defer s.mu.Unlock()
s.readOnly = v
}
// Size returns the number of bytes in the data file.
func (s *TranslateStore) Size() int64 {
if s.db == nil {
return 0
}
tx, err := s.db.Begin(false)
if err != nil {
return 0
}
defer func() { _ = tx.Rollback() }()
return tx.Size()
}
// TranslateKeys converts a string key to an integer ID.
// If key does not have an associated id then one is created.
func (s *TranslateStore) TranslateKey(key string) (id uint64, _ error) {
// Find id by key under read lock.
if err := s.db.View(func(tx *bolt.Tx) error {
id = findIDByKey(tx.Bucket([]byte("keys")), key)
return nil
}); err != nil {
return 0, err
} else if id != 0 {
return id, nil
}
if s.ReadOnly() {
return 0, pilosa.ErrTranslateStoreReadOnly
}
// Find or create id under write lock.
var written bool
if err := s.db.Update(func(tx *bolt.Tx) (err error) {
bkt := tx.Bucket([]byte("keys"))
if id = findIDByKey(bkt, key); id != 0 {
return nil
} else if id, err = bkt.NextSequence(); err != nil {
return err
} else if err := bkt.Put([]byte(key), u64tob(id)); err != nil {
return err
} else if err := tx.Bucket([]byte("ids")).Put(u64tob(id), []byte(key)); err != nil {
return err
}
written = true
return nil
}); err != nil {
return 0, err
}
if written {
s.notifyWrite()
}
return id, nil
}
// TranslateKeys converts a string key to an integer ID.
// If key does not have an associated id then one is created.
func (s *TranslateStore) TranslateKeys(keys []string) (ids []uint64, _ error) {
if len(keys) == 0 {
return nil, nil
}
// Allocate slice for ID mapping.
ids = make([]uint64, len(keys))
// Find ids by key under read lock.
var found int
if err := s.db.View(func(tx *bolt.Tx) error {
bkt := tx.Bucket([]byte("keys"))
for i, key := range keys {
if id := findIDByKey(bkt, key); id != 0 {
ids[i] = id
found++
}
}
return nil
}); err != nil {
return nil, err
} else if found == len(keys) {
return ids, nil
}
if s.ReadOnly() {
return ids, pilosa.ErrTranslateStoreReadOnly
}
// Find or create ids under write lock if any keys were not found.
var written bool
if err := s.db.Update(func(tx *bolt.Tx) (err error) {
bkt := tx.Bucket([]byte("keys"))
for i, key := range keys {
if ids[i] != 0 {
continue
}
if ids[i] = findIDByKey(bkt, key); ids[i] != 0 {
continue
} else if ids[i], err = bkt.NextSequence(); err != nil {
return err
} else if err := bkt.Put([]byte(key), u64tob(ids[i])); err != nil {
return err
} else if err := tx.Bucket([]byte("ids")).Put(u64tob(ids[i]), []byte(key)); err != nil {
return err
}
written = true
}
return nil
}); err != nil {
return nil, err
}
if written {
s.notifyWrite()
}
return ids, nil
}
// TranslateID converts an integer ID to a string key.
// Returns a blank string if ID does not exist.
func (s *TranslateStore) TranslateID(id uint64) (string, error) {
tx, err := s.db.Begin(false)
if err != nil {
return "", err
}
defer func() { _ = tx.Rollback() }()
return findKeyByID(tx.Bucket([]byte("ids")), id), nil
}
// TranslateIDs converts a list of integer IDs to a list of string keys.
func (s *TranslateStore) TranslateIDs(ids []uint64) ([]string, error) {
if len(ids) == 0 {
return nil, nil
}
tx, err := s.db.Begin(false)
if err != nil {
return nil, err
}
defer func() { _ = tx.Rollback() }()
keys := make([]string, len(ids))
for i, id := range ids {
keys[i] = findKeyByID(tx.Bucket([]byte("ids")), id)
}
return keys, nil
}
// ForceSet writes the id/key pair to the store even if read only. Used by replication.
func (s *TranslateStore) ForceSet(id uint64, key string) error {
if err := s.db.Update(func(tx *bolt.Tx) (err error) {
if err := tx.Bucket([]byte("keys")).Put([]byte(key), u64tob(id)); err != nil {
return err
} else if err := tx.Bucket([]byte("ids")).Put(u64tob(id), []byte(key)); err != nil {
return err
}
return nil
}); err != nil {
return err
}
s.notifyWrite()
return nil
}
// Reader returns a reader that streams the underlying data file.
func (s *TranslateStore) EntryReader(ctx context.Context, offset uint64) (pilosa.TranslateEntryReader, error) {
ctx, cancel := context.WithCancel(ctx)
return &TranslateEntryReader{ctx: ctx, cancel: cancel, store: s, offset: offset}, nil
}
// WriteNotify returns a channel that is closed when a new entry is written.
func (s *TranslateStore) WriteNotify() <-chan struct{} {
s.mu.RLock()
ch := s.writeNotify
s.mu.RUnlock()
return ch
}
// notifyWrite sends a write notification under write lock.
func (s *TranslateStore) notifyWrite() {
s.mu.Lock()
defer s.mu.Unlock()
close(s.writeNotify)
s.writeNotify = make(chan struct{})
}
// MaxID returns the highest id in the store.
func (s *TranslateStore) MaxID() (max uint64, err error) {
if err := s.db.View(func(tx *bolt.Tx) error {
if key, _ := tx.Bucket([]byte("ids")).Cursor().Last(); key != nil {
max = btou64(key)
}
return nil
}); err != nil {
return 0, err
}
return max, nil
}
type TranslateEntryReader struct {
ctx context.Context
store *TranslateStore
offset uint64
cancel func()
}
// Close closes the reader.
func (r *TranslateEntryReader) Close() error {
r.cancel()
return nil
}
// ReadEntry reads the next entry from the underlying translate store.
func (r *TranslateEntryReader) ReadEntry(entry *pilosa.TranslateEntry) error {
// Ensure reader has not been closed before read.
select {
case <-r.ctx.Done():
return r.ctx.Err()
case <-r.store.closing:
return ErrTranslateStoreClosed
default:
}
for {
// Obtain notification channel before read to ensure concurrency issues.
writeNotify := r.store.WriteNotify()
// Find next ID/key pair in transaction.
var found bool
if err := r.store.db.View(func(tx *bolt.Tx) error {
// Find ID/key lookup at offset or later.
cur := tx.Bucket([]byte("ids")).Cursor()
key, value := cur.Seek(u64tob(r.offset))
if key == nil {
return nil
}
// Copy ID & key to entry and mark as found.
found = true
entry.Index = r.store.index
entry.Field = r.store.field
entry.ID = btou64(key)
entry.Key = string(value)
// Update offset position.
r.offset = entry.ID + 1
return nil
}); err != nil {
return err
} else if found {
return nil
}
// If no entry found, wait for new write or reader close.
select {
case <-r.ctx.Done():
return r.ctx.Err()
case <-r.store.closing:
return ErrTranslateStoreClosed
case <-writeNotify:
}
}
}
func findIDByKey(bkt *bolt.Bucket, key string) uint64 {
if value := bkt.Get([]byte(key)); value != nil {
return btou64(value)
}
return 0
}
func findKeyByID(bkt *bolt.Bucket, id uint64) string {
return string(bkt.Get(u64tob(id)))
}

326
boltdb/translate_test.go Normal file
View file

@ -0,0 +1,326 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package boltdb_test
import (
"context"
"io/ioutil"
"os"
"testing"
"time"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/boltdb"
)
func TestTranslateStore_TranslateKey(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Ensure initial key translates to ID 1.
if id, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(1); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
// Ensure next key autoincrements.
if id, err := s.TranslateKey("bar"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(2); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
// Ensure retranslating existing key returns original ID.
if id, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(1); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
}
func TestTranslateStore_TranslateKeys(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Ensure initial keys translate to incrementing IDs.
if ids, err := s.TranslateKeys([]string{"foo", "bar"}); err != nil {
t.Fatal(err)
} else if got, want := ids[0], uint64(1); got != want {
t.Fatalf("TranslateKeys()[0]=%d, want %d", got, want)
} else if got, want := ids[1], uint64(2); got != want {
t.Fatalf("TranslateKeys()[1]=%d, want %d", got, want)
}
// Ensure retranslation returns original IDs.
if ids, err := s.TranslateKeys([]string{"foo", "bar"}); err != nil {
t.Fatal(err)
} else if got, want := ids[0], uint64(1); got != want {
t.Fatalf("TranslateKeys()[0]=%d, want %d", got, want)
} else if got, want := ids[1], uint64(2); got != want {
t.Fatalf("TranslateKeys()[1]=%d, want %d", got, want)
}
// Ensure retranslating with existing and non-existing keys returns correctly.
if ids, err := s.TranslateKeys([]string{"foo", "baz", "bar"}); err != nil {
t.Fatal(err)
} else if got, want := ids[0], uint64(1); got != want {
t.Fatalf("TranslateKeys()[0]=%d, want %d", got, want)
} else if got, want := ids[1], uint64(3); got != want {
t.Fatalf("TranslateKeys()[1]=%d, want %d", got, want)
} else if got, want := ids[2], uint64(2); got != want {
t.Fatalf("TranslateKeys()[2]=%d, want %d", got, want)
}
}
func TestTranslateStore_TranslateID(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Setup initial keys.
if _, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if _, err := s.TranslateKey("bar"); err != nil {
t.Fatal(err)
}
// Ensure IDs can be translated back to keys.
if key, err := s.TranslateID(1); err != nil {
t.Fatal(err)
} else if got, want := key, "foo"; got != want {
t.Fatalf("TranslateID()=%s, want %s", got, want)
}
if key, err := s.TranslateID(2); err != nil {
t.Fatal(err)
} else if got, want := key, "bar"; got != want {
t.Fatalf("TranslateID()=%s, want %s", got, want)
}
}
func TestTranslateStore_TranslateIDs(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Setup initial keys.
if _, err := s.TranslateKeys([]string{"foo", "bar"}); err != nil {
t.Fatal(err)
}
// Ensure IDs can be translated back to keys.
if keys, err := s.TranslateIDs([]uint64{1, 2, 3}); err != nil {
t.Fatal(err)
} else if got, want := keys[0], "foo"; got != want {
t.Fatalf("TranslateIDs()[0]=%s, want %s", got, want)
} else if got, want := keys[1], "bar"; got != want {
t.Fatalf("TranslateIDs()[1]=%s, want %s", got, want)
} else if got, want := keys[2], ""; got != want {
t.Fatalf("TranslateIDs()[2]=%s, want %s", got, want)
}
}
func TestTranslateStore_EntryReader(t *testing.T) {
t.Run("OK", func(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Create multiple new keys.
if _, err := s.TranslateKeys([]string{"foo", "bar"}); err != nil {
t.Fatal(err)
}
// Start reader from initial position.
var entry pilosa.TranslateEntry
r, err := s.EntryReader(context.Background(), 0)
if err != nil {
t.Fatal(err)
}
defer r.Close()
// Read first entry.
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(1); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "foo"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
}
// Read next entry.
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(2); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "bar"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
}
// Insert next key while reader is open.
if _, err := s.TranslateKey("baz"); err != nil {
t.Fatal(err)
}
// Read newly created entry.
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(3); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "baz"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
}
// Ensure reader closes cleanly.
if err := r.Close(); err != nil {
t.Fatal(err)
}
})
// Ensure reader will read as soon as a new write comes in using WriteNotify().
t.Run("WriteNotify", func(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Start reader from initial position.
r, err := s.EntryReader(context.Background(), 0)
if err != nil {
t.Fatal(err)
}
defer r.Close()
// Insert key in separate goroutine.
// Sleep momentarily to reader hangs.
translateErr := make(chan error)
go func() {
time.Sleep(100 * time.Millisecond)
if _, err := s.TranslateKey("foo"); err != nil {
translateErr <- err
}
}()
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if got, want := entry.ID, uint64(1); got != want {
t.Fatalf("ReadEntry() ID=%d, want %d", got, want)
} else if got, want := entry.Key, "foo"; got != want {
t.Fatalf("ReadEntry() Key=%s, want %s", got, want)
}
select {
case err := <-translateErr:
t.Fatalf("translate error: %s", err)
default:
}
})
// Ensure exits read on close.
t.Run("Close", func(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Start reader from initial position.
r, err := s.EntryReader(context.Background(), 0)
if err != nil {
t.Fatal(err)
}
defer r.Close()
// Insert key in separate goroutine.
// Sleep momentarily to reader hangs.
closeErr := make(chan error)
go func() {
time.Sleep(100 * time.Millisecond)
if err := r.Close(); err != nil {
closeErr <- err
}
}()
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != context.Canceled {
t.Fatalf("unexpected error: %#v", err)
}
select {
case err := <-closeErr:
t.Fatalf("close error: %s", err)
default:
}
})
// Ensure exits read on store close.
t.Run("StoreClose", func(t *testing.T) {
s := MustOpenNewTranslateStore()
defer MustCloseTranslateStore(s)
// Start reader from initial position.
r, err := s.EntryReader(context.Background(), 0)
if err != nil {
t.Fatal(err)
}
defer r.Close()
// Insert key in separate goroutine.
// Sleep momentarily to reader hangs.
closeErr := make(chan error)
go func() {
time.Sleep(100 * time.Millisecond)
if err := s.Close(); err != nil {
closeErr <- err
}
}()
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != boltdb.ErrTranslateStoreClosed {
t.Fatalf("unexpected error: %#v", err)
}
select {
case err := <-closeErr:
t.Fatalf("close error: %s", err)
default:
}
})
}
// MustNewTranslateStore returns a new TranslateStore with a temporary path.
func MustNewTranslateStore() *boltdb.TranslateStore {
f, err := ioutil.TempFile("", "")
if err != nil {
panic(err)
} else if err := f.Close(); err != nil {
panic(err)
}
s := boltdb.NewTranslateStore("I", "F")
s.Path = f.Name()
return s
}
// MustOpenNewTranslateStore returns a new, opened TranslateStore.
func MustOpenNewTranslateStore() *boltdb.TranslateStore {
s := MustNewTranslateStore()
if err := s.Open(); err != nil {
panic(err)
}
return s
}
// MustCloseTranslateStore closes s and removes the underlying data file.
func MustCloseTranslateStore(s *boltdb.TranslateStore) {
if err := s.Close(); err != nil {
panic(err)
} else if err := os.Remove(s.Path); err != nil {
panic(err)
}
}

View file

@ -75,6 +75,14 @@ type Node struct {
State string `json:"state"`
}
func (n *Node) Clone() *Node {
if n == nil {
return nil
}
other := *n
return &other
}
func (n Node) String() string {
return fmt.Sprintf("Node:%s:%s:%s", n.URI, n.State, n.ID)
}
@ -382,6 +390,14 @@ func (c *cluster) addNode(node *Node) error {
return nil
}
// If the cluster membership has changed, reset the primary for
// translate store replication.
if c.holder != nil {
if err := c.holder.setPrimaryTranslateStore(c.unprotectedPrimaryReplicaNode()); err != nil {
return err
}
}
// add to topology
if c.Topology == nil {
return fmt.Errorf("Cluster.Topology is nil")
@ -401,6 +417,14 @@ func (c *cluster) removeNode(nodeID string) error {
// remove from cluster
c.removeNodeBasicSorted(nodeID)
// If the cluster membership has changed, reset the primary for
// translate store replication.
if c.holder != nil {
if err := c.holder.setPrimaryTranslateStore(c.unprotectedPrimaryReplicaNode()); err != nil {
return err
}
}
// remove from topology
if c.Topology == nil {
return fmt.Errorf("Cluster.Topology is nil")
@ -876,6 +900,7 @@ func (c *cluster) ownsShard(nodeID string, index string, shard uint64) bool {
// partitionNodes returns a list of nodes that own a partition. unprotected.
func (c *cluster) partitionNodes(partitionID int) []*Node {
// Default replica count to between one and the number of nodes.
// The replica count can be zero if there are no nodes.
replicaN := c.ReplicaN
@ -1918,7 +1943,7 @@ func (c *cluster) nodeStatus() *NodeStatus {
func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error {
c.mu.Lock()
defer c.mu.Unlock()
c.logger.Printf("merge cluster status: %v", cs)
c.logger.Printf("merge cluster status: node=%s cluster=%v", c.Node.ID, cs)
// Ignore status updates from self (coordinator).
if c.unprotectedIsCoordinator() {
return nil
@ -1966,10 +1991,6 @@ func (c *cluster) mergeClusterStatus(cs *ClusterStatus) error {
}
}
// If the cluster membership has changed, reset the primary for
// translate store replication.
c.holder.setPrimaryTranslateStore(c.unprotectedPreviousNode())
c.unprotectedSetState(cs.State)
c.markAsJoined()
@ -1995,6 +2016,22 @@ func (c *cluster) unprotectedPreviousNode() *Node {
}
}
// PrimaryReplicaNode returns the node listed before the current node in c.Nodes.
// This is different than "previous node" as the first node always returns nil.
func (c *cluster) PrimaryReplicaNode() *Node {
c.mu.RLock()
defer c.mu.RUnlock()
return c.unprotectedPrimaryReplicaNode()
}
func (c *cluster) unprotectedPrimaryReplicaNode() *Node {
pos := c.nodePositionByID(c.Node.ID)
if pos <= 0 {
return nil
}
return c.nodes[pos-1]
}
// setStatic is unprotected, but only called before the cluster has been started
// (and therefore not concurrently).
func (c *cluster) setStatic(hosts []string) error {

View file

@ -54,9 +54,6 @@ type executor struct {
// Maximum number of Set() or Clear() commands per request.
MaxWritesPerRequest int
// Stores key/id translation data.
TranslateStore TranslateStore
workersWG sync.WaitGroup
workerPoolSize int
work chan job
@ -184,7 +181,7 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
// Translate column attributes, if necessary.
if idx.Keys() {
for _, col := range columnAttrSets {
v, err := e.Holder.translateFile.TranslateColumnToString(index, col.ID)
v, err := idx.translateStore.TranslateID(col.ID)
if err != nil {
return resp, err
}
@ -2645,17 +2642,18 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
fieldName = callArgString(c, "field")
rowKey = "row"
}
// Translate column key.
if idx.Keys() {
if c.Args[colKey] != nil && !isString(c.Args[colKey]) {
return errors.New("column value must be a string when index 'keys' option enabled")
}
if value := callArgString(c, colKey); value != "" {
ids, err := e.TranslateStore.TranslateColumnsToUint64(index, []string{value})
id, err := idx.translateStore.TranslateKey(value)
if err != nil {
return err
}
c.Args[colKey] = ids[0]
c.Args[colKey] = id
}
} else {
if isString(c.Args[colKey]) {
@ -2692,11 +2690,11 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
return errors.New("row value must be a string when field 'keys' option enabled")
}
if value := callArgString(c, rowKey); value != "" {
ids, err := e.TranslateStore.TranslateRowsToUint64(index, fieldName, []string{value})
id, err := field.translateStore.TranslateKey(value)
if err != nil {
return err
}
c.Args[rowKey] = ids[0]
c.Args[rowKey] = id
}
} else {
if isString(c.Args[rowKey]) {
@ -2765,11 +2763,11 @@ func (e *executor) translateGroupByCall(index string, idx *Index, c *pql.Call) e
if !ok {
return errors.New("prev value must be a string when field 'keys' option enabled")
}
ids, err := e.TranslateStore.TranslateRowsToUint64(index, field.Name(), []string{prevStr})
id, err := field.translateStore.TranslateKey(prevStr)
if err != nil {
return errors.Wrapf(err, "translating row key '%s'", prevStr)
}
previous[i] = ids[0]
previous[i] = id
} else {
if prevStr, ok := prev.(string); ok {
return errors.Errorf("got string row val '%s' in 'previous' for field %s which doesn't use string keys", prevStr, field.Name())
@ -2800,7 +2798,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
other := &Row{Attrs: result.Attrs}
for _, segment := range result.Segments() {
for _, col := range segment.Columns() {
key, err := e.TranslateStore.TranslateColumnToString(index, col)
key, err := idx.translateStore.TranslateID(col)
if err != nil {
return nil, err
}
@ -2817,7 +2815,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 := e.TranslateStore.TranslateRowToString(index, fieldName, result.ID)
key, err := field.translateStore.TranslateID(result.ID)
if err != nil {
return nil, err
}
@ -2838,7 +2836,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
if field.keys() {
other := make([]Pair, len(result))
for i := range result {
key, err := e.TranslateStore.TranslateRowToString(index, fieldName, result[i].ID)
key, err := field.translateStore.TranslateID(result[i].ID)
if err != nil {
return nil, err
}
@ -2862,7 +2860,7 @@ func (e *executor) translateResult(index string, idx *Index, call *pql.Call, res
return nil, ErrFieldNotFound
}
if field.keys() {
key, err := e.TranslateStore.TranslateRowToString(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")
}
@ -2890,7 +2888,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.TranslateStore.TranslateRowToString(index, fieldName, id)
key, err := field.translateStore.TranslateID(id)
if err != nil {
return nil, errors.Wrap(err, "translating row ID")
}

View file

@ -34,14 +34,6 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
t.Fatalf("opening holder: %v", err)
}
e.TranslateStore = e.Holder.translateFile
tf, _ := ioutil.TempFile("", "")
e.Holder.translateFile.Path = tf.Name()
err = e.Holder.translateFile.Open()
if err != nil {
t.Fatalf("opening translateFile: %v", err)
}
idx, err := e.Holder.CreateIndex("i", IndexOptions{})
if err != nil {
t.Fatalf("creating index: %v", err)
@ -54,12 +46,6 @@ func TestExecutor_TranslateGroupByCall(t *testing.T) {
t.Fatalf("creating fields %v, %v, %v", erra, errb, errc)
}
_, erra = e.TranslateStore.TranslateRowsToUint64("i", "ak", []string{"la"})
_, errb = e.TranslateStore.TranslateRowsToUint64("i", "ck", []string{"ha"})
if erra != nil || errb != nil {
t.Fatalf("translating rows %v, %v", erra, errb)
}
query, err := pql.ParseString(`GroupBy(Rows(ak), Rows(b), Rows(ck), previous=["la", 0, "ha"])`)
if err != nil {
t.Fatalf("parsing query: %v", err)

View file

@ -32,6 +32,8 @@ import (
"github.com/google/go-cmp/cmp"
"github.com/google/go-cmp/cmp/cmpopts"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/boltdb"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
"github.com/pkg/errors"
@ -2712,7 +2714,12 @@ func TestExecutor_ExecuteOptions(t *testing.T) {
// Ensure an existence field is maintained.
func TestExecutor_Execute_Existence(t *testing.T) {
t.Run("Row", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
c := test.MustRunCluster(t, 1, []server.CommandOption{
server.OptCommandServerOptions(
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.OpenTranslateReader),
),
})
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{TrackExistence: true})

View file

@ -73,6 +73,9 @@ type Field struct {
// Row attribute storage and cache
rowAttrStore AttrStore
// Key/ID translation store.
translateStore TranslateStore
broadcaster broadcaster
Stats stats.StatsClient
@ -87,6 +90,9 @@ type Field struct {
logger logger.Logger
snapshotQueue chan *fragment
// Instantiates new translation store on open.
OpenTranslateStore OpenTranslateStoreFunc
}
// FieldOption is a functional option type for pilosa.fieldOptions.
@ -232,6 +238,8 @@ func newField(path, index, name string, opts FieldOption) (*Field, error) {
remoteAvailableShards: roaring.NewBitmap(),
logger: logger.NopLogger,
OpenTranslateStore: OpenInMemTranslateStore,
}
return f, nil
}
@ -248,6 +256,9 @@ func (f *Field) Path() string { return f.path }
// RowAttrStore returns the attribute storage.
func (f *Field) RowAttrStore() AttrStore { return f.rowAttrStore }
// TranslateStore returns the underlying translation store for the field.
func (f *Field) TranslateStore() TranslateStore { return f.translateStore }
// AvailableShards returns a bitmap of shards that contain data.
func (f *Field) AvailableShards() *roaring.Bitmap {
f.mu.RLock()
@ -390,7 +401,7 @@ func (f *Field) Options() FieldOptions {
// Open opens and initializes the field.
func (f *Field) Open() error {
if err := func() error {
if err := func() (err error) {
// Ensure the field's path exists.
f.logger.Debugf("ensure field path exists: %s", f.path)
if err := os.MkdirAll(f.path, 0777); err != nil {
@ -423,6 +434,11 @@ func (f *Field) Open() error {
return errors.Wrap(err, "opening attrstore")
}
// Instantiate & open translation store.
if f.translateStore, err = f.OpenTranslateStore(filepath.Join(f.path, "keys"), f.index, f.name); err != nil {
return errors.Wrap(err, "opening translate store")
}
return nil
}(); err != nil {
f.Close()
@ -669,6 +685,12 @@ func (f *Field) Close() error {
}
f.viewMap = make(map[string]*view)
if f.translateStore != nil {
if err := f.translateStore.Close(); err != nil {
return err
}
}
return nil
}

6
go.mod
View file

@ -6,7 +6,7 @@ require (
github.com/BurntSushi/toml v0.3.1 // indirect
github.com/CAFxX/gcnotifier v0.0.0-20190112062741-224a280d589d
github.com/DataDog/datadog-go v0.0.0-20180822151419-281ae9f2d895
github.com/StackExchange/wmi v0.0.0-20181212234831-e0a55b97c705 // indirect
github.com/StackExchange/wmi v0.0.0-20190523213315-cbe66965904d // indirect
github.com/boltdb/bolt v1.3.1
github.com/cespare/xxhash v1.1.0
github.com/codahale/hdrhistogram v0.0.0-20161010025455-3a0bb77429bd // indirect
@ -24,7 +24,7 @@ require (
github.com/pkg/errors v0.8.1
github.com/prometheus/client_golang v0.9.3
github.com/prometheus/client_model v0.0.0-20190129233127-fd36f4220a90
github.com/remyoudompheng/bigfft v0.0.0-20190321074620-2f0d2b0e0001 // indirect
github.com/remyoudompheng/bigfft v0.0.0-20190728182440-6a916e37a237 // indirect
github.com/satori/go.uuid v1.2.0
github.com/shirou/gopsutil v2.18.12+incompatible
github.com/shirou/w32 v0.0.0-20160930032740-bb4de0191aa4 // indirect
@ -33,7 +33,7 @@ require (
github.com/spf13/viper v1.3.1
github.com/uber-go/atomic v1.4.0 // indirect
github.com/uber/jaeger-client-go v2.16.0+incompatible
github.com/uber/jaeger-lib v2.0.0+incompatible // indirect
github.com/uber/jaeger-lib v2.2.0+incompatible // indirect
go.uber.org/atomic v1.4.0 // indirect
golang.org/x/crypto v0.0.0-20190426145343-a29dc8fdc734 // indirect
golang.org/x/net v0.0.0-20190424112056-4829fb13d2c6 // indirect

13
go.sum
View file

@ -6,8 +6,8 @@ github.com/DataDog/datadog-go v0.0.0-20180822151419-281ae9f2d895 h1:dmc/C8bpE5Vk
github.com/DataDog/datadog-go v0.0.0-20180822151419-281ae9f2d895/go.mod h1:LButxg5PwREeZtORoXG3tL4fMGNddJ+vMq1mwgfaqoQ=
github.com/OneOfOne/xxhash v1.2.2 h1:KMrpdQIwFcEqXDklaen+P1axHaj9BSKzvpUUfnHldSE=
github.com/OneOfOne/xxhash v1.2.2/go.mod h1:HSdplMjZKSmBqAxg5vPj2TmRDmfkzw+cTzAElWljhcU=
github.com/StackExchange/wmi v0.0.0-20181212234831-e0a55b97c705 h1:UUppSQnhf4Yc6xGxSkoQpPhb7RVzuv5Nb1mwJ5VId9s=
github.com/StackExchange/wmi v0.0.0-20181212234831-e0a55b97c705/go.mod h1:3eOhrUMpNV+6aFIbp5/iudMxNCF27Vw2OZgy4xEx0Fg=
github.com/StackExchange/wmi v0.0.0-20190523213315-cbe66965904d h1:G0m3OIz70MZUWq3EgK3CesDbo8upS2Vm9/P3FtgI+Jk=
github.com/StackExchange/wmi v0.0.0-20190523213315-cbe66965904d/go.mod h1:3eOhrUMpNV+6aFIbp5/iudMxNCF27Vw2OZgy4xEx0Fg=
github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc=
github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0=
github.com/armon/consul-api v0.0.0-20180202201655-eb2c6b5be1b6/go.mod h1:grANhF5doyWs3UAsr3K4I6qtAmlQcZDesFNEHPZAzj8=
@ -108,8 +108,8 @@ github.com/prometheus/procfs v0.0.0-20181005140218-185b4288413d/go.mod h1:c3At6R
github.com/prometheus/procfs v0.0.0-20190507164030-5867b95ac084 h1:sofwID9zm4tzrgykg80hfFph1mryUeLRsUfoocVVmRY=
github.com/prometheus/procfs v0.0.0-20190507164030-5867b95ac084/go.mod h1:TjEm7ze935MbeOT/UhFTIMYKhuLP4wbCsTZCD3I8kEA=
github.com/prometheus/tsdb v0.7.1/go.mod h1:qhTCs0VvXwvX/y3TZrWD7rabWM+ijKTux40TwIPHuXU=
github.com/remyoudompheng/bigfft v0.0.0-20190321074620-2f0d2b0e0001 h1:YDeskXpkNDhPdWN3REluVa46HQOVuVkjkd2sWnrABNQ=
github.com/remyoudompheng/bigfft v0.0.0-20190321074620-2f0d2b0e0001/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
github.com/remyoudompheng/bigfft v0.0.0-20190728182440-6a916e37a237 h1:HQagqIiBmr8YXawX/le3+O26N+vPPC1PtjaF3mwnook=
github.com/remyoudompheng/bigfft v0.0.0-20190728182440-6a916e37a237/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo=
github.com/satori/go.uuid v1.2.0 h1:0uYX9dsZ2yD7q2RtLRtPSdGDWzjeM3TbMJP9utgA0ww=
github.com/satori/go.uuid v1.2.0/go.mod h1:dA0hQrYB0VpLJoorglMZABFdXlWrHn1NEOzdhQKdks0=
github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529 h1:nn5Wsu0esKSJiIVhscUtVbo7ada43DJhG55ua/hjS5I=
@ -140,8 +140,8 @@ github.com/uber-go/atomic v1.4.0 h1:yOuPqEq4ovnhEjpHmfFwsqBXDYbQeT6Nb0bwD6XnD5o=
github.com/uber-go/atomic v1.4.0/go.mod h1:/Ct5t2lcmbJ4OSe/waGBoaVvVqtO0bmtfVNex1PFV8g=
github.com/uber/jaeger-client-go v2.16.0+incompatible h1:Q2Pp6v3QYiocMxomCaJuwQGFt7E53bPYqEgug/AoBtY=
github.com/uber/jaeger-client-go v2.16.0+incompatible/go.mod h1:WVhlPFC8FDjOFMMWRy2pZqQJSXxYSwNYOkTr/Z6d3Kk=
github.com/uber/jaeger-lib v2.0.0+incompatible h1:iMSCV0rmXEogjNWPh2D0xk9YVKvrtGoHJNe9ebLu/pw=
github.com/uber/jaeger-lib v2.0.0+incompatible/go.mod h1:ComeNDZlWwrWnDv8aPp0Ba6+uUTzImX/AauajbLI56U=
github.com/uber/jaeger-lib v2.2.0+incompatible h1:MxZXOiR2JuoANZ3J6DE/U0kSFv/eJ/GfSYVCjK7dyaw=
github.com/uber/jaeger-lib v2.2.0+incompatible/go.mod h1:ComeNDZlWwrWnDv8aPp0Ba6+uUTzImX/AauajbLI56U=
github.com/ugorji/go/codec v0.0.0-20181204163529-d75b2dcb6bc8/go.mod h1:VFNgLljTbGfSG7qAOspJ7OScBnGdDN/yBr0sguwnwf0=
github.com/xordataexchange/crypt v0.0.3-0.20170626215501-b2862e3d0a77/go.mod h1:aYKd//L2LvnjZzWKhF00oedf4jCCReLcmhLdhm1A27Q=
go.uber.org/atomic v1.4.0 h1:cxzIVoETapQEqDhQu3QfnvXAV4AlzcvUCxkVUFw3+EU=
@ -155,6 +155,7 @@ golang.org/x/crypto v0.0.0-20190426145343-a29dc8fdc734 h1:p/H982KKEjUnLJkM3tt/Le
golang.org/x/crypto v0.0.0-20190426145343-a29dc8fdc734/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/net v0.0.0-20181023162649-9b4f9f5ad519 h1:x6rhz8Y9CjbgQkccRGmELH6K+LJj7tOoh3XWeC1yaQM=
golang.org/x/net v0.0.0-20181023162649-9b4f9f5ad519/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
golang.org/x/net v0.0.0-20181114220301-adae6a3d119a h1:gOpx8G595UYyvj8UK4+OFyY4rx037g3fmfhe5SasG3U=
golang.org/x/net v0.0.0-20181114220301-adae6a3d119a/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
golang.org/x/net v0.0.0-20190424112056-4829fb13d2c6 h1:FP8hkuE6yUEaJnK7O2eTuejKWwW+Rhfj80dQ2JcKxCU=

293
holder.go
View file

@ -53,10 +53,6 @@ type Holder struct {
// Indexes by name.
indexes map[string]*Index
// Key/ID translation
translateFile *TranslateFile
NewPrimaryTranslateStore func(interface{}) TranslateStore
// opened channel is closed once Open() completes.
opened lockedChan
@ -80,6 +76,14 @@ type Holder struct {
Logger logger.Logger
snapshotQueue chan *fragment
// Manages replication from the primary node.
primaryTranslateNode *Node
translateStoreReplicator *holderTranslateStoreReplicator
// Instantiates new translation stores for indexes & fields.
OpenTranslateStore OpenTranslateStoreFunc // local store
OpenTranslateReader OpenTranslateReaderFunc // replication
}
// lockedChan looks a little ridiculous admittedly, but exists for good reason.
@ -116,9 +120,6 @@ func NewHolder() *Holder {
opened: lockedChan{ch: make(chan struct{})},
translateFile: NewTranslateFile(),
NewPrimaryTranslateStore: newNopTranslateStore,
broadcaster: NopBroadcaster,
Stats: stats.NopStatsClient,
@ -127,6 +128,8 @@ func NewHolder() *Holder {
cacheFlushInterval: defaultCacheFlushInterval,
Logger: logger.NopLogger,
OpenTranslateStore: OpenInMemTranslateStore,
}
}
@ -142,6 +145,13 @@ func (h *Holder) Open() error {
return errors.Wrap(err, "creating directory")
}
// Verify that we are not trying to open with v1 translation data.
if ok, err := h.hasV1TranslateKeysFile(); err != nil {
return errors.Wrap(err, "verify v1 translation file")
} else if !ok {
return ErrCannotOpenV1TranslateFile
}
// Open path to read all index directories.
f, err := os.Open(h.Path)
if err != nil {
@ -174,6 +184,7 @@ func (h *Holder) Open() error {
} else if err != nil {
return errors.Wrap(err, "opening index")
}
if err := index.Open(); err != nil {
if err == ErrName {
h.Logger.Printf("ERROR opening index: %s, err=%s", index.Name(), err)
@ -216,17 +227,18 @@ func (h *Holder) Close() error {
h.snapshotQueue = nil
}
if h.translateFile != nil {
if err := h.translateFile.Close(); err != nil {
return err
}
}
// Reset opened in case Holder needs to be reopened.
h.opened.mu.Lock()
h.opened.ch = make(chan struct{})
h.opened.mu.Unlock()
h.mu.Lock()
if h.translateStoreReplicator != nil {
h.translateStoreReplicator.Close()
h.translateStoreReplicator = nil
}
h.mu.Unlock()
return nil
}
@ -266,6 +278,16 @@ func (h *Holder) HasData() (bool, error) {
return false, nil
}
// hasV1TranslateKeysFile returns true if a v1 translation data file exists on disk.
func (h *Holder) hasV1TranslateKeysFile() (bool, error) {
if _, err := os.Stat(filepath.Join(h.Path, ".keys")); os.IsNotExist(err) {
return true, nil
} else if err != nil {
return false, err
}
return false, nil
}
// availableShardsByIndex returns a bitmap of all shards by indexes.
func (h *Holder) availableShardsByIndex() map[string]*roaring.Bitmap {
m := make(map[string]*roaring.Bitmap)
@ -430,6 +452,9 @@ func (h *Holder) createIndex(name string, opt IndexOptions) (*Index, error) {
// Update options.
h.indexes[index.Name()] = index
// Restart replication.
go h.refreshTranslateStoreReplicator()
return index, nil
}
@ -444,6 +469,8 @@ func (h *Holder) newIndex(path, name string) (*Index, error) {
index.newAttrStore = h.NewAttrStore
index.columnAttrs = h.NewAttrStore(filepath.Join(index.path, ".data"))
index.snapshotQueue = h.snapshotQueue
index.holder = h
index.OpenTranslateStore = h.OpenTranslateStore
return index, nil
}
@ -640,13 +667,241 @@ func (h *Holder) logStartup() error {
return nil
}
func (h *Holder) setPrimaryTranslateStore(node *Node) {
var nodeID string
if node != nil {
nodeID = node.ID
// TranslateStore returns store for the given index or field.
func (h *Holder) TranslateStore(index, field string) (TranslateStore, error) {
if field == "" {
idx := h.Index(index)
if idx == nil {
return nil, ErrIndexNotFound
}
return idx.TranslateStore(), nil
}
f := h.Field(index, field)
if f == nil {
return nil, ErrFieldNotFound
}
return f.TranslateStore(), nil
}
// TranslateOffsetMap returns a map of offsets for all indexes & fields.
func (h *Holder) TranslateOffsetMap() (TranslateOffsetMap, error) {
m := make(TranslateOffsetMap)
for _, idx := range h.Indexes() {
id, err := idx.TranslateStore().MaxID()
if err != nil {
return nil, err
}
m.SetIndexOffset(idx.Name(), id+1)
for _, field := range idx.Fields() {
id, err := field.TranslateStore().MaxID()
if err != nil {
return nil, err
}
m.SetFieldOffset(idx.Name(), field.Name(), id+1)
}
}
return m, nil
}
func (h *Holder) setTranslateStoreReadOnly(v bool) {
for _, idx := range h.Indexes() {
idx.TranslateStore().SetReadOnly(v)
for _, field := range idx.Fields() {
field.TranslateStore().SetReadOnly(v)
}
}
}
func (h *Holder) setPrimaryTranslateStore(node *Node) error {
if node != nil && h.OpenTranslateReader == nil {
return nil
}
h.mu.Lock()
h.primaryTranslateNode = node.Clone()
h.mu.Unlock()
go h.refreshTranslateStoreReplicator()
return nil
}
func (h *Holder) refreshTranslateStoreReplicator() {
h.mu.RLock()
node := h.primaryTranslateNode
h.mu.RUnlock()
var nodeURL string
if node != nil {
u := node.URI.URL()
nodeURL = u.String()
}
// Stop existing replication, if running.
h.mu.Lock()
if h.translateStoreReplicator != nil {
h.translateStoreReplicator.Close()
h.translateStoreReplicator = nil
}
h.mu.Unlock()
// Set all stores read only mode based on if we have a primary.
h.setTranslateStoreReadOnly(node != nil)
// Start replication monitor, if needed.
h.mu.Lock()
defer h.mu.Unlock()
if nodeURL != "" {
h.translateStoreReplicator = newHolderTranslateStoreReplicator(h, nodeURL)
h.translateStoreReplicator.logger = h.Logger
if err := h.translateStoreReplicator.Open(); err != nil {
h.Logger.Printf("cannot open translate store replicator: %s", err)
}
}
}
// TranslateEntryReader returns a reader that merges all index & field reader
// that are specified in the offsets map.
func (h *Holder) 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 readers.
for indexName, m := range offsets {
for fieldName, offset := range m {
var store TranslateStore
idx := h.Index(indexName)
if idx == nil {
return nil, ErrIndexNotFound
}
// Fetch from index or field store.
if fieldName == "" {
store = idx.TranslateStore()
} else {
f := idx.Field(fieldName)
if f == nil {
return nil, ErrFieldNotFound
}
store = f.TranslateStore()
}
// Generate reader and append to multireader.
r, err := store.EntryReader(ctx, uint64(offset))
if err != nil {
return nil, errors.Wrap(err, "translate reader")
}
a = append(a, r)
}
}
return NewMultiTranslateEntryReader(ctx, a), nil
}
// holderTranslateStoreReplicator manages the replication of translation store
// data from a primary store to the local replica. Continually tries to
// reconnect on disconnect.
type holderTranslateStoreReplicator struct {
ctx context.Context
cancel func()
wg sync.WaitGroup
holder *Holder
nodeURL string
logger logger.Logger
}
func newHolderTranslateStoreReplicator(h *Holder, nodeURL string) *holderTranslateStoreReplicator {
r := &holderTranslateStoreReplicator{
holder: h,
nodeURL: nodeURL,
logger: logger.NopLogger,
}
r.ctx, r.cancel = context.WithCancel(context.Background())
return r
}
// Open starts the background monitoring goroutine.
func (r *holderTranslateStoreReplicator) Open() error {
r.wg.Add(1)
go func() { defer r.wg.Done(); r.monitor() }()
return nil
}
// Close stops the replicator.
func (r *holderTranslateStoreReplicator) Close() error {
r.cancel()
return nil
}
// monitor runs in a background goroutine and continually tries to connect and
// stream translate changes from the primary store.
func (r *holderTranslateStoreReplicator) monitor() {
for {
select {
case <-r.ctx.Done():
return
default:
if err := r.replicate(); err != nil {
r.logger.Printf("cannot replicate: nodeURL=%s err=%s", r.nodeURL, err)
}
time.Sleep(1 * time.Second)
}
}
}
func (r *holderTranslateStoreReplicator) replicate() error {
// Determine the offsets of every index & field store.
offsets, err := r.holder.TranslateOffsetMap()
if err != nil {
return err
} else if len(offsets) == 0 {
return nil
}
// Begin streaming from remote primary.
rd, err := r.holder.OpenTranslateReader(r.ctx, r.nodeURL, offsets)
if err != nil {
return err
}
defer rd.Close()
for {
var entry TranslateEntry
if err := rd.ReadEntry(&entry); err != nil {
return err
}
// Find appropriate store.
var store TranslateStore
if entry.Field == "" {
idx := r.holder.Index(entry.Index)
if idx == nil {
return ErrIndexNotFound
}
store = idx.TranslateStore()
} else {
f := r.holder.Field(entry.Index, entry.Field)
if f == nil {
return ErrFieldNotFound
}
store = f.TranslateStore()
}
// Apply replication to store.
if err := store.ForceSet(entry.ID, entry.Key); err != nil {
return err
}
}
ts := h.NewPrimaryTranslateStore(node)
h.translateFile.SetPrimaryStore(nodeID, ts)
}
// holderSyncer is an active anti-entropy tool that compares the local holder

View file

@ -688,7 +688,7 @@ func TestClient_ImportKeys(t *testing.T) {
{Key: "blue", Count: 2},
{Key: "purple", Count: 1},
}) {
t.Fatalf("unexpected topn result: %v", pairs)
t.Fatalf("unexpected topn result: %#v", pairs)
}
})
})

View file

@ -181,6 +181,8 @@ func (h *Handler) populateValidators() {
h.validators["GetIndex"] = queryValidationSpecRequired()
h.validators["PostIndex"] = queryValidationSpecRequired()
h.validators["DeleteIndex"] = queryValidationSpecRequired()
h.validators["GetTranslateData"] = queryValidationSpecRequired("offset")
h.validators["PostTranslateKeys"] = queryValidationSpecRequired()
h.validators["PostField"] = queryValidationSpecRequired()
h.validators["DeleteField"] = queryValidationSpecRequired()
h.validators["PostImport"] = queryValidationSpecRequired().Optional("clear", "ignoreKeyCheck")
@ -201,8 +203,6 @@ func (h *Handler) populateValidators() {
h.validators["PostFieldAttrDiff"] = queryValidationSpecRequired()
h.validators["GetNodes"] = queryValidationSpecRequired()
h.validators["GetShardMax"] = queryValidationSpecRequired()
h.validators["GetTranslateData"] = queryValidationSpecRequired("offset")
h.validators["PostTranslateKeys"] = queryValidationSpecRequired()
}
func (h *Handler) queryArgValidator(next http.Handler) http.Handler {
@ -306,12 +306,12 @@ func newRouter(handler *Handler) *mux.Router {
router.HandleFunc("/internal/fragment/data", handler.handleGetFragmentData).Methods("GET").Name("GetFragmentData")
router.HandleFunc("/internal/fragment/nodes", handler.handleGetFragmentNodes).Methods("GET").Name("GetFragmentNodes")
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/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")
router.HandleFunc("/internal/shards/max", handler.handleGetShardsMax).Methods("GET").Name("GetShardsMax") // TODO: deprecate, but it's being used by the client
router.HandleFunc("/internal/translate/data", handler.handleGetTranslateData).Methods("GET").Name("GetTranslateData")
router.HandleFunc("/internal/translate/keys", handler.handlePostTranslateKeys).Methods("POST").Name("PostTranslateKeys")
router.Use(handler.queryArgValidator)
router.Use(handler.extractTracing)
@ -508,9 +508,14 @@ func (h *Handler) handlePostQuery(w http.ResponseWriter, r *http.Request) {
resp, err := h.api.Query(r.Context(), req)
if err != nil {
switch errors.Cause(resp.Err) {
switch errors.Cause(err) {
case pilosa.ErrTooManyWrites:
w.WriteHeader(http.StatusRequestEntityTooLarge)
case pilosa.ErrTranslateStoreReadOnly:
u := h.api.PrimaryReplicaNodeURL()
u.Path, u.RawQuery = r.URL.Path, r.URL.RawQuery
http.Redirect(w, r, u.String(), http.StatusFound)
return
default:
w.WriteHeader(http.StatusBadRequest)
}
@ -1495,54 +1500,46 @@ func (h *Handler) handlePostClusterMessage(w http.ResponseWriter, r *http.Reques
type defaultClusterMessageResponse struct{}
// translateStoreBufferSize is the buffer size used for streaming data.
const translateStoreBufferSize = 1 << 16 // 64k
func (h *Handler) handleGetTranslateData(w http.ResponseWriter, r *http.Request) {
q := r.URL.Query()
offset, _ := strconv.ParseInt(q.Get("offset"), 10, 64)
rdr, err := h.api.GetTranslateData(r.Context(), offset)
if err != nil {
if errors.Cause(err) == pilosa.ErrNotImplemented {
http.Error(w, err.Error(), http.StatusNotImplemented)
} else {
http.Error(w, err.Error(), http.StatusInternalServerError)
}
func (h *Handler) handlePostTranslateData(w http.ResponseWriter, r *http.Request) {
// Parse offsets for all indexes and fields from POST body.
offsets := make(pilosa.TranslateOffsetMap)
if err := json.NewDecoder(r.Body).Decode(&offsets); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Stream all translation data.
rd, err := h.api.GetTranslateEntryReader(r.Context(), offsets)
if errors.Cause(err) == pilosa.ErrNotImplemented {
http.Error(w, err.Error(), http.StatusNotImplemented)
return
} else if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
defer rd.Close()
// Flush header so client can continue.
w.WriteHeader(http.StatusOK)
if w, ok := w.(http.Flusher); ok {
w.Flush()
}
w.(http.Flusher).Flush()
// Copy from reader to client until store or client disconnect.
useBufferSize := translateStoreBufferSize
buf := make([]byte, useBufferSize)
enc := json.NewEncoder(w)
for {
// Read from store.
n, err := rdr.Read(buf)
if err == io.EOF {
var entry pilosa.TranslateEntry
if err := rd.ReadEntry(&entry); err == io.EOF {
return
} else if err != nil {
h.logger.Printf("http: translate store read error: %s", err)
return
} else if n == 0 {
// Reset the default buffer size.
useBufferSize = translateStoreBufferSize
buf = make([]byte, useBufferSize)
continue
}
// Write to response & flush.
if _, err := w.Write(buf[:n]); err != nil {
h.logger.Printf("http: translate store response write error: %s", err)
if err := enc.Encode(&entry); err != nil {
return
} else if w, ok := w.(http.Flusher); ok {
w.Flush()
}
w.(http.Flusher).Flush()
}
}
@ -1676,6 +1673,7 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request
_, err = w.Write(buf)
if err != nil {
h.logger.Printf("writing import-roaring response: %v", err)
return
}
}

View file

@ -17,113 +17,103 @@ package http
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"io/ioutil"
"net/http"
"net/url"
"strconv"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/logger"
)
// Ensure implementation implements inteface.
var _ pilosa.TranslateStore = (*translateStore)(nil)
// translateStore represents an implementation of pilosa.TranslateStore that
// communicates over HTTP. This is used with the TranslateHandler.
type translateStore struct {
node *pilosa.Node
httpClient *http.Client
}
// NewTranslateStore returns a new instance of TranslateStore based on node.
// DEPRECATED: Providing a string url to this function is being deprecated. Instead,
// provide a *pilosa.Node.
// TODO (2.0) Refactor to avoid panic
func NewTranslateStore(node interface{}) pilosa.TranslateStore {
var n *pilosa.Node
switch v := node.(type) {
case string:
if uri, err := pilosa.NewURIFromAddress(v); err != nil {
panic("bad uri for translatestore in deprecated api")
} else {
n = &pilosa.Node{
ID: v,
URI: *uri,
}
}
case *pilosa.Node:
n = v
default:
panic("*pilosa.Node is the only type supported by NewTranslateStore().")
}
return &translateStore{node: n, httpClient: http.DefaultClient}
}
func NewTranslateStoreWithHTTPClient(httpClient *http.Client) func(interface{}) pilosa.TranslateStore {
return func(node interface{}) pilosa.TranslateStore {
store := NewTranslateStore(node)
store.(*translateStore).httpClient = httpClient
return store
}
}
// TranslateColumnsToUint64 is not currently implemented.
func (s *translateStore) TranslateColumnsToUint64(index string, values []string) ([]uint64, error) {
return nil, pilosa.ErrNotImplemented
}
// TranslateColumnToString is not currently implemented.
func (s *translateStore) TranslateColumnToString(index string, values uint64) (string, error) {
return "", pilosa.ErrNotImplemented
}
// TranslateRowsToUint64 is not currently implemented.
func (s *translateStore) TranslateRowsToUint64(index, frame string, values []string) ([]uint64, error) {
return nil, pilosa.ErrNotImplemented
}
// TranslateRowToString is not currently implemented.
func (s *translateStore) TranslateRowToString(index, frame string, values uint64) (string, error) {
return "", pilosa.ErrNotImplemented
}
// Reader returns a reader that can stream data from a remote store.
func (s *translateStore) Reader(ctx context.Context, off int64) (io.ReadCloser, error) {
// Generate remote URL.
u, err := url.Parse(s.node.URI.String())
if err != nil {
func OpenTranslateReader(ctx context.Context, nodeURL string, offsets pilosa.TranslateOffsetMap) (pilosa.TranslateEntryReader, error) {
r := NewTranslateEntryReader(ctx)
r.URL = nodeURL + "/internal/translate/data"
r.Offsets = offsets
if err := r.Open(); err != nil {
return nil, err
}
u.Path = "/internal/translate/data"
u.RawQuery = (url.Values{
"offset": {strconv.FormatInt(off, 10)},
}).Encode()
return r, nil
}
// Connect a stream to the remote server.
req, err := http.NewRequest("GET", u.String(), nil)
if err != nil {
return nil, err
}
req = req.WithContext(ctx)
// TranslateEntryReader represents an implementation of pilosa.TranslateEntryReader.
// It consolidates all index & field translate entries into a single reader.
type TranslateEntryReader struct {
ctx context.Context
cancel func()
// Connect a stream to the remote server.
resp, err := s.httpClient.Do(req)
body io.ReadCloser
dec *json.Decoder
// Lookup of offsets for each index & field.
// Must be set before calling Open().
Offsets pilosa.TranslateOffsetMap
// URL to stream entries from.
// Must be set before calling Open().
URL string
HTTPClient *http.Client
Logger logger.Logger
}
// NewTranslateEntryReader returns a new instance of TranslateEntryReader.
func NewTranslateEntryReader(ctx context.Context) *TranslateEntryReader {
r := &TranslateEntryReader{HTTPClient: http.DefaultClient, Logger: logger.NopLogger}
r.ctx, r.cancel = context.WithCancel(ctx)
return r
}
// Open initiates the reader.
func (r *TranslateEntryReader) Open() error {
// Serialize map of offsets to request body.
requestBody, err := json.Marshal(r.Offsets)
if err != nil {
return nil, fmt.Errorf("http: cannot connect to translate store endpoint: %s", err)
return err
}
// Handle error codes or return body as stream.
switch resp.StatusCode {
case http.StatusOK:
return resp.Body, nil
case http.StatusNotImplemented:
resp.Body.Close()
return nil, pilosa.ErrNotImplemented
default:
// Connect a stream to the remote server.
req, err := http.NewRequest("POST", r.URL, bytes.NewReader(requestBody))
if err != nil {
return err
}
req = req.WithContext(r.ctx)
// Connect a stream to the remote server.
resp, err := r.HTTPClient.Do(req)
if err != nil {
return fmt.Errorf("http: cannot connect to translate store endpoint: url=%s err=%s", r.URL, err)
}
r.body = resp.Body
r.dec = json.NewDecoder(r.body)
// Handle error codes.
if resp.StatusCode == http.StatusNotImplemented {
r.body.Close()
return pilosa.ErrNotImplemented
} else if resp.StatusCode != http.StatusOK {
body, _ := ioutil.ReadAll(resp.Body)
resp.Body.Close()
return nil, fmt.Errorf("http: invalid translate store endpoint status: code=%d url=%s body=%q", resp.StatusCode, u.String(), bytes.TrimSpace(body))
r.body.Close()
return fmt.Errorf("http: invalid translate store endpoint status: code=%d url=%s body=%q", resp.StatusCode, r.URL, bytes.TrimSpace(body))
}
return nil
}
// Close stops the reader.
func (r *TranslateEntryReader) Close() error {
if r.cancel != nil {
r.cancel()
}
if r.body != nil {
return r.body.Close()
}
return nil
}
// ReadEntry reads the next entry from the stream into entry.
// Returns io.EOF at the end of the stream.
func (r *TranslateEntryReader) ReadEntry(entry *pilosa.TranslateEntry) error {
return r.dec.Decode(&entry)
}

View file

@ -17,19 +17,15 @@ package http_test
import (
"context"
"fmt"
"io"
"io/ioutil"
"testing"
"time"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/mock"
"github.com/pilosa/pilosa/v2/server"
"github.com/pilosa/pilosa/v2/test"
)
func TestTranslateStore_Reader(t *testing.T) {
func TestTranslateStore_EntryReader(t *testing.T) {
// Ensure client can connect and stream the translate store data.
t.Run("OK", func(t *testing.T) {
t.Run("ServerDisconnect", func(t *testing.T) {
@ -58,50 +54,88 @@ func TestTranslateStore_Reader(t *testing.T) {
}
// Connect to server and stream all available data.
store := http.NewTranslateStore(primary.URL())
r := http.NewTranslateEntryReader(context.Background())
r.URL = primary.URL()
// Wait to ensure writes make it to translate store
time.Sleep(500 * time.Millisecond)
rc, err := store.Reader(context.Background(), 11) // offset=11 skips the first entry: \n\x01\x01i\x00\x01\x01\x03foo
// Close the primary to disconnect reader.
primary.Close()
if err != nil {
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if data, err := ioutil.ReadAll(rc); err != nil {
t.Fatal(err)
} else if string(data) != "\n\x01\x01i\x00\x01\x02\x03bar\n\x01\x01i\x00\x01\x03\x03baz" {
t.Fatalf("unexpected data: %q", data)
} else if err := rc.Close(); err != nil {
} else if got, want := entry.ID, uint64(1); got != want {
t.Fatalf("entry.ID=%v, want %v", got, want)
} else if got, want := entry.Key, "-"; got != want {
t.Fatalf("entry.Key=%v, want %v", got, want)
}
if err := r.Close(); err != nil {
t.Fatal(err)
}
})
// Ensure server closes store reader if client disconnects.
t.Run("ClientDisconnect", func(t *testing.T) {
/*
// Ensure server closes store reader if client disconnects.
t.Run("ClientDisconnect", func(t *testing.T) {
t.Skip() // can't mock server from http package
// Setup mock so that Read() hangs.
done := make(chan struct{})
var mrc mock.ReadCloser
mrc.ReadFunc = func(p []byte) (int, error) {
<-done
return 0, io.EOF
}
closeInvoked := make(chan struct{})
mrc.CloseFunc = func() error {
close(closeInvoked)
return nil
}
var translateStore mock.TranslateStore
translateStore.ReaderFunc = func(ctx context.Context, off int64) (io.ReadCloser, error) {
return &mrc, nil
}
opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore))
cluster := test.MustRunCluster(t, 1, []server.CommandOption{opts})
defer cluster.Close()
primary := cluster[0]
defer close(done)
// Connect to server and begin streaming.
ctx, cancel := context.WithCancel(context.Background())
store := http.NewTranslateStore(primary.URL())
if _, err := store.Reader(ctx, 0); err != nil {
t.Fatal(err)
}
// Cancel the context and check if server is closed.
cancel()
select {
case <-time.NewTimer(time.Millisecond * 100).C:
t.Fatal("expected server close")
case <-closeInvoked:
return
}
})
*/
})
/*
// Ensure client is notified if the server doesn't support streaming replication.
t.Run("ErrNotImplemented", func(t *testing.T) {
t.Skip() // can't mock server from http package
// Setup mock so that Read() hangs.
done := make(chan struct{})
var mrc mock.ReadCloser
mrc.ReadFunc = func(p []byte) (int, error) {
<-done
return 0, io.EOF
}
closeInvoked := make(chan struct{})
mrc.CloseFunc = func() error {
close(closeInvoked)
return nil
}
var translateStore mock.TranslateStore
translateStore.ReaderFunc = func(ctx context.Context, off int64) (io.ReadCloser, error) {
return &mrc, nil
return nil, pilosa.ErrNotImplemented
}
opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore))
@ -109,43 +143,11 @@ func TestTranslateStore_Reader(t *testing.T) {
defer cluster.Close()
primary := cluster[0]
defer close(done)
// Connect to server and begin streaming.
ctx, cancel := context.WithCancel(context.Background())
store := http.NewTranslateStore(primary.URL())
if _, err := store.Reader(ctx, 0); err != nil {
t.Fatal(err)
}
// Cancel the context and check if server is closed.
cancel()
select {
case <-time.NewTimer(time.Millisecond * 100).C:
t.Fatal("expected server close")
case <-closeInvoked:
return
ts := http.NewTranslateStore(primary.URL())
_, err := ts.Reader(context.Background(), 0)
if err != pilosa.ErrNotImplemented {
t.Fatalf("unexpected error: %s", err)
}
})
})
// Ensure client is notified if the server doesn't support streaming replication.
t.Run("ErrNotImplemented", func(t *testing.T) {
t.Skip() // can't mock server from http package
var translateStore mock.TranslateStore
translateStore.ReaderFunc = func(ctx context.Context, off int64) (io.ReadCloser, error) {
return nil, pilosa.ErrNotImplemented
}
opts := server.OptCommandServerOptions(pilosa.OptServerPrimaryTranslateStore(translateStore))
cluster := test.MustRunCluster(t, 1, []server.CommandOption{opts})
defer cluster.Close()
primary := cluster[0]
ts := http.NewTranslateStore(primary.URL())
_, err := ts.Reader(context.Background(), 0)
if err != pilosa.ErrNotImplemented {
t.Fatalf("unexpected error: %s", err)
}
})
*/
}

View file

@ -52,11 +52,19 @@ type Index struct {
// Column attribute storage and cache.
columnAttrs AttrStore
translateStore TranslateStore
broadcaster broadcaster
Stats stats.StatsClient
logger logger.Logger
snapshotQueue chan *fragment
// Used for notifying holder when a field is added.
holder *Holder
// Instantiates new translation stores for fields.
OpenTranslateStore OpenTranslateStoreFunc
}
// NewIndex returns a new instance of Index.
@ -78,6 +86,8 @@ func NewIndex(path, name string) (*Index, error) {
Stats: stats.NopStatsClient,
logger: logger.NopLogger,
trackExistence: true,
OpenTranslateStore: OpenInMemTranslateStore,
}, nil
}
@ -93,6 +103,9 @@ func (i *Index) Keys() bool { return i.keys }
// ColumnAttrStore returns the storage for column attributes.
func (i *Index) ColumnAttrStore() AttrStore { return i.columnAttrs }
// TranslateStore returns the underlying translation store for the index.
func (i *Index) TranslateStore() TranslateStore { return i.translateStore }
// Options returns all options for this index.
func (i *Index) Options() IndexOptions {
i.mu.RLock()
@ -108,7 +121,7 @@ func (i *Index) options() IndexOptions {
}
// Open opens and initializes the index.
func (i *Index) Open() error {
func (i *Index) Open() (err error) {
// Ensure the path exists.
i.logger.Debugf("ensure index path exists: %s", i.path)
if err := os.MkdirAll(i.path, 0777); err != nil {
@ -136,6 +149,11 @@ func (i *Index) Open() error {
return errors.Wrap(err, "opening attrstore")
}
// Instantiate & open translation store.
if i.translateStore, err = i.OpenTranslateStore(filepath.Join(i.path, "keys"), i.name, ""); err != nil {
return errors.Wrap(err, "opening translate store")
}
return nil
}
@ -260,6 +278,12 @@ func (i *Index) Close() error {
}
i.fields = make(map[string]*Field)
if i.translateStore != nil {
if err := i.translateStore.Close(); err != nil {
return err
}
}
return nil
}
@ -420,6 +444,11 @@ func (i *Index) createField(name string, opt FieldOptions) (*Field, error) {
// Add to index's field lookup.
i.fields[name] = f
// Update replication, if needed.
if i.holder != nil {
go i.holder.refreshTranslateStoreReplicator()
}
return f, nil
}
@ -433,6 +462,7 @@ func (i *Index) newField(path, name string) (*Field, error) {
f.broadcaster = i.broadcaster
f.rowAttrStore = i.newAttrStore(filepath.Join(f.path, ".data"))
f.snapshotQueue = i.snapshotQueue
f.OpenTranslateStore = i.OpenTranslateStore
return f, nil
}

View file

@ -1,229 +0,0 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package inmem
import (
"context"
"io"
"sync"
"github.com/pilosa/pilosa/v2"
)
// Ensure type implements interface.
var _ pilosa.TranslateStore = &translateStore{}
// translateStore is an in-memory storage engine for translating string-to-uint64 values.
type translateStore struct {
mu sync.RWMutex
cols map[string]*translateIndex
rows map[frameKey]*translateIndex
}
// NewTranslateStore returns a new instance of TranslateStore.
func NewTranslateStore() *translateStore {
return &translateStore{
cols: make(map[string]*translateIndex),
rows: make(map[frameKey]*translateIndex),
}
}
// Reader returns an error because it is not supported by the inmem store.
func (s *translateStore) Reader(ctx context.Context, offset int64) (io.ReadCloser, error) {
return nil, pilosa.ErrReplicationNotSupported
}
// TranslateColumnsToUint64 converts value to a uint64 id.
// If value does not have an associated id then one is created.
func (s *translateStore) TranslateColumnsToUint64(index string, values []string) ([]uint64, error) {
ret := make([]uint64, len(values))
// Read value under read lock.
s.mu.RLock()
if idx := s.cols[index]; idx != nil {
var writeRequired bool
for i := range values {
v, ok := idx.lookup[values[i]]
if !ok {
writeRequired = true
}
ret[i] = v
}
if !writeRequired {
s.mu.RUnlock()
return ret, nil
}
}
s.mu.RUnlock()
// If any values not found then recheck and then add under a write lock.
s.mu.Lock()
defer s.mu.Unlock()
// Recheck if value was created between the read lock and write lock.
idx := s.cols[index]
if idx != nil {
var writeRequired bool
for i := range values {
if ret[i] != 0 {
continue
}
v, ok := idx.lookup[values[i]]
if !ok {
writeRequired = true
continue
}
ret[i] = v
}
if !writeRequired {
return ret, nil
}
}
// Create index map if it doesn't exists.
if idx == nil {
idx = newTranslateIndex()
s.cols[index] = idx
}
// Add new identifiers.
for i := range values {
if ret[i] != 0 {
continue
}
idx.seq++
v := idx.seq
ret[i] = v
idx.lookup[values[i]] = v
idx.reverse[v] = values[i]
}
return ret, nil
}
// TranslateColumnToString converts a uint64 id to its associated string value.
// If the id is not associated with a string value then a blank string is returned.
func (s *translateStore) TranslateColumnToString(index string, value uint64) (string, error) {
s.mu.RLock()
if idx := s.cols[index]; idx != nil {
if ret, ok := idx.reverse[value]; ok {
s.mu.RUnlock()
return ret, nil
}
}
s.mu.RUnlock()
return "", nil
}
func (s *translateStore) TranslateRowsToUint64(index, frame string, values []string) ([]uint64, error) {
key := frameKey{index, frame}
ret := make([]uint64, len(values))
// Read value under read lock.
s.mu.RLock()
if idx := s.rows[key]; idx != nil {
var writeRequired bool
for i := range values {
v, ok := idx.lookup[values[i]]
if !ok {
writeRequired = true
}
ret[i] = v
}
if !writeRequired {
s.mu.RUnlock()
return ret, nil
}
}
s.mu.RUnlock()
// If any values not found then recheck and then add under a write lock.
s.mu.Lock()
defer s.mu.Unlock()
// Recheck if value was created between the read lock and write lock.
idx := s.rows[key]
if idx != nil {
var writeRequired bool
for i := range values {
if ret[i] != 0 {
continue
}
v, ok := idx.lookup[values[i]]
if !ok {
writeRequired = true
continue
}
ret[i] = v
}
if !writeRequired {
return ret, nil
}
}
// Create map if it doesn't exists.
if idx == nil {
idx = newTranslateIndex()
s.rows[key] = idx
}
// Add new identifiers.
for i := range values {
if ret[i] != 0 {
continue
}
idx.seq++
v := idx.seq
ret[i] = v
idx.lookup[values[i]] = v
idx.reverse[v] = values[i]
}
return ret, nil
}
func (s *translateStore) TranslateRowToString(index, frame string, value uint64) (string, error) {
s.mu.RLock()
if idx := s.rows[frameKey{index, frame}]; idx != nil {
if ret, ok := idx.reverse[value]; ok {
s.mu.RUnlock()
return ret, nil
}
}
s.mu.RUnlock()
return "", nil
}
type frameKey struct {
index string
frame string
}
type translateIndex struct {
seq uint64
lookup map[string]uint64
reverse map[uint64]string
}
func newTranslateIndex() *translateIndex {
return &translateIndex{
lookup: make(map[string]uint64),
reverse: make(map[uint64]string),
}
}

View file

@ -1,146 +0,0 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package inmem_test
import (
"fmt"
"math/rand"
"reflect"
"testing"
"github.com/pilosa/pilosa/v2/inmem"
)
func TestTranslateStore_TranslateColumn(t *testing.T) {
s := inmem.NewTranslateStore()
// First translation should start id at zero.
if ids, err := s.TranslateColumnsToUint64("IDX0", []string{"foo"}); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(ids, []uint64{1}) {
t.Fatalf("unexpected id: %#v", ids)
}
// Next translation on the same index should move to one.
if ids, err := s.TranslateColumnsToUint64("IDX0", []string{"bar"}); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(ids, []uint64{2}) {
t.Fatalf("unexpected id: %#v", ids)
}
// Translation on a different index restarts at 0.
if ids, err := s.TranslateColumnsToUint64("IDX1", []string{"bar"}); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(ids, []uint64{1}) {
t.Fatalf("unexpected id: %#v", ids)
}
// Ensure that string values can be looked up by ID.
if value, err := s.TranslateColumnToString("IDX0", 2); err != nil {
t.Fatal(err)
} else if value != "bar" {
t.Fatalf("unexpected value: %s", value)
}
}
func TestTranslateStore_TranslateRow(t *testing.T) {
s := inmem.NewTranslateStore()
// First translation should start id at zero.
if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME0", []string{"foo"}); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(ids, []uint64{1}) {
t.Fatalf("unexpected id: %#v", ids)
}
// Next translation on the same index should move to one.
if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME0", []string{"bar"}); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(ids, []uint64{2}) {
t.Fatalf("unexpected id: %#v", ids)
}
// Translation on a different index restarts at 0.
if ids, err := s.TranslateRowsToUint64("IDX1", "FRAME0", []string{"bar"}); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(ids, []uint64{1}) {
t.Fatalf("unexpected id: %#v", ids)
}
// Translation on a different frame restarts at 0.
if ids, err := s.TranslateRowsToUint64("IDX0", "FRAME1", []string{"bar"}); err != nil {
t.Fatal(err)
} else if !reflect.DeepEqual(ids, []uint64{1}) {
t.Fatalf("unexpected id: %#v", ids)
}
// Ensure that string values can be looked up by ID.
if value, err := s.TranslateRowToString("IDX0", "FRAME0", 2); err != nil {
t.Fatal(err)
} else if value != "bar" {
t.Fatalf("unexpected value: %s", value)
}
}
func BenchmarkTranslateStore_TranslateColumnsToUint64(b *testing.B) {
const batchSize = 1000
s := inmem.NewTranslateStore()
// Generate keys before benchmark begins
keySets := make([][]string, b.N/1000)
for i := range keySets {
keySets[i] = make([]string, batchSize)
for j, jv := range rand.New(rand.NewSource(0)).Perm(batchSize) {
keySets[i][j] = fmt.Sprintf("%08d%08d", jv, i)
}
}
b.ResetTimer()
for _, keySet := range keySets {
if _, err := s.TranslateColumnsToUint64("IDX0", keySet); err != nil {
b.Fatal(err)
}
}
}
func BenchmarkTranslateStore_TranslateColumnToString(b *testing.B) {
const batchSize = 1000
s := inmem.NewTranslateStore()
// Generate keys before benchmark begins
for i := 0; i < b.N; i += batchSize {
keySet := make([]string, batchSize)
for j, jv := range rand.New(rand.NewSource(0)).Perm(batchSize) {
keySet[j] = fmt.Sprintf("%08d%08d", jv, i)
}
if _, err := s.TranslateColumnsToUint64("IDX0", keySet); err != nil {
b.Fatal(err)
}
}
// Generate random key access.
perm := rand.New(rand.NewSource(0)).Perm(b.N)
b.ResetTimer()
for i := 0; i < b.N; i++ {
if _, err := s.TranslateColumnToString("IDX0", uint64(perm[i])); err != nil {
b.Fatal(err)
}
}
}

View file

@ -16,7 +16,6 @@ package mock
import (
"context"
"io"
"github.com/pilosa/pilosa/v2"
)
@ -24,29 +23,69 @@ import (
var _ pilosa.TranslateStore = (*TranslateStore)(nil)
type TranslateStore struct {
TranslateColumnsToUint64Func func(index string, values []string) ([]uint64, error)
TranslateColumnToStringFunc func(index string, values uint64) (string, error)
TranslateRowsToUint64Func func(index, field string, values []string) ([]uint64, error)
TranslateRowToStringFunc func(index, field string, values uint64) (string, error)
ReaderFunc func(ctx context.Context, off int64) (io.ReadCloser, error)
CloseFunc func() error
MaxIDFunc func() (uint64, error)
ReadOnlyFunc func() bool
SetReadOnlyFunc func(v bool)
TranslateKeyFunc func(key string) (uint64, error)
TranslateKeysFunc func(keys []string) ([]uint64, error)
TranslateIDFunc func(id uint64) (string, error)
TranslateIDsFunc func(ids []uint64) ([]string, error)
ForceSetFunc func(id uint64, key string) error
EntryReaderFunc func(ctx context.Context, offset uint64) (pilosa.TranslateEntryReader, error)
}
func (s TranslateStore) TranslateColumnsToUint64(index string, values []string) ([]uint64, error) {
return s.TranslateColumnsToUint64Func(index, values)
func (s *TranslateStore) Close() error {
return s.CloseFunc()
}
func (s TranslateStore) TranslateColumnToString(index string, values uint64) (string, error) {
return s.TranslateColumnToStringFunc(index, values)
func (s *TranslateStore) MaxID() (uint64, error) {
return s.MaxIDFunc()
}
func (s TranslateStore) TranslateRowsToUint64(index, field string, values []string) ([]uint64, error) {
return s.TranslateRowsToUint64Func(index, field, values)
func (s *TranslateStore) ReadOnly() bool {
return s.ReadOnlyFunc()
}
func (s TranslateStore) TranslateRowToString(index, field string, value uint64) (string, error) {
return s.TranslateRowToStringFunc(index, field, value)
func (s *TranslateStore) SetReadOnly(v bool) {
s.SetReadOnlyFunc(v)
}
func (s TranslateStore) Reader(ctx context.Context, off int64) (io.ReadCloser, error) {
return s.ReaderFunc(ctx, off)
func (s *TranslateStore) TranslateKey(key string) (uint64, error) {
return s.TranslateKeyFunc(key)
}
func (s *TranslateStore) TranslateKeys(keys []string) ([]uint64, error) {
return s.TranslateKeysFunc(keys)
}
func (s *TranslateStore) TranslateID(id uint64) (string, error) {
return s.TranslateIDFunc(id)
}
func (s *TranslateStore) TranslateIDs(ids []uint64) ([]string, error) {
return s.TranslateIDsFunc(ids)
}
func (s *TranslateStore) ForceSet(id uint64, key string) error {
return s.ForceSetFunc(id, key)
}
func (s *TranslateStore) EntryReader(ctx context.Context, offset uint64) (pilosa.TranslateEntryReader, error) {
return s.EntryReaderFunc(ctx, offset)
}
var _ pilosa.TranslateEntryReader = (*TranslateEntryReader)(nil)
type TranslateEntryReader struct {
CloseFunc func() error
ReadEntryFunc func(entry *pilosa.TranslateEntry) error
}
func (r *TranslateEntryReader) Close() error {
return r.CloseFunc()
}
func (r *TranslateEntryReader) ReadEntry(entry *pilosa.TranslateEntry) error {
return r.ReadEntryFunc(entry)
}

View file

@ -201,17 +201,6 @@ func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption {
}
}
// OptServerPrimaryTranslateStoreFunc is a functional option on Server
// used to specify the function used to create a new primary translate
// store.
func OptServerPrimaryTranslateStoreFunc(tf func(interface{}) TranslateStore) ServerOption {
return func(s *Server) error {
s.holder.NewPrimaryTranslateStore = tf
return nil
}
}
// OptServerStatsClient is a functional option on Server
// used to specify the stats client.
func OptServerStatsClient(sc stats.StatsClient) ServerOption {
@ -286,11 +275,20 @@ func OptServerClusterHasher(h Hasher) ServerOption {
}
}
// OptServerTranslateFileMapSize is a functional option on Server
// used to specify the size of the translate file.
func OptServerTranslateFileMapSize(mapSize int) ServerOption {
// OptServerOpenTranslateStore is a functional option on Server
// used to specify the translation data store type.
func OptServerOpenTranslateStore(fn OpenTranslateStoreFunc) ServerOption {
return func(s *Server) error {
s.holder.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize))
s.holder.OpenTranslateStore = fn
return nil
}
}
// OptServerOpenTranslateReader is a functional option on Server
// used to specify the remote translation data reader.
func OptServerOpenTranslateReader(fn OpenTranslateReaderFunc) ServerOption {
return func(s *Server) error {
s.holder.OpenTranslateReader = fn
return nil
}
}
@ -331,7 +329,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
}
s.executor = newExecutor(executorOpts...)
s.holder.translateFile.logger = s.logger
// s.holder.translateFile.logger = s.logger
path, err := expandDirName(s.dataDir)
if err != nil {
@ -339,7 +337,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
}
s.holder.Path = path
s.holder.translateFile.Path = filepath.Join(path, ".keys")
// s.holder.translateFile.Path = filepath.Join(path, ".keys")
s.holder.Logger = s.logger
s.holder.Stats.SetLogger(s.logger)
@ -374,7 +372,6 @@ func NewServer(opts ...ServerOption) (*Server, error) {
s.executor.Holder = s.holder
s.executor.Node = node
s.executor.Cluster = s.cluster
s.executor.TranslateStore = s.holder.translateFile
s.executor.MaxWritesPerRequest = s.maxWritesPerRequest
s.cluster.broadcaster = s
s.cluster.maxWritesPerRequest = s.maxWritesPerRequest
@ -399,11 +396,6 @@ func (s *Server) UpAndDown() error {
log.Println(errors.Wrap(err, "logging startup"))
}
// Initialize id-key storage.
if err := s.holder.translateFile.Open(); err != nil {
return errors.Wrap(err, "opening TranslateFile")
}
// Open holder.
if err := s.holder.Open(); err != nil {
return errors.Wrap(err, "opening Holder")
@ -427,11 +419,6 @@ func (s *Server) Open() error {
log.Println(errors.Wrap(err, "logging startup"))
}
// Initialize id-key storage.
if err := s.holder.translateFile.Open(); err != nil {
return errors.Wrap(err, "opening TranslateFile")
}
// Open Cluster management.
if err := s.cluster.waitForStarted(); err != nil {
return errors.Wrap(err, "opening Cluster")

View file

@ -31,6 +31,7 @@ import (
"time"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/boltdb"
"github.com/pilosa/pilosa/v2/encoding/proto"
"github.com/pilosa/pilosa/v2/http"
"github.com/pilosa/pilosa/v2/server"
@ -1033,10 +1034,10 @@ func TestClusterTranslator(t *testing.T) {
if err != nil {
t.Fatalf("starting cluster 1: %v", err)
}
httpTranslateStore := http.NewTranslateStore(cluster[0].URL())
cluster[1] = test.NewCommandNode(false,
server.OptCommandServerOptions(
pilosa.OptServerPrimaryTranslateStore(httpTranslateStore),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.OpenTranslateReader),
),
)
cluster[1].Config.Gossip.Port = "0"
@ -1051,14 +1052,16 @@ func TestClusterTranslator(t *testing.T) {
test.MustDo("POST", cluster[0].URL()+"/index/i0/query", "Set(\"foo\", f0=\"bar\")")
// wait for key to replicate to second node
time.Sleep(500 * time.Millisecond)
result0 := test.MustDo("POST", cluster[0].URL()+"/index/i0/query", "Row(f0=\"bar\")").Body
result1 := test.MustDo("POST", cluster[1].URL()+"/index/i0/query", "Row(f0=\"bar\")").Body
if result0 != result1 {
t.Fatalf("`%s` != `%s`", result0, result1)
var result0, result1 string
if err := test.RetryUntil(2*time.Second, func() error {
result0 = test.MustDo("POST", cluster[0].URL()+"/index/i0/query", "Row(f0=\"bar\")").Body
result1 = test.MustDo("POST", cluster[1].URL()+"/index/i0/query", "Row(f0=\"bar\")").Body
if result0 != result1 {
return fmt.Errorf("`%s` != `%s`", result0, result1)
}
return nil
}); err != nil {
t.Fatal(err)
}
for _, i := range []string{result0, result1} {

View file

@ -307,6 +307,8 @@ func (m *Command) SetupServer() error {
pilosa.OptServerMetricInterval(time.Duration(m.Config.Metric.PollInterval)),
pilosa.OptServerDiagnosticsInterval(diagnosticsInterval),
pilosa.OptServerExecutorPoolSize(m.Config.WorkerPoolSize),
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.OpenTranslateReader),
pilosa.OptServerLogger(m.logger),
pilosa.OptServerAttrStoreFunc(boltdb.NewAttrStore),
pilosa.OptServerSystemInfo(gopsutil.NewSystemInfo()),
@ -314,20 +316,11 @@ func (m *Command) SetupServer() error {
pilosa.OptServerStatsClient(statsClient),
pilosa.OptServerURI(advertiseURI),
pilosa.OptServerInternalClient(http.NewInternalClientFromURI(uri, c)),
pilosa.OptServerPrimaryTranslateStoreFunc(http.NewTranslateStoreWithHTTPClient(c)),
pilosa.OptServerClusterDisabled(m.Config.Cluster.Disabled, m.Config.Cluster.Hosts),
pilosa.OptServerSerializer(proto.Serializer{}),
coordinatorOpt,
}
if m.Config.Translation.MapSize > 0 {
serverOptions = append(
serverOptions,
pilosa.OptServerTranslateFileMapSize(
m.Config.Translation.MapSize,
),
)
}
serverOptions = append(serverOptions, m.serverOptions...)
m.Server, err = pilosa.NewServer(serverOptions...)

View file

@ -442,3 +442,23 @@ type httpResponse struct {
*gohttp.Response
Body string
}
// RetryUntil repeatedly executes fn until it returns nil or timeout occurs.
func RetryUntil(timeout time.Duration, fn func() error) (err error) {
timer := time.NewTimer(timeout)
defer timer.Stop()
ticker := time.NewTicker(10 * time.Millisecond)
defer ticker.Stop()
for {
if err = fn(); err == nil {
return nil
}
select {
case <-timer.C:
return err
case <-ticker.C:
}
}
}

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

179
translator_test.go Normal file
View file

@ -0,0 +1,179 @@
// Copyright 2017 Pilosa Corp.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package pilosa_test
import (
"context"
"errors"
"io"
"testing"
"github.com/google/go-cmp/cmp"
"github.com/pilosa/pilosa/v2"
"github.com/pilosa/pilosa/v2/mock"
)
func TestInMemTranslateStore_TranslateKey(t *testing.T) {
s := pilosa.NewInMemTranslateStore("IDX", "FLD")
// Ensure initial key translates to ID 1.
if id, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(1); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
// Ensure next key autoincrements.
if id, err := s.TranslateKey("bar"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(2); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
// Ensure retranslating existing key returns original ID.
if id, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if got, want := id, uint64(1); got != want {
t.Fatalf("TranslateKey()=%d, want %d", got, want)
}
}
func TestInMemTranslateStore_TranslateID(t *testing.T) {
s := pilosa.NewInMemTranslateStore("IDX", "FLD")
// Setup initial keys.
if _, err := s.TranslateKey("foo"); err != nil {
t.Fatal(err)
} else if _, err := s.TranslateKey("bar"); err != nil {
t.Fatal(err)
}
// Ensure IDs can be translated back to keys.
if key, err := s.TranslateID(1); err != nil {
t.Fatal(err)
} else if got, want := key, "foo"; got != want {
t.Fatalf("TranslateID()=%s, want %s", got, want)
}
if key, err := s.TranslateID(2); err != nil {
t.Fatal(err)
} else if got, want := key, "bar"; got != want {
t.Fatalf("TranslateID()=%s, want %s", got, want)
}
}
func TestMultiTranslateEntryReader(t *testing.T) {
t.Run("None", func(t *testing.T) {
r := pilosa.NewMultiTranslateEntryReader(context.Background(), nil)
defer r.Close()
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != io.EOF {
t.Fatal(err)
}
})
t.Run("Single", func(t *testing.T) {
var r0 mock.TranslateEntryReader
r0.CloseFunc = func() error { return nil }
r0.ReadEntryFunc = func(entry *pilosa.TranslateEntry) error {
*entry = pilosa.TranslateEntry{Index: "i", Field: "f", ID: 1, Key: "foo"}
return nil
}
r := pilosa.NewMultiTranslateEntryReader(context.Background(), []pilosa.TranslateEntryReader{&r0})
defer r.Close()
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if diff := cmp.Diff(entry, pilosa.TranslateEntry{Index: "i", Field: "f", ID: 1, Key: "foo"}); diff != "" {
t.Fatal(diff)
}
if err := r.Close(); err != nil {
t.Fatal(err)
}
})
t.Run("Multi", func(t *testing.T) {
ready0, ready1 := make(chan struct{}), make(chan struct{})
var r0 mock.TranslateEntryReader
r0.CloseFunc = func() error { return nil }
r0.ReadEntryFunc = func(entry *pilosa.TranslateEntry) error {
if _, ok := <-ready0; !ok {
return io.EOF
}
*entry = pilosa.TranslateEntry{Index: "i0", Field: "f0", ID: 1, Key: "foo"}
return nil
}
var r1 mock.TranslateEntryReader
r1.CloseFunc = func() error { return nil }
r1.ReadEntryFunc = func(entry *pilosa.TranslateEntry) error {
if _, ok := <-ready1; !ok {
return io.EOF
}
*entry = pilosa.TranslateEntry{Index: "i1", Field: "f1", ID: 2, Key: "bar"}
return nil
}
r := pilosa.NewMultiTranslateEntryReader(context.Background(), []pilosa.TranslateEntryReader{&r1, &r0})
defer r.Close()
// Ensure r0 is read first
ready0 <- struct{}{}
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if diff := cmp.Diff(entry, pilosa.TranslateEntry{Index: "i0", Field: "f0", ID: 1, Key: "foo"}); diff != "" {
t.Fatal(diff)
}
// Unblock r1.
ready1 <- struct{}{}
// Read from r1 next.
if err := r.ReadEntry(&entry); err != nil {
t.Fatal(err)
} else if diff := cmp.Diff(entry, pilosa.TranslateEntry{Index: "i1", Field: "f1", ID: 2, Key: "bar"}); diff != "" {
t.Fatal(diff)
}
// Close both readers.
close(ready0)
close(ready1)
if err := r.Close(); err != nil {
t.Fatal(err)
}
})
t.Run("Error", func(t *testing.T) {
var r0 mock.TranslateEntryReader
r0.CloseFunc = func() error { return nil }
r0.ReadEntryFunc = func(entry *pilosa.TranslateEntry) error {
return errors.New("marker")
}
r := pilosa.NewMultiTranslateEntryReader(context.Background(), []pilosa.TranslateEntryReader{&r0})
defer r.Close()
var entry pilosa.TranslateEntry
if err := r.ReadEntry(&entry); err == nil || err.Error() != `marker` {
t.Fatalf("unexpected error: %s", err)
}
if err := r.Close(); err != nil {
t.Fatal(err)
}
})
}

7
uri.go
View file

@ -17,6 +17,8 @@ package pilosa
import (
"encoding/json"
"fmt"
"net"
"net/url"
"regexp"
"strconv"
"strings"
@ -47,6 +49,11 @@ type URI struct {
Port uint16 `json:"port"`
}
// URL returns a url.URL representation of the URI.
func (u *URI) URL() url.URL {
return url.URL{Scheme: u.Scheme, Host: net.JoinHostPort(u.Host, strconv.Itoa(int(u.Port)))}
}
// defaultURI creates and returns the default URI.
func defaultURI() *URI {
return &URI{