diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index c70b02bc8..9887f6bce 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -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 diff --git a/weed/server/filer_grpc_server_create_route_test.go b/weed/server/filer_grpc_server_create_route_test.go index f5f4dc99b..0c266c539 100644 --- a/weed/server/filer_grpc_server_create_route_test.go +++ b/weed/server/filer_grpc_server_create_route_test.go @@ -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) diff --git a/weed/server/filer_grpc_server_dlm.go b/weed/server/filer_grpc_server_dlm.go index 3ff9b5f44..08427350a 100644 --- a/weed/server/filer_grpc_server_dlm.go +++ b/weed/server/filer_grpc_server_dlm.go @@ -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{ diff --git a/weed/server/filer_grpc_server_object_txn_test.go b/weed/server/filer_grpc_server_object_txn_test.go index 2aaa47e7a..ad2c56129 100644 --- a/weed/server/filer_grpc_server_object_txn_test.go +++ b/weed/server/filer_grpc_server_object_txn_test.go @@ -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, diff --git a/weed/server/filer_grpc_server_posix_lock.go b/weed/server/filer_grpc_server_posix_lock.go index c7024a883..869606ac8 100644 --- a/weed/server/filer_grpc_server_posix_lock.go +++ b/weed/server/filer_grpc_server_posix_lock.go @@ -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 + } } } diff --git a/weed/server/filer_grpc_server_route.go b/weed/server/filer_grpc_server_route.go index d60254bfb..18126eaff 100644 --- a/weed/server/filer_grpc_server_route.go +++ b/weed/server/filer_grpc_server_route.go @@ -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) +} diff --git a/weed/server/filer_server.go b/weed/server/filer_server.go index 959e1ac8c..b1c859d11 100644 --- a/weed/server/filer_server.go +++ b/weed/server/filer_server.go @@ -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