Merge pull request #110 from travisturner/translatestore-fixes

WIP: Thread OpenTranslateStore through Holder to Index
This commit is contained in:
Travis Turner 2020-02-12 11:56:49 -06:00 • committed by GitHub
commit 4713ccd0c8
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
7 changed files with 159 additions and 79 deletions

View file

@ -144,7 +144,7 @@ func (s *TranslateStore) Size() int64 {
return tx.Size()
}
// TranslateKeys converts a string key to an integer ID.
// TranslateKey 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.
@ -188,8 +188,8 @@ func (s *TranslateStore) TranslateKey(key string) (id uint64, _ error) {
return id, nil
}
// TranslateKeys converts a string key to an integer ID.
// If key does not have an associated id then one is created.
// TranslateKeys converts a slice of string keys to a slice of integer IDs.
// If a 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

View file

@ -138,6 +138,8 @@ func NewHolder(partitionN int) *Holder {
cacheFlushInterval: defaultCacheFlushInterval,
OpenTranslateStore: OpenInMemTranslateStore,
Logger: logger.NopLogger,
}
}
@ -496,6 +498,7 @@ 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.OpenTranslateStore = h.OpenTranslateStore
index.holder = h
return index, nil
}
@ -1011,7 +1014,7 @@ func (s *holderSyncer) initializeIndexTranslateReplication() error {
continue
}
// Connect to remote not and begin streaming.
// Connect to remote node and begin streaming.
rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m)
if err != nil {
return err

View file

@ -107,7 +107,7 @@ func (i *Index) Path() string { return i.path }
// TranslateStorePath returns the translation database path for a partition.
func (i *Index) TranslateStorePath(partitionID int) string {
return filepath.Join(i.path, "keys", strconv.Itoa(partitionID))
return filepath.Join(i.path, translateStoreDir, strconv.Itoa(partitionID))
}
// TranslateStore returns the translation store for a given partition.
@ -165,12 +165,26 @@ func (i *Index) Open() (err error) {
}
i.logger.Debugf("open translate store for index: %s", i.name)
var g errgroup.Group
var mu sync.Mutex
for partitionID := 0; partitionID < i.partitionN; partitionID++ {
store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.partitionN)
if err != nil {
return errors.Wrap(err, "opening index translate store")
}
i.translateStores[partitionID] = store
partitionID := partitionID
g.Go(func() error {
store, err := i.OpenTranslateStore(i.TranslateStorePath(partitionID), i.name, "", partitionID, i.partitionN)
if err != nil {
return errors.Wrapf(err, "opening index translate store: partition=%d", partitionID)
}
mu.Lock()
defer mu.Unlock()
i.translateStores[partitionID] = store
return nil
})
}
if err := g.Wait(); err != nil {
return err
}
return nil

View file

@ -1031,13 +1031,34 @@ func TestHandler_Endpoints(t *testing.T) {
})
}
func TestClusterTranslator(t *testing.T) {
cluster := make(test.Cluster, 2)
cluster[0] = test.NewCommandNode(true)
func TestCluster_TranslateStore(t *testing.T) {
cluster := make(test.Cluster, 1)
cluster[0] = test.NewCommandNode(true,
server.OptCommandServerOptions(
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
pilosa.OptServerOpenTranslateReader(http.GetOpenTranslateReaderFunc(nil)),
),
)
cluster[0].Config.Gossip.Port = "0"
err := cluster[0].Start()
if err != nil {
t.Fatalf("starting cluster 1: %v", err)
t.Fatalf("starting cluster 0: %v", err)
}
test.MustDo("POST", cluster[0].URL()+"/index/i0", "{\"options\": {\"keys\": true}}")
}
func TestClusterTranslator(t *testing.T) {
cluster := make(test.Cluster, 2)
cluster[0] = test.NewCommandNode(true,
server.OptCommandServerOptions(
pilosa.OptServerOpenTranslateStore(boltdb.OpenTranslateStore),
),
)
cluster[0].Config.Gossip.Port = "0"
err := cluster[0].Start()
if err != nil {
t.Fatalf("starting cluster 0: %v", err)
}
cluster[1] = test.NewCommandNode(false,
server.OptCommandServerOptions(

View file

@ -27,7 +27,6 @@ import (
"strconv"
"strings"
"testing"
"testing/quick"
"time"
"github.com/pelletier/go-toml"
@ -51,78 +50,77 @@ func TestMain_Set_Quick(t *testing.T) {
t.Skip("short")
}
if err := quick.Check(func(cmds []SetCommand) bool {
m := test.MustRunCommand()
defer m.Close()
for i := 0; i < 100; i++ {
t.Run(fmt.Sprint(i), func(t *testing.T) {
t.Parallel()
// Create client.
client, err := http.NewInternalClient(m.API.Node().URI.HostPort(), http.GetHTTPClient(nil))
if err != nil {
t.Fatal(err)
}
rand := rand.New(rand.NewSource(int64(i)))
cmds := GenerateSetCommands(1000, rand)
// Execute Set() commands.
for _, cmd := range cmds {
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
m := test.MustRunCommand()
defer m.Close()
// Create client.
client, err := http.NewInternalClient(m.API.Node().URI.HostPort(), http.GetHTTPClient(nil))
if err != nil {
t.Fatal(err)
}
if err := client.CreateField(context.Background(), "i", cmd.Field); err != nil && err != pilosa.ErrFieldExists {
t.Fatal(err)
}
if _, err := m.Query("i", "", fmt.Sprintf(`Set(%d, %s=%d)`, cmd.ColumnID, cmd.Field, cmd.ID)); err != nil {
t.Fatal(err)
}
}
// Validate data.
for field, fieldSet := range SetCommands(cmds).Fields() {
for id, columnIDs := range fieldSet {
exp := MustMarshalJSON(map[string]interface{}{
"results": []interface{}{
map[string]interface{}{
"columns": columnIDs,
"attrs": map[string]interface{}{},
},
},
}) + "\n"
if res, err := m.Query("i", "", fmt.Sprintf(`Row(%s=%d)`, field, id)); err != nil {
// Execute Set() commands.
for _, cmd := range cmds {
if err := client.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}); err != nil && err != pilosa.ErrIndexExists {
t.Fatal(err)
}
if err := client.CreateField(context.Background(), "i", cmd.Field); err != nil && err != pilosa.ErrFieldExists {
t.Fatal(err)
}
if _, err := m.Query("i", "", fmt.Sprintf(`Set(%d, %s=%d)`, cmd.ColumnID, cmd.Field, cmd.ID)); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result:\n\ngot=%s\n\nexp=%s\n\n", res, exp)
}
}
}
if err := m.Reopen(); err != nil {
t.Fatal(err)
}
// Validate data after reopening.
for field, fieldSet := range SetCommands(cmds).Fields() {
for id, columnIDs := range fieldSet {
exp := MustMarshalJSON(map[string]interface{}{
"results": []interface{}{
map[string]interface{}{
"columns": columnIDs,
"attrs": map[string]interface{}{},
// Validate data.
for field, fieldSet := range SetCommands(cmds).Fields() {
for id, columnIDs := range fieldSet {
exp := MustMarshalJSON(map[string]interface{}{
"results": []interface{}{
map[string]interface{}{
"columns": columnIDs,
"attrs": map[string]interface{}{},
},
},
},
}) + "\n"
if res, err := m.Query("i", "", fmt.Sprintf(`Row(%s=%d)`, field, id)); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result (reopen):\n\ngot=%s\n\nexp=%s\n\n", res, exp)
}) + "\n"
if res, err := m.Query("i", "", fmt.Sprintf(`Row(%s=%d)`, field, id)); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result:\n\ngot=%s\n\nexp=%s\n\n", res, exp)
}
}
}
}
return true
}, &quick.Config{
Values: func(values []reflect.Value, rand *rand.Rand) {
values[0] = reflect.ValueOf(GenerateSetCommands(1000, rand))
},
}); err != nil {
t.Fatal(err)
if err := m.Reopen(); err != nil {
t.Fatal(err)
}
// Validate data after reopening.
for field, fieldSet := range SetCommands(cmds).Fields() {
for id, columnIDs := range fieldSet {
exp := MustMarshalJSON(map[string]interface{}{
"results": []interface{}{
map[string]interface{}{
"columns": columnIDs,
"attrs": map[string]interface{}{},
},
},
}) + "\n"
if res, err := m.Query("i", "", fmt.Sprintf(`Row(%s=%d)`, field, id)); err != nil {
t.Fatal(err)
} else if res != exp {
t.Fatalf("unexpected result (reopen):\n\ngot=%s\n\nexp=%s\n\n", res, exp)
}
}
}
})
}
}
@ -497,7 +495,7 @@ func TestClusteringNodesReplica1(t *testing.T) {
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[2].Command.GossipTransport().URI.Port))
cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr)
cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cluster[2].Command.Config = config
// Run new program.
@ -571,7 +569,7 @@ func TestClusteringNodesReplica2(t *testing.T) {
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[2].Command.GossipTransport().URI.Port))
cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr)
cluster[2].Command = server.NewCommand(cluster[2].Stdin, cluster[2].Stdout, cluster[2].Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cluster[2].Command.Config = config
// Run new program.
@ -591,7 +589,7 @@ func TestClusteringNodesReplica2(t *testing.T) {
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[1].Command.GossipTransport().URI.Port))
cluster[1].Command = server.NewCommand(cluster[1].Stdin, cluster[1].Stdout, cluster[1].Stderr)
cluster[1].Command = server.NewCommand(cluster[1].Stdin, cluster[1].Stdout, cluster[1].Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cluster[1].Command.Config = config
// Run new program.
@ -864,7 +862,7 @@ func TestClusterQueriesAfterRestart(t *testing.T) {
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cmd1.Command.GossipTransport().URI.Port))
cmd1.Command = server.NewCommand(cmd1.Stdin, cmd1.Stdout, cmd1.Stderr)
cmd1.Command = server.NewCommand(cmd1.Stdin, cmd1.Stdout, cmd1.Stderr, server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
cmd1.Command.Config = config
err = cmd1.Start()
if err != nil {
@ -885,6 +883,7 @@ func TestClusterQueriesAfterRestart(t *testing.T) {
if results.Results[0].(uint64) != 100 {
t.Fatalf("Count should be 100, but got %v of type %[1]T", results.Results[0])
}
}
// TODO: confirm that things keep working if a node is hard-closed (no nodeLeave event) and immediately restarted with a different address.

View file

@ -89,6 +89,10 @@ func newCommand(opts ...server.CommandOption) *Command {
// NewCommandNode returns a new instance of Command with clustering enabled.
func NewCommandNode(isCoordinator bool, opts ...server.CommandOption) *Command {
// We want tests to default to using the in-memory translate store, so we
// prepend opts with that functional option. If a different translate store
// has been specified, it will override this one.
opts = prependWithMemStore(opts)
m := newCommand(opts...)
m.Config.Cluster.Disabled = false
m.Config.Cluster.Coordinator = isCoordinator
@ -97,7 +101,7 @@ func NewCommandNode(isCoordinator bool, opts ...server.CommandOption) *Command {
// MustRunCommand returns a new, running Main. Panic on error.
func MustRunCommand() *Command {
m := newCommand()
m := newCommand(server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)))
m.Config.Metric.Diagnostics = false // Disable diagnostics.
m.Config.Gossip.Port = "0"
if err := m.Start(); err != nil {
@ -398,6 +402,11 @@ func runCluster(size int, opts ...[]server.CommandOption) (Cluster, error) {
// MustRunCluster creates and starts a new cluster
func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Cluster {
// We want tests to default to using the in-memory translate store, so we
// prepend opts with that functional option. If a different translate store
// has been specified, it will override this one.
opts = prependOpts(opts)
tb.Helper()
c, err := runCluster(size, opts...)
if err != nil {
@ -406,6 +415,29 @@ func MustRunCluster(tb testing.TB, size int, opts ...[]server.CommandOption) Clu
return c
}
// prependOpts applies prependWithMemStore to each of the ops (one per
// node, or one for the entire cluser).
func prependOpts(opts [][]server.CommandOption) [][]server.CommandOption {
if len(opts) == 0 {
opts = [][]server.CommandOption{
prependWithMemStore([]server.CommandOption{}),
}
} else {
for i := range opts {
opts[i] = prependWithMemStore(opts[i])
}
}
return opts
}
// prependWithMemStore prepends opts with the OpenInMemTranslateStore.
func prependWithMemStore(opts []server.CommandOption) []server.CommandOption {
defaultOpts := []server.CommandOption{
server.OptCommandServerOptions(pilosa.OptServerOpenTranslateStore(pilosa.OpenInMemTranslateStore)),
}
return append(defaultOpts, opts...)
}
////////////////////////////////////////////////////////////////////////////////////
// MustDo executes http.Do() with an http.NewRequest(). Panic on error.

View file

@ -22,6 +22,12 @@ import (
"github.com/pkg/errors"
)
const (
// translateStoreDir is the subdirctory into which the partitioned
// translate store data is stored.
translateStoreDir = "_keys"
)
// Translate store errors.
var (
ErrTranslateStoreClosed = errors.New("translate store closed")
@ -70,6 +76,11 @@ type OpenTranslateStoreFunc func(path, index, field string, partitionID, partiti
// GenerateNextPartitionedID returns the next ID within the same partition.
func GenerateNextPartitionedID(index string, prev uint64, partitionID, partitionN int) uint64 {
// If the translation store is not partitioned, just return
// the next ID.
if partitionID == -1 {
return prev + 1
}
// Try to use the next ID if it is in the same partition.
// Otherwise find ID in next shard that has a matching partition.
for id := prev + 1; ; id += ShardWidth {