diff --git a/weed/server/filer_server_format.go b/weed/server/filer_server_format.go index 830570cfa..8a156390d 100644 --- a/weed/server/filer_server_format.go +++ b/weed/server/filer_server_format.go @@ -58,6 +58,19 @@ func formatChunkIdentity(chunks []*filer_pb.FileChunk) []byte { return digest.Sum(nil) } +// repackSourceIdentity digests every entry field that influenced a repack's +// output: the chunk list it read, the size the layout was validated against, +// the TTL and expiry anchors its new chunks were assigned with, and the +// hard-link and remote state its guards evaluated. +func repackSourceIdentity(entry *filer.Entry) []byte { + digest := md5.New() + digest.Write(formatChunkIdentity(entry.GetChunks())) + fmt.Fprintf(digest, "%d:%d:%d:%d:%t:%x:%t", + entry.FileSize, entry.TtlSec, entry.Crtime.UnixNano(), entry.Mtime.UnixNano(), + entry.IsExpireS3Enabled(), []byte(entry.HardLinkId), entry.Remote != nil) + return digest.Sum(nil) +} + // roundUpToVolumeTTL returns the smallest volume-TTL-representable seconds // value not below the argument. A volume TTL is at most 255 of one unit and // SecondsToTTL truncates anything else downward, which would let chunks @@ -307,7 +320,7 @@ func (fs *FilerServer) formatRepack(ctx context.Context, w http.ResponseWriter, return } oldChunks := entry.GetChunks() - oldIdentity := formatChunkIdentity(oldChunks) + sourceIdentity := repackSourceIdentity(entry) if len(oldChunks) == 0 { writeJsonError(w, r, http.StatusBadRequest, errors.New("entry has no chunks to repack")) return @@ -408,9 +421,12 @@ func (fs *FilerServer) formatRepack(ctx context.Context, w http.ResponseWriter, } // The entry lock is filer-local, so a writer on another filer is not - // blocked by it. Re-read from the store and revalidate right before the - // swap: the unguarded window shrinks from the whole repack to this - // commit. Full enforcement needs owner routing. + // blocked by it. Re-read from the store and conflict on any change to + // state that influenced this repack - chunks, size, TTL and expiry + // anchors, hard-link and remote state - so the swap can never pair fresh + // metadata with chunks built from stale inputs. The unguarded window + // shrinks from the whole repack to this commit; within it repack races + // like any ordinary writer. Closing it needs owner routing. current, err := fs.filer.FindEntry(ctx, fullPath) if err != nil { cleanup() @@ -421,7 +437,7 @@ func (fs *FilerServer) formatRepack(ctx context.Context, w http.ResponseWriter, } return } - if !bytes.Equal(formatChunkIdentity(current.GetChunks()), oldIdentity) { + if !bytes.Equal(repackSourceIdentity(current), sourceIdentity) { cleanup() writeJsonError(w, r, http.StatusConflict, errors.New("entry changed during repack")) return diff --git a/weed/server/filer_server_format_test.go b/weed/server/filer_server_format_test.go index f95d0e596..77cae537b 100644 --- a/weed/server/filer_server_format_test.go +++ b/weed/server/filer_server_format_test.go @@ -2,8 +2,10 @@ package weed_server import ( "bytes" + "github.com/seaweedfs/seaweedfs/weed/filer" "math" "testing" + "time" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" "github.com/seaweedfs/seaweedfs/weed/storage/needle" @@ -22,7 +24,7 @@ func TestRoundUpToVolumeTTL(t *testing.T) { {3600, 3600}, {3601, 3660}, {255 * 60, 255 * 60}, - {255*60 + 1, 5 * 3600}, // minutes overflow 255, ceil to hours + {255*60 + 1, 5 * 3600}, // minutes overflow 255, ceil to hours {20_000_000, 232 * 24 * 3600}, // ~231.5 days, ceil to days {int64(math.MaxInt32), 0}, // beyond every unit's 255 cap: no TTL, never a shortened one } @@ -45,6 +47,39 @@ func TestRoundUpToVolumeTTL(t *testing.T) { } } +func TestRepackSourceIdentity(t *testing.T) { + base := func() *filer.Entry { + return &filer.Entry{ + FullPath: "/videos/movie.ts", + Attr: filer.Attr{ + FileSize: 30, TtlSec: 600, + Crtime: time.Unix(1000, 0), Mtime: time.Unix(2000, 0), + }, + Chunks: []*filer_pb.FileChunk{{FileId: "1,ab", Offset: 0, Size: 30}}, + } + } + identity := repackSourceIdentity(base()) + if !bytes.Equal(identity, repackSourceIdentity(base())) { + t.Fatalf("identity is not deterministic") + } + mutations := map[string]func(*filer.Entry){ + "ttl cleared": func(e *filer.Entry) { e.TtlSec = 0 }, + "size changed": func(e *filer.Entry) { e.FileSize = 31 }, + "mtime moved": func(e *filer.Entry) { e.Mtime = time.Unix(3000, 0) }, + "crtime moved": func(e *filer.Entry) { e.Crtime = time.Unix(1001, 0) }, + "hard linked": func(e *filer.Entry) { e.HardLinkId = []byte{1} }, + "went remote": func(e *filer.Entry) { e.Remote = &filer_pb.RemoteEntry{} }, + "chunk moved": func(e *filer.Entry) { e.Chunks[0].Offset = 1 }, + } + for name, mutate := range mutations { + changed := base() + mutate(changed) + if bytes.Equal(identity, repackSourceIdentity(changed)) { + t.Fatalf("identity ignored: %s", name) + } + } +} + func TestFormatChunkIdentity(t *testing.T) { chunks := []*filer_pb.FileChunk{ {FileId: "1,ab", Offset: 0, Size: 10},