fix(offsite): 数据库包与桶内已有份数合并排名,只传会留下的,不再每小时重传再剪掉
This commit is contained in:
2 files changed
+65
-20
No files matched your search
@@ -392,6 +392,33 @@ func TestSyncDBKeepsNewest(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestSyncDBRanksWithTheBucket: an old local bundle (kept here by its label's
|
||||||
|
// own retention) that newer bundles in the bucket outrank is not sent, so a
|
||||||
|
// pass does not upload what it then prunes, and the next pass the same again.
|
||||||
|
func TestSyncDBRanksWithTheBucket(t *testing.T) {
|
||||||
|
s, b := newSyncer(t, &fakeCatalog{})
|
||||||
|
writeFile(t, s.DBDir, "felis-db-20260910T030000Z-pre-migrate.tar", 50)
|
||||||
|
writeFile(t, s.DBDir, "felis-db-20260924T030000Z-daily.tar", 50)
|
||||||
|
b.objs["db/felis-db-20260923T030000Z-daily.tar.fenc"] = []byte("gone here, kept there")
|
||||||
|
|
||||||
|
for pass := 1; pass <= 2; pass++ {
|
||||||
|
res, err := s.Run(context.Background())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
want := 0
|
||||||
|
if pass == 1 {
|
||||||
|
want = 1
|
||||||
|
}
|
||||||
|
if res.DBUploaded != want || res.DBPruned != 0 || res.RemoteDB != 2 {
|
||||||
|
t.Fatalf("pass %d: result = %+v, want %d uploaded, none pruned, 2 held", pass, res, want)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if _, ok := b.objs["db/felis-db-20260910T030000Z-pre-migrate.tar.fenc"]; ok {
|
||||||
|
t.Fatal("the outranked bundle was sent")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// TestFetchWorldsRestoresVolume: after a rebuild, every present archive the
|
// TestFetchWorldsRestoresVolume: after a rebuild, every present archive the
|
||||||
// volume lacks comes back byte for byte; ones already there are left alone.
|
// volume lacks comes back byte for byte; ones already there are left alone.
|
||||||
func TestFetchWorldsRestoresVolume(t *testing.T) {
|
func TestFetchWorldsRestoresVolume(t *testing.T) {
|
||||||
|
|||||||
+38
-20
@@ -17,9 +17,11 @@ import (
|
|||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
|
"maps"
|
||||||
"os"
|
"os"
|
||||||
"path"
|
"path"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"slices"
|
||||||
"sort"
|
"sort"
|
||||||
"strings"
|
"strings"
|
||||||
"syscall"
|
"syscall"
|
||||||
@@ -240,12 +242,32 @@ func (s *Syncer) syncDB(ctx context.Context, res *Result, fail func(string, ...a
|
|||||||
fail("list %s in the bucket: %v", dbDir, err)
|
fail("list %s in the bucket: %v", dbDir, err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
// Only the newest keep bundles are worth sending: older ones would be
|
// The bucket keeps the newest keep bundles of what it holds and what is
|
||||||
// pruned again at the end of this very pass.
|
// here together. Local retention is per label, so an old pre-migrate
|
||||||
if len(local) > keep {
|
// bundle can outlive newer dailies here that the bucket still has:
|
||||||
local = local[:keep]
|
// sending it would only see it pruned again at the end of this pass, and
|
||||||
|
// sent again on the next.
|
||||||
|
names := map[string]bool{}
|
||||||
|
for key := range remote {
|
||||||
|
name := strings.TrimSuffix(strings.TrimPrefix(key, dbDir), objExt)
|
||||||
|
if _, _, ok := dbbackup.ParseBundleName(name); ok && strings.HasSuffix(key, objExt) {
|
||||||
|
names[name] = true
|
||||||
|
}
|
||||||
}
|
}
|
||||||
for _, b := range local {
|
for _, b := range local {
|
||||||
|
names[b.Name] = true
|
||||||
|
}
|
||||||
|
// Bundle names start with their UTC stamp, so reversed order is newest first.
|
||||||
|
ranked := slices.Sorted(maps.Keys(names))
|
||||||
|
slices.Reverse(ranked)
|
||||||
|
kept := map[string]bool{}
|
||||||
|
for _, name := range ranked[:min(keep, len(ranked))] {
|
||||||
|
kept[name] = true
|
||||||
|
}
|
||||||
|
for _, b := range local {
|
||||||
|
if !kept[b.Name] {
|
||||||
|
continue
|
||||||
|
}
|
||||||
key := DBKey(b.Name)
|
key := DBKey(b.Name)
|
||||||
want := SealedSize(b.Size)
|
want := SealedSize(b.Size)
|
||||||
if size, ok := remote[key]; ok && size == want {
|
if size, ok := remote[key]; ok && size == want {
|
||||||
@@ -260,29 +282,25 @@ func (s *Syncer) syncDB(ctx context.Context, res *Result, fail func(string, ...a
|
|||||||
res.BytesUploaded += b.Size
|
res.BytesUploaded += b.Size
|
||||||
s.logf("copied database bundle %s (%s)", b.Name, HumanBytes(b.Size))
|
s.logf("copied database bundle %s (%s)", b.Name, HumanBytes(b.Size))
|
||||||
}
|
}
|
||||||
|
for _, name := range ranked {
|
||||||
var names []string
|
key := DBKey(name)
|
||||||
for key := range remote {
|
if _, ok := remote[key]; !ok || kept[name] {
|
||||||
name := strings.TrimSuffix(strings.TrimPrefix(key, dbDir), objExt)
|
|
||||||
if _, _, ok := dbbackup.ParseBundleName(name); ok && strings.HasSuffix(key, objExt) {
|
|
||||||
names = append(names, name)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
// Bundle names start with their UTC stamp, so newest sorts last.
|
|
||||||
sort.Sort(sort.Reverse(sort.StringSlice(names)))
|
|
||||||
for i, name := range names {
|
|
||||||
if i < keep {
|
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if err := s.Bucket.Remove(ctx, DBKey(name)); err != nil {
|
if err := s.Bucket.Remove(ctx, key); err != nil {
|
||||||
fail("prune database bundle %s: %v", name, err)
|
fail("prune database bundle %s: %v", name, err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
delete(remote, key)
|
||||||
res.DBPruned++
|
res.DBPruned++
|
||||||
}
|
}
|
||||||
res.RemoteDB = min(len(names), keep)
|
for _, name := range ranked {
|
||||||
if len(names) > 0 {
|
if _, ok := remote[DBKey(name)]; ok {
|
||||||
res.NewestDB = names[0]
|
res.RemoteDB++
|
||||||
|
if res.NewestDB == "" {
|
||||||
|
res.NewestDB = name
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user