From 9a73e222ca84d4d58afdd6f91d743eaab120e3a6 Mon Sep 17 00:00:00 2001 From: Ben McClelland Date: Fri, 28 Aug 2026 09:55:12 -0700 Subject: [PATCH] fix: harden posix conditional publish locking Fail closed on unexpected advisory-lock errors instead of silently reducing cross-process exclusion to a local mutex. Make local publish-slot waits honor request cancellation. Move version snapshots under the publish lock so concurrent versioned PUTs preserve publication order. Add regression coverage for canceled lock waiters. Also add an option to disable flock files and only rely on in process locking. --- backend/posix/objlock.go | 74 ++++++++---- backend/posix/objlock_unix.go | 5 + backend/posix/objlock_windows.go | 5 + backend/posix/posix.go | 67 ++++++----- backend/posix/posix_conditional_put_test.go | 124 ++++++++++++++++++++ backend/scoutfs/scoutfs_compat.go | 1 + cmd/internal/gwcli/posix.go | 8 ++ cmd/versitygw/main.go | 9 +- extra/example.conf | 10 ++ 9 files changed, 250 insertions(+), 53 deletions(-) diff --git a/backend/posix/objlock.go b/backend/posix/objlock.go index a1b2b255..fc926ff8 100644 --- a/backend/posix/objlock.go +++ b/backend/posix/objlock.go @@ -38,13 +38,15 @@ import ( // contention within one process is resolved cheaply and each process presents // at most one waiter to the filesystem lock. // -// Lock identity: bucket/.sgwtmp/objlock/, where shard is the first -// byte of sha256(object key) rendered as two hex characters. Hashing the key -// gives a stable, traversal-safe, fixed-length name; sharding (256 slots per -// bucket) keeps the number of lock files bounded while still letting writes -// to unrelated keys proceed concurrently in the common case. Lock files are -// empty, created on demand, and never unlinked: unlinking an flock file opens -// a classic race where a waiter holds a lock on an unlinked inode while a new +// Lock identity: /.vgwlocks//, where shard is the +// first byte of sha256(object key) rendered as two hex characters. Hashing the +// bucket and key gives stable, traversal-safe, fixed-length names; sharding +// (256 slots per bucket) keeps the number of lock files bounded while still +// letting writes to unrelated keys proceed concurrently in the common case. +// Lock files are outside the bucket tree so Windows bucket deletion cannot +// fail merely because an active request has the lock file open. They are empty, +// created on demand, and never unlinked: unlinking an flock file opens a +// classic race where a waiter holds a lock on an unlinked inode while a new // file takes its place, which is unsafe to detect reliably on NFS due to // attribute caching. // @@ -63,8 +65,8 @@ import ( // process-local exclusion and logs a warning once. const ( - // objLockDir is the per-bucket directory holding object publish lock files - objLockDir = MetaTmpDir + "/objlock" + // objLockDir is the root directory holding object publish lock files. + objLockDir = ".vgwlocks" // objLockShards is the number of lock shards per bucket objLockShards = 256 ) @@ -83,12 +85,24 @@ func objLockShard(object string) uint8 { func (p *Posix) lockObjectPublish(ctx context.Context, bucket, object string) (func(), error) { shard := objLockShard(object) - mu := &p.objLockMus[shard] - mu.Lock() + slot := p.objLockSlots[shard] + select { + case <-ctx.Done(): + return nil, ctx.Err() + case <-slot: + } + releaseLocal := func() { slot <- struct{}{} } + if err := ctx.Err(); err != nil { + releaseLocal() + return nil, err + } + if p.forceNoObjLockFile { + return releaseLocal, nil + } f, err := p.openObjLockFile(bucket, shard) if err != nil { - mu.Unlock() + releaseLocal() return nil, err } @@ -96,33 +110,47 @@ func (p *Posix) lockObjectPublish(ctx context.Context, bucket, object string) (f if err != nil { f.Close() if ctx.Err() != nil { - mu.Unlock() + releaseLocal() return nil, ctx.Err() } - // The filesystem does not support advisory locking (e.g. NFS - // mounted with -o nolock). Fall back to process-local exclusion - // and warn once: conditional writes are then only atomic within - // this gateway process. + if !isAdvisoryLockUnsupported(err) { + releaseLocal() + return nil, fmt.Errorf("lock object publish: %w", err) + } + // The filesystem does not support advisory locking (e.g. NFS mounted + // with -o nolock). Fall back to process-local exclusion and warn once: + // conditional writes are then only atomic within this gateway process. p.objLockWarn.Do(func() { debuglogger.Logf("object lock file locking unavailable (%v): "+ "conditional write atomicity limited to this process", err) }) - return mu.Unlock, nil + return releaseLocal, nil } return func() { // closing the file releases the advisory lock f.Close() - mu.Unlock() + releaseLocal() }, nil } +func newObjLockSlots() [objLockShards]chan struct{} { + var slots [objLockShards]chan struct{} + for i := range slots { + slots[i] = make(chan struct{}, 1) + slots[i] <- struct{}{} + } + return slots +} + // openObjLockFile opens (creating as needed) the lock file for the shard in // the given bucket. func (p *Posix) openObjLockFile(bucket string, shard uint8) (*os.File, error) { - name := filepath.Join(bucket, objLockDir, fmt.Sprintf("%02x", shard)) + bucketHash := sha256.Sum256([]byte(bucket)) + lockDir := filepath.Join(p.rootdir, objLockDir, fmt.Sprintf("%x", bucketHash)) + name := filepath.Join(lockDir, fmt.Sprintf("%02x", shard)) - f, err := os.OpenFile(name, os.O_RDWR|os.O_CREATE, os.FileMode(defaultNewFilePerm)) + f, err := os.OpenFile(name, os.O_RDWR|os.O_CREATE, defaultNewFilePerm) if err == nil { return f, nil } @@ -131,11 +159,11 @@ func (p *Posix) openObjLockFile(bucket string, shard uint8) (*os.File, error) { } // lock dir not created yet - err = backend.MkdirAll(filepath.Join(bucket, objLockDir), 0, 0, false, p.newDirPerm) + err = backend.MkdirAll(lockDir, 0, 0, false, p.newDirPerm) if err != nil { return nil, fmt.Errorf("make object lock dir: %w", err) } - f, err = os.OpenFile(name, os.O_RDWR|os.O_CREATE, os.FileMode(defaultNewFilePerm)) + f, err = os.OpenFile(name, os.O_RDWR|os.O_CREATE, defaultNewFilePerm) if err != nil { return nil, fmt.Errorf("open object lock file: %w", err) } diff --git a/backend/posix/objlock_unix.go b/backend/posix/objlock_unix.go index fe667723..3078dd63 100644 --- a/backend/posix/objlock_unix.go +++ b/backend/posix/objlock_unix.go @@ -49,3 +49,8 @@ func lockFileExclusive(ctx context.Context, f *os.File) error { backoff = min(backoff*2, objLockMaxBackoff) } } + +func isAdvisoryLockUnsupported(err error) bool { + return errors.Is(err, unix.ENOTSUP) || errors.Is(err, unix.EOPNOTSUPP) || + errors.Is(err, unix.ENOSYS) +} diff --git a/backend/posix/objlock_windows.go b/backend/posix/objlock_windows.go index 5a6740ee..c5e4a306 100644 --- a/backend/posix/objlock_windows.go +++ b/backend/posix/objlock_windows.go @@ -50,3 +50,8 @@ func lockFileExclusive(ctx context.Context, f *os.File) error { backoff = min(backoff*2, objLockMaxBackoff) } } + +func isAdvisoryLockUnsupported(err error) bool { + return errors.Is(err, windows.ERROR_INVALID_FUNCTION) || + errors.Is(err, windows.ERROR_NOT_SUPPORTED) +} diff --git a/backend/posix/posix.go b/backend/posix/posix.go index e4a6e945..834b4561 100644 --- a/backend/posix/posix.go +++ b/backend/posix/posix.go @@ -92,6 +92,10 @@ type Posix struct { // support copy_file_range is mounted over NFSv4.2. forceNoCopyFileRange bool + // forceNoObjLockFile disables the shared advisory lock file used for + // conditional object publishes. Process-local serialization remains active. + forceNoObjLockFile bool + // enableODirect is a flag to open object data files with O_DIRECT. // This is best-effort and falls back to buffered I/O when unsupported. enableODirect bool @@ -132,11 +136,10 @@ type Posix struct { // multipart composite: "ALGO--" dataIntegrityEtag bool - // objLockMus are process-local striped mutexes taken alongside the - // shared per-bucket lock files by lockObjectPublish so that in-process - // contention on an object's publish lock is resolved without filesystem - // lock thrash. See objlock.go. - objLockMus [objLockShards]sync.Mutex + // objLockSlots are process-local striped slots taken alongside the shared + // per-bucket lock files by lockObjectPublish. Channels make local lock + // acquisition cancelable while avoiding filesystem lock thrash. + objLockSlots [objLockShards]chan struct{} // objLockWarn ensures the advisory-locking-unavailable warning is // logged at most once objLockWarn sync.Once @@ -234,6 +237,10 @@ type PosixOpts struct { ForceNoTmpFile bool // ForceNoCopyFileRange disables the use of io.Copy for multipart uploads parts ForceNoCopyFileRange bool + // ForceNoObjLockFile disables the shared advisory lock file used for + // conditional object publishes. Conditional writes are only atomic within a + // single gateway process when enabled. + ForceNoObjLockFile bool // EnableODirect enables best-effort O_DIRECT for object data reads/writes. // Disabled by default. EnableODirect bool @@ -273,26 +280,25 @@ type PosixOpts struct { func New(rootdir string, meta meta.MetadataStorer, opts PosixOpts) (*Posix, error) { ioBufferSize := ioBufferSizeOrDefault(opts.IOBufferSize) + rootdirAbs, err := filepath.Abs(rootdir) + if err != nil { + return nil, fmt.Errorf("get absolute path of %v: %w", rootdir, err) + } if opts.SideCarDir != "" && strings.HasPrefix(opts.SideCarDir, rootdir) { return nil, fmt.Errorf("sidecar directory cannot be inside the gateway root directory") } - err := os.Chdir(rootdir) + err = os.Chdir(rootdirAbs) if err != nil { return nil, fmt.Errorf("chdir %v: %w", rootdir, err) } - f, err := os.Open(rootdir) + f, err := os.Open(rootdirAbs) if err != nil { return nil, fmt.Errorf("open %v: %w", rootdir, err) } - rootdirAbs, err := filepath.Abs(rootdir) - if err != nil { - return nil, fmt.Errorf("get absolute path of %v: %w", rootdir, err) - } - var versioningdirAbs string // Ensure the versioning directory isn't within the root directory if opts.VersioningDir != "" { @@ -345,7 +351,7 @@ func New(rootdir string, meta meta.MetadataStorer, opts PosixOpts) (*Posix, erro return &Posix{ meta: meta, rootfd: f, - rootdir: rootdir, + rootdir: rootdirAbs, euid: euid, egid: egid, chownuid: opts.ChownUID, @@ -356,6 +362,7 @@ func New(rootdir string, meta meta.MetadataStorer, opts PosixOpts) (*Posix, erro newFilePerm: newFilePerm, forceNoTmpFile: opts.ForceNoTmpFile, forceNoCopyFileRange: opts.ForceNoCopyFileRange, + forceNoObjLockFile: opts.ForceNoObjLockFile, enableODirect: opts.EnableODirect, validateBucketName: opts.ValidateBucketNames, actionLimiter: semaphore.NewWeighted(int64(concurrencyOrDefault(opts.Concurrency))), @@ -367,6 +374,7 @@ func New(rootdir string, meta meta.MetadataStorer, opts PosixOpts) (*Posix, erro return &b }}, dataIntegrityEtag: opts.DataIntegrityEtag, + objLockSlots: newObjLockSlots(), }, nil } @@ -602,6 +610,10 @@ func (p *Posix) ListBuckets(ctx context.Context, input s3response.ListBucketsInp } func (p *Posix) isBucketValid(bucket string) bool { + if bucket == objLockDir { + return false + } + if !p.validateBucketName { return true } @@ -4253,17 +4265,9 @@ func (p *Posix) PutObjectWithPostFunc(ctx context.Context, po s3response.PutObje return s3response.PutObjectOutput{}, s3err.GetAPIError(s3err.ErrExistingObjectIsDirectory) } - // If versioning is enabled, first create the file object version. - // Conditional writes defer the snapshot until the precondition has been - // confirmed under the object publish lock, so that a losing conditional - // write does not create a spurious version snapshot. + // Version snapshots are taken under the object publish lock below. This + // keeps concurrent versioned overwrites ordered with their publications. conditional := po.IfMatch != nil || po.IfNoneMatch != nil - if err == nil && !conditional { - verr := p.snapshotObjVersion(*po.Bucket, *po.Key, vStatus, acct) - if verr != nil { - return s3response.PutObjectOutput{}, verr - } - } if isErrNameTooLong(err) { return s3response.PutObjectOutput{}, s3err.GetKeyTooLongErr(int64(len(*po.Key)), 1024) } @@ -4386,12 +4390,13 @@ func (p *Posix) PutObjectWithPostFunc(ctx context.Context, po s3response.PutObje return s3response.PutObjectOutput{}, err } - // snapshot the object version that is about to be replaced, now - // that this writer is known to win the conditional race - verr := p.snapshotObjVersion(*po.Bucket, *po.Key, vStatus, acct) - if verr != nil { - return s3response.PutObjectOutput{}, verr - } + } + + // Snapshot the object that is about to be replaced only after the + // conditional check has succeeded and while holding the publish lock. + verr := p.snapshotObjVersion(*po.Bucket, *po.Key, vStatus, acct) + if verr != nil { + return s3response.PutObjectOutput{}, verr } // Before finalizing the object creation remove @@ -7201,6 +7206,10 @@ func listBucketFileInfos(bucketlinks bool) ([]fs.FileInfo, error) { var fis []fs.FileInfo for _, entry := range entries { + if entry.Name() == objLockDir { + continue + } + fi, err := entry.Info() if err != nil { continue diff --git a/backend/posix/posix_conditional_put_test.go b/backend/posix/posix_conditional_put_test.go index ad94362e..6617db6d 100644 --- a/backend/posix/posix_conditional_put_test.go +++ b/backend/posix/posix_conditional_put_test.go @@ -27,6 +27,7 @@ import ( "strings" "sync" "testing" + "time" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/service/s3" @@ -79,6 +80,129 @@ func createTestBucket(t *testing.T, p *Posix, bucket string) { } } +func TestObjectPublishLockHonorsContextWhileWaiting(t *testing.T) { + p := newTestPosix(t, metaModes(t)["xattr"]) + bucket := "testbucket" + createTestBucket(t, p, bucket) + + shard := objLockShard("cancel-wait") + <-p.objLockSlots[shard] + defer func() { p.objLockSlots[shard] <- struct{}{} }() + if _, err := os.Stat(filepath.Join(bucket, objLockDir)); !errors.Is(err, fs.ErrNotExist) { + t.Fatalf("bucket contains publish lock directory: %v", err) + } + + ctx, cancel := context.WithCancel(context.Background()) + result := make(chan error, 1) + go func() { + unlock, err := p.lockObjectPublish(ctx, bucket, "cancel-wait") + if unlock != nil { + unlock() + } + result <- err + }() + + select { + case err := <-result: + t.Fatalf("lockObjectPublish returned before cancellation: %v", err) + case <-time.After(10 * time.Millisecond): + } + cancel() + + var err error + select { + case err = <-result: + case <-time.After(time.Second): + t.Fatal("lockObjectPublish did not honor cancellation") + } + if !errors.Is(err, context.Canceled) { + t.Fatalf("lockObjectPublish error = %v, want context.Canceled", err) + } +} + +func TestObjectPublishLockHonorsCancellationAfterSlotAcquired(t *testing.T) { + p := newTestPosix(t, metaModes(t)["xattr"]) + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + unlock, err := p.lockObjectPublish(ctx, "testbucket", "cancel-ready-slot") + if unlock != nil { + unlock() + } + if !errors.Is(err, context.Canceled) { + t.Fatalf("lockObjectPublish error = %v, want context.Canceled", err) + } +} + +func TestObjectPublishLockCanDisableSharedLockFile(t *testing.T) { + root := t.TempDir() + p, err := New(root, meta.XattrMeta{}, PosixOpts{ + NewDirPerm: 0755, + ForceNoObjLockFile: true, + }) + if err != nil { + t.Fatalf("new posix: %v", err) + } + defer p.Shutdown() + + unlock, err := p.lockObjectPublish(context.Background(), "testbucket", "object") + if err != nil { + t.Fatalf("lock object publish: %v", err) + } + unlock() + + if _, err := os.Stat(filepath.Join(root, objLockDir)); !errors.Is(err, fs.ErrNotExist) { + t.Fatalf("shared lock directory exists or stat failed: %v", err) + } +} + +func TestObjectPublishLockUsesRootForRelativeBackendPath(t *testing.T) { + parent := t.TempDir() + rootName := "gateway-root" + root := filepath.Join(parent, rootName) + if err := os.Mkdir(root, 0755); err != nil { + t.Fatalf("make root: %v", err) + } + t.Chdir(parent) + + p, err := New(rootName, meta.XattrMeta{}, PosixOpts{NewDirPerm: 0755}) + if err != nil { + t.Fatalf("new posix: %v", err) + } + defer p.Shutdown() + + unlock, err := p.lockObjectPublish(context.Background(), "testbucket", "relative-root") + if err != nil { + t.Fatalf("lock object publish: %v", err) + } + unlock() + + if _, err := os.Stat(filepath.Join(root, objLockDir)); err != nil { + t.Fatalf("root lock directory: %v", err) + } + if _, err := os.Stat(filepath.Join(root, rootName, objLockDir)); !errors.Is(err, fs.ErrNotExist) { + t.Fatalf("nested lock directory exists or stat failed: %v", err) + } +} + +func TestListBucketsAndOwnersExcludesObjectLockDirectory(t *testing.T) { + p := newTestPosix(t, metaModes(t)["xattr"]) + + unlock, err := p.lockObjectPublish(context.Background(), "testbucket", "object") + if err != nil { + t.Fatalf("lock object publish: %v", err) + } + unlock() + + buckets, err := p.ListBucketsAndOwners(context.Background()) + if err != nil { + t.Fatalf("list buckets and owners: %v", err) + } + if len(buckets) != 0 { + t.Fatalf("buckets = %#v, want no internal lock directory", buckets) + } +} + func testPut(p *Posix, bucket, key string, body []byte, ifMatch, ifNoneMatch *string) (s3response.PutObjectOutput, error) { return p.PutObject(context.Background(), s3response.PutObjectInput{ Bucket: &bucket, diff --git a/backend/scoutfs/scoutfs_compat.go b/backend/scoutfs/scoutfs_compat.go index 34b70b40..a281c6c4 100644 --- a/backend/scoutfs/scoutfs_compat.go +++ b/backend/scoutfs/scoutfs_compat.go @@ -84,6 +84,7 @@ func New(rootdir string, opts ScoutfsOpts) (*ScoutFS, error) { CopyObjectThreshold: opts.CopyObjectThreshold, DefaultEtag: opts.DefaultEtag, DataIntegrityEtag: opts.DataIntegrityEtag, + ForceNoObjLockFile: true, // scoutfs flock not cluster consistent } if opts.newDirPermSet { posixOpts.SetNewDirPerm(opts.NewDirPerm) diff --git a/cmd/internal/gwcli/posix.go b/cmd/internal/gwcli/posix.go index 782180cf..fdca5f89 100644 --- a/cmd/internal/gwcli/posix.go +++ b/cmd/internal/gwcli/posix.go @@ -39,6 +39,7 @@ var ( nometa bool forceNoTmpFile bool forceNoCopyFileRange bool + forceNoObjLockFile bool enableODirect bool actionsConcurrency int ioBufferSize int @@ -142,6 +143,12 @@ will be translated into the file /mnt/fs/gwroot/mybucket/a/b/c/myobject`, EnvVars: []string{"VGW_DISABLE_COPY_FILE_RANGE"}, Destination: &forceNoCopyFileRange, }, + &cli.BoolFlag{ + Name: "disable-object-lock-file", + Usage: "disable shared advisory lock files for conditional object publishes (unsafe with multiple gateway processes)", + EnvVars: []string{"VGW_DISABLE_OBJECT_LOCK_FILE"}, + Destination: &forceNoObjLockFile, + }, &cli.BoolFlag{ Name: "enable-odirect", Usage: "enable best-effort O_DIRECT for object data reads/writes", @@ -194,6 +201,7 @@ func runPosix(ctx *cli.Context) error { VersioningDir: versioningDir, ForceNoTmpFile: forceNoTmpFile, ForceNoCopyFileRange: forceNoCopyFileRange, + ForceNoObjLockFile: forceNoObjLockFile, EnableODirect: enableODirect, ValidateBucketNames: DisableStrictBucketNames, Concurrency: actionsConcurrency, diff --git a/cmd/versitygw/main.go b/cmd/versitygw/main.go index ea716362..458da43c 100644 --- a/cmd/versitygw/main.go +++ b/cmd/versitygw/main.go @@ -123,9 +123,16 @@ func main() { gwcli.RunIAM = runIAM app := initApp() + posixCommand := gwcli.PosixCommand() + posixCommand.Before = func(ctx *cli.Context) error { + if ctx.Bool("disable-object-lock-file") { + fmt.Println("Warning: shared object publish locking disabled; conditional write atomicity is limited to this gateway process") + } + return nil + } app.Commands = []*cli.Command{ - gwcli.PosixCommand(), + posixCommand, gwcli.ScoutfsCommand(), gwcli.S3Command(), gwcli.AzureCommand(), diff --git a/extra/example.conf b/extra/example.conf index 93cea2b3..ace350c7 100644 --- a/extra/example.conf +++ b/extra/example.conf @@ -655,6 +655,16 @@ ROOT_SECRET_ACCESS_KEY= # NFS servers that may hang on this call. #VGW_DISABLE_COPY_FILE_RANGE=false +# The VGW_DISABLE_OBJECT_LOCK_FILE option disables shared advisory lock files +# used for conditional object publishes. This can be necessary on filesystems +# that do not support advisory file locking. When enabled, conditional writes +# are atomic only within a single versitygw process and are unsafe with +# multiple gateway processes sharing the same backend. +# This can be set to true to disable the .vgwlocks directory if running in +# single versitygw instance mode or conditional writes are not important +# to this deployment. +#VGW_DISABLE_OBJECT_LOCK_FILE=false + # The VGW_ENABLE_O_DIRECT option enables best-effort O_DIRECT for object data # reads and writes. This is disabled by default. When enabled, versitygw # attempts to open object data files with O_DIRECT and automatically falls back