feat: add best-effort O_DIRECT support for posix put/get-object put-part

This change introduces an opt-in O_DIRECT mode for POSIX object data paths while
keeping behavior safe and predictable across different filesystems and kernel
constraints. We now use direct I/O when available and beneficial, but preserve
correctness by falling back to buffered I/O when runtime read behavior indicates
alignment or capability mismatches.

The implementation keeps fast paths available for full-object reads and
descriptor-to-descriptor copy operations so kernel-level optimizations can still
be used where possible. At the same time, it avoids global assumptions from
single runtime failures and performs fallback at the stream level so requests
can continue successfully without broad feature disablement.
This commit is contained in:
Ben McClelland
2026-07-31 13:29:48 -07:00
parent 029a002777
commit 062bf0bf47
14 changed files with 651 additions and 62 deletions
+192
View File
@@ -0,0 +1,192 @@
// Copyright 2026 Versity Software
// This file is licensed under the Apache License, Version 2.0
// (the "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
//go:build linux
package posix
import (
"fmt"
"io"
"os"
"path/filepath"
"strconv"
"sync"
"syscall"
"github.com/versity/versitygw/backend"
)
// openDataRead opens object data for reading and applies O_DIRECT when
// requested. O_DIRECT is best-effort and falls back to buffered I/O when
// unsupported by the filesystem.
func openDataRead(name string, useODirect bool) (*os.File, error) {
if !useODirect {
return os.Open(name)
}
f, err := os.OpenFile(name, syscall.O_RDONLY|syscall.O_DIRECT, 0)
if err != nil {
if isODirectUnsupportedOpenErr(err) {
warnODirectUnsupportedOnce("openDataRead", err)
return os.Open(name)
}
return nil, err
}
return f, nil
}
func buildGetObjectBody(f *os.File, path string, startOffset, length, objSize int64, useODirect bool, readBufferSize int) (io.ReadCloser, error) {
if startOffset == 0 && length == objSize {
if useODirect {
return newODirectReadFallbackFile(path, f, true), nil
}
return f, nil
}
if !useODirect {
rdr := io.NewSectionReader(f, startOffset, length)
return withReadBufferSize(&backend.FileSectionReadCloser{R: rdr, F: f}, readBufferSize), nil
}
rf := newODirectReadFallbackFile(path, f, true)
if _, err := rf.Seek(startOffset, io.SeekStart); err != nil {
_ = rf.Close()
return nil, fmt.Errorf("seek range start: %w", err)
}
return withReadBufferSize(&readerWithCloser{r: io.LimitReader(rf, length), c: rf}, readBufferSize), nil
}
type readerWithCloser struct {
r io.Reader
c io.Closer
}
func (r *readerWithCloser) Read(p []byte) (int, error) {
return r.r.Read(p)
}
func (r *readerWithCloser) Close() error {
return r.c.Close()
}
type odirectReadFallbackFile struct {
mu sync.Mutex
path string
f *os.File
useODirect bool
}
func newODirectReadFallbackFile(path string, f *os.File, useODirect bool) *odirectReadFallbackFile {
return &odirectReadFallbackFile{
path: path,
f: f,
useODirect: useODirect,
}
}
func (r *odirectReadFallbackFile) Read(p []byte) (int, error) {
r.mu.Lock()
defer r.mu.Unlock()
n, err := r.f.Read(p)
if err != nil && n == 0 && r.useODirect && isODirectRuntimeFallbackErr(err) {
if fallbackErr := r.switchToBufferedAtCurrentOffsetLocked(); fallbackErr != nil {
return 0, fallbackErr
}
return r.f.Read(p)
}
return n, err
}
func (r *odirectReadFallbackFile) WriteTo(w io.Writer) (int64, error) {
r.mu.Lock()
useODirect := r.useODirect
f := r.f
r.mu.Unlock()
if !useODirect {
return io.Copy(w, &onlyRead{r})
}
writerTo, ok := interface{}(f).(io.WriterTo)
if !ok {
return io.Copy(w, &onlyRead{r})
}
n, err := writerTo.WriteTo(w)
if err == nil || !isODirectRuntimeFallbackErr(err) {
return n, err
}
r.mu.Lock()
fallbackErr := r.switchToBufferedAtCurrentOffsetLocked()
r.mu.Unlock()
if fallbackErr != nil {
return n, fallbackErr
}
m, err := io.Copy(w, &onlyRead{r})
return n + m, err
}
func (r *odirectReadFallbackFile) Seek(offset int64, whence int) (int64, error) {
r.mu.Lock()
defer r.mu.Unlock()
return r.f.Seek(offset, whence)
}
func (r *odirectReadFallbackFile) Close() error {
r.mu.Lock()
defer r.mu.Unlock()
return r.f.Close()
}
func (r *odirectReadFallbackFile) switchToBufferedAtOffset(offset int64) error {
fd := strconv.Itoa(int(r.f.Fd()))
bf, openErr := os.Open(filepath.Join(procfddir, fd))
if openErr != nil {
return openErr
}
if _, seekErr := bf.Seek(offset, io.SeekStart); seekErr != nil {
_ = bf.Close()
return seekErr
}
if closeErr := r.f.Close(); closeErr != nil {
_ = bf.Close()
return closeErr
}
r.f = bf
r.useODirect = false
return nil
}
func (r *odirectReadFallbackFile) switchToBufferedAtCurrentOffsetLocked() error {
offset, seekErr := r.f.Seek(0, io.SeekCurrent)
if seekErr != nil {
return seekErr
}
return r.switchToBufferedAtOffset(offset)
}
+33
View File
@@ -0,0 +1,33 @@
// Copyright 2026 Versity Software
// This file is licensed under the Apache License, Version 2.0
// (the "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
//go:build !linux
package posix
import (
"io"
"os"
"github.com/versity/versitygw/backend"
)
func buildGetObjectBody(f *os.File, _ string, startOffset, length, objSize int64, _ bool, readBufferSize int) (io.ReadCloser, error) {
if startOffset == 0 && length == objSize {
return f, nil
}
rdr := io.NewSectionReader(f, startOffset, length)
return withReadBufferSize(&backend.FileSectionReadCloser{R: rdr, F: f}, readBufferSize), nil
}
+26
View File
@@ -0,0 +1,26 @@
// Copyright 2026 Versity Software
// This file is licensed under the Apache License, Version 2.0
// (the "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
//go:build !linux && !windows
package posix
import "os"
func openDataRead(name string, useODirect bool) (*os.File, error) {
if useODirect {
warnODirectUnsupportedOnce("openDataRead-nonlinux", os.ErrInvalid)
}
return os.Open(name)
}
+46
View File
@@ -0,0 +1,46 @@
// Copyright 2026 Versity Software
// This file is licensed under the Apache License, Version 2.0
// (the "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
//go:build windows
package posix
import (
"os"
"syscall"
)
func openDataRead(name string, useODirect bool) (*os.File, error) {
if useODirect {
warnODirectUnsupportedOnce("openDataRead-windows", os.ErrInvalid)
}
ptr, err := syscall.UTF16PtrFromString(name)
if err != nil {
return nil, &os.PathError{Op: "open", Path: name, Err: err}
}
h, err := syscall.CreateFile(
ptr,
syscall.GENERIC_READ,
syscall.FILE_SHARE_READ|syscall.FILE_SHARE_WRITE|syscall.FILE_SHARE_DELETE,
nil,
syscall.OPEN_EXISTING,
syscall.FILE_ATTRIBUTE_NORMAL,
0,
)
if err != nil {
return nil, &os.PathError{Op: "open", Path: name, Err: err}
}
return os.NewFile(uintptr(h), name), nil
}
+3 -5
View File
@@ -46,9 +46,7 @@ func isErrDirNotEmpty(err error) bool {
return errors.Is(err, syscall.ENOTEMPTY)
}
// openForRead opens a file for reading. On non-Windows systems, os.Open is
// sufficient because POSIX allows removing (unlinking) a file that is still
// open by another process.
func openForRead(name string) (*os.File, error) {
return os.Open(name)
// openForRead opens an object data file for reading.
func openForRead(name string, useODirect bool) (*os.File, error) {
return openDataRead(name, useODirect)
}
+2 -18
View File
@@ -110,22 +110,6 @@ func isErrNotDir(err error) bool {
// is held open for streaming the GET response body. Without this flag,
// Windows returns "The process cannot access the file because it is being
// used by another process" on the Remove call.
func openForRead(name string) (*os.File, error) {
ptr, err := syscall.UTF16PtrFromString(name)
if err != nil {
return nil, &os.PathError{Op: "open", Path: name, Err: err}
}
h, err := syscall.CreateFile(
ptr,
syscall.GENERIC_READ,
syscall.FILE_SHARE_READ|syscall.FILE_SHARE_WRITE|syscall.FILE_SHARE_DELETE,
nil,
syscall.OPEN_EXISTING,
syscall.FILE_ATTRIBUTE_NORMAL,
0,
)
if err != nil {
return nil, &os.PathError{Op: "open", Path: name, Err: err}
}
return os.NewFile(uintptr(h), name), nil
func openForRead(name string, useODirect bool) (*os.File, error) {
return openDataRead(name, useODirect)
}
+29
View File
@@ -0,0 +1,29 @@
// Copyright 2026 Versity Software
// This file is licensed under the Apache License, Version 2.0
// (the "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
//go:build linux
package posix
import (
"errors"
"syscall"
)
// Shared by Linux O_DIRECT open paths in data reads and tmpfile creation.
func isODirectUnsupportedOpenErr(err error) bool {
return errors.Is(err, syscall.EINVAL) ||
errors.Is(err, syscall.EOPNOTSUPP) ||
errors.Is(err, syscall.ENOTSUP)
}
+61
View File
@@ -15,22 +15,83 @@
package posix
import (
"errors"
"fmt"
"io"
"math/rand"
"os"
"strconv"
"syscall"
"time"
)
const odirectMinWriteAlign = 512
func (tmp *tmpfile) Write(b []byte) (int, error) {
if int64(len(b)) > tmp.size {
return 0, fmt.Errorf("write exceeds content length %v", tmp.size)
}
if tmp.useODirect && !isODirectLenAligned(len(b)) {
if err := tmp.switchToBufferedAtCurrentOffset(fmt.Sprintf("unaligned write length: len=%d len%%512=%d", len(b), len(b)%odirectMinWriteAlign)); err != nil {
return 0, err
}
}
n, err := tmp.f.Write(b)
if err != nil && n == 0 && tmp.useODirect && isODirectRuntimeFallbackErr(err) {
warnODirectUnsupportedOnce("tmpfile.Write", err)
if fallbackErr := tmp.switchToBufferedAtCurrentOffset("O_DIRECT write failure"); fallbackErr != nil {
return 0, fallbackErr
}
n, err = tmp.f.Write(b)
}
tmp.size -= int64(n)
return n, err
}
func (tmp *tmpfile) switchToBufferedAtCurrentOffset(reason string) error {
offset, seekErr := tmp.f.Seek(0, io.SeekCurrent)
if seekErr != nil {
return fmt.Errorf("capture write offset before fallback reopen: %w", seekErr)
}
name := tmp.f.Name()
f, openErr := os.OpenFile(name, os.O_RDWR, 0)
if openErr != nil {
return fmt.Errorf("reopen temp file in buffered mode after O_DIRECT fallback (%s): %w", reason, openErr)
}
if _, seekErr = f.Seek(offset, io.SeekStart); seekErr != nil {
f.Close()
return fmt.Errorf("restore write offset after fallback reopen: %w", seekErr)
}
if closeErr := tmp.f.Close(); closeErr != nil {
f.Close()
return fmt.Errorf("close O_DIRECT temp file after fallback reopen: %w", closeErr)
}
tmp.f = f
if tmp.isOTmp {
tmp.procFDName = strconv.Itoa(int(f.Fd()))
}
tmp.useODirect = false
return nil
}
func isODirectLenAligned(n int) bool {
return n%odirectMinWriteAlign == 0
}
func isODirectRuntimeFallbackErr(err error) bool {
return errors.Is(err, syscall.EINVAL) ||
errors.Is(err, syscall.EOPNOTSUPP) ||
errors.Is(err, syscall.ENOTSUP)
}
func (tmp *tmpfile) File() *os.File {
return tmp.f
}
+87 -20
View File
@@ -29,6 +29,7 @@ import (
"sort"
"strconv"
"strings"
"sync"
"syscall"
"time"
@@ -87,6 +88,10 @@ type Posix struct {
// support copy_file_range is mounted over NFSv4.2.
forceNoCopyFileRange 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
// enable posix level bucket name validations, not needed if the
// frontend handlers are already validating bucket names
validateBucketName bool
@@ -108,6 +113,10 @@ type Posix struct {
// rejected with an 'InvalidRequest' to comply with the S3 limit
// of 5 GiB.
copyObjectThreshold int64
// ioBufferSize is the buffer size used by buffered copy/read paths.
ioBufferSize int
ioBufferPool sync.Pool
}
var _ backend.Backend = &Posix{}
@@ -149,10 +158,17 @@ const (
doFalloc = true
skipFalloc = false
odirectAllowed odirectPolicy = true
odirectNotAllowed odirectPolicy = false
// defaultConcurrency is the default limit for concurrent POSIX actions.
defaultConcurrency = 5000
// defaultIOBufferSize is the default buffer size used by io.CopyBuffer paths.
defaultIOBufferSize = 1024 * 1024
)
type odirectPolicy bool
// PosixOpts are the options for the Posix backend
type PosixOpts struct {
// ChownUID sets the UID of the object to the UID of the user on PUT
@@ -172,6 +188,9 @@ type PosixOpts struct {
ForceNoTmpFile bool
// ForceNoCopyFileRange disables the use of io.Copy for multipart uploads parts
ForceNoCopyFileRange bool
// EnableODirect enables best-effort O_DIRECT for object data reads/writes.
// Disabled by default.
EnableODirect bool
// ValidateBucketNames enables minimal bucket name validation to prevent
// incorrect access to the filesystem. This is only needed if the
// frontend is not already validating bucket names.
@@ -194,9 +213,14 @@ type PosixOpts struct {
// attribute (e.g. files placed on the filesystem outside of versitygw).
// When empty, such objects are served with an empty ETag.
DefaultEtag string
// IOBufferSize sets the buffer size (in bytes) for copy/read paths that use
// io.CopyBuffer or buffered readers. Defaults to 1MiB when unset or invalid.
IOBufferSize int
}
func New(rootdir string, meta meta.MetadataStorer, opts PosixOpts) (*Posix, error) {
ioBufferSize := ioBufferSizeOrDefault(opts.IOBufferSize)
if opts.SideCarDir != "" && strings.HasPrefix(opts.SideCarDir, rootdir) {
return nil, fmt.Errorf("sidecar directory cannot be inside the gateway root directory")
}
@@ -255,10 +279,16 @@ func New(rootdir string, meta meta.MetadataStorer, opts PosixOpts) (*Posix, erro
newDirPerm: opts.NewDirPerm,
forceNoTmpFile: opts.ForceNoTmpFile,
forceNoCopyFileRange: opts.ForceNoCopyFileRange,
enableODirect: opts.EnableODirect,
validateBucketName: opts.ValidateBucketNames,
actionLimiter: semaphore.NewWeighted(int64(concurrencyOrDefault(opts.Concurrency))),
copyObjectThreshold: opts.CopyObjectThreshold,
defaultEtag: opts.DefaultEtag,
ioBufferSize: ioBufferSize,
ioBufferPool: sync.Pool{New: func() any {
b := make([]byte, ioBufferSize)
return &b
}},
}, nil
}
@@ -270,6 +300,32 @@ func concurrencyOrDefault(n int) int {
return defaultConcurrency
}
func ioBufferSizeOrDefault(n int) int {
if n > 0 {
return n
}
return defaultIOBufferSize
}
func (p *Posix) getIOBuffer() []byte {
bp, ok := p.ioBufferPool.Get().(*[]byte)
if !ok || bp == nil || cap(*bp) < p.ioBufferSize {
return make([]byte, p.ioBufferSize)
}
return (*bp)[:p.ioBufferSize]
}
func (p *Posix) putIOBuffer(b []byte) {
if b == nil {
return
}
if cap(b) < p.ioBufferSize {
return
}
b = b[:p.ioBufferSize]
p.ioBufferPool.Put(&b)
}
func validateSubDir(root, dir string) (string, error) {
absDir, err := filepath.Abs(dir)
if err != nil {
@@ -945,12 +1001,13 @@ func (p *Posix) createObjVersion(bucket, key string, size int64, acc auth.Accoun
versioningKey := filepath.Join(genObjVersionKey(key), versionId)
versionTmpPath := filepath.Join(versionBucketPath, MetaTmpDir)
f, err := p.openTmpFile(versionTmpPath, versionBucketPath, versioningKey,
size, acc, doFalloc, p.forceNoTmpFile)
size, acc, doFalloc, p.forceNoTmpFile, odirectNotAllowed)
if err != nil {
return versionPath, err
}
defer f.cleanup()
// Prioritize copy_file_range for internal file-to-file version copies.
_, err = io.Copy(f.File(), sf)
if err != nil {
return versionPath, err
@@ -2092,7 +2149,7 @@ func (p *Posix) CompleteMultipartUploadWithCopy(ctx context.Context, input *s3.C
}
f, err := p.openTmpFile(filepath.Join(bucket, MetaTmpDir), bucket, object,
totalsize, acct, skipFalloc, p.forceNoTmpFile)
totalsize, acct, skipFalloc, p.forceNoTmpFile, odirectNotAllowed)
if err != nil {
if errors.Is(err, syscall.EDQUOT) {
return res, "", s3err.GetAPIError(s3err.ErrQuotaExceeded)
@@ -2119,12 +2176,13 @@ func (p *Posix) CompleteMultipartUploadWithCopy(ctx context.Context, input *s3.C
// Fail back to standard copy
debuglogger.Logf("custom data block move failed (%q/%q): %v, failing back to io.Copy()",
bucket, object, err)
fw := f.File()
fw.Seek(0, io.SeekEnd)
_, _ = f.File().Seek(0, io.SeekEnd)
if p.forceNoCopyFileRange {
_, err = io.Copy(fw, &onlyRead{pf})
_, err = io.Copy(f, &onlyRead{pf})
} else {
_, err = io.Copy(fw, pf)
// Keep both endpoints as *os.File here so
// io.Copy can use copy_file_range.
_, err = io.Copy(f.File(), pf)
}
}
if !idemp && err == nil {
@@ -2145,8 +2203,10 @@ func (p *Posix) CompleteMultipartUploadWithCopy(ctx context.Context, input *s3.C
}
} else {
if p.forceNoCopyFileRange {
_, err = io.Copy(f.File(), &onlyRead{pf})
_, err = io.Copy(f, &onlyRead{pf})
} else {
// Keep both endpoints as *os.File here so
// io.Copy can use copy_file_range.
_, err = io.Copy(f.File(), pf)
}
}
@@ -3030,7 +3090,7 @@ func (p *Posix) UploadPartWithPostFunc(ctx context.Context, input *s3.UploadPart
partPath := filepath.Join(mpPath, fmt.Sprintf("%v", *part))
f, err := p.openTmpFile(filepath.Join(bucket, objdir),
bucket, partPath, length, acct, doFalloc, p.forceNoTmpFile)
bucket, partPath, length, acct, doFalloc, p.forceNoTmpFile, odirectAllowed)
if err != nil {
if errors.Is(err, syscall.EDQUOT) {
drainBody(r)
@@ -3152,7 +3212,10 @@ func (p *Posix) UploadPartWithPostFunc(ctx context.Context, input *s3.UploadPart
}
}
_, err = io.Copy(f, tr)
buf := p.getIOBuffer()
defer p.putIOBuffer(buf)
_, err = io.CopyBuffer(f, tr, buf)
if err != nil {
if errors.Is(err, syscall.EDQUOT) {
drainBody(tr)
@@ -3415,7 +3478,7 @@ func (p *Posix) UploadPartCopy(ctx context.Context, upi *s3.UploadPartCopyInput)
}
f, err := p.openTmpFile(filepath.Join(*upi.Bucket, objdir),
*upi.Bucket, partPath, length, acct, doFalloc, p.forceNoTmpFile)
*upi.Bucket, partPath, length, acct, doFalloc, p.forceNoTmpFile, odirectNotAllowed)
if err != nil {
if errors.Is(err, syscall.EDQUOT) {
return s3response.CopyPartResult{}, s3err.GetAPIError(s3err.ErrQuotaExceeded)
@@ -3803,7 +3866,7 @@ func (p *Posix) PutObjectWithPostFunc(ctx context.Context, po s3response.PutObje
}
f, err := p.openTmpFile(filepath.Join(*po.Bucket, MetaTmpDir),
*po.Bucket, *po.Key, contentLength, acct, doFalloc, p.forceNoTmpFile)
*po.Bucket, *po.Key, contentLength, acct, doFalloc, p.forceNoTmpFile, odirectAllowed)
if err != nil {
if errors.Is(err, syscall.EDQUOT) {
drainBody(po.Body)
@@ -3832,7 +3895,10 @@ func (p *Posix) PutObjectWithPostFunc(ctx context.Context, po s3response.PutObje
rdr = hashRdr
}
_, err = io.Copy(f, rdr)
buf := p.getIOBuffer()
defer p.putIOBuffer(buf)
_, err = io.CopyBuffer(f, rdr, buf)
if err != nil {
if errors.Is(err, syscall.EDQUOT) {
drainBody(rdr)
@@ -4234,13 +4300,14 @@ func (p *Posix) DeleteObject(ctx context.Context, input *s3.DeleteObjectInput) (
f, err := p.openTmpFile(filepath.Join(bucket, MetaTmpDir),
bucket, object, srcObjVersion.Size(), acct, doFalloc,
p.forceNoTmpFile)
p.forceNoTmpFile, odirectNotAllowed)
if err != nil {
return nil, fmt.Errorf("open tmp file: %w", err)
}
defer f.cleanup()
_, err = io.Copy(f, sf)
// Prioritize copy_file_range for internal file-to-file version restores.
_, err = io.Copy(f.File(), sf)
if err != nil {
_ = sf.Close()
return nil, fmt.Errorf("copy object %w", err)
@@ -4678,7 +4745,7 @@ func (p *Posix) GetObject(ctx context.Context, input *s3.GetObjectInput) (*s3.Ge
// openForRead opens with FILE_SHARE_DELETE on Windows so that a concurrent
// DeleteObject can call os.Remove on this file while the GET response body
// is still being streamed. On POSIX, os.Open is sufficient.
f, err := openForRead(objPath)
f, err := openForRead(objPath, p.enableODirect)
if errors.Is(err, fs.ErrNotExist) {
return nil, s3err.GetAPIError(s3err.ErrNoSuchKey)
}
@@ -4781,11 +4848,11 @@ func (p *Posix) GetObject(ctx context.Context, input *s3.GetObjectInput) (*s3.Ge
}
}
// using an os.File allows zero-copy sendfile via io.Copy(os.File, net.Conn)
var body io.ReadCloser = f
if startOffset != 0 || length != objSize {
rdr := io.NewSectionReader(f, startOffset, length)
body = &backend.FileSectionReadCloser{R: rdr, F: f}
// Full-object responses can keep the underlying *os.File for sendfile.
// Linux range reads on O_DIRECT may need runtime fallback to buffered I/O.
body, err := buildGetObjectBody(f, objPath, startOffset, length, objSize, p.enableODirect, p.ioBufferSize)
if err != nil {
return nil, fmt.Errorf("build get object body: %w", err)
}
return &s3.GetObjectOutput{
+55
View File
@@ -0,0 +1,55 @@
// Copyright 2026 Versity Software
// This file is licensed under the Apache License, Version 2.0
// (the "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.
package posix
import (
"bufio"
"io"
"log"
"sync"
)
func withReadBufferSize(r io.ReadCloser, size int) io.ReadCloser {
if size <= 0 {
return r
}
return &bufferedReadCloser{
r: bufio.NewReaderSize(r, size),
c: r,
}
}
type bufferedReadCloser struct {
r *bufio.Reader
c io.Closer
}
func (b *bufferedReadCloser) Read(p []byte) (int, error) {
return b.r.Read(p)
}
func (b *bufferedReadCloser) Close() error {
return b.c.Close()
}
var odirectUnsupportedWarnByOp sync.Map
func warnODirectUnsupportedOnce(op string, err error) {
v, _ := odirectUnsupportedWarnByOp.LoadOrStore(op, &sync.Once{})
v.(*sync.Once).Do(func() {
log.Printf("WARNING: O_DIRECT is enabled but unsupported (%s: %v); falling back to buffered I/O. This warning is shown once per operation.", op, err)
})
}
+70 -15
View File
@@ -42,6 +42,8 @@ type tmpfile struct {
bucket string
objname string
isOTmp bool
procFDName string
useODirect bool
size int64
doChown bool
uid int
@@ -54,11 +56,11 @@ var (
defaultFilePerm uint32 = 0644
)
func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Account, dofalloc bool, forceNoTmpFile bool) (*tmpfile, error) {
func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Account, dofalloc bool, forceNoTmpFile bool, allowODirect odirectPolicy) (*tmpfile, error) {
uid, gid, doChown := p.getChownIDs(acct)
if forceNoTmpFile {
return p.openMkTemp(dir, bucket, obj, size, dofalloc, uid, gid, doChown)
return p.openMkTemp(dir, bucket, obj, size, dofalloc, uid, gid, doChown, allowODirect)
}
// O_TMPFILE allows for a file handle to an unnamed file in the filesystem.
@@ -67,14 +69,34 @@ func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Accou
// file descriptor into the namespace.
// Not all filesystems support this, so fallback to CreateTemp for when
// this is not supported.
fd, err := unix.Open(dir, unix.O_RDWR|unix.O_TMPFILE|unix.O_CLOEXEC, defaultFilePerm)
openFlags := unix.O_RDWR | unix.O_TMPFILE | unix.O_CLOEXEC
useODirect := false
if p.enableODirect && bool(allowODirect) {
openFlags |= unix.O_DIRECT
useODirect = true
}
fd, err := unix.Open(dir, openFlags, defaultFilePerm)
if err != nil {
if errors.Is(err, syscall.EROFS) {
return nil, s3err.GetAPIError(s3err.ErrMethodNotAllowed)
}
// O_TMPFILE not supported, try fallback
return p.openMkTemp(dir, bucket, obj, size, dofalloc, uid, gid, doChown)
if p.enableODirect && bool(allowODirect) && isODirectUnsupportedOpenErr(err) {
warnODirectUnsupportedOnce("openTmpFile", err)
fd, err = unix.Open(dir, unix.O_RDWR|unix.O_TMPFILE|unix.O_CLOEXEC, defaultFilePerm)
if err == nil {
useODirect = false
} else if errors.Is(err, syscall.EROFS) {
return nil, s3err.GetAPIError(s3err.ErrMethodNotAllowed)
}
}
if err != nil {
// O_TMPFILE not supported, try fallback
return p.openMkTemp(dir, bucket, obj, size, dofalloc, uid, gid, doChown, allowODirect)
}
}
// for O_TMPFILE, filename is /proc/self/fd/<fd> to be used
@@ -86,6 +108,8 @@ func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Accou
bucket: bucket,
objname: obj,
isOTmp: true,
procFDName: strconv.Itoa(fd),
useODirect: useODirect,
size: size,
doChown: doChown,
uid: uid,
@@ -109,7 +133,7 @@ func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Accou
return tmp, nil
}
func (p *Posix) openMkTemp(dir, bucket, obj string, size int64, dofalloc bool, uid, gid int, doChown bool) (*tmpfile, error) {
func (p *Posix) openMkTemp(dir, bucket, obj string, size int64, dofalloc bool, uid, gid int, doChown bool, allowODirect odirectPolicy) (*tmpfile, error) {
err := backend.MkdirAll(dir, uid, gid, doChown, p.newDirPerm)
if err != nil {
if errors.Is(err, syscall.EROFS) {
@@ -125,14 +149,41 @@ func (p *Posix) openMkTemp(dir, bucket, obj string, size int64, dofalloc bool, u
}
return nil, err
}
useODirect := false
if p.enableODirect && bool(allowODirect) {
name := f.Name()
if err := f.Close(); err != nil {
os.Remove(name)
return nil, fmt.Errorf("close temp file before O_DIRECT reopen: %w", err)
}
fd, err := unix.Open(name, unix.O_RDWR|unix.O_CLOEXEC|unix.O_DIRECT, defaultFilePerm)
if err == nil {
f = os.NewFile(uintptr(fd), name)
useODirect = true
} else if isODirectUnsupportedOpenErr(err) {
warnODirectUnsupportedOnce("openMkTemp", err)
f, err = os.OpenFile(name, os.O_RDWR, 0)
if err != nil {
os.Remove(name)
return nil, fmt.Errorf("reopen temp file after O_DIRECT fallback: %w", err)
}
} else {
os.Remove(name)
return nil, fmt.Errorf("open temp file with O_DIRECT: %w", err)
}
}
tmp := &tmpfile{
f: f,
bucket: bucket,
objname: obj,
size: size,
doChown: doChown,
uid: uid,
gid: gid,
f: f,
bucket: bucket,
objname: obj,
useODirect: useODirect,
size: size,
doChown: doChown,
uid: uid,
gid: gid,
}
// falloc is best effort, its fine if this fails
if size > 0 && dofalloc {
@@ -237,8 +288,12 @@ func (tmp *tmpfile) link() error {
if err != nil {
return fmt.Errorf("open parent dir: %w", err)
}
srcFDName := tmp.procFDName
if srcFDName == "" {
srcFDName = filepath.Base(tmp.f.Name())
}
err = linkatOTmpfile(int(procdir.Fd()), int(dirf.Fd()),
filepath.Base(tmp.f.Name()), filepath.Base(objPath))
srcFDName, filepath.Base(objPath))
dirf.Close()
if errors.Is(err, syscall.ENOENT) {
// The directory was removed between open and linkat; backoff and retry.
@@ -254,7 +309,7 @@ func (tmp *tmpfile) link() error {
}
if err != nil {
return fmt.Errorf("link tmpfile (fd %q as %q): %w",
filepath.Base(tmp.f.Name()), objPath, err)
srcFDName, objPath, err)
}
break
}
+15 -4
View File
@@ -36,9 +36,13 @@ const (
)
type tmpfile struct {
f *os.File
bucket string
objname string
f *os.File
bucket string
objname string
// Retained for compatibility with shared tmpfile methods in otmpfile_common.
isOTmp bool
procFDName string
useODirect bool
size int64
newDirPerm fs.FileMode
uid int
@@ -46,9 +50,13 @@ type tmpfile struct {
doChown bool
}
func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Account, _ bool, _ bool) (*tmpfile, error) {
func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Account, _ bool, _ bool, allowODirect odirectPolicy) (*tmpfile, error) {
uid, gid, doChown := p.getChownIDs(acct)
if p.enableODirect && bool(allowODirect) {
warnODirectUnsupportedOnce("openTmpFile-nonlinux", os.ErrInvalid)
}
// Create a temp file for upload while in progress (see link comments below).
var err error
err = backend.MkdirAll(dir, uid, gid, doChown, p.newDirPerm)
@@ -80,6 +88,9 @@ func (p *Posix) openTmpFile(dir, bucket, obj string, size int64, acct auth.Accou
f: f,
bucket: bucket,
objname: obj,
isOTmp: false,
procFDName: "",
useODirect: false,
size: size,
newDirPerm: p.newDirPerm,
uid: uid,
+17
View File
@@ -33,7 +33,9 @@ var (
nometa bool
forceNoTmpFile bool
forceNoCopyFileRange bool
enableODirect bool
actionsConcurrency int
ioBufferSize int
defaultEtag string
)
@@ -98,6 +100,13 @@ will be translated into the file /mnt/fs/gwroot/mybucket/a/b/c/myobject`,
Value: 5000,
Destination: &actionsConcurrency,
},
&cli.IntFlag{
Name: "io-buffer-size",
Usage: "buffer size in bytes used by POSIX put/get/part read and write paths (<=0 uses backend default 1MiB)",
EnvVars: []string{"VGW_POSIX_IO_BUFFER_SIZE"},
Value: 1024 * 1024,
Destination: &ioBufferSize,
},
&cli.BoolFlag{
Name: "nometa",
Usage: "disable metadata storage",
@@ -116,6 +125,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: "enable-odirect",
Usage: "enable best-effort O_DIRECT for object data reads/writes",
EnvVars: []string{"VGW_ENABLE_O_DIRECT"},
Destination: &enableODirect,
},
&cli.StringFlag{
Name: "default-etag",
Usage: "default ETag value returned for objects that do not have a stored etag attribute (e.g. files placed on the filesystem outside of versitygw)",
@@ -153,8 +168,10 @@ func runPosix(ctx *cli.Context) error {
NewDirPerm: fs.FileMode(dirPerms),
ForceNoTmpFile: forceNoTmpFile,
ForceNoCopyFileRange: forceNoCopyFileRange,
EnableODirect: enableODirect,
ValidateBucketNames: disableStrictBucketNames,
Concurrency: actionsConcurrency,
IOBufferSize: ioBufferSize,
CopyObjectThreshold: copyObjectThreshold,
DefaultEtag: defaultEtag,
}
+15
View File
@@ -602,6 +602,12 @@ ROOT_SECRET_ACCESS_KEY=
# memory use.
#VGW_POSIX_CONCURRENCY=5000
# The VGW_POSIX_IO_BUFFER_SIZE option sets the data transfer buffer size used
# by the posix backend for PutObject, UploadPart, and GetObject read/write
# paths. The default is 1048576 bytes (1 MiB). When set to 0 or a negative
# value, the backend falls back to the same default.
#VGW_POSIX_IO_BUFFER_SIZE=1048576
# The gateway will use O_TMPFILE for writing objects while uploading and
# link the file to the final object name when the upload is complete if the
# filesystem supports O_TMPFILE. This creates an atomic object creation
@@ -623,6 +629,15 @@ ROOT_SECRET_ACCESS_KEY=
# NFS servers that may hang on this call.
#VGW_DISABLE_COPY_FILE_RANGE=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
# to buffered I/O when O_DIRECT is not supported by the filesystem/open call.
# Some systems may still enforce alignment constraints for direct I/O. The
# VGW_POSIX_IO_BUFFER_SIZE setting can be used to tune buffered read/write
# chunk sizing while keeping O_DIRECT best-effort behavior enabled.
#VGW_ENABLE_O_DIRECT=false
# The VGW_DEFAULT_ETAG option sets the ETag value returned for objects that do
# not have a stored etag attribute. This applies to files that were placed on
# the filesystem outside of versitygw and therefore lack S3 metadata. Some S3