filer: honor is_moved only from ring member connections (#11456)

* filer.remote.sync: stamp entries with IF_CHUNKS_EQUAL so a stale write-back cannot delete live chunks

updateLocalEntry records the RemoteEntry stamp after an upload by writing the
event's entry back with UpdateEntry. The filer deletes every stored chunk
absent from an updated entry, so when the file was rewritten while its upload
was in flight (or the event is a replay), the stale snapshot deletes the
rewrite's chunks: the entry then points at the new fid with no needle behind
it, and the rewrite's own upload fails and is skipped as superseded.

The stamp write now carries WriteCondition IF_CHUNKS_EQUAL over the event's
chunk fids, evaluated by the filer under the path lock. A refused stamp means
the filer moved past this event; the superseding event follows in the log and
stamps the current entry, so the refusal is logged and skipped like a
superseded upload.

Reproduction: weed server -filer plus a weed server -s3 remote, remote.mount,
filer.remote.sync; hold the remote (docker pause) so one upload stays in
flight, rewrite the file through the filer, unpause. Before: the entry's chunk
is 404 on every volume server. After: the stale stamp is refused, the rewrite's
chunk stays live and reads back after a vacuum.

* filer.remote.sync: stamp entries with IF_ENTRY_EQUAL so stale inline content or metadata cannot be restored

The IF_CHUNKS_EQUAL guard compared only the chunk fid multiset, so a
rewrite that touched inline content or metadata alone still compared
equal and the stale snapshot overwrote the live entry. The new clause
compares the whole stored entry against the event's entry under the
same path lock.

* filer: route conditional UpdateEntry to the entry's owner filer

Two filers locking the same path locally could still pass a stale
condition on the non-owner while the owner's entry had moved on. When a
condition or expected_extended precondition is set, forward the request
to the entry's owner the same way conditional CreateEntry does, with
is_moved bounding the hop.

* filer: compare IF_ENTRY_EQUAL against the normalized expected entry

FindEntry grows FileSize to the chunk extent, so a raw event entry with
FileSize still zero failed the condition on an unchanged file and the
stamp was skipped, letting a replay upload the object again.

* filer.remote.sync: classify refused stamps by gRPC status only

A FailedPrecondition substring in an unrelated error would have been
swallowed as a skipped stamp; status.FromError already unwraps.

* remote sync: keep the event entry intact for IF_ENTRY_EQUAL

* filer: honor is_moved only from ring member connections

is_moved is caller-controlled, so a request could set it to skip owner
routing and run a conditional check under a non-owner's lock. Verify the
marker against the peer's connection address and the lock ring members;
an unverified marker is ignored and the request routes like a fresh one.

* filer: refuse unverifiable is_moved at a non-owner, cache ring IPs

Follow-up fixes from review on the is_moved provenance check:

- checkMovedMarker replaces "ignore and re-forward" for markers that did
  not arrive on a ring member's connection. Re-forwarding a claimed hop
  could cycle while rings disagree; instead the request is refused with
  FailedPrecondition unless this filer is the key's owner, in which case
  applying locally is correct anyway.
- ringMemberIPs caches resolved member addresses per ring membership so
  hostname-advertising deployments do not pay a DNS lookup per forwarded
  request; failed lookups are not cached so a DNS blip self-heals.
- DistributedUnlock no longer dereferences the nil response of a failed
  next-hop RPC.

* filer: refuse unverifiable is_moved with PermissionDenied, not FailedPrecondition

A routing refusal is different in kind from a write-condition mismatch:
remote sync treats FailedPrecondition as a stale stamp and skips it, so
reusing that code let a routing failure pass as synced. Owner checks now
also run before the peer-IP lookup so the common accept path does no DNS.

* filer: expire resolved ring member IPs after 5 minutes

A member's hostname can re-resolve to a new IP while its ring address
stays unchanged; caching forever would reject its genuine forwards until
a membership change or restart.

* filer: deduplicate concurrent ring member DNS lookups

At cache expiry, parallel forwarded requests would each resolve every
member hostname serially; singleflight collapses them into one lookup
per ring membership.

* filer: detach the shared ring lookup from the caller's context

The singleflight winner's ctx is cancelled when its request ends; the
shared result would then be an incomplete member list and genuine
forwards denied. The lookup now runs on a detached context with its
own deadline so a canceled caller cannot poison it.

* filer: resolve ring member hostnames in parallel

The shared lookup gave every member one serial budget, so a few slow
resolutions could leave later members out of the cached list and reject
their genuine forwards. Each member now resolves concurrently under its
own detached deadline.

* filer: gather literal member IPs before spawning lookups

A ring mixing IP literals and hostnames raced: the literal appends ran
unlocked alongside the resolver goroutines' locked appends. Split into
two passes so only hostname results share the mutex.

---------

Co-authored-by: jsas <1351492+jsas@users.noreply.github.com>
This commit is contained in:
Chris Lu
2026-09-26 19:42:23 +08:00
committed by GitHub
co-authored by jsas
parent 80fd3635d2
commit 2f641a63d6
7 changed files with 246 additions and 28 deletions
+17 -2
View File
@@ -205,7 +205,8 @@ func (fs *FilerServer) CreateEntry(ctx context.Context, req *filer_pb.CreateEntr
// upsert. Route it to the entry's ring owner so one filer's lock arbitrates
// every creator cluster-wide; is_moved bounds this to one hop. Plain creates
// are upserts either way and stay local.
if !req.IsMoved && (req.OExcl || conditionIsSet(req.Condition)) {
routed := req.OExcl || conditionIsSet(req.Condition)
if !req.IsMoved && routed {
fullpath := util.NewFullPath(req.Directory, req.Entry.Name)
// Held apart from the named resp, which the local path below writes into:
// a failed forward must not leave it nil.
@@ -228,6 +229,10 @@ func (fs *FilerServer) CreateEntry(ctx context.Context, req *filer_pb.CreateEntr
}
return ownerResp, nil
}
} else if req.IsMoved && routed {
if err := fs.checkMovedMarker(ctx, req.IsMoved, fs.writeOwner(entryRouteKey(util.NewFullPath(req.Directory, req.Entry.Name)))); err != nil {
return &filer_pb.CreateEntryResponse{}, err
}
}
chunks, garbage, err2 := fs.cleanupChunks(ctx, util.Join(req.Directory, req.Entry.Name), nil, req.Entry)
@@ -321,6 +326,11 @@ func (fs *FilerServer) ObjectTransaction(ctx context.Context, req *filer_pb.Obje
// serialization point — even when the caller's ring view was stale. is_moved
// bounds this to one hop: a forwarded transaction is applied locally, so two
// filers that disagree on the owner during a ring change cannot loop.
if req.RouteKey != "" {
if err := fs.checkMovedMarker(ctx, req.IsMoved, fs.writeOwner(req.RouteKey)); err != nil {
return &filer_pb.ObjectTransactionResponse{Error: err.Error()}, nil
}
}
if req.RouteKey != "" && !req.IsMoved {
// Rebuild rather than copy the request struct (it carries a mutex); the
// pointer/slice fields are shared since the original is not mutated.
@@ -666,7 +676,8 @@ func (fs *FilerServer) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntr
// per-path lock below only makes atomic on this filer. Route it to the
// entry's ring owner so one filer's lock arbitrates every writer
// cluster-wide; is_moved bounds this to one hop.
if !req.IsMoved && (conditionIsSet(req.Condition) || len(req.ExpectedExtended) > 0) {
routed := conditionIsSet(req.Condition) || len(req.ExpectedExtended) > 0
if !req.IsMoved && routed {
var ownerResp *filer_pb.UpdateEntryResponse
handled, forwardErr := fs.forwardToWriteOwner(ctx, entryRouteKey(util.FullPath(fullpath)), func(owner pb.ServerAddress) error {
glog.V(2).InfofCtx(ctx, "UpdateEntry %s: forwarding to owner %s", fullpath, owner)
@@ -686,6 +697,10 @@ func (fs *FilerServer) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntr
}
return ownerResp, nil
}
} else if req.IsMoved && routed {
if err := fs.checkMovedMarker(ctx, req.IsMoved, fs.writeOwner(entryRouteKey(util.FullPath(fullpath)))); err != nil {
return &filer_pb.UpdateEntryResponse{}, err
}
}
// Serialize concurrent mutations to the same path on this filer so the
@@ -3,10 +3,14 @@ package weed_server
import (
"context"
"fmt"
"net"
"testing"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/peer"
"google.golang.org/grpc/status"
"github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager"
"github.com/seaweedfs/seaweedfs/weed/pb"
@@ -89,6 +93,14 @@ func TestCreateEntryPlainStaysLocal(t *testing.T) {
}
}
// ringPeerCtx simulates a request arriving over a connection from a ring
// member, the provenance an is_moved marker needs before it is trusted.
func ringPeerCtx() context.Context {
return peer.NewContext(context.Background(), &peer.Peer{
Addr: &net.TCPAddr{IP: net.ParseIP("127.0.0.1"), Port: 49999},
})
}
// is_moved bounds forwarding to one hop: the owner applies it even if its own
// ring view says someone else owns the key.
func TestCreateEntryMovedAppliesLocally(t *testing.T) {
@@ -97,7 +109,7 @@ func TestCreateEntryMovedAppliesLocally(t *testing.T) {
req := createReq(name, true)
req.IsMoved = true
resp, err := fs.CreateEntry(context.Background(), req)
resp, err := fs.CreateEntry(ringPeerCtx(), req)
if err != nil || resp.Error != "" {
t.Fatalf("forwarded create must be applied locally, got err=%v resp=%v", err, resp.Error)
}
@@ -106,6 +118,53 @@ func TestCreateEntryMovedAppliesLocally(t *testing.T) {
}
}
// A caller can set is_moved, but without a ring member's connection it is not
// trusted: on a non-owner the request is refused rather than evaluated under
// the wrong lock or forwarded again, where it could cycle while rings disagree.
func TestCreateEntryMovedFromClientIsNotTrusted(t *testing.T) {
fs, store := createRouteServer(t)
name := peerOwnedName(t, fs)
req := createReq(name, true)
req.IsMoved = true
_, err := fs.CreateEntry(context.Background(), req)
if status.Code(err) != codes.PermissionDenied {
t.Fatalf("a forged is_moved must be refused, got err=%v", err)
}
if _, found := store.entries[string(util.NewFullPath("/test", name))]; found {
t.Fatal("a forged is_moved must not cause a local write")
}
}
// Forwarding an unlock to an unreachable owner must surface the RPC error, not
// panic dereferencing the nil response of the failed call.
func TestDistributedUnlockFailedForwardReturnsError(t *testing.T) {
fs, _ := createRouteServer(t)
var name string
for i := 0; i < 4000; i++ {
candidate := fmt.Sprintf("lock-%d", i)
if fs.filer.Dlm.LockRing.GetPrimary(candidate) != fs.option.Host {
name = candidate
break
}
}
if name == "" {
t.Skip("no lock owned by the peer")
}
resp, err := fs.DistributedUnlock(context.Background(), &filer_pb.UnlockRequest{
Name: name,
RenewToken: "tok",
})
if err != nil {
t.Fatalf("unlock forward should report via resp.Error, not err: %v", err)
}
if resp.Error == "" {
t.Fatal("a failed forward must surface an error")
}
}
// An exclusive create this filer owns is applied locally, under its per-path lock.
func TestCreateEntryExclusiveAppliedByOwner(t *testing.T) {
fs, store := createRouteServer(t)
+13 -3
View File
@@ -29,7 +29,9 @@ func (fs *FilerServer) DistributedLock(ctx context.Context, req *filer_pb.LockRe
glog.V(4).Infof("FILER LOCK: LockWithTimeout result - name=%s lockOwner=%s renewToken=%s movedTo=%s err=%v",
req.Name, resp.LockOwner, resp.RenewToken, movedTo, err)
glog.V(4).Infof("lock %s %v %v %v, isMoved=%v %v", req.Name, req.SecondsToLock, req.RenewToken, req.Owner, req.IsMoved, movedTo)
if movedTo != "" && movedTo != fs.option.Host && !req.IsMoved {
if mErr := fs.checkMovedMarker(ctx, req.IsMoved, movedTo); mErr != nil {
err = mErr
} else if !req.IsMoved && movedTo != "" && movedTo != fs.option.Host {
glog.V(0).Infof("FILER LOCK: Forwarding to correct filer - from=%s to=%s", fs.option.Host, movedTo)
err = pb.WithFilerClient(false, 0, movedTo, fs.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
secondResp, err := client.DistributedLock(ctx, &filer_pb.LockRequest{
@@ -76,15 +78,20 @@ func (fs *FilerServer) DistributedUnlock(ctx context.Context, req *filer_pb.Unlo
var movedTo pb.ServerAddress
movedTo, err = fs.filer.Dlm.Unlock(req.Name, req.RenewToken)
if !req.IsMoved && movedTo != "" {
if mErr := fs.checkMovedMarker(ctx, req.IsMoved, movedTo); mErr != nil {
err = mErr
} else if !req.IsMoved && movedTo != "" {
err = pb.WithFilerClient(false, 0, movedTo, fs.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
secondResp, err := client.DistributedUnlock(ctx, &filer_pb.UnlockRequest{
Name: req.Name,
RenewToken: req.RenewToken,
IsMoved: true,
})
if err != nil {
return err
}
resp.Error = secondResp.Error
return err
return nil
})
}
@@ -101,6 +108,9 @@ func (fs *FilerServer) DistributedUnlock(ctx context.Context, req *filer_pb.Unlo
func (fs *FilerServer) FindLockOwner(ctx context.Context, req *filer_pb.FindLockOwnerRequest) (*filer_pb.FindLockOwnerResponse, error) {
owner, movedTo, err := fs.filer.Dlm.FindLockOwner(req.Name)
if mErr := fs.checkMovedMarker(ctx, req.IsMoved, movedTo); mErr != nil {
return nil, mErr
}
if !req.IsMoved && movedTo != "" {
err = pb.WithFilerClient(false, 0, movedTo, fs.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
secondResp, err := client.FindLockOwner(ctx, &filer_pb.FindLockOwnerRequest{
@@ -769,10 +769,11 @@ func TestObjectTransactionRouteKeyOwnerAppliesLocally(t *testing.T) {
}
}
// A forwarded transaction (is_moved) applies locally even when the ring names a
// different owner: is_moved bounds forwarding to a single hop, so two filers that
// disagree on the owner during a ring change cannot loop. If is_moved were
// ignored, this would attempt to dial the bogus owner instead of applying.
// A forwarded transaction (is_moved over a ring member's connection) applies
// locally even when the ring names a different owner: is_moved bounds
// forwarding to a single hop, so two filers that disagree on the owner during
// a ring change cannot loop. Without the marker this would attempt to dial
// the bogus owner instead of applying.
func TestObjectTransactionIsMovedSkipsForward(t *testing.T) {
self := pb.ServerAddress("localhost:1")
other := pb.ServerAddress("localhost:2")
@@ -781,7 +782,7 @@ func TestObjectTransactionIsMovedSkipsForward(t *testing.T) {
})
withRing(fs, self, other) // ring owner is "other", not self
resp, err := fs.ObjectTransaction(context.Background(), &filer_pb.ObjectTransactionRequest{
resp, err := fs.ObjectTransaction(ringPeerCtx(), &filer_pb.ObjectTransactionRequest{
LockKey: "/buckets/b/obj",
RouteKey: "s3.object.write:/buckets/b/obj",
IsMoved: true,
+22 -17
View File
@@ -63,26 +63,31 @@ func (fs *FilerServer) PosixLock(ctx context.Context, req *filer_pb.PosixLockReq
return &filer_pb.PosixLockResponse{}, fmt.Errorf("lock is required")
}
if !req.IsMoved && fs.filer.Dlm != nil {
if fs.filer.Dlm != nil {
if owner := fs.filer.Dlm.LockRing.GetPrimary(req.Key); owner != "" && owner != fs.option.Host {
forwarded := &filer_pb.PosixLockRequest{
Key: req.Key,
IsMoved: true,
Op: req.Op,
Lock: req.Lock,
Locks: req.Locks,
}
glog.V(4).InfofCtx(ctx, "PosixLock %s op=%v: forwarding to owner %s", req.Key, req.Op, owner)
var resp *filer_pb.PosixLockResponse
err := pb.WithFilerClient(false, 0, owner, fs.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
var e error
resp, e = client.PosixLock(ctx, forwarded)
return e
})
if err != nil {
if err := fs.checkMovedMarker(ctx, req.IsMoved, owner); err != nil {
return &filer_pb.PosixLockResponse{}, err
}
return resp, nil
if !req.IsMoved {
forwarded := &filer_pb.PosixLockRequest{
Key: req.Key,
IsMoved: true,
Op: req.Op,
Lock: req.Lock,
Locks: req.Locks,
}
glog.V(4).InfofCtx(ctx, "PosixLock %s op=%v: forwarding to owner %s", req.Key, req.Op, owner)
var resp *filer_pb.PosixLockResponse
err := pb.WithFilerClient(false, 0, owner, fs.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
var e error
resp, e = client.PosixLock(ctx, forwarded)
return e
})
if err != nil {
return &filer_pb.PosixLockResponse{}, err
}
return resp, nil
}
}
}
+123
View File
@@ -2,11 +2,18 @@ package weed_server
import (
"context"
"net"
"strings"
"sync"
"time"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
"github.com/seaweedfs/seaweedfs/weed/util"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/peer"
"google.golang.org/grpc/status"
)
// writeOwner returns the filer that serializes writes to key, or "" when this
@@ -55,3 +62,119 @@ func (fs *FilerServer) forwardToWriteOwner(ctx context.Context, key string, send
func entryRouteKey(fullpath util.FullPath) string {
return s3_constants.ObjectWriteRouteKeyPrefix + string(fullpath)
}
// movedFromPeer reports whether an is_moved marker arrived on a connection
// from a ring member, i.e. it marks a genuine forwarded hop.
func (fs *FilerServer) movedFromPeer(ctx context.Context, isMoved bool) bool {
if !isMoved || fs.filer.Dlm == nil {
return false
}
p, ok := peer.FromContext(ctx)
if !ok {
return false
}
peerHost, _, err := net.SplitHostPort(p.Addr.String())
if err != nil {
return false
}
peerIP := net.ParseIP(peerHost)
if peerIP == nil {
return false
}
for _, ip := range fs.ringMemberIPs(ctx) {
if ip.Equal(peerIP) {
return true
}
}
return false
}
// ringPeerIPs caches resolved member addresses of one ring membership. A
// member's hostname may re-resolve under the same ring address, so the cache
// expires rather than trusting the resolution forever.
type ringPeerIPs struct {
members string
ips []net.IP
expires time.Time
}
const ringPeerIPTTL = 5 * time.Minute
// ringMemberIPs returns the ring members' addresses as IPs. Members can
// advertise hostnames, so resolution is cached per membership to keep DNS off
// each forwarded request; a failed lookup is not cached, so a DNS blip does
// not keep rejecting genuine forwards until the next ring change.
func (fs *FilerServer) ringMemberIPs(ctx context.Context) []net.IP {
members := fs.filer.Dlm.LockRing.GetSnapshot()
var sb strings.Builder
for _, member := range members {
sb.WriteString(string(member))
sb.WriteByte(' ')
}
key := sb.String()
if cached := fs.ringPeerIPs.Load(); cached != nil && cached.members == key && time.Now().Before(cached.expires) {
return cached.ips
}
resolved, _, _ := fs.ringResolveGroup.Do(key, func() (any, error) {
if cached := fs.ringPeerIPs.Load(); cached != nil && cached.members == key && time.Now().Before(cached.expires) {
return cached.ips, nil
}
// Shared by every caller waiting on this key: the lookups outlive the
// first request's cancellation, and run in parallel so one slow member
// cannot starve the rest of the shared deadline.
var ips []net.IP
var hosts []string
for _, member := range members {
host, _, err := net.SplitHostPort(string(member))
if err != nil {
continue
}
if ip := net.ParseIP(host); ip != nil {
ips = append(ips, ip)
continue
}
hosts = append(hosts, host)
}
var wg sync.WaitGroup
var mu sync.Mutex
failed := false
for _, host := range hosts {
wg.Add(1)
go func(host string) {
defer wg.Done()
lookupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 3*time.Second)
defer cancel()
found, err := net.DefaultResolver.LookupIP(lookupCtx, "ip", host)
mu.Lock()
defer mu.Unlock()
if err != nil {
failed = true
return
}
ips = append(ips, found...)
}(host)
}
wg.Wait()
if !failed {
fs.ringPeerIPs.Store(&ringPeerIPs{members: key, ips: ips, expires: time.Now().Add(ringPeerIPTTL)})
}
return ips, nil
})
return resolved.([]net.IP)
}
// checkMovedMarker refuses a request whose is_moved marker did not arrive from
// a ring member while this filer is not the key's owner. The marker is
// caller-controlled: applying it would evaluate a conditional mutation under a
// non-owner's lock, and re-forwarding a claimed hop can cycle while rings
// disagree — so an unverifiable marker on a non-owner is refused instead.
// PermissionDenied keeps the refusal distinct from a write condition's
// FailedPrecondition, which callers use to detect a stale stamp.
// owner=="" means this filer is the serialization point and the request can
// be applied locally.
func (fs *FilerServer) checkMovedMarker(ctx context.Context, isMoved bool, owner pb.ServerAddress) error {
if !isMoved || owner == "" || owner == fs.option.Host || fs.movedFromPeer(ctx, isMoved) {
return nil
}
return status.Errorf(codes.PermissionDenied, "is_moved not sent by a ring member; the key's owner is %s", owner)
}
+5
View File
@@ -136,6 +136,11 @@ type FilerServer struct {
// chunks behind the first request's back.
tusActiveUploads sync.Map
// ringPeerIPs caches resolved ring member addresses per ring version so
// verifying a forwarded request's peer does not pay a DNS lookup per hop.
ringPeerIPs atomic.Pointer[ringPeerIPs]
ringResolveGroup singleflight.Group
// entryLockTable serializes mutations to the same entry path on this filer.
// CreateEntry takes it today; UpdateEntry and DeleteEntry are intended to take
// it too as their callers route a key's writes to this node, making it the