Manual buffer management in restore-tar (#2249)

* manual buffer management
reuse buffer on backuptar as well
test coverage for backuptar
This commit is contained in:
tgruben 2023-02-16 14:57:00 -06:00 • committed by GitHub
parent 6a9843b4a0
commit 4172976e6a
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
3 changed files with 157 additions and 92 deletions

View file

@ -164,16 +164,17 @@ func (cmd *BackupTarCommand) Run(ctx context.Context) (err error) {
tw := tar.NewWriter(w)
defer tw.Close()
buf := new(bytes.Buffer)
// Backup schema.
if err := cmd.backupTarSchema(ctx, tw, schema); err != nil {
return fmt.Errorf("cannot back up schema: %w", err)
} else if err := cmd.backupTarIDAllocData(ctx, tw); err != nil {
} else if err := cmd.backupTarIDAllocData(ctx, tw, buf); err != nil {
return fmt.Errorf("cannot back up id alloc data: %w", err)
}
// Backup data for each index.
for _, ii := range schema.Indexes {
if err := cmd.backupTarIndex(ctx, tw, ii); err != nil {
if err := cmd.backupTarIndex(ctx, tw, ii, buf); err != nil {
return err
}
}
@ -219,7 +220,8 @@ func (cmd *BackupTarCommand) backupTarSchema(ctx context.Context, tw *tar.Writer
return nil
}
func (cmd *BackupTarCommand) backupTarIDAllocData(ctx context.Context, tw *tar.Writer) error {
func (cmd *BackupTarCommand) backupTarIDAllocData(ctx context.Context, tw *tar.Writer, buf *bytes.Buffer) error {
buf.Reset()
logger := cmd.Logger()
logger.Printf("backing up id alloc data")
@ -230,7 +232,6 @@ func (cmd *BackupTarCommand) backupTarIDAllocData(ctx context.Context, tw *tar.W
defer rc.Close()
// Read to buffer to determine size.
var buf bytes.Buffer
if _, err := buf.ReadFrom(rc); err != nil {
return fmt.Errorf("copying id alloc data to memory: %w", err)
}
@ -243,7 +244,7 @@ func (cmd *BackupTarCommand) backupTarIDAllocData(ctx context.Context, tw *tar.W
ModTime: time.Now(),
}); err != nil {
return err
} else if _, err := io.Copy(tw, &buf); err != nil {
} else if _, err := io.Copy(tw, buf); err != nil {
return fmt.Errorf("copying id alloc data to archive: %w", err)
}
@ -251,7 +252,7 @@ func (cmd *BackupTarCommand) backupTarIDAllocData(ctx context.Context, tw *tar.W
}
// backupTarIndex backs up all shards for a given index.
func (cmd *BackupTarCommand) backupTarIndex(ctx context.Context, tw *tar.Writer, ii *pilosa.IndexInfo) error {
func (cmd *BackupTarCommand) backupTarIndex(ctx context.Context, tw *tar.Writer, ii *pilosa.IndexInfo, buf *bytes.Buffer) error {
logger := cmd.Logger()
logger.Printf("backing up index: %q", ii.Name)
@ -262,14 +263,14 @@ func (cmd *BackupTarCommand) backupTarIndex(ctx context.Context, tw *tar.Writer,
// Back up all bitmap data for the index.
for _, shard := range shards {
if err := cmd.backupTarShard(ctx, tw, ii.Name, shard); err != nil {
if err := cmd.backupTarShard(ctx, tw, ii.Name, shard, buf); err != nil {
return fmt.Errorf("cannot backup shard %d on index %q: %w", shard, ii.Name, err)
}
}
if ii.Options.Keys {
// Back up translation data after bitmap data so we ensure we can translate all data.
if err := cmd.backupTarIndexTranslateData(ctx, tw, ii.Name); err != nil {
// Back up translation data after bitmap data so we ensurean translate all data.
if err := cmd.backupTarIndexTranslateData(ctx, tw, ii.Name, buf); err != nil {
return err
}
}
@ -279,7 +280,7 @@ func (cmd *BackupTarCommand) backupTarIndex(ctx context.Context, tw *tar.Writer,
if !fi.Options.Keys {
continue
}
if err := cmd.backupTarFieldTranslateData(ctx, tw, ii.Name, fi.Name); err != nil {
if err := cmd.backupTarFieldTranslateData(ctx, tw, ii.Name, fi.Name, buf); err != nil {
return fmt.Errorf("cannot backup field translation data for field %q on index %q: %w", fi.Name, ii.Name, err)
}
}
@ -288,7 +289,7 @@ func (cmd *BackupTarCommand) backupTarIndex(ctx context.Context, tw *tar.Writer,
}
// backupTarShard backs up a single shard from a single index.
func (cmd *BackupTarCommand) backupTarShard(ctx context.Context, tw *tar.Writer, indexName string, shard uint64) (err error) {
func (cmd *BackupTarCommand) backupTarShard(ctx context.Context, tw *tar.Writer, indexName string, shard uint64, buf *bytes.Buffer) (err error) {
nodes, err := cmd.client.FragmentNodes(ctx, indexName, shard)
if err != nil {
return fmt.Errorf("cannot determine fragment nodes: %w", err)
@ -297,7 +298,7 @@ func (cmd *BackupTarCommand) backupTarShard(ctx context.Context, tw *tar.Writer,
}
for _, node := range nodes {
if e := cmd.backupTarShardNode(ctx, tw, indexName, shard, node); e == nil {
if e := cmd.backupTarShardNode(ctx, tw, indexName, shard, node, buf); e == nil {
break
} else if err == nil {
err = e // save first error, try next node
@ -305,7 +306,7 @@ func (cmd *BackupTarCommand) backupTarShard(ctx context.Context, tw *tar.Writer,
}
for _, node := range nodes {
if e := cmd.backupTarShardDataframe(ctx, tw, indexName, shard, node); e == nil {
if e := cmd.backupTarShardDataframe(ctx, tw, indexName, shard, node, buf); e == nil {
break
} else if err == nil {
err = e // save first error, try next node
@ -315,7 +316,8 @@ func (cmd *BackupTarCommand) backupTarShard(ctx context.Context, tw *tar.Writer,
}
// backupTarShardNode backs up a single shard from a single index on a specific node.
func (cmd *BackupTarCommand) backupTarShardNode(ctx context.Context, tw *tar.Writer, indexName string, shard uint64, node *disco.Node) error {
func (cmd *BackupTarCommand) backupTarShardNode(ctx context.Context, tw *tar.Writer, indexName string, shard uint64, node *disco.Node, buf *bytes.Buffer) error {
buf.Reset()
logger := cmd.Logger()
logger.Printf("backing up shard: index=%q id=%d", indexName, shard)
@ -333,7 +335,6 @@ func (cmd *BackupTarCommand) backupTarShardNode(ctx context.Context, tw *tar.Wri
// Read to buffer to determine size.
// TODO: Provide size via the reader itself.
var buf bytes.Buffer
if _, err := buf.ReadFrom(rc); err != nil {
return fmt.Errorf("copying shard data to memory: %w", err)
}
@ -346,14 +347,15 @@ func (cmd *BackupTarCommand) backupTarShardNode(ctx context.Context, tw *tar.Wri
ModTime: time.Now(),
}); err != nil {
return err
} else if _, err := io.Copy(tw, &buf); err != nil {
} else if _, err := io.Copy(tw, buf); err != nil {
return fmt.Errorf("copying shard data to archive: %w", err)
}
return nil
}
func (cmd *BackupTarCommand) backupTarShardDataframe(ctx context.Context, tw *tar.Writer, indexName string, shard uint64, node *disco.Node) error {
func (cmd *BackupTarCommand) backupTarShardDataframe(ctx context.Context, tw *tar.Writer, indexName string, shard uint64, node *disco.Node, buf *bytes.Buffer) error {
buf.Reset()
logger := cmd.Logger()
logger.Printf("backing up dataframe shard: index=%q shard=%d", indexName, shard)
@ -374,7 +376,6 @@ func (cmd *BackupTarCommand) backupTarShardDataframe(ctx context.Context, tw *ta
filename := filepath.Join("indexes", indexName, "dataframe", fmt.Sprintf("%04d", shard))
logger.Printf("writing %v", filename)
var buf bytes.Buffer
if _, err := buf.ReadFrom(resp.Body); err != nil {
return fmt.Errorf("copying shard data to memory: %w", err)
}
@ -387,24 +388,25 @@ func (cmd *BackupTarCommand) backupTarShardDataframe(ctx context.Context, tw *ta
ModTime: time.Now(),
}); err != nil {
return err
} else if _, err := io.Copy(tw, &buf); err != nil {
} else if _, err := io.Copy(tw, buf); err != nil {
return fmt.Errorf("copying shard data to archive: %w", err)
}
return nil
}
func (cmd *BackupTarCommand) backupTarIndexTranslateData(ctx context.Context, tw *tar.Writer, name string) error {
func (cmd *BackupTarCommand) backupTarIndexTranslateData(ctx context.Context, tw *tar.Writer, name string, buf *bytes.Buffer) error {
// TODO: Fetch holder partition count.
partitionN := disco.DefaultPartitionN
for partitionID := 0; partitionID < partitionN; partitionID++ {
if err := cmd.backupTarIndexPartitionTranslateData(ctx, tw, name, partitionID); err != nil {
if err := cmd.backupTarIndexPartitionTranslateData(ctx, tw, name, partitionID, buf); err != nil {
return fmt.Errorf("cannot backup index translation data for partition %d on %q: %w", partitionID, name, err)
}
}
return nil
}
func (cmd *BackupTarCommand) backupTarIndexPartitionTranslateData(ctx context.Context, tw *tar.Writer, name string, partitionID int) error {
func (cmd *BackupTarCommand) backupTarIndexPartitionTranslateData(ctx context.Context, tw *tar.Writer, name string, partitionID int, buf *bytes.Buffer) error {
buf.Reset()
logger := cmd.Logger()
logger.Printf("backing up index translation data: %s/%d", name, partitionID)
@ -417,7 +419,6 @@ func (cmd *BackupTarCommand) backupTarIndexPartitionTranslateData(ctx context.Co
defer rc.Close()
// Read to buffer to determine size.
var buf bytes.Buffer
if _, err := buf.ReadFrom(rc); err != nil {
return fmt.Errorf("copying translate data to memory: %w", err)
}
@ -430,14 +431,15 @@ func (cmd *BackupTarCommand) backupTarIndexPartitionTranslateData(ctx context.Co
ModTime: time.Now(),
}); err != nil {
return err
} else if _, err := io.Copy(tw, &buf); err != nil {
} else if _, err := io.Copy(tw, buf); err != nil {
return fmt.Errorf("copying translate data to archive: %w", err)
}
return nil
}
func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw *tar.Writer, indexName, fieldName string) error {
func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw *tar.Writer, indexName, fieldName string, buf *bytes.Buffer) error {
buf.Reset()
logger := cmd.Logger()
logger.Printf("backing up field translation data: %s/%s", indexName, fieldName)
@ -450,7 +452,6 @@ func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw
defer rc.Close()
// Read to buffer to determine size.
var buf bytes.Buffer
if _, err := buf.ReadFrom(rc); err != nil {
return fmt.Errorf("copying translate data to memory: %w", err)
}
@ -463,7 +464,7 @@ func (cmd *BackupTarCommand) backupTarFieldTranslateData(ctx context.Context, tw
ModTime: time.Now(),
}); err != nil {
return err
} else if _, err := io.Copy(tw, &buf); err != nil {
} else if _, err := io.Copy(tw, buf); err != nil {
return fmt.Errorf("copying translate data to archive: %w", err)
}
return nil

View file

@ -134,7 +134,9 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
return errors.New("no primary")
}
c := &gohttp.Client{}
buf := new(bytes.Buffer)
for {
buf.Reset()
header, err := tarReader.Next()
if err == io.EOF {
break
@ -176,16 +178,17 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
return fmt.Errorf("no fragmentNodes available")
}
shardBytes, err := io.ReadAll(tarReader) // this feels wrong but works for now
_, err = io.Copy(buf, tarReader)
if err != nil {
return err
return errors.Wrap(err, "copying")
}
g, _ := errgroup.WithContext(ctx)
for _, node := range fragmentNodes {
node := node
g.Go(func() error {
client := &gohttp.Client{}
rd := bytes.NewReader(shardBytes)
rd := bytes.NewReader(buf.Bytes())
logger.Printf("shard %v %v", shard, indexName)
url := node.URI.Path(fmt.Sprintf("/internal/restore/%v/%v", indexName, shard))
_, err = client.Post(url, "application/octet-stream", rd)
@ -207,16 +210,17 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
return fmt.Errorf("no fragmentNodes available")
}
shardBytes, err := io.ReadAll(tarReader) // this feels wrong but works for now
_, err = io.Copy(buf, tarReader)
if err != nil {
return err
return errors.Wrap(err, "copying")
}
g, _ := errgroup.WithContext(ctx)
for _, node := range fragmentNodes {
node := node
g.Go(func() error {
client := &gohttp.Client{}
rd := bytes.NewReader(shardBytes)
rd := bytes.NewReader(buf.Bytes())
logger.Printf("dataframe shard %v %v", shard, indexName)
url := node.URI.Path(fmt.Sprintf("/internal/dataframe/restore/%v/%v", indexName, shard))
_, err = client.Post(url, "application/octet-stream", rd)
@ -236,17 +240,19 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
if err != nil {
return err
}
shardBytes, err := io.ReadAll(tarReader) // this feels wrong but works for now
_, err = io.Copy(buf, tarReader)
if err != nil {
return err
return errors.Wrap(err, "copying")
}
g, _ := errgroup.WithContext(ctx)
for _, node := range partitionNodes {
node := node
g.Go(func() error {
// rd := bytes.NewReader(shardBytes)
rd := func() (io.Reader, error) {
return bytes.NewReader(shardBytes), nil
return bytes.NewReader(buf.Bytes()), nil
}
return cmd.client.ImportIndexKeys(ctx, &node.URI, indexName, partitionID, false, rd)
@ -264,17 +270,19 @@ func (cmd *RestoreTarCommand) Run(ctx context.Context) (err error) {
case "translate":
logger.Printf("field keys %v %v", indexName, fieldName)
// needs to go to all nodes
shardBytes, err := io.ReadAll(tarReader) // this feels wrong but works for now
_, err = io.Copy(buf, tarReader)
if err != nil {
return err
return errors.Wrap(err, "copying")
}
g, _ := errgroup.WithContext(ctx)
for _, node := range nodes {
node := node
g.Go(func() error {
// rd := bytes.NewReader(shardBytes)
rd := func() (io.Reader, error) {
return bytes.NewReader(shardBytes), nil
return bytes.NewReader(buf.Bytes()), nil
}
return cmd.client.ImportFieldKeys(ctx, &node.URI, indexName, fieldName, false, rd)

View file

@ -1369,7 +1369,6 @@ func TestExecutor_Execute_Count(t *testing.T) {
t.Fatalf("unexpected n: %d", res.Results[0])
}
})
}
// Ensure a set query can be executed.
@ -1565,7 +1564,6 @@ func TestExecutor_Execute_SetBool(t *testing.T) {
if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: c.Idx(), Query: `Set(100, f=1)`}); err == nil {
t.Fatalf("expected invalid bool type error")
}
})
}
@ -1754,7 +1752,6 @@ func TestExecutor_Execute_SetValue(t *testing.T) {
t.Fatalf("unexpected value: %v", value)
}
})
}
func TestExecutor_ExecuteTopK(t *testing.T) {
@ -2919,7 +2916,6 @@ func TestExecutor_Execute_Sum(t *testing.T) {
t.Fatalf("unexpected result: %s", spew.Sdump(result))
}
})
})
})
@ -3176,7 +3172,6 @@ func TestExecutor_Execute_Row_BSIGroup(t *testing.T) {
}
})
}
})
// Ensure that the NotNull code path gets run.
@ -3883,10 +3878,8 @@ func TestExecutor_Time_Clear_Quantums(t *testing.T) {
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, tt.expected) {
t.Fatalf("unexpected columns: %+v", columns)
}
})
}
}
func TestReopenCluster(t *testing.T) {
@ -4003,7 +3996,6 @@ func TestExecutor_Execute_Existence(t *testing.T) {
// Ensure a not query can be executed.
func TestExecutor_Execute_Not(t *testing.T) {
}
// Ensure an all query can be executed.
@ -4237,7 +4229,6 @@ func TestExecutor_Execute_Limit(t *testing.T) {
t.Errorf("expected %v but got %v", expect, got)
}
})
}
func TestExecutor_Sort(t *testing.T) {
@ -4344,7 +4335,6 @@ func TestExecutor_Sort(t *testing.T) {
t.Errorf("expected %v but got %v", expect[i], resp.Results[0])
}
}
})
}
@ -4644,10 +4634,11 @@ func TestExecutor_Execute_SetRow(t *testing.T) {
}
// Set bits.
if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: c.Idx(), Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10),
if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{
Index: c.Idx(), Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10),
}); err != nil {
t.Fatal(err)
}
@ -4697,10 +4688,11 @@ func TestExecutor_Execute_SetRow(t *testing.T) {
}
// Set bits.
if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: c.Idx(), Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10),
if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{
Index: c.Idx(), Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10),
}); err != nil {
t.Fatal(err)
}
@ -4750,12 +4742,13 @@ func TestExecutor_Execute_SetRow(t *testing.T) {
}
// Set bits.
if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: c.Idx(), Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", 1, 20) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 20),
if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{
Index: c.Idx(), Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", 1, 20) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 20),
}); err != nil {
t.Fatal(err)
}
@ -5685,9 +5678,7 @@ func TestExecutor_Execute_Rows_Keys(t *testing.T) {
for i := shard; i < shard+10; i++ {
for row := i; row >= 0 && row > i-3; row-- {
query.WriteString(fmt.Sprintf("Set(\"%d\", f=\"%d\")", shard*pilosa.ShardWidth+i, row))
}
}
}
_, err = c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{
@ -5823,7 +5814,6 @@ func TestExecutor_Execute_Rows_Keys(t *testing.T) {
}
})
}
}
func TestExecutor_ForeignIndex(t *testing.T) {
@ -6130,7 +6120,6 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
results := c.Query(t, c.Idx(), `GroupBy(Rows(general, previous=10), limit=1)`).Results[0].(*pilosa.GroupCounts).Groups()
test.CheckGroupBy(t, expected, results)
})
c.CreateField(t, c.Idx(), pilosa.IndexOptions{}, "a")
@ -6164,15 +6153,21 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{3, 3}, // no overlap
})
c.ImportBits(t, c.Idx(), "wb", [][2]uint64{
{0, 0}, {0, 1}, {0, 2},
{0, 0},
{0, 1},
{0, 2},
{1, 1},
{2, 0}, {2, 2},
{2, 0},
{2, 2},
{3, 3},
})
c.ImportBits(t, c.Idx(), "wc", [][2]uint64{
{0, 0}, {0, 1}, {0, 2},
{0, 0},
{0, 1},
{0, 2},
{1, 1},
{2, 0}, {2, 2},
{2, 0},
{2, 2},
{3, 3},
})
@ -6274,7 +6269,6 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{Group: []pilosa.FieldRow{{Field: "na", RowID: 1}, {Field: "nb", RowID: 1}}, Count: 2},
}
test.CheckGroupBy(t, expected, results)
})
// test paging over results using previous. set the same bits in three
@ -6286,19 +6280,31 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
{0, 0},
{1, 0},
{2, 0},
{3, 0}, {3, 91000}, {3, ShardWidth}, {3, ShardWidth * 2}, {3, ShardWidth * 3},
{3, 0},
{3, 91000},
{3, ShardWidth},
{3, ShardWidth * 2},
{3, ShardWidth * 3},
})
c.ImportBits(t, c.Idx(), "ppb", [][2]uint64{
{0, 0},
{1, 0},
{2, 0},
{3, 0}, {3, 91000}, {3, ShardWidth}, {3, ShardWidth * 2}, {3, ShardWidth * 3},
{3, 0},
{3, 91000},
{3, ShardWidth},
{3, ShardWidth * 2},
{3, ShardWidth * 3},
})
c.ImportBits(t, c.Idx(), "ppc", [][2]uint64{
{0, 0},
{1, 0},
{2, 0},
{3, 0}, {3, 91000}, {3, ShardWidth}, {3, ShardWidth * 2}, {3, ShardWidth * 3},
{3, 0},
{3, 91000},
{3, ShardWidth},
{3, ShardWidth * 2},
{3, ShardWidth * 3},
})
t.Run("test wrapping with previous", func(t *testing.T) {
@ -6352,7 +6358,6 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
results := c.Query(t, c.Idx(), `GroupBy(Rows(generalk), Rows(subk))`).Results[0].(*pilosa.GroupCounts).Groups()
test.CheckGroupBy(t, expected, results)
})
// Foreign Index
@ -6420,7 +6425,6 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
},
results,
)
})
// Create some time-quantum data:
c.Query(t, c.Idx(), "Set(0, tq=1, 2022-01-01T01:01)")
@ -6451,7 +6455,6 @@ func TestExecutor_Execute_GroupBy(t *testing.T) {
results := c.Query(t, c.Idx(), query).Results[0].(*pilosa.GroupCounts).Groups()
test.CheckGroupBy(t, want, results)
}
})
}
for _, size := range []int{1, 3} {
@ -6521,7 +6524,6 @@ func BenchmarkGroupBy(b *testing.B) {
// TODO benchmark over multiple shards
// TODO benchmark paging over large numbers of rows
}
// NOTE: The shift function in its current state is unsupported.
@ -6597,10 +6599,10 @@ func TestExecutor_Execute_Shift(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := c.GetHolder(0)
hldr.SetBit(c.Idx(), "general", 10, ShardWidth-2) //shardwidth -1
hldr.SetBit(c.Idx(), "general", 10, ShardWidth-1) //shardwidth
hldr.SetBit(c.Idx(), "general", 10, ShardWidth) //shardwidth +1
hldr.SetBit(c.Idx(), "general", 10, ShardWidth+2) //shardwidth +3
hldr.SetBit(c.Idx(), "general", 10, ShardWidth-2) // shardwidth -1
hldr.SetBit(c.Idx(), "general", 10, ShardWidth-1) // shardwidth
hldr.SetBit(c.Idx(), "general", 10, ShardWidth) // shardwidth +1
hldr.SetBit(c.Idx(), "general", 10, ShardWidth+2) // shardwidth +3
exp := []uint64{ShardWidth - 1, ShardWidth, ShardWidth + 1, ShardWidth + 3}
if res, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: c.Idx(), Query: `Shift(Row(general=10), n=1)`}); err != nil {
@ -6978,7 +6980,6 @@ func TestExecutor_Execute_CountDistinct(t *testing.T) {
}
if r.Pos.Count() != 1 {
t.Fatalf("invalid pilosa.SignedRow.Pos.Count, expected: 1, got: %v", r.Pos.Count())
}
if r.Pos.Columns()[0] != 100 {
t.Fatalf("invalid pilosa.SignedRow.Pos.Columns, expected: [100], got: %v", r.Pos.Columns())
@ -7395,6 +7396,9 @@ func TestVariousQueries(t *testing.T) {
t.Run("backup-full", func(t *testing.T) {
backupTest(t, c, "") // test backup/restore of all indexes
})
t.Run("backuptar-full", func(t *testing.T) {
backupTarTest(t, c, "") // test backup/restore of all indexes
})
})
}
}
@ -7420,6 +7424,27 @@ func backupTest(t *testing.T, c *test.Cluster, index string) {
}
}
func backupTarTest(t *testing.T, c *test.Cluster, index string) {
// should this really be in executor? No. But all these
// integration-y query tests probably shouldn't be either. My goal
// putting this here is to take advantage of already-existing
// clusters and data.
sum := chkSumCluster(t, c)
backupDir := backupClusterTar(t, c, index)
cnew := test.MustRunUnsharedCluster(t, 3) // this way we test 1->3 3->3 5->3
defer cnew.Close()
restoreClusterTar(t, backupDir, cnew)
sumNew := chkSumCluster(t, cnew)
if sum != sumNew {
t.Fatalf("old/new checksum mismatch, old:\n%s\nnew:\n%s", sum, sumNew)
}
}
func chkSumCluster(t *testing.T, c *test.Cluster) string {
t.Helper()
buf := &bytes.Buffer{}
@ -7466,6 +7491,39 @@ func restoreCluster(t *testing.T, backupDir string, c *test.Cluster) {
}
}
func backupClusterTar(t *testing.T, c *test.Cluster, index string) (backupFileName string) {
td, err := testhook.TempDir(t, "backupTestTar")
if err != nil {
t.Fatalf("can't even get a temp dir, what a ripoff: %v", err)
}
td = td + "/backupTestTar.tar"
buf := &bytes.Buffer{}
// backupLog := logger.NewStandardLogger(buf)
backupTarCommand := ctl.NewBackupTarCommand(buf)
backupTarCommand.Host = c.Nodes[len(c.Nodes)-1].URL() // don't pick node 0 so we don't always get primary (better code coverage)
backupTarCommand.Index = index
backupTarCommand.OutputPath = td
if err := backupTarCommand.Run(context.Background()); err != nil {
t.Log(buf.String())
t.Fatalf("running backuptar: %v", err)
}
return td
}
func restoreClusterTar(t *testing.T, backupFileDir string, c *test.Cluster) {
buf := &bytes.Buffer{}
restoreLog := logger.NewStandardLogger(buf)
restore := ctl.NewRestoreTarCommand(restoreLog)
restore.Host = c.Nodes[len(c.Nodes)-1].URL()
restore.Path = backupFileDir
if err := restore.Run(context.Background()); err != nil {
t.Fatalf("restoring: %v", err)
}
}
// tests for abbreviating time values in queries
func variousQueriesOnPercentiles(t *testing.T, c *test.Cluster) {
// todo, make rand more random, 42 isnt the answer to everything
@ -8089,8 +8147,8 @@ func variousQueriesOnLargeEpoch(t *testing.T, c *test.Cluster) {
{Val: 0, Key: "userA"}, // 9999-12-31
// {Val: 1, Key: "userE"},
// {Val: -1, Key: "userD"},
{Val: -maxSec, Key: "userB"}, //1970....
{Val: -maxSec + minSec, Key: "userC"}, //0001
{Val: -maxSec, Key: "userB"}, // 1970....
{Val: -maxSec + minSec, Key: "userC"}, // 0001
})
c.CreateField(t, index, pilosa.IndexOptions{Keys: true, TrackExistence: true}, "unix_milli_min", pilosa.OptFieldTypeTimestamp(minTime, "ms"))
@ -9367,6 +9425,7 @@ func TestExternalLookup(t *testing.T) {
}
})
}
func TestToRows(t *testing.T) {
ids := &pilosa.RowIdentifiers{
Rows: []uint64{1, 2, 3},
@ -9475,7 +9534,6 @@ func TestToRows(t *testing.T) {
if e != nil {
t.Fatal("Shouldn't be err ", e)
}
}
// TestMinMaxTimestampVariableNode tests Min() and Max() queries on
@ -9596,7 +9654,7 @@ func TestExecutor_Execute_ExtractWithTime(t *testing.T) {
t.Run("Extract With From Time", func(t *testing.T) {
resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment,from=2022-07-03T00:00))")
//resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment, from=))")
// resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment, from=))")
if len(resp.Results) != 1 {
t.Fatalf("expected 1 result but got %d", len(resp.Results))
}
@ -9629,7 +9687,7 @@ func TestExecutor_Execute_ExtractWithTime(t *testing.T) {
})
t.Run("Extract With Time No Opt", func(t *testing.T) {
resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment))")
//resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment, from=))")
// resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment, from=))")
if len(resp.Results) != 1 {
t.Fatalf("expected 1 result but got %d", len(resp.Results))
}
@ -9663,7 +9721,7 @@ func TestExecutor_Execute_ExtractWithTime(t *testing.T) {
t.Run("Extract With ToTime ", func(t *testing.T) {
resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment,to=2022-07-02T00:00))")
//resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment, from=))")
// resp := c.Query(t, c.Idx(), "Extract(All(), Rows(segment, from=))")
if len(resp.Results) != 1 {
t.Fatalf("expected 1 result but got %d", len(resp.Results))
}
@ -9772,7 +9830,6 @@ func TestExecutorTimeRange(t *testing.T) {
t.Fatalf("unexpected columns: %+v for query: %+v", columns, readQueries[i])
}
}
})
// test standard view enabled
@ -9808,5 +9865,4 @@ func TestExecutorTimeRange(t *testing.T) {
}
}
})
}