mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-05 16:15:56 +00:00
A few things were going wrong here. First, we take a "RetryPeriod" option on backup and restore which is meant to be roughly the total amount of time we spend retrying any given request before failing. However we were incorrectly passing that as the RetryMaxWait which is the maximum amount of time to sleep between any two attempts. We now do some fuzzy math to figure out approximately how many attempts we should make given a minimum sleep of 100ms and the fact that we double the sleep time every attempt. Second, during the backup test, if a host was totally stopped when we started the request, it would fail immediately and then retry, but if the host was stopped during the request (after DNS had resolved), then the request would wait for the DialTimeout which we default to 30s, so turning off the cluster for 5 seconds and turning it back on resulted in the backup completing rather than failing. Because of this, we change the commandClient to have a default dial timeout of 1 second. I was tempted to change the global default to 1s which I think would be fine, but didn't want to break anything too badly.
442 lines
11 KiB
Go
442 lines
11 KiB
Go
// Copyright 2021 Molecula Corp. All rights reserved.
|
|
package ctl
|
|
|
|
import (
|
|
"context"
|
|
"crypto/tls"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"math"
|
|
"net/http"
|
|
"os"
|
|
"path/filepath"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/hashicorp/go-retryablehttp"
|
|
|
|
pilosa "github.com/molecula/featurebase/v2"
|
|
fb_http "github.com/molecula/featurebase/v2/http"
|
|
"github.com/molecula/featurebase/v2/server"
|
|
"github.com/molecula/featurebase/v2/topology"
|
|
"github.com/pkg/errors"
|
|
"golang.org/x/sync/errgroup"
|
|
)
|
|
|
|
// RestoreCommand represents a command for restoring a backup to
|
|
type RestoreCommand struct {
|
|
tlsConfig *tls.Config
|
|
Host string
|
|
|
|
Concurrency int
|
|
|
|
// Filepath to the backup file.
|
|
Path string
|
|
|
|
// Amount of time after first failed request to continue retrying.
|
|
RetryPeriod time.Duration `json:"retry-period"`
|
|
|
|
// Host:port on which to listen for pprof.
|
|
Pprof string `json:"pprof"`
|
|
|
|
// Reusable client.
|
|
client pilosa.InternalClient
|
|
|
|
// Standard input/output
|
|
*pilosa.CmdIO
|
|
TLS server.TLSConfig
|
|
}
|
|
|
|
// NewRestoreCommand returns a new instance of RestoreCommand.
|
|
func NewRestoreCommand(stdin io.Reader, stdout, stderr io.Writer) *RestoreCommand {
|
|
return &RestoreCommand{
|
|
CmdIO: pilosa.NewCmdIO(stdin, stdout, stderr),
|
|
RetryPeriod: time.Second * 30,
|
|
Concurrency: 1,
|
|
Pprof: "localhost:43809",
|
|
}
|
|
}
|
|
|
|
// Run executes the restore.
|
|
func (cmd *RestoreCommand) Run(ctx context.Context) (err error) {
|
|
logger := cmd.Logger()
|
|
close, err := startProfilingServer(cmd.Pprof, logger)
|
|
if err != nil {
|
|
return errors.Wrap(err, "starting profiling server")
|
|
}
|
|
defer close()
|
|
|
|
// Validate arguments.
|
|
if cmd.Path == "" {
|
|
return fmt.Errorf("-s flag required")
|
|
} else if cmd.Concurrency <= 0 {
|
|
return fmt.Errorf("concurrency must be at least one")
|
|
}
|
|
|
|
// Parse TLS configuration for node-specific clients.
|
|
tls := cmd.TLSConfiguration()
|
|
if cmd.tlsConfig, err = server.GetTLSConfig(&tls, logger); err != nil {
|
|
return fmt.Errorf("parsing tls config: %w", err)
|
|
}
|
|
// Create a client to the server.
|
|
client, err := commandClient(cmd, fb_http.WithClientRetryPeriod(cmd.RetryPeriod))
|
|
if err != nil {
|
|
return fmt.Errorf("creating client: %w", err)
|
|
}
|
|
cmd.client = client
|
|
|
|
nodes, err := cmd.client.Nodes(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
var primary *topology.Node
|
|
for _, node := range nodes {
|
|
if node.IsPrimary {
|
|
primary = node
|
|
break
|
|
}
|
|
|
|
}
|
|
if primary == nil {
|
|
return errors.New("no primary")
|
|
}
|
|
|
|
if err := cmd.restoreSchema(ctx, primary); err != nil {
|
|
return fmt.Errorf("cannot restore schema: %w", err)
|
|
} else if err := cmd.restoreIDAlloc(ctx, primary); err != nil {
|
|
return fmt.Errorf("cannot restore idalloc: %w", err)
|
|
}
|
|
if err := cmd.restoreShards(ctx); err != nil {
|
|
return fmt.Errorf("cannot restore shards: %w", err)
|
|
} else if err := cmd.restoreIndexTranslation(ctx); err != nil {
|
|
return fmt.Errorf("cannot restore index translation: %w", err)
|
|
} else if err := cmd.restoreFieldTranslation(ctx, nodes); err != nil {
|
|
return fmt.Errorf("cannot restore field translation: %w", err)
|
|
}
|
|
|
|
/* Fetch the cluster nodes from the target host.
|
|
For each index:
|
|
Upload the RBF snapshot for each shard to the nodes that own the shard.
|
|
Upload the index & field translation BoltDB snapshots to each node.
|
|
If possible, trigger the node to reload itself. Otherwise a restart would be required.
|
|
*/
|
|
|
|
return nil
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreSchema(ctx context.Context, primary *topology.Node) error {
|
|
f, err := os.Open(filepath.Join(cmd.Path, "schema"))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer f.Close()
|
|
|
|
existingSchema, err := cmd.client.Schema(ctx)
|
|
if len(existingSchema) == 0 {
|
|
cmd.Logger().Printf("Load Schema")
|
|
url := primary.URI.Path("/schema")
|
|
client := cmd.newClient()
|
|
_, err = client.Post(url, "application/json", f)
|
|
} else {
|
|
schema := &pilosa.Schema{}
|
|
if err := json.NewDecoder(f).Decode(schema); err != nil {
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
exists := func(indexName string) bool {
|
|
for _, i := range existingSchema {
|
|
if i.Name == indexName {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
logger := cmd.Logger()
|
|
//NOTE SHOULD ONLY BE ONE
|
|
for _, index := range schema.Indexes {
|
|
if exists(index.Name) {
|
|
return fmt.Errorf("index Exists %v", index.Name)
|
|
}
|
|
logger.Printf("Create INDEX %v", index.Name)
|
|
err = cmd.client.CreateIndex(ctx, index.Name, index.Options)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, field := range index.Fields {
|
|
logger.Printf("Create Field %v", field.Name)
|
|
err = cmd.client.CreateFieldWithOptions(ctx, index.Name, field.Name, field.Options)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
}
|
|
return err
|
|
}
|
|
|
|
func retryWith400(ctx context.Context, resp *http.Response, err error) (bool, error) {
|
|
if resp != nil && resp.StatusCode >= 400 { // we have some dumb status codes
|
|
return true, nil
|
|
}
|
|
return retryablehttp.DefaultRetryPolicy(ctx, resp, err)
|
|
}
|
|
|
|
// This logic is taken from featurebase/http/client.go If this logic
|
|
// is not the same as what's there, that could be a problem. Ideally
|
|
// all network calls from restore would go through the client and this
|
|
// would not longer be needed.
|
|
func (cmd *RestoreCommand) newClient() *retryablehttp.Client {
|
|
min := time.Millisecond * 100
|
|
|
|
// do some math to figure out how many attempts we need to get our
|
|
// total sleep time close to the period
|
|
attempts := math.Log2(float64(cmd.RetryPeriod)) - math.Log2(float64(min))
|
|
attempts += 0.3 // mmmm, fudge
|
|
if attempts < 1 {
|
|
attempts = 1
|
|
}
|
|
client := retryablehttp.NewClient()
|
|
client.RetryWaitMin = min
|
|
client.RetryMax = int(attempts)
|
|
client.CheckRetry = retryWith400
|
|
return client
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreIDAlloc(ctx context.Context, primary *topology.Node) error {
|
|
logger := cmd.Logger()
|
|
|
|
f, err := os.Open(filepath.Join(cmd.Path, "idalloc"))
|
|
if os.IsNotExist(err) {
|
|
logger.Printf("No idalloc, skipping")
|
|
return nil
|
|
} else if err != nil {
|
|
return err
|
|
}
|
|
defer f.Close()
|
|
|
|
logger.Printf("Load idalloc")
|
|
url := primary.URI.Path("/internal/idalloc/restore")
|
|
|
|
client := cmd.newClient()
|
|
_, err = client.Post(url, "application/octet-stream", f)
|
|
return err
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreShards(ctx context.Context) error {
|
|
filenames, err := filepath.Glob(filepath.Join(cmd.Path, "indexes", "*", "shards", "*"))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
ch := make(chan string, len(filenames))
|
|
for _, filename := range filenames {
|
|
ch <- filename
|
|
}
|
|
close(ch)
|
|
|
|
g, ctx := errgroup.WithContext(ctx)
|
|
for i := 0; i < cmd.Concurrency; i++ {
|
|
g.Go(func() error {
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case filename, ok := <-ch:
|
|
if !ok {
|
|
return nil
|
|
} else if err := cmd.restoreShard(ctx, filename); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
})
|
|
}
|
|
return g.Wait()
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreShard(ctx context.Context, filename string) error {
|
|
logger := cmd.Logger()
|
|
|
|
rel, err := filepath.Rel(cmd.Path, filename)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Parse filename.
|
|
record := strings.Split(rel, string(os.PathSeparator))
|
|
indexName := record[1]
|
|
shard, err := strconv.ParseUint(record[3], 10, 64)
|
|
if err != nil {
|
|
return nil // not a shard file
|
|
}
|
|
|
|
nodes, err := cmd.client.FragmentNodes(ctx, indexName, shard)
|
|
if err != nil {
|
|
return fmt.Errorf("cannot determine fragment nodes: %w", err)
|
|
} else if len(nodes) == 0 {
|
|
return fmt.Errorf("no nodes available")
|
|
}
|
|
|
|
for _, node := range nodes {
|
|
logger.Printf("shard %v %v", shard, indexName)
|
|
|
|
f, err := os.Open(filename)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer f.Close()
|
|
|
|
url := node.URI.Path(fmt.Sprintf("/internal/restore/%v/%v", indexName, shard))
|
|
req, err := retryablehttp.NewRequest("POST", url, f)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
req = req.WithContext(ctx)
|
|
req.Header.Set("Content-Type", "application/octet-stream")
|
|
|
|
client := cmd.newClient()
|
|
resp, err := client.Do(req)
|
|
if err != nil {
|
|
return err
|
|
} else if err := resp.Body.Close(); err != nil {
|
|
return err
|
|
} else if resp.StatusCode != http.StatusOK {
|
|
return fmt.Errorf("unexpected status code: %d", resp.StatusCode)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreIndexTranslation(ctx context.Context) error {
|
|
filenames, err := filepath.Glob(filepath.Join(cmd.Path, "indexes", "*", "translate", "*"))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
ch := make(chan string, len(filenames))
|
|
for _, filename := range filenames {
|
|
ch <- filename
|
|
}
|
|
close(ch)
|
|
|
|
g, ctx := errgroup.WithContext(ctx)
|
|
for i := 0; i < cmd.Concurrency; i++ {
|
|
g.Go(func() error {
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case filename, ok := <-ch:
|
|
if !ok {
|
|
return nil
|
|
} else if err := cmd.restoreIndexTranslationFile(ctx, filename); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
})
|
|
}
|
|
return g.Wait()
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreIndexTranslationFile(ctx context.Context, filename string) error {
|
|
logger := cmd.Logger()
|
|
|
|
rel, err := filepath.Rel(cmd.Path, filename)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
record := strings.Split(rel, string(os.PathSeparator))
|
|
indexName := record[1]
|
|
partitionID, err := strconv.Atoi(record[3])
|
|
if err != nil {
|
|
return err
|
|
}
|
|
logger.Printf("column keys %v (%v)", indexName, partitionID)
|
|
|
|
nodes, err := cmd.client.PartitionNodes(ctx, partitionID)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, node := range nodes {
|
|
if err := func() error {
|
|
readerFunc := func() (io.Reader, error) {
|
|
return os.Open(filename) // gets used as an HTTP request body and closed by http library
|
|
}
|
|
|
|
return cmd.client.ImportIndexKeys(ctx, &node.URI, indexName, partitionID, false, readerFunc)
|
|
}(); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreFieldTranslation(ctx context.Context, nodes []*topology.Node) error {
|
|
filenames, err := filepath.Glob(filepath.Join(cmd.Path, "indexes", "*", "fields", "*", "translate"))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
ch := make(chan string, len(filenames))
|
|
for _, filename := range filenames {
|
|
ch <- filename
|
|
}
|
|
close(ch)
|
|
|
|
g, ctx := errgroup.WithContext(ctx)
|
|
for i := 0; i < cmd.Concurrency; i++ {
|
|
g.Go(func() error {
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return ctx.Err()
|
|
case filename, ok := <-ch:
|
|
if !ok {
|
|
return nil
|
|
} else if err := cmd.restoreFieldTranslationFile(ctx, nodes, filename); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
})
|
|
}
|
|
return g.Wait()
|
|
}
|
|
|
|
func (cmd *RestoreCommand) restoreFieldTranslationFile(ctx context.Context, nodes []*topology.Node, filename string) error {
|
|
logger := cmd.Logger()
|
|
|
|
rel, err := filepath.Rel(cmd.Path, filename)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
record := strings.Split(rel, string(os.PathSeparator))
|
|
indexName, fieldName := record[1], record[3]
|
|
|
|
logger.Printf("field keys %v %v", indexName, fieldName)
|
|
|
|
for _, node := range nodes {
|
|
if err := func() error {
|
|
readerFunc := func() (io.Reader, error) {
|
|
return os.Open(filename)
|
|
}
|
|
|
|
return cmd.client.ImportFieldKeys(ctx, &node.URI, indexName, fieldName, false, readerFunc)
|
|
}(); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (cmd *RestoreCommand) TLSHost() string { return cmd.Host }
|
|
|
|
func (cmd *RestoreCommand) TLSConfiguration() server.TLSConfig { return cmd.TLS }
|