diff --git a/Makefile b/Makefile index 6fb679f91..67912d2f3 100644 --- a/Makefile +++ b/Makefile @@ -137,6 +137,13 @@ clustertests: vendor docker-compose -f $(DOCKER_COMPOSE) up --exit-code-from=client1 +DOCKER_COMPOSE_INDEX_KEY_REPLICATION=internal/clustertests/docker-compose-index-key-replication.yml +# Check clustertests target for more info +clustertests-index-key-replication: vendor + docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) down + docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) build client1 + docker-compose -f $(DOCKER_COMPOSE_INDEX_KEY_REPLICATION) up --exit-code-from=client1 + # Like clustertests, but rebuilds all images. clustertests-build: vendor docker-compose -f $(DOCKER_COMPOSE) down -v diff --git a/cluster.go b/cluster.go index e86883c1c..6419f3681 100644 --- a/cluster.go +++ b/cluster.go @@ -1165,6 +1165,9 @@ func (c *cluster) followResizeInstruction(ctx context.Context, instr *ResizeInst } } + // fire off translation sync + _ = c.translationSyncer.Reset() + return nil } @@ -1179,6 +1182,8 @@ func (c *cluster) resizeAbort() error { if c.resizeCancel != nil { c.resizeCancel() } + // fire off translation sync + _ = c.translationSyncer.Reset() return nil } diff --git a/holder.go b/holder.go index 0ec881add..0daf8416d 100644 --- a/holder.go +++ b/holder.go @@ -1472,7 +1472,10 @@ type holderSyncer struct { Cluster *cluster // Translation sync handling. - readers []TranslateEntryReader + readers []TranslateEntryReader + readersMu sync.Mutex + pendingReaders int + stopInitializeReplicationCh chan struct{} syncers errgroup.Group @@ -1596,6 +1599,23 @@ func (s *holderSyncer) syncFragment(index, field, view string, shard uint64) err // resetTranslationSync reinitializes streaming sync of translation data. func (s *holderSyncer) resetTranslationSync() error { + if s.stopInitializeReplicationCh == nil { + // suppose stopTranslationSync[S] holds the lock s.readersMu and tries + // to send a signal to s.stopInitializeRepliationCh. If + // s.stopInitializeReplicationCh is unbufferd and + // s.initializeReplication[I] is trying to acquire s.readersMu, this + // will result in a deadlock. + // Hence, the channel is buffered to prevent this scenario. Once + // [S] releases the lock, [I] will acquire it, however, the value + // of s.pendingReaders will be -1, at this point [I] should not + // attempt to add any more readers since those readers will be 'stale' + // [I] should also drain the channel to prevent the next invocation + // of [I] receiving a stop signal that was meant for the current one + // This is needlessly complicated and is the result of me running + // into various deadlocks while trying to fix handling of column + // key replication. + s.stopInitializeReplicationCh = make(chan struct{}) + } // Stop existing streams. if err := s.stopTranslationSync(); err != nil { return errors.Wrap(err, "stop translation sync") @@ -1655,7 +1675,15 @@ func newActiveTranslationSyncer(ch chan struct{}) *activeTranslationSyncer { // Reset resets the server's translation syncer. func (a *activeTranslationSyncer) Reset() error { - a.ch <- struct{}{} + // just in case some other part of the code has fired + // off a translation sync and it hasn't been received yet + // therefore we don't want to block on send since a.ch is + // (for now) unbuffered. One translationSync is as good as + // another + select { + case a.ch <- struct{}{}: + default: + } return nil } @@ -1665,6 +1693,17 @@ func (a *activeTranslationSyncer) Reset() error { // to complete. This should be called before reconnecting to the cluster in case // of a cluster resize or schema change. func (s *holderSyncer) stopTranslationSync() error { + s.readersMu.Lock() + defer func() { + s.readers = nil // will be populated by initializeReplication + s.readersMu.Unlock() + }() + // send signal to stop initializing more readers + if s.pendingReaders > 0 { + s.pendingReaders = -1 + close(s.stopInitializeReplicationCh) + s.stopInitializeReplicationCh = make(chan struct{}) + } var g errgroup.Group for i := range s.readers { rd := s.readers[i] @@ -1724,7 +1763,6 @@ func (s *holderSyncer) setTranslateReadOnlyFlags(snap *topology.ClusterSnapshot) // replicate these from whichever node is primary for that partition). func (s *holderSyncer) initializeReplication(snap *topology.ClusterSnapshot) error { nodeMaps := make(map[string]TranslateOffsetMap) - if snap.ReplicaN > 1 { if err := s.populateIndexReplication(nodeMaps, snap); err != nil { return err @@ -1734,26 +1772,99 @@ func (s *holderSyncer) initializeReplication(snap *topology.ClusterSnapshot) err return err } + // filter out empty nodes + nodes := make(map[*topology.Node]bool) for _, node := range snap.Nodes { m := nodeMaps[node.ID] - if m.Empty() { - continue + if !m.Empty() { + nodes[node] = true + } + } + + // connect to remote nodes and set up readers + readersCh := make(chan TranslateEntryReader, len(nodes)) + ctx, cancelAddingMoreReaders := context.WithCancel(context.Background()) + defer cancelAddingMoreReaders() + s.readersMu.Lock() + s.pendingReaders = len(nodes) + s.readersMu.Unlock() + go func() { + for { + for node := range nodes { + // check if ctx cancelled + // this means there was a signal sent to stop further init + // of readers + select { + case <-s.Closing: + return + case <-ctx.Done(): + close(readersCh) + return + default: + } + // connect to remote node + m := nodeMaps[node.ID] + rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m) + if err != nil { + continue + } + readersCh <- rd + delete(nodes, node) + } + if len(nodes) == 0 { + close(readersCh) + return + } + time.Sleep(10 * time.Second) + + } + }() + + for { + select { + case <-s.Closing: + cancelAddingMoreReaders() + return nil + case <-s.stopInitializeReplicationCh: + return nil + case rd, ok := <-readersCh: + // all translate readers have been launched, hence channel is + // closed + if !ok { + return nil + } + s.readersMu.Lock() + // [S] has been initiated and acquired the lock first + // at this point we should close the reader we've recieved rather + // than start replication on it. + // [S] should have already closed all the rest + // of the reads if they were still in action. + // we are also draining the channel since the signal for stopping + // further replication was meant for us + if s.pendingReaders == -1 { + rd.Close() + cancelAddingMoreReaders() + drain: + for { + select { + case <-s.stopInitializeReplicationCh: + default: + break drain + } + } + s.readersMu.Unlock() + return nil + } + s.pendingReaders-- + s.readers = append(s.readers, rd) + s.syncers.Go(func() error { + defer rd.Close() + s.readBothTranslateReader(rd, snap) + return nil + }) + s.readersMu.Unlock() } - - // Connect to remote node and begin streaming. - rd, err := s.Holder.OpenTranslateReader(context.Background(), node.URI.String(), m) - if err != nil { - return err - } - s.readers = append(s.readers, rd) - - s.syncers.Go(func() error { - defer rd.Close() - s.readBothTranslateReader(rd, snap) - return nil - }) } - return nil } // populateFieldReplication populates a map from node IDs to TranslateOffsetMaps diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index 892e66383..faad15337 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -21,7 +21,7 @@ import ( "testing" "time" - "github.com/molecula/featurebase/v2" + pilosa "github.com/molecula/featurebase/v2" "github.com/molecula/featurebase/v2/disco" picli "github.com/molecula/featurebase/v2/http" ) diff --git a/internal/clustertests/docker-compose-index-key-replication.yml b/internal/clustertests/docker-compose-index-key-replication.yml new file mode 100644 index 000000000..6f84d5947 --- /dev/null +++ b/internal/clustertests/docker-compose-index-key-replication.yml @@ -0,0 +1,73 @@ +version: '2' +services: + pilosa1: + build: + context: ../.. + dockerfile: Dockerfile-clustertests + image: ptest + ports: + - "33455:10101" + environment: + - PILOSA_NAME=pilosa1 + - PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201 + - PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa1:10201 + - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 + - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa1:10301 + - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 + - PILOSA_CLUSTER_REPLICAS=3 + networks: + - pilosanet + command: + - "/featurebase server --bind pilosa1:10101" + pilosa2: + build: + context: ../.. + dockerfile: Dockerfile-clustertests + image: ptest + ports: + - "33456:10101" + environment: + - PILOSA_NAME=pilosa2 + - PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201 + - PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa2:10201 + - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 + - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa2:10301 + - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 + - PILOSA_CLUSTER_REPLICAS=3 + networks: + - pilosanet + command: + - "/featurebase server --bind pilosa2:10101" + pilosa3: + build: + context: ../.. + dockerfile: Dockerfile-clustertests + image: ptest + ports: + - "33457:10101" + environment: + - PILOSA_NAME=pilosa3 + - PILOSA_ETCD_LISTEN_CLIENT_ADDRESS=http://0.0.0.0:10201 + - PILOSA_ETCD_ADVERTISE_CLIENT_ADDRESS=http://pilosa3:10201 + - PILOSA_ETCD_LISTEN_PEER_ADDRESS=http://0.0.0.0:10301 + - PILOSA_ETCD_ADVERTISE_PEER_ADDRESS=http://pilosa3:10301 + - PILOSA_ETCD_INITIAL_CLUSTER=pilosa1=http://pilosa1:10301,pilosa2=http://pilosa2:10301,pilosa3=http://pilosa3:10301 + - PILOSA_CLUSTER_REPLICAS=3 + networks: + - pilosanet + command: + - "/featurebase server --bind pilosa3:10101" + client1: + build: + context: . + environment: + - ENABLE_PILOSA_CLUSTER_TESTS_FOR_INDEX_KEY_REPLICATION=1 + - GO111MODULE=on + networks: + - pilosanet + volumes: + - /var/run/docker.sock:/var/run/docker.sock + command: + - "cd /go/src/github.com/molecula/featurebase/ && go test -mod=vendor -v -run=IndexKey -count=1 github.com/molecula/featurebase/v2/internal/clustertests" +networks: + pilosanet: diff --git a/internal/clustertests/index_key_replication_test.go b/internal/clustertests/index_key_replication_test.go new file mode 100644 index 000000000..54c3dfa63 --- /dev/null +++ b/internal/clustertests/index_key_replication_test.go @@ -0,0 +1,375 @@ +// 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 clustertest + +import ( + "context" + "crypto/tls" + "fmt" + "math/rand" + "os" + "os/exec" + "sync" + "testing" + "time" + + pilosa "github.com/molecula/featurebase/v2" + "github.com/molecula/featurebase/v2/disco" + "github.com/molecula/featurebase/v2/http" + picli "github.com/molecula/featurebase/v2/http" + "github.com/molecula/featurebase/v2/net" + "github.com/molecula/featurebase/v2/topology" +) + +// index -> key -> ids from all replicas +type translationRes map[string]map[string][]uint64 + +func defaultTranslationResults(indexes []string) translationRes { + res := make(translationRes) + for _, index := range indexes { + res[index] = make(map[string][]uint64) + } + return res +} + +func verify(allRes translationRes, indexes []string, replicasN, count int) error { + received := 0 + allIDsSame := func(ids []uint64) bool { + if len(ids) <= 1 { + // trivially true + return true + } + id := ids[0] + for _, other := range ids { + if id != other { + return false + } + } + return true + } + for _, index := range indexes { + res, ok := allRes[index] + if !ok { + return fmt.Errorf("Expected index '%s' but not present", index) + } + for key, ids := range res { + // first id is after successful translation, rest are from + // translate readers + replicationCount := len(ids) - 1 + received += replicationCount + if replicationCount != replicasN { + return fmt.Errorf("Count for ids for key '%s'(%d), index '%s' not equal than replicasN(%d)", key, replicationCount, index, replicasN) + } + if !allIDsSame(ids) { + return fmt.Errorf("Expected all ids for key '%s', index '%s' to be the same across the cluster %+v", key, index, ids) + } + } + } + + if received != count { + return fmt.Errorf("Expected %d count of keys, received %d", count, received) + } + + return nil +} +func getURIsFromAddresses(addrs []string) ([]*net.URI, error) { + uris := make([]*net.URI, 0, len(addrs)) + for _, addr := range addrs { + uri, err := net.NewURIFromAddress(addr) + if err != nil { + return nil, err + } + uris = append(uris, uri) + } + return uris, nil +} + +func getClients(addrs []string) ([]*http.InternalClient, error) { + clients := make([]*http.InternalClient, 0, len(addrs)) + for _, addr := range addrs { + c, err := picli.NewInternalClient(addr, picli.GetHTTPClient(nil)) + if err != nil { + return nil, err + } + clients = append(clients, c) + } + return clients, nil +} + +func genIndexNames(indexCount int) []string { + indexNames := make([]string, 0, indexCount) + for i := 1; i <= indexCount; i++ { + indexNames = append(indexNames, fmt.Sprintf("idx-%d", i)) + } + return indexNames +} + +func parseDuration(t *testing.T, s string) time.Duration { + parsed, err := time.ParseDuration(s) + if err != nil { + t.Fatal(err) + } + return parsed +} + +func durationToSeconds(d time.Duration) string { + s := int(d.Seconds()) + return fmt.Sprintf("%ds", s) +} + +func TestIndexKeyReplication(t *testing.T) { + if os.Getenv("ENABLE_PILOSA_CLUSTER_TESTS_FOR_INDEX_KEY_REPLICATION") != "1" { + t.Skip("pilosa cluster tests for index key replication are not enabled") + } + // configurations for test + replicasN := 3 + indexCount := 4 + intervalDurationArg := "100ms" + totalInsertionDurationArg := "10s" + coolOffDurationArg := "5s" + numKeysToInsertPerDuration := 100 + addresses := []string{"pilosa1:10101", "pilosa2:10101", "pilosa3:10101"} + + intervalDuration := parseDuration(t, intervalDurationArg) + totalInsertionDuration := parseDuration(t, totalInsertionDurationArg) + coolOffDuration := parseDuration(t, coolOffDurationArg) + + indexes := genIndexNames(indexCount) + clients, err := getClients(addresses) + cli := clients[0] + if err != nil { + t.Fatalf("on init clients from addresses: %v, %v", addresses, err) + } + uris, err := getURIsFromAddresses(addresses) + if err != nil { + t.Fatalf("on init clients from addresses: %v, %v", addresses, err) + } + ctx := context.Background() + + // create index keyed + for _, index := range indexes { + err = cli.EnsureIndex(ctx, index, pilosa.IndexOptions{ + Keys: true, + }) + if err != nil { + t.Fatalf("creating/asserting index: %v", err) + } + } + + // set up index keys to insert + keysInserted := 0 + allRes := defaultTranslationResults(indexes) + { + ctxForIndexKeyCreation, cancelFurtherInsertions := context.WithCancel(ctx) + var wg sync.WaitGroup + var mu = &sync.Mutex{} + wg.Add(len(indexes)) + seed := time.Now().UnixNano() + t.Logf("start inserting index keys for %v, seed(%v)", totalInsertionDuration, seed) + rng := rand.New(rand.NewSource(seed)) + for _, index := range indexes { + translations := allRes[index] + go func(index string, translations map[string][]uint64) { + keys := make([]string, numKeysToInsertPerDuration) + offset := 1 + ticker := time.NewTicker(intervalDuration) + defer func() { + ticker.Stop() + wg.Done() + }() + for { + select { + case <-ticker.C: + for i := 0; i < numKeysToInsertPerDuration; i++ { + keys[i] = fmt.Sprintf("key-%d", i+offset) + } + // pick random node to send key creation to + r := rng.Intn(len(addresses)) + cli, uri := clients[r], uris[r] + + // insert index keys + transmap, err := cli.CreateIndexKeysNode(ctxForIndexKeyCreation, uri, index, keys...) + if err != nil { + if err == context.Canceled { + return + } else { + t.Logf("creating index keys for index(%s) send to node(%s): %v", index, uri.String(), err) + continue + } + } + for key, id := range transmap { + translations[key] = append(translations[key], id) + } + mu.Lock() + keysInserted += len(transmap) + mu.Unlock() + offset += numKeysToInsertPerDuration + case <-ctxForIndexKeyCreation.Done(): + return + } + } + }(index, translations) + } + // inject fault + pcmd := exec.Command("/pumba", "netem", "--duration", + durationToSeconds(totalInsertionDuration), + "loss", "--percent", "50", "--correlation", "60", + "clustertests_pilosa3_1") + pcmd.Stdout = os.Stdout + pcmd.Stderr = os.Stderr + t.Logf("sending pumba fault injection cmd: %v", pcmd.String()) + err = pcmd.Start() + if err != nil { + t.Fatalf("starting pumba command: %v", err) + } + err = pcmd.Wait() + if err != nil { + t.Fatalf("waiting on pumba pause cmd: %v", err) + } + + // wait for index keys to be created + t.Logf("start wait to complete index key creation") + time.Sleep(totalInsertionDuration) + cancelFurtherInsertions() + wg.Wait() + t.Logf("done with inserting index keys. Total keys inserted: %d", keysInserted) + } + + t.Logf("start cool off period: %v\n", coolOffDuration) + time.Sleep(coolOffDuration) + t.Log("done with cool off period, waiting for stability") + waitForStatus(t, clients[0].Status, string(disco.ClusterStateNormal), 30, 1*time.Second) + t.Log("done with waiting for stability, starting verifying persistence of index keys") + + // get all nodes + nodes, err := cli.Nodes(ctx) + if err != nil { + t.Fatal(err) + } + + // prepare translate offset maps + nodeMaps := make(map[string]pilosa.TranslateOffsetMap) + for _, n := range nodes { + nodeMaps[n.ID] = make(pilosa.TranslateOffsetMap) + } + schema, err := cli.Schema(ctx) + if err != nil { + t.Fatal(err) + } + for _, indexInfo := range schema { + index := indexInfo.Name + isKeyed := indexInfo.Options.Keys + if !isKeyed { + continue + } + + partitionN := topology.DefaultPartitionN + for partition := 0; partition < partitionN; partition++ { + nodes, err := cli.PartitionNodes(ctx, partition) + if err != nil { + t.Fatal(err) + } + for _, n := range nodes { + m := nodeMaps[n.ID] + m.SetIndexPartitionOffset(index, partition, 0) + } + } + } + + // open translate reader + readers := make([]pilosa.TranslateEntryReader, 0, len(nodes)) + closeAllTranslateReaders := func() []error { + var errs []error + for _, tr := range readers { + err := tr.Close() + if err != nil { + errs = append(errs, err) + } + } + return errs + } + var tlsConfig *tls.Config = nil + var wg sync.WaitGroup + entries := make(chan *pilosa.TranslateEntry) + for _, n := range nodes { + client := http.GetHTTPClient(tlsConfig) + nodeURL := n.URI.String() + offsets := nodeMaps[n.ID] + openTranslateReader := http.GetOpenTranslateReaderFunc(client) + tr, err := openTranslateReader(ctx, nodeURL, offsets) + if err != nil { + t.Fatal(err) + } + wg.Add(1) + readers = append(readers, tr) + go func(node *topology.Node, tr pilosa.TranslateEntryReader) { + defer func() { + wg.Done() + }() + for { + var entry pilosa.TranslateEntry + err := tr.ReadEntry(&entry) + if err != nil { + // TODO also ignore transient http: read on clsoed response body errors + // for now just print error + if err != context.Canceled { + fmt.Printf("node(%s) On read from translate entry reader: %v", node.URI.String(), err) + } + return + } + // ignore field keys + if entry.Field != "" { + continue + } + entries <- &entry + } + }(n, tr) + } + + // receive all translation entries made + count := replicasN * keysInserted + i := 0 + for { + entry := <-entries + res, ok := allRes[entry.Index] + if !ok { + // ignore indexes we did not create for this test + continue + } + ids, ok := res[entry.Key] + if !ok { + // ignore keys we did not insert for the test + continue + } + res[entry.Key] = append(ids, entry.ID) + i++ + if i == count { + break + } + } + errs := closeAllTranslateReaders() + wg.Wait() + for _, err := range errs { + if err != nil { + t.Errorf("on close translate readers: %v", err) + } + } + close(entries) + + err = verify(allRes, indexes, replicasN, count) + if err != nil { + t.Fatal(err) + } +} diff --git a/server.go b/server.go index c9b4d8d5e..93d9947a0 100644 --- a/server.go +++ b/server.go @@ -437,7 +437,7 @@ func NewServer(opts ...ServerOption) (*Server, error) { confirmDownRetries: defaultConfirmDownRetries, confirmDownSleep: defaultConfirmDownSleep, - resetTranslationSyncCh: make(chan struct{}), + resetTranslationSyncCh: make(chan struct{}, 1), logger: logger.NopLogger, } @@ -569,11 +569,6 @@ func (s *Server) Open() error { log.Println(errors.Wrap(err, "logging startup")) } - // Start background process listening for translation - // sync resets. - s.wg.Add(1) - go func() { defer s.wg.Done(); s.monitorResetTranslationSync() }() - // Start DisCo. ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second) defer cancel() @@ -612,6 +607,12 @@ func (s *Server) Open() error { s.syncer.Closing = s.closing s.syncer.Stats = s.holder.Stats.WithTags("component:HolderSyncer") + // Start background process listening for translation + // sync resets. + s.wg.Add(1) + go func() { defer s.wg.Done(); s.monitorResetTranslationSync() }() + go func() { _ = s.translationSyncer.Reset() }() + // Open holder. func() { s.holder.startMsgsMu.Lock()