fix(offsite): DB 包先同步,每个对象按大小单独限时,大归档传不完不再拖住整轮同步

This commit is contained in:
Lemon-miaow committed 2026-09-27 00:55:44 +08:00
1 parent 3f04171424
commit 5c58104e09
7 files changed
+140 -13

No files matched your search

+3 -1
View File
@@ -532,7 +532,9 @@ func (s *Syncer) putBytes(ctx context.Context, key string, plain []byte) error {
if err := Encrypt(&sealed, bytes.NewReader(plain), s.Key); err != nil {
return err
}
return s.Bucket.Put(ctx, key, &sealed, int64(sealed.Len()))
putCtx, cancel := context.WithTimeout(ctx, s.uploadBudget(int64(len(plain))))
defer cancel()
return s.Bucket.Put(putCtx, key, &sealed, int64(sealed.Len()))
}
func (s *Syncer) putImageIndex(ctx context.Context, stamp string, x *ImageIndex) error {
+74 -1
View File
@@ -118,13 +118,24 @@ type memBucket struct {
putErr map[string]error
puts int
removed []string
// stall holds keys whose upload never finishes: Put waits out its context.
stall map[string]bool
// started lists every Put in the order it began.
started []string
}
func newMemBucket() *memBucket {
return &memBucket{objs: map[string][]byte{}, putErr: map[string]error{}}
}
func (b *memBucket) Put(_ context.Context, key string, r io.Reader, size int64) error {
func (b *memBucket) Put(ctx context.Context, key string, r io.Reader, size int64) error {
b.mu.Lock()
b.started = append(b.started, key)
b.mu.Unlock()
if b.stall[key] {
<-ctx.Done()
return ctx.Err()
}
if err := b.putErr[key]; err != nil {
io.Copy(io.Discard, r)
return err
@@ -471,3 +482,65 @@ func TestWorldKeyRejectsOddRefs(t *testing.T) {
}
}
}
// TestSyncBigArchiveDoesNotStallThePass: an archive the uplink cannot send in
// time fails at its own deadline and stays pending, while the database bundle
// (sent before any world) and the archives after it still reach the bucket.
// One deadline for the whole pass used to go to the big archive, every hour,
// with the bundles queued behind it.
func TestSyncBigArchiveDoesNotStallThePass(t *testing.T) {
cat := &fakeCatalog{rows: []*row{
{WorldBackup: WorldBackup{ID: "b1", Server: "alpha", Ref: "/a/alpha-1.tar.gz"}, status: "present"},
{WorldBackup: WorldBackup{ID: "b2", Server: "beta", Ref: "/a/beta-2.tar.gz"}, status: "present"},
}}
s, b := newSyncer(t, cat)
s.UploadGrace = 50 * time.Millisecond
writeFile(t, s.ArchiveDir, "alpha-1.tar.gz", 100)
writeFile(t, s.ArchiveDir, "beta-2.tar.gz", 100)
writeFile(t, s.DBDir, "felis-db-20260924T030000Z-daily.tar", 50)
b.stall = map[string]bool{"worlds/alpha-1.tar.gz.fenc": true}
type ran struct {
res Result
err error
}
done := make(chan ran, 1)
go func() {
res, err := s.Run(context.Background())
done <- ran{res, err}
}()
var r ran
select {
case r = <-done:
case <-time.After(10 * time.Second):
t.Fatal("the pass is still waiting on the archive that cannot be sent")
}
if r.err == nil || !strings.Contains(r.err.Error(), "alpha-1.tar.gz") || !strings.Contains(r.err.Error(), "the next run tries again") {
t.Fatalf("run error = %v, want the stalled archive named as retried next run", r.err)
}
if r.res.DBUploaded != 1 || r.res.WorldsUploaded != 1 || r.res.WorldsPending != 1 {
t.Fatalf("result = %+v, want the bundle and beta copied, alpha pending", r.res)
}
if !cat.rows[0].offsite.IsZero() || cat.rows[1].offsite.IsZero() {
t.Fatalf("offsite_at alpha=%v beta=%v, want only beta recorded", cat.rows[0].offsite, cat.rows[1].offsite)
}
if len(b.started) == 0 || !strings.HasPrefix(b.started[0], dbDir) {
t.Fatalf("uploads began in the order %v, want the database bundle first", b.started)
}
}
// TestUploadBudgetFitsTheUplink: by default a 10 GiB archive may take well over
// what a 20 Mbit/s uplink needs to send it, and the smallest object still gets
// the grace.
func TestUploadBudgetFitsTheUplink(t *testing.T) {
s := &Syncer{}
const big = 10 << 30
need := time.Duration(float64(big) / (20e6 / 8) * float64(time.Second))
if got := s.uploadBudget(big); got < 3*need {
t.Errorf("budget for 10 GiB = %s, want at least %s (3× a 20 Mbit/s uplink)", got, 3*need)
}
if got := s.uploadBudget(0); got < defaultUploadGrace {
t.Errorf("budget for an empty object = %s, want at least %s", got, defaultUploadGrace)
}
}
+47 -8
View File
@@ -102,8 +102,38 @@ type Syncer struct {
// UploadsDir is the host directory of the uploads volume, whose submission
// contexts are copied (uploads.go); empty copies none.
UploadsDir string
Now func() time.Time
Log io.Writer
// UploadGrace and MinRate bound one object's upload: UploadGrace plus the
// time the object takes at MinRate bytes a second. Zero takes the defaults.
UploadGrace time.Duration
MinRate int64
Now func() time.Time
Log io.Writer
}
// Each object gets its own deadline, scaled to its size, so an archive too big
// for the uplink fails alone and the rest of the pass still goes; one deadline
// for the whole pass let a 10 GiB world use all of it, hour after hour, with
// every bundle queued behind it. 512 KiB/s (about 4 Mbit/s) gives a 10 GiB
// archive close to six hours, several times what a 20 Mbit/s uplink needs.
const (
defaultUploadGrace = 10 * time.Minute
defaultMinRate = 512 << 10
)
// uploadBudget is how long an object of size plaintext bytes may take.
func (s *Syncer) uploadBudget(size int64) time.Duration {
grace := s.UploadGrace
if grace <= 0 {
grace = defaultUploadGrace
}
return grace + time.Duration(float64(SealedSize(size))/float64(s.rate())*float64(time.Second))
}
func (s *Syncer) rate() int64 {
if s.MinRate > 0 {
return s.MinRate
}
return defaultMinRate
}
// Result is what one Run did and found.
@@ -158,9 +188,10 @@ func (s *Syncer) logf(format string, args ...any) {
}
}
// Run does one pass: world archives, database bundles, registry images,
// submission uploads, then expiry. A failure on one item is recorded and the pass carries on; the
// returned error is non-nil when anything failed.
// Run does one pass: database bundles, world archives, registry images,
// submission uploads, then expiry. The bundles go first: they are small, and
// every restore starts from one. A failure on one item is recorded and the
// pass carries on; the returned error is non-nil when anything failed.
func (s *Syncer) Run(ctx context.Context) (Result, error) {
var res Result
fail := func(format string, args ...any) {
@@ -173,8 +204,8 @@ func (s *Syncer) Run(ctx context.Context) (Result, error) {
if err != nil {
return res, fmt.Errorf("list %s in the bucket: %w", worldsDir, err)
}
s.syncWorlds(ctx, remoteWorlds, &res, fail)
s.syncDB(ctx, &res, fail)
s.syncWorlds(ctx, remoteWorlds, &res, fail)
s.syncImages(ctx, &res, fail)
s.syncUploads(ctx, &res, fail)
s.expireWorlds(ctx, remoteWorlds, &res, fail)
@@ -376,14 +407,22 @@ func (s *Syncer) putFile(ctx context.Context, key, p string, size int64) error {
// putStream encrypts size bytes of src into key. The sealed size is known in
// advance, so the upload streams: nothing larger than one part is buffered.
// An error from src fails the upload.
// An error from src fails the upload, and so does running past the object's
// uploadBudget.
func (s *Syncer) putStream(ctx context.Context, key string, src io.Reader, size int64) error {
budget := s.uploadBudget(size)
putCtx, cancel := context.WithTimeout(ctx, budget)
defer cancel()
pr, pw := io.Pipe()
go func() {
pw.CloseWithError(Encrypt(pw, src, s.Key))
}()
err := s.Bucket.Put(ctx, key, pr, SealedSize(size))
err := s.Bucket.Put(putCtx, key, pr, SealedSize(size))
pr.CloseWithError(errors.New("upload finished"))
if err != nil && ctx.Err() == nil && errors.Is(putCtx.Err(), context.DeadlineExceeded) {
return fmt.Errorf("not finished within %s (%s at under %s/s); the next run tries again: %w",
budget.Round(time.Minute), HumanBytes(SealedSize(size)), HumanBytes(s.rate()), err)
}
return err
}