Files
seaweedfs/weed/server/filer_server.go
T
2f641a63d6 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>
2026-09-26 19:42:23 +08:00

388 lines
16 KiB
Go

package weed_server
import (
"context"
"fmt"
"net/http"
"os"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/seaweedfs/seaweedfs/weed/credential"
"github.com/seaweedfs/seaweedfs/weed/stats"
"golang.org/x/sync/singleflight"
"google.golang.org/grpc"
"github.com/seaweedfs/seaweedfs/weed/util/grace"
"github.com/seaweedfs/seaweedfs/weed/operation"
"github.com/seaweedfs/seaweedfs/weed/pb"
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/remote_pb"
"github.com/seaweedfs/seaweedfs/weed/util"
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
"github.com/seaweedfs/seaweedfs/weed/filer"
_ "github.com/seaweedfs/seaweedfs/weed/filer/arangodb"
_ "github.com/seaweedfs/seaweedfs/weed/filer/cassandra"
_ "github.com/seaweedfs/seaweedfs/weed/filer/cassandra2"
_ "github.com/seaweedfs/seaweedfs/weed/filer/elastic/v7"
_ "github.com/seaweedfs/seaweedfs/weed/filer/etcd"
_ "github.com/seaweedfs/seaweedfs/weed/filer/foundationdb"
_ "github.com/seaweedfs/seaweedfs/weed/filer/hbase"
_ "github.com/seaweedfs/seaweedfs/weed/filer/leveldb"
_ "github.com/seaweedfs/seaweedfs/weed/filer/leveldb2"
_ "github.com/seaweedfs/seaweedfs/weed/filer/leveldb3"
_ "github.com/seaweedfs/seaweedfs/weed/filer/mongodb"
_ "github.com/seaweedfs/seaweedfs/weed/filer/mysql"
_ "github.com/seaweedfs/seaweedfs/weed/filer/mysql2"
"github.com/seaweedfs/seaweedfs/weed/filer/posixlock"
_ "github.com/seaweedfs/seaweedfs/weed/filer/postgres"
_ "github.com/seaweedfs/seaweedfs/weed/filer/postgres2"
_ "github.com/seaweedfs/seaweedfs/weed/filer/redis"
_ "github.com/seaweedfs/seaweedfs/weed/filer/redis2"
_ "github.com/seaweedfs/seaweedfs/weed/filer/redis3"
_ "github.com/seaweedfs/seaweedfs/weed/filer/sqlite"
_ "github.com/seaweedfs/seaweedfs/weed/filer/tarantool"
_ "github.com/seaweedfs/seaweedfs/weed/filer/ydb"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/notification"
_ "github.com/seaweedfs/seaweedfs/weed/notification/aws_sqs"
_ "github.com/seaweedfs/seaweedfs/weed/notification/gocdk_pub_sub"
_ "github.com/seaweedfs/seaweedfs/weed/notification/google_pub_sub"
_ "github.com/seaweedfs/seaweedfs/weed/notification/kafka"
_ "github.com/seaweedfs/seaweedfs/weed/notification/log"
_ "github.com/seaweedfs/seaweedfs/weed/notification/webhook"
"github.com/seaweedfs/seaweedfs/weed/security"
)
type FilerOption struct {
Masters *pb.ServerDiscovery
FilerGroup string
Collection string
DefaultReplication string
DisableDirListing bool
MaxMB int
DirListingLimit int
DataCenter string
Rack string
DataNode string
DefaultLevelDbDir string
DisableHttp bool
Host pb.ServerAddress
recursiveDelete bool
Cipher bool
SaveToFilerLimit int64
ConcurrentUploadLimit int64
ConcurrentFileUploadLimit int64
ShowUIDirectoryDelete bool
DownloadMaxBytesPs int64
DiskType string
AllowedOrigins []string
ExposeDirectoryData bool
TusBasePath string
TusMaxSize int64
TusSessionExpiry time.Duration
S3ConfigFile string // optional path to static S3 identity config file
CredentialManager *credential.CredentialManager
// AllowUntrustedRemoteEndpoints lets a read of a remote-only entry dial a
// mounted endpoint that resolves to a loopback / private / metadata host.
AllowUntrustedRemoteEndpoints bool
}
type FilerServer struct {
inFlightDataSize int64
inFlightUploads int64
inFlightDataLimitCond *sync.Cond
filer_pb.UnimplementedSeaweedFilerServer
option *FilerOption
filer *filer.Filer
filerGuard *security.Guard
volumeGuard *security.Guard
grpcDialOption grpc.DialOption
// metrics read from the master
metricsAddress string
metricsIntervalSec int
// track known metadata listeners
knownListenersLock sync.Mutex
knownListeners map[int32]int32
// live metadata subscribers (FUSE mounts, S3, peer filers, ...) keyed by
// clientId, guarded by knownListenersLock. Exposed via ListMetadataSubscribers.
subscribers map[int32]*metadataSubscriber
// deduplicates concurrent remote object caching operations
remoteCacheGroup singleflight.Group
recentCopyRequestsMu sync.Mutex
recentCopyRequests map[string]recentCopyRequest
// credential manager for IAM operations
CredentialManager *credential.CredentialManager
// mountPeerRegistry backs the MountRegister / MountList RPCs for peer
// chunk sharing (tier 1). Always populated.
mountPeerRegistry *filer.MountPeerRegistry
// tusActiveUploads marks TUS sessions with a mutating request in flight, so
// a concurrent PATCH or DELETE is refused instead of recording duplicate
// 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
// local serialization point for read-modify-write operations that replaces
// the distributed lock for that key. Idle keys are evicted automatically, so
// the table stays bounded.
entryLockTable *util.LockTable[util.FullPath]
// posixLocks is the in-memory authority for cross-mount POSIX advisory locks
// on inodes this filer owns (per the route-by-key ring). Lock state is kept
// here rather than in replicated metadata: it is transient coordination, so
// keeping it off the meta-log avoids churn.
posixLocks *posixlock.Manager
// posixLockSweeperStop stops the lease-reaping sweeper goroutine on Shutdown.
posixLockSweeperStop chan struct{}
// posixLockReadyAt is the unix-nanos when this filer began serving POSIX
// locks. For posixLockWarmup after it, the owner defers would-be grants while
// mounts re-assert, so a (re)started owner does not double-grant from empty
// state. Atomic so the handler reads it without locking; 0 means "not warming
// up" (e.g. in tests).
posixLockReadyAt atomic.Int64
}
func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption) (fs *FilerServer, err error) {
v := util.GetViper()
signingKey := v.GetString("jwt.filer_signing.key")
v.SetDefault("jwt.filer_signing.expires_after_seconds", 10)
expiresAfterSec := v.GetInt("jwt.filer_signing.expires_after_seconds")
readSigningKey := v.GetString("jwt.filer_signing.read.key")
v.SetDefault("jwt.filer_signing.read.expires_after_seconds", 60)
readExpiresAfterSec := v.GetInt("jwt.filer_signing.read.expires_after_seconds")
volumeSigningKey := v.GetString("jwt.signing.key")
v.SetDefault("jwt.signing.expires_after_seconds", 10)
volumeExpiresAfterSec := v.GetInt("jwt.signing.expires_after_seconds")
volumeReadSigningKey := v.GetString("jwt.signing.read.key")
v.SetDefault("jwt.signing.read.expires_after_seconds", 60)
volumeReadExpiresAfterSec := v.GetInt("jwt.signing.read.expires_after_seconds")
v.SetDefault("cors.allowed_origins.values", "*")
allowedOrigins := v.GetString("cors.allowed_origins.values")
domains := strings.Split(allowedOrigins, ",")
option.AllowedOrigins = domains
// -exposeDirectoryData and filer.expose_directory_metadata both default to
// on, and either one turning it off has to hold: this is what keeps the
// directory listing off a filer whose reads are otherwise unauthenticated.
v.SetDefault("filer.expose_directory_metadata.enabled", true)
option.ExposeDirectoryData = option.ExposeDirectoryData && v.GetBool("filer.expose_directory_metadata.enabled")
fs = &FilerServer{
option: option,
grpcDialOption: security.LoadClientTLS(util.GetViper(), "grpc.filer"),
knownListeners: make(map[int32]int32),
subscribers: make(map[int32]*metadataSubscriber),
inFlightDataLimitCond: sync.NewCond(new(sync.Mutex)),
recentCopyRequests: make(map[string]recentCopyRequest),
CredentialManager: option.CredentialManager,
entryLockTable: util.NewLockTable[util.FullPath](),
posixLocks: posixlock.NewManager(),
}
fs.startPosixLockSweeper()
fs.mountPeerRegistry = filer.NewMountPeerRegistry()
go fs.runMountPeerRegistrySweeper()
option.Masters.RefreshBySrvIfAvailable()
if len(option.Masters.GetInstances()) == 0 {
glog.Fatal("master list is required!")
}
if !util.LoadConfiguration("filer", false) {
v.SetDefault("leveldb2.enabled", true)
v.SetDefault("leveldb2.dir", option.DefaultLevelDbDir)
_, err := os.Stat(option.DefaultLevelDbDir)
if os.IsNotExist(err) {
os.MkdirAll(option.DefaultLevelDbDir, 0755)
}
glog.V(0).Infof("default to create filer store dir in %s", option.DefaultLevelDbDir)
} else {
glog.Warningf("skipping default store dir in %s", option.DefaultLevelDbDir)
}
util.LoadConfiguration("notification", false)
v.SetDefault("filer.options.max_file_name_length", 255)
maxFilenameLength := v.GetUint32("filer.options.max_file_name_length")
glog.V(0).Infof("max_file_name_length %d", maxFilenameLength)
fs.filer = filer.NewFiler(*option.Masters, fs.grpcDialOption, option.Host, option.FilerGroup, option.Collection, option.DefaultReplication, option.DataCenter, maxFilenameLength, nil)
fs.filer.Cipher = option.Cipher
fs.filer.DefaultDiskType = option.DiskType
fs.filer.BuildGuardedRemoteClient = BuildGuardedRemoteStorageClient
fs.filer.AllowUntrustedRemoteEndpoints = option.AllowUntrustedRemoteEndpoints
fs.filer.RemoteStorage.SetConfValidator(func(ctx context.Context, conf *remote_pb.RemoteConf) error {
return ValidateRemoteConfForLoad(ctx, conf, option.AllowUntrustedRemoteEndpoints)
})
// we do not support IP whitelist right now https://github.com/seaweedfs/seaweedfs/issues/7094
if v.GetString("guard.white_list") != "" {
glog.Warningf("filer: guard.white_list is configured but the IP whitelist feature is currently disabled. See https://github.com/seaweedfs/seaweedfs/issues/7094")
}
fs.filerGuard = security.NewGuard([]string{}, signingKey, expiresAfterSec, readSigningKey, readExpiresAfterSec)
fs.volumeGuard = security.NewGuard([]string{}, volumeSigningKey, volumeExpiresAfterSec, volumeReadSigningKey, volumeReadExpiresAfterSec)
fs.checkWithMaster()
go stats.LoopPushingMetric("filer", string(fs.option.Host), fs.metricsAddress, fs.metricsIntervalSec)
go fs.filer.MasterClient.KeepConnectedToMaster(context.Background())
fs.option.recursiveDelete = v.GetBool("filer.options.recursive_delete")
v.SetDefault("filer.options.buckets_folder", "/buckets")
fs.filer.DirBucketsPath = v.GetString("filer.options.buckets_folder")
// TODO deprecated, will be removed after 2020-12-31
// replaced by https://github.com/seaweedfs/seaweedfs/wiki/Path-Specific-Configuration
// fs.filer.FsyncBuckets = v.GetStringSlice("filer.options.buckets_fsync")
isFresh := fs.filer.LoadConfiguration(v)
notification.LoadConfiguration(v, "notification.")
handleStaticResources(defaultMux)
if !option.DisableHttp {
defaultMux.HandleFunc("/healthz", requestIDMiddleware(fs.filerHealthzHandler))
defaultMux.HandleFunc("/readyz", requestIDMiddleware(fs.filerHealthzHandler))
// TUS resumable upload protocol handler
if option.TusBasePath != "" {
// Normalize TusPath to always have a leading slash and no trailing slash
if !strings.HasPrefix(option.TusBasePath, "/") {
option.TusBasePath = "/" + option.TusBasePath
}
option.TusBasePath = strings.TrimRight(option.TusBasePath, "/")
// Disallow using "/" as TUS base to avoid hijacking all filer routes
if option.TusBasePath == "" {
glog.Warningf("Invalid TUS base path; TUS disabled (must not be root '/')")
} else {
if option.TusMaxSize <= 0 {
option.TusMaxSize = TusDefaultMaxSize
}
if option.TusSessionExpiry <= 0 {
option.TusSessionExpiry = TusDefaultSessionExpiry
}
handlePath := option.TusBasePath + "/"
defaultMux.HandleFunc(handlePath, fs.filerGuard.WhiteList(requestIDMiddleware(fs.tusHandler)))
// Start background cleanup of expired TUS sessions (every hour)
fs.StartTusSessionCleanup(1 * time.Hour)
}
}
defaultMux.HandleFunc("/", fs.filerGuard.WhiteList(requestIDMiddleware(fs.filerHandler)))
}
if defaultMux != readonlyMux {
handleStaticResources(readonlyMux)
readonlyMux.HandleFunc("/healthz", requestIDMiddleware(fs.filerHealthzHandler))
readonlyMux.HandleFunc("/readyz", requestIDMiddleware(fs.filerHealthzHandler))
readonlyMux.HandleFunc("/", fs.filerGuard.WhiteList(requestIDMiddleware(fs.readonlyFilerHandler)))
}
existingNodes := fs.filer.ListExistingPeerUpdates(context.Background())
startFromTime := time.Now().Add(-filer.LogFlushInterval)
if isFresh {
glog.V(0).Infof("%s bootstrap from peers %+v", option.Host, existingNodes)
if err := fs.filer.MaybeBootstrapFromOnePeer(option.Host, existingNodes, startFromTime); err != nil {
glog.Fatalf("%s bootstrap from %+v: %v", option.Host, existingNodes, err)
}
}
v.SetDefault("filer.options.s3.empty_folder_cleanup_delay", "2m")
if d, err := time.ParseDuration(v.GetString("filer.options.s3.empty_folder_cleanup_delay")); err == nil {
fs.filer.EmptyFolderCleanupDelay = d
}
fs.filer.AggregateFromPeers(option.Host, existingNodes, startFromTime)
fs.filer.LoadFilerConf()
fs.filer.LoadRemoteStorageConfAndMapping()
fs.filer.RebuildRemoteDeletionTombstones(context.Background())
grace.OnReload(fs.Reload)
fs.SetupDlmReplication()
fs.filer.Dlm.LockRing.SetTakeSnapshotCallback(fs.OnDlmChangeSnapshot)
if fs.CredentialManager != nil {
fs.CredentialManager.SetFilerAddressFunc(func() pb.ServerAddress {
return fs.option.Host
}, fs.grpcDialOption)
fs.CredentialManager.SetMasterClient(fs.filer.MasterClient, fs.grpcDialOption)
}
return fs, nil
}
func (fs *FilerServer) checkWithMaster() {
isConnected := false
for !isConnected {
fs.option.Masters.RefreshBySrvIfAvailable()
for _, master := range fs.option.Masters.GetInstances() {
readErr := operation.WithMasterServerClient(context.Background(), false, master, fs.grpcDialOption, func(masterClient master_pb.SeaweedClient) error {
resp, err := masterClient.GetMasterConfiguration(context.Background(), &master_pb.GetMasterConfigurationRequest{})
if err != nil {
return fmt.Errorf("get master %s configuration: %v", master, err)
}
fs.metricsAddress, fs.metricsIntervalSec = resp.MetricsAddress, int(resp.MetricsIntervalSeconds)
return nil
})
if readErr == nil {
isConnected = true
} else {
time.Sleep(7 * time.Second)
}
}
}
}
// Shutdown gracefully shuts down the filer server by waiting for in-flight uploads to complete.
// This prevents data corruption when the process receives SIGTERM during active uploads.
func (fs *FilerServer) Shutdown() {
glog.V(0).Infof("Shutting down filer")
if fs.posixLockSweeperStop != nil {
close(fs.posixLockSweeperStop)
}
fs.filer.Shutdown()
}
func (fs *FilerServer) Reload() {
glog.V(0).Infoln("Reload filer server...")
util.LoadConfiguration("security", false)
v := util.GetViper()
fs.filerGuard.UpdateSigningKeys(
v.GetString("jwt.filer_signing.key"),
v.GetInt("jwt.filer_signing.expires_after_seconds"),
v.GetString("jwt.filer_signing.read.key"),
v.GetInt("jwt.filer_signing.read.expires_after_seconds"),
)
fs.volumeGuard.UpdateSigningKeys(
v.GetString("jwt.signing.key"),
v.GetInt("jwt.signing.expires_after_seconds"),
v.GetString("jwt.signing.read.key"),
v.GetInt("jwt.signing.read.expires_after_seconds"),
)
util_http.ReloadJwtSigningReadConfig()
}