mirror of
https://github.com/netbirdio/gvisor.git
synced 2026-05-22 17:12:49 -07:00
@@ -184,6 +184,13 @@ const (
|
||||
FUSE_KERNEL_MINOR_VERSION = 31
|
||||
)
|
||||
|
||||
// Constants relevant to FUSE operations.
|
||||
const (
|
||||
FUSE_NAME_MAX = 1024
|
||||
FUSE_PAGE_SIZE = 4096
|
||||
FUSE_DIRENT_ALIGN = 8
|
||||
)
|
||||
|
||||
// FUSEInitIn is the request sent by the kernel to the daemon,
|
||||
// to negotiate the version and flags.
|
||||
//
|
||||
@@ -392,6 +399,36 @@ type FUSEOpenOut struct {
|
||||
|
||||
// OpenFlag for the opened file.
|
||||
OpenFlag uint32
|
||||
}
|
||||
|
||||
// FUSE_READ flags, consistent with the ones in include/uapi/linux/fuse.h.
|
||||
const (
|
||||
FUSE_READ_LOCKOWNER = 1 << 1
|
||||
)
|
||||
|
||||
// FUSEReadIn is the request sent by the kernel to the daemon
|
||||
// for FUSE_READ.
|
||||
//
|
||||
// +marshal
|
||||
type FUSEReadIn struct {
|
||||
// Fh is the file handle in userspace.
|
||||
Fh uint64
|
||||
|
||||
// Offset is the read offset.
|
||||
Offset uint64
|
||||
|
||||
// Size is the number of bytes to read.
|
||||
Size uint32
|
||||
|
||||
// ReadFlags for this FUSE_READ request.
|
||||
// Currently only contains FUSE_READ_LOCKOWNER.
|
||||
ReadFlags uint32
|
||||
|
||||
// LockOwner is the id of the lock owner if there is one.
|
||||
LockOwner uint64
|
||||
|
||||
// Flags for the underlying file.
|
||||
Flags uint32
|
||||
|
||||
_ uint32
|
||||
}
|
||||
|
||||
@@ -36,7 +36,9 @@ go_library(
|
||||
"fusefs.go",
|
||||
"init.go",
|
||||
"inode_refs.go",
|
||||
"read_write.go",
|
||||
"register.go",
|
||||
"regular_file.go",
|
||||
"request_list.go",
|
||||
],
|
||||
visibility = ["//pkg/sentry:internal"],
|
||||
@@ -46,6 +48,7 @@ go_library(
|
||||
"//pkg/log",
|
||||
"//pkg/marshal",
|
||||
"//pkg/refs",
|
||||
"//pkg/safemem",
|
||||
"//pkg/sentry/fsimpl/devtmpfs",
|
||||
"//pkg/sentry/fsimpl/kernfs",
|
||||
"//pkg/sentry/kernel",
|
||||
|
||||
@@ -161,6 +161,7 @@ type connection struct {
|
||||
bgLock sync.Mutex
|
||||
|
||||
// maxRead is the maximum size of a read buffer in in bytes.
|
||||
// Initialized from a fuse fs parameter.
|
||||
maxRead uint32
|
||||
|
||||
// maxWrite is the maximum size of a write buffer in bytes.
|
||||
@@ -206,7 +207,7 @@ type connection struct {
|
||||
}
|
||||
|
||||
// newFUSEConnection creates a FUSE connection to fd.
|
||||
func newFUSEConnection(_ context.Context, fd *vfs.FileDescription, maxInFlightRequests uint64) (*connection, error) {
|
||||
func newFUSEConnection(_ context.Context, fd *vfs.FileDescription, opts *filesystemOptions) (*connection, error) {
|
||||
// Mark the device as ready so it can be used. /dev/fuse can only be used if the FD was used to
|
||||
// mount a FUSE filesystem.
|
||||
fuseFD := fd.Impl().(*DeviceFD)
|
||||
@@ -216,13 +217,14 @@ func newFUSEConnection(_ context.Context, fd *vfs.FileDescription, maxInFlightRe
|
||||
hdrLen := uint32((*linux.FUSEHeaderOut)(nil).SizeBytes())
|
||||
fuseFD.writeBuf = make([]byte, hdrLen)
|
||||
fuseFD.completions = make(map[linux.FUSEOpID]*futureResponse)
|
||||
fuseFD.fullQueueCh = make(chan struct{}, maxInFlightRequests)
|
||||
fuseFD.fullQueueCh = make(chan struct{}, opts.maxActiveRequests)
|
||||
fuseFD.writeCursor = 0
|
||||
|
||||
return &connection{
|
||||
fd: fuseFD,
|
||||
maxBackground: fuseDefaultMaxBackground,
|
||||
congestionThreshold: fuseDefaultCongestionThreshold,
|
||||
maxRead: opts.maxRead,
|
||||
maxPages: fuseDefaultMaxPagesPerReq,
|
||||
initializedChan: make(chan struct{}),
|
||||
connected: true,
|
||||
|
||||
@@ -401,10 +401,12 @@ func (fd *DeviceFD) sendError(ctx context.Context, errno int32, req *Request) er
|
||||
// receiver is going to be waiting on the future channel. This is to be used by:
|
||||
// FUSE_INIT.
|
||||
func (fd *DeviceFD) noReceiverAction(ctx context.Context, r *Response) error {
|
||||
if r.opcode == linux.FUSE_INIT {
|
||||
switch r.opcode {
|
||||
case linux.FUSE_INIT:
|
||||
creds := auth.CredentialsFromContext(ctx)
|
||||
rootUserNs := kernel.KernelFromContext(ctx).RootUserNamespace()
|
||||
return fd.fs.conn.InitRecv(r, creds.HasCapabilityIn(linux.CAP_SYS_ADMIN, rootUserNs))
|
||||
// TODO(gvisor.dev/issue/3247): support async read: correctly process the response using information from r.options.
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -16,6 +16,7 @@
|
||||
package fuse
|
||||
|
||||
import (
|
||||
"math"
|
||||
"strconv"
|
||||
"sync/atomic"
|
||||
|
||||
@@ -58,6 +59,11 @@ type filesystemOptions struct {
|
||||
// exist at any time. Any further requests will block when trying to
|
||||
// Call the server.
|
||||
maxActiveRequests uint64
|
||||
|
||||
// maxRead is the max number of bytes to read,
|
||||
// specified as "max_read" in fs parameters.
|
||||
// If not specified by user, use math.MaxUint32 as default value.
|
||||
maxRead uint32
|
||||
}
|
||||
|
||||
// filesystem implements vfs.FilesystemImpl.
|
||||
@@ -144,6 +150,21 @@ func (fsType FilesystemType) GetFilesystem(ctx context.Context, vfsObj *vfs.Virt
|
||||
// Set the maxInFlightRequests option.
|
||||
fsopts.maxActiveRequests = maxActiveRequestsDefault
|
||||
|
||||
if maxReadStr, ok := mopts["max_read"]; ok {
|
||||
delete(mopts, "max_read")
|
||||
maxRead, err := strconv.ParseUint(maxReadStr, 10, 32)
|
||||
if err != nil {
|
||||
log.Warningf("%s.GetFilesystem: invalid max_read: max_read=%s", fsType.Name(), maxReadStr)
|
||||
return nil, nil, syserror.EINVAL
|
||||
}
|
||||
if maxRead < fuseMinMaxRead {
|
||||
maxRead = fuseMinMaxRead
|
||||
}
|
||||
fsopts.maxRead = uint32(maxRead)
|
||||
} else {
|
||||
fsopts.maxRead = math.MaxUint32
|
||||
}
|
||||
|
||||
// Check for unparsed options.
|
||||
if len(mopts) != 0 {
|
||||
log.Warningf("%s.GetFilesystem: unknown options: %v", fsType.Name(), mopts)
|
||||
@@ -179,7 +200,7 @@ func NewFUSEFilesystem(ctx context.Context, devMinor uint32, opts *filesystemOpt
|
||||
opts: opts,
|
||||
}
|
||||
|
||||
conn, err := newFUSEConnection(ctx, device, opts.maxActiveRequests)
|
||||
conn, err := newFUSEConnection(ctx, device, opts)
|
||||
if err != nil {
|
||||
log.Warningf("fuse.NewFUSEFilesystem: NewFUSEConnection failed with error: %v", err)
|
||||
return nil, syserror.EINVAL
|
||||
@@ -244,6 +265,7 @@ func (fs *filesystem) newInode(nodeID uint64, attr linux.FUSEAttr) *kernfs.Dentr
|
||||
i := &inode{fs: fs, NodeID: nodeID}
|
||||
creds := auth.Credentials{EffectiveKGID: auth.KGID(attr.UID), EffectiveKUID: auth.KUID(attr.UID)}
|
||||
i.InodeAttrs.Init(&creds, linux.UNNAMED_MAJOR, fs.devMinor, fs.NextIno(), linux.FileMode(attr.Mode))
|
||||
atomic.StoreUint64(&i.size, attr.Size)
|
||||
i.OrderedChildren.Init(kernfs.OrderedChildrenOptions{})
|
||||
i.EnableLeakCheck()
|
||||
i.dentry.Init(i)
|
||||
@@ -269,10 +291,13 @@ func (i *inode) Open(ctx context.Context, rp *vfs.ResolvingPath, vfsd *vfs.Dentr
|
||||
fd = &(directoryFD.fileDescription)
|
||||
fdImpl = directoryFD
|
||||
} else {
|
||||
// FOPEN_KEEP_CACHE is the defualt flag for noOpen.
|
||||
fd = &fileDescription{OpenFlag: linux.FOPEN_KEEP_CACHE}
|
||||
fdImpl = fd
|
||||
regularFD := ®ularFileFD{}
|
||||
fd = &(regularFD.fileDescription)
|
||||
fdImpl = regularFD
|
||||
}
|
||||
// FOPEN_KEEP_CACHE is the defualt flag for noOpen.
|
||||
fd.OpenFlag = linux.FOPEN_KEEP_CACHE
|
||||
|
||||
// Only send open request when FUSE server support open or is opening a directory.
|
||||
if !i.fs.conn.noOpen || isDir {
|
||||
kernelTask := kernel.TaskFromContext(ctx)
|
||||
@@ -281,21 +306,25 @@ func (i *inode) Open(ctx context.Context, rp *vfs.ResolvingPath, vfsd *vfs.Dentr
|
||||
return nil, syserror.EINVAL
|
||||
}
|
||||
|
||||
// Build the request.
|
||||
var opcode linux.FUSEOpcode
|
||||
if isDir {
|
||||
opcode = linux.FUSE_OPENDIR
|
||||
} else {
|
||||
opcode = linux.FUSE_OPEN
|
||||
}
|
||||
|
||||
in := linux.FUSEOpenIn{Flags: opts.Flags & ^uint32(linux.O_CREAT|linux.O_EXCL|linux.O_NOCTTY)}
|
||||
if !i.fs.conn.atomicOTrunc {
|
||||
in.Flags &= ^uint32(linux.O_TRUNC)
|
||||
}
|
||||
|
||||
req, err := i.fs.conn.NewRequest(auth.CredentialsFromContext(ctx), uint32(kernelTask.ThreadID()), i.NodeID, opcode, &in)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Send the request and receive the reply.
|
||||
res, err := i.fs.conn.Call(kernelTask, req)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -309,15 +338,17 @@ func (i *inode) Open(ctx context.Context, rp *vfs.ResolvingPath, vfsd *vfs.Dentr
|
||||
if err := res.UnmarshalPayload(&out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Process the reply.
|
||||
fd.OpenFlag = out.OpenFlag
|
||||
if isDir {
|
||||
fd.OpenFlag &= ^uint32(linux.FOPEN_DIRECT_IO)
|
||||
}
|
||||
|
||||
fd.Fh = out.Fh
|
||||
}
|
||||
}
|
||||
|
||||
if isDir {
|
||||
fd.OpenFlag &= ^uint32(linux.FOPEN_DIRECT_IO)
|
||||
}
|
||||
|
||||
// TODO(gvisor.dev/issue/3234): invalidate mmap after implemented it for FUSE Inode
|
||||
fd.DirectIO = fd.OpenFlag&linux.FOPEN_DIRECT_IO != 0
|
||||
fdOptions := &vfs.FileDescriptionOptions{}
|
||||
@@ -457,6 +488,16 @@ func (i *inode) Readlink(ctx context.Context, mnt *vfs.Mount) (string, error) {
|
||||
return i.link, nil
|
||||
}
|
||||
|
||||
// getFUSEAttr returns a linux.FUSEAttr of this inode stored in local cache.
|
||||
// TODO(gvisor.dev/issue/3679): Add support for other fields.
|
||||
func (i *inode) getFUSEAttr() linux.FUSEAttr {
|
||||
return linux.FUSEAttr{
|
||||
Ino: i.Ino(),
|
||||
Size: atomic.LoadUint64(&i.size),
|
||||
Mode: uint32(i.Mode()),
|
||||
}
|
||||
}
|
||||
|
||||
// statFromFUSEAttr makes attributes from linux.FUSEAttr to linux.Statx. The
|
||||
// opts.Sync attribute is ignored since the synchronization is handled by the
|
||||
// FUSE server.
|
||||
@@ -510,47 +551,90 @@ func statFromFUSEAttr(attr linux.FUSEAttr, mask, devMinor uint32) linux.Statx {
|
||||
return stat
|
||||
}
|
||||
|
||||
// Stat implements kernfs.Inode.Stat.
|
||||
func (i *inode) Stat(ctx context.Context, fs *vfs.Filesystem, opts vfs.StatOptions) (linux.Statx, error) {
|
||||
fusefs := fs.Impl().(*filesystem)
|
||||
conn := fusefs.conn
|
||||
task, creds := kernel.TaskFromContext(ctx), auth.CredentialsFromContext(ctx)
|
||||
// getAttr gets the attribute of this inode by issuing a FUSE_GETATTR request
|
||||
// or read from local cache.
|
||||
// It updates the corresponding attributes if necessary.
|
||||
func (i *inode) getAttr(ctx context.Context, fs *vfs.Filesystem, opts vfs.StatOptions) (linux.FUSEAttr, error) {
|
||||
attributeVersion := atomic.LoadUint64(&i.fs.conn.attributeVersion)
|
||||
|
||||
// TODO(gvisor.dev/issue/3679): send the request only if
|
||||
// - invalid local cache for fields specified in the opts.Mask
|
||||
// - forced update
|
||||
// - i.attributeTime expired
|
||||
// If local cache is still valid, return local cache.
|
||||
// Currently we always send a request,
|
||||
// and we always set the metadata with the new result,
|
||||
// unless attributeVersion has changed.
|
||||
|
||||
task := kernel.TaskFromContext(ctx)
|
||||
if task == nil {
|
||||
log.Warningf("couldn't get kernel task from context")
|
||||
return linux.Statx{}, syserror.EINVAL
|
||||
return linux.FUSEAttr{}, syserror.EINVAL
|
||||
}
|
||||
|
||||
creds := auth.CredentialsFromContext(ctx)
|
||||
|
||||
var in linux.FUSEGetAttrIn
|
||||
// We don't set any attribute in the request, because in VFS2 fstat(2) will
|
||||
// finally be translated into vfs.FilesystemImpl.StatAt() (see
|
||||
// pkg/sentry/syscalls/linux/vfs2/stat.go), resulting in the same flow
|
||||
// as stat(2). Thus GetAttrFlags and Fh variable will never be used in VFS2.
|
||||
req, err := conn.NewRequest(creds, uint32(task.ThreadID()), i.NodeID, linux.FUSE_GETATTR, &in)
|
||||
req, err := i.fs.conn.NewRequest(creds, uint32(task.ThreadID()), i.NodeID, linux.FUSE_GETATTR, &in)
|
||||
if err != nil {
|
||||
return linux.Statx{}, err
|
||||
return linux.FUSEAttr{}, err
|
||||
}
|
||||
|
||||
res, err := conn.Call(task, req)
|
||||
res, err := i.fs.conn.Call(task, req)
|
||||
if err != nil {
|
||||
return linux.Statx{}, err
|
||||
return linux.FUSEAttr{}, err
|
||||
}
|
||||
if err := res.Error(); err != nil {
|
||||
return linux.Statx{}, err
|
||||
return linux.FUSEAttr{}, err
|
||||
}
|
||||
|
||||
var out linux.FUSEGetAttrOut
|
||||
if err := res.UnmarshalPayload(&out); err != nil {
|
||||
return linux.Statx{}, err
|
||||
return linux.FUSEAttr{}, err
|
||||
}
|
||||
|
||||
// Set all metadata into kernfs.InodeAttrs.
|
||||
// Local version is newer, return the local one.
|
||||
// Skip the update.
|
||||
if attributeVersion != 0 && atomic.LoadUint64(&i.attributeVersion) > attributeVersion {
|
||||
return i.getFUSEAttr(), nil
|
||||
}
|
||||
|
||||
// Set the metadata of kernfs.InodeAttrs.
|
||||
if err := i.SetStat(ctx, fs, creds, vfs.SetStatOptions{
|
||||
Stat: statFromFUSEAttr(out.Attr, linux.STATX_ALL, fusefs.devMinor),
|
||||
Stat: statFromFUSEAttr(out.Attr, linux.STATX_ALL, i.fs.devMinor),
|
||||
}); err != nil {
|
||||
return linux.FUSEAttr{}, err
|
||||
}
|
||||
|
||||
// Set the size if no error (after SetStat() check).
|
||||
atomic.StoreUint64(&i.size, out.Attr.Size)
|
||||
|
||||
return out.Attr, nil
|
||||
}
|
||||
|
||||
// reviseAttr attempts to update the attributes for internal purposes
|
||||
// by calling getAttr with a pre-specified mask.
|
||||
// Used by read, write, lseek.
|
||||
func (i *inode) reviseAttr(ctx context.Context) error {
|
||||
// Never need atime for internal purposes.
|
||||
_, err := i.getAttr(ctx, i.fs.VFSFilesystem(), vfs.StatOptions{
|
||||
Mask: linux.STATX_BASIC_STATS &^ linux.STATX_ATIME,
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// Stat implements kernfs.Inode.Stat.
|
||||
func (i *inode) Stat(ctx context.Context, fs *vfs.Filesystem, opts vfs.StatOptions) (linux.Statx, error) {
|
||||
attr, err := i.getAttr(ctx, fs, opts)
|
||||
if err != nil {
|
||||
return linux.Statx{}, err
|
||||
}
|
||||
|
||||
return statFromFUSEAttr(out.Attr, opts.Mask, fusefs.devMinor), nil
|
||||
return statFromFUSEAttr(attr, opts.Mask, i.fs.devMinor), nil
|
||||
}
|
||||
|
||||
// DecRef implements kernfs.Inode.
|
||||
|
||||
@@ -29,9 +29,10 @@ const (
|
||||
// Follow the same behavior as unix fuse implementation.
|
||||
fuseMaxTimeGranNs = 1000000000
|
||||
|
||||
// Minimum value for MaxWrite.
|
||||
// Minimum value for MaxWrite and MaxRead.
|
||||
// Follow the same behavior as unix fuse implementation.
|
||||
fuseMinMaxWrite = 4096
|
||||
fuseMinMaxRead = 4096
|
||||
|
||||
// Temporary default value for max readahead, 128kb.
|
||||
fuseDefaultMaxReadahead = 131072
|
||||
|
||||
@@ -0,0 +1,152 @@
|
||||
// Copyright 2020 The gVisor Authors.
|
||||
//
|
||||
// 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 fuse
|
||||
|
||||
import (
|
||||
"io"
|
||||
"sync/atomic"
|
||||
|
||||
"gvisor.dev/gvisor/pkg/abi/linux"
|
||||
"gvisor.dev/gvisor/pkg/context"
|
||||
"gvisor.dev/gvisor/pkg/log"
|
||||
"gvisor.dev/gvisor/pkg/sentry/kernel"
|
||||
"gvisor.dev/gvisor/pkg/sentry/kernel/auth"
|
||||
"gvisor.dev/gvisor/pkg/syserror"
|
||||
"gvisor.dev/gvisor/pkg/usermem"
|
||||
)
|
||||
|
||||
// ReadInPages sends FUSE_READ requests for the size after round it up to
|
||||
// a multiple of page size, blocks on it for reply, processes the reply
|
||||
// and returns the payload (or joined payloads) as a byte slice.
|
||||
// This is used for the general purpose reading.
|
||||
// We do not support direct IO (which read the exact number of bytes)
|
||||
// at this moment.
|
||||
func (fs *filesystem) ReadInPages(ctx context.Context, fd *regularFileFD, off uint64, size uint32) ([][]byte, uint32, error) {
|
||||
attributeVersion := atomic.LoadUint64(&fs.conn.attributeVersion)
|
||||
|
||||
t := kernel.TaskFromContext(ctx)
|
||||
if t == nil {
|
||||
log.Warningf("fusefs.Read: couldn't get kernel task from context")
|
||||
return nil, 0, syserror.EINVAL
|
||||
}
|
||||
|
||||
// Round up to a multiple of page size.
|
||||
readSize, _ := usermem.PageRoundUp(uint64(size))
|
||||
|
||||
// One request cannnot exceed either maxRead or maxPages.
|
||||
maxPages := fs.conn.maxRead >> usermem.PageShift
|
||||
if maxPages > uint32(fs.conn.maxPages) {
|
||||
maxPages = uint32(fs.conn.maxPages)
|
||||
}
|
||||
|
||||
var outs [][]byte
|
||||
var sizeRead uint32
|
||||
|
||||
// readSize is a multiple of usermem.PageSize.
|
||||
// Always request bytes as a multiple of pages.
|
||||
pagesRead, pagesToRead := uint32(0), uint32(readSize>>usermem.PageShift)
|
||||
|
||||
// Reuse the same struct for unmarshalling to avoid unnecessary memory allocation.
|
||||
in := linux.FUSEReadIn{
|
||||
Fh: fd.Fh,
|
||||
LockOwner: 0, // TODO(gvisor.dev/issue/3245): file lock
|
||||
ReadFlags: 0, // TODO(gvisor.dev/issue/3245): |= linux.FUSE_READ_LOCKOWNER
|
||||
Flags: fd.statusFlags(),
|
||||
}
|
||||
|
||||
// This loop is intended for fragmented read where the bytes to read is
|
||||
// larger than either the maxPages or maxRead.
|
||||
// For the majority of reads with normal size, this loop should only
|
||||
// execute once.
|
||||
for pagesRead < pagesToRead {
|
||||
pagesCanRead := pagesToRead - pagesRead
|
||||
if pagesCanRead > maxPages {
|
||||
pagesCanRead = maxPages
|
||||
}
|
||||
|
||||
in.Offset = off + (uint64(pagesRead) << usermem.PageShift)
|
||||
in.Size = pagesCanRead << usermem.PageShift
|
||||
|
||||
req, err := fs.conn.NewRequest(auth.CredentialsFromContext(ctx), uint32(t.ThreadID()), fd.inode().NodeID, linux.FUSE_READ, &in)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
|
||||
// TODO(gvisor.dev/issue/3247): support async read.
|
||||
|
||||
res, err := fs.conn.Call(t, req)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
if err := res.Error(); err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
|
||||
// Not enough bytes in response,
|
||||
// either we reached EOF,
|
||||
// or the FUSE server sends back a response
|
||||
// that cannot even fit the hdr.
|
||||
if len(res.data) <= res.hdr.SizeBytes() {
|
||||
// We treat both case as EOF here for now
|
||||
// since there is no reliable way to detect
|
||||
// the over-short hdr case.
|
||||
break
|
||||
}
|
||||
|
||||
// Directly using the slice to avoid extra copy.
|
||||
out := res.data[res.hdr.SizeBytes():]
|
||||
|
||||
outs = append(outs, out)
|
||||
sizeRead += uint32(len(out))
|
||||
|
||||
pagesRead += pagesCanRead
|
||||
}
|
||||
|
||||
defer fs.ReadCallback(ctx, fd, off, size, sizeRead, attributeVersion)
|
||||
|
||||
// No bytes returned: offset >= EOF.
|
||||
if len(outs) == 0 {
|
||||
return nil, 0, io.EOF
|
||||
}
|
||||
|
||||
return outs, sizeRead, nil
|
||||
}
|
||||
|
||||
// ReadCallback updates several information after receiving a read response.
|
||||
// Due to readahead, sizeRead can be larger than size.
|
||||
func (fs *filesystem) ReadCallback(ctx context.Context, fd *regularFileFD, off uint64, size uint32, sizeRead uint32, attributeVersion uint64) {
|
||||
// TODO(gvisor.dev/issue/3247): support async read.
|
||||
// If this is called by an async read, correctly process it.
|
||||
// May need to update the signature.
|
||||
|
||||
i := fd.inode()
|
||||
// TODO(gvisor.dev/issue/1193): Invalidate or update atime.
|
||||
|
||||
// Reached EOF.
|
||||
if sizeRead < size {
|
||||
// TODO(gvisor.dev/issue/3630): If we have writeback cache, then we need to fill this hole.
|
||||
// Might need to update the buf to be returned from the Read().
|
||||
|
||||
// Update existing size.
|
||||
newSize := off + uint64(sizeRead)
|
||||
fs.conn.mu.Lock()
|
||||
if attributeVersion == i.attributeVersion && newSize < atomic.LoadUint64(&i.size) {
|
||||
fs.conn.attributeVersion++
|
||||
i.attributeVersion = i.fs.conn.attributeVersion
|
||||
atomic.StoreUint64(&i.size, newSize)
|
||||
}
|
||||
fs.conn.mu.Unlock()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,125 @@
|
||||
// Copyright 2020 The gVisor Authors.
|
||||
//
|
||||
// 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 fuse
|
||||
|
||||
import (
|
||||
"io"
|
||||
"math"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
|
||||
"gvisor.dev/gvisor/pkg/abi/linux"
|
||||
"gvisor.dev/gvisor/pkg/context"
|
||||
"gvisor.dev/gvisor/pkg/sentry/vfs"
|
||||
"gvisor.dev/gvisor/pkg/syserror"
|
||||
"gvisor.dev/gvisor/pkg/usermem"
|
||||
)
|
||||
|
||||
type regularFileFD struct {
|
||||
fileDescription
|
||||
|
||||
// off is the file offset.
|
||||
off int64
|
||||
// offMu protects off.
|
||||
offMu sync.Mutex
|
||||
}
|
||||
|
||||
// PRead implements vfs.FileDescriptionImpl.PRead.
|
||||
func (fd *regularFileFD) PRead(ctx context.Context, dst usermem.IOSequence, offset int64, opts vfs.ReadOptions) (int64, error) {
|
||||
if offset < 0 {
|
||||
return 0, syserror.EINVAL
|
||||
}
|
||||
|
||||
// Check that flags are supported.
|
||||
//
|
||||
// TODO(gvisor.dev/issue/2601): Support select preadv2 flags.
|
||||
if opts.Flags&^linux.RWF_HIPRI != 0 {
|
||||
return 0, syserror.EOPNOTSUPP
|
||||
}
|
||||
|
||||
size := dst.NumBytes()
|
||||
if size == 0 {
|
||||
// Early return if count is 0.
|
||||
return 0, nil
|
||||
} else if size > math.MaxUint32 {
|
||||
// FUSE only supports uint32 for size.
|
||||
// Overflow.
|
||||
return 0, syserror.EINVAL
|
||||
}
|
||||
|
||||
// TODO(gvisor.dev/issue/3678): Add direct IO support.
|
||||
|
||||
inode := fd.inode()
|
||||
|
||||
// Reading beyond EOF, update file size if outdated.
|
||||
if uint64(offset+size) > atomic.LoadUint64(&inode.size) {
|
||||
if err := inode.reviseAttr(ctx); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
// If the offset after update is still too large, return error.
|
||||
if uint64(offset) >= atomic.LoadUint64(&inode.size) {
|
||||
return 0, io.EOF
|
||||
}
|
||||
}
|
||||
|
||||
// Truncate the read with updated file size.
|
||||
fileSize := atomic.LoadUint64(&inode.size)
|
||||
if uint64(offset+size) > fileSize {
|
||||
size = int64(fileSize) - offset
|
||||
}
|
||||
|
||||
buffers, n, err := inode.fs.ReadInPages(ctx, fd, uint64(offset), uint32(size))
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
// TODO(gvisor.dev/issue/3237): support indirect IO (e.g. caching),
|
||||
// store the bytes that were read ahead.
|
||||
|
||||
// Update the number of bytes to copy for short read.
|
||||
if n < uint32(size) {
|
||||
size = int64(n)
|
||||
}
|
||||
|
||||
// Copy the bytes read to the dst.
|
||||
// This loop is intended for fragmented reads.
|
||||
// For the majority of reads, this loop only execute once.
|
||||
var copied int64
|
||||
for _, buffer := range buffers {
|
||||
toCopy := int64(len(buffer))
|
||||
if copied+toCopy > size {
|
||||
toCopy = size - copied
|
||||
}
|
||||
cp, err := dst.DropFirst64(copied).CopyOut(ctx, buffer[:toCopy])
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
if int64(cp) != toCopy {
|
||||
return 0, syserror.EIO
|
||||
}
|
||||
copied += toCopy
|
||||
}
|
||||
|
||||
return copied, nil
|
||||
}
|
||||
|
||||
// Read implements vfs.FileDescriptionImpl.Read.
|
||||
func (fd *regularFileFD) Read(ctx context.Context, dst usermem.IOSequence, opts vfs.ReadOptions) (int64, error) {
|
||||
fd.offMu.Lock()
|
||||
n, err := fd.PRead(ctx, dst, fd.off, opts)
|
||||
fd.off += n
|
||||
fd.offMu.Unlock()
|
||||
return n, err
|
||||
}
|
||||
@@ -36,3 +36,8 @@ syscall_test(
|
||||
fuse = "True",
|
||||
test = "//test/fuse/linux:mkdir_test",
|
||||
)
|
||||
|
||||
syscall_test(
|
||||
fuse = "True",
|
||||
test = "//test/fuse/linux:read_test",
|
||||
)
|
||||
|
||||
@@ -112,3 +112,16 @@ cc_library(
|
||||
"@com_google_absl//absl/strings:str_format",
|
||||
],
|
||||
)
|
||||
|
||||
cc_binary(
|
||||
name = "read_test",
|
||||
testonly = 1,
|
||||
srcs = ["read_test.cc"],
|
||||
deps = [
|
||||
gtest,
|
||||
":fuse_base",
|
||||
"//test/util:fuse_util",
|
||||
"//test/util:test_main",
|
||||
"//test/util:test_util",
|
||||
],
|
||||
)
|
||||
@@ -129,7 +129,8 @@ void FuseTest::SkipServerActualRequest() {
|
||||
|
||||
// Sends the `kSetInodeLookup` command, expected mode, and the path of the
|
||||
// inode to create under the mount point.
|
||||
void FuseTest::SetServerInodeLookup(const std::string& path, mode_t mode) {
|
||||
void FuseTest::SetServerInodeLookup(const std::string& path, mode_t mode,
|
||||
uint64_t size) {
|
||||
uint32_t cmd = static_cast<uint32_t>(FuseTestCmd::kSetInodeLookup);
|
||||
EXPECT_THAT(RetryEINTR(write)(sock_[0], &cmd, sizeof(cmd)),
|
||||
SyscallSucceedsWithValue(sizeof(cmd)));
|
||||
@@ -137,6 +138,9 @@ void FuseTest::SetServerInodeLookup(const std::string& path, mode_t mode) {
|
||||
EXPECT_THAT(RetryEINTR(write)(sock_[0], &mode, sizeof(mode)),
|
||||
SyscallSucceedsWithValue(sizeof(mode)));
|
||||
|
||||
EXPECT_THAT(RetryEINTR(write)(sock_[0], &size, sizeof(size)),
|
||||
SyscallSucceedsWithValue(sizeof(size)));
|
||||
|
||||
// Pad 1 byte for null-terminate c-string.
|
||||
EXPECT_THAT(RetryEINTR(write)(sock_[0], path.c_str(), path.size() + 1),
|
||||
SyscallSucceedsWithValue(path.size() + 1));
|
||||
@@ -144,10 +148,10 @@ void FuseTest::SetServerInodeLookup(const std::string& path, mode_t mode) {
|
||||
WaitServerComplete();
|
||||
}
|
||||
|
||||
void FuseTest::MountFuse() {
|
||||
void FuseTest::MountFuse(const char* mountOpts) {
|
||||
EXPECT_THAT(dev_fd_ = open("/dev/fuse", O_RDWR), SyscallSucceeds());
|
||||
|
||||
std::string mount_opts = absl::StrFormat("fd=%d,%s", dev_fd_, kMountOpts);
|
||||
std::string mount_opts = absl::StrFormat("fd=%d,%s", dev_fd_, mountOpts);
|
||||
mount_point_ = ASSERT_NO_ERRNO_AND_VALUE(TempPath::CreateDir());
|
||||
EXPECT_THAT(mount("fuse", mount_point_.path().c_str(), "fuse",
|
||||
MS_NODEV | MS_NOSUID, mount_opts.c_str()),
|
||||
@@ -311,11 +315,15 @@ void FuseTest::ServerHandleCommand() {
|
||||
// request with this specific path comes in.
|
||||
void FuseTest::ServerReceiveInodeLookup() {
|
||||
mode_t mode;
|
||||
uint64_t size;
|
||||
std::vector<char> buf(FUSE_MIN_READ_BUFFER);
|
||||
|
||||
EXPECT_THAT(RetryEINTR(read)(sock_[1], &mode, sizeof(mode)),
|
||||
SyscallSucceedsWithValue(sizeof(mode)));
|
||||
|
||||
EXPECT_THAT(RetryEINTR(read)(sock_[1], &size, sizeof(size)),
|
||||
SyscallSucceedsWithValue(sizeof(size)));
|
||||
|
||||
EXPECT_THAT(RetryEINTR(read)(sock_[1], buf.data(), buf.size()),
|
||||
SyscallSucceeds());
|
||||
|
||||
@@ -332,6 +340,9 @@ void FuseTest::ServerReceiveInodeLookup() {
|
||||
// comply with the unqiueness of different path.
|
||||
++nodeid_;
|
||||
|
||||
// Set the size.
|
||||
out_payload.attr.size = size;
|
||||
|
||||
memcpy(buf.data(), &out_header, sizeof(out_header));
|
||||
memcpy(buf.data() + sizeof(out_header), &out_payload, sizeof(out_payload));
|
||||
lookups_.AddMemBlock(FUSE_LOOKUP, buf.data(), out_len);
|
||||
|
||||
@@ -137,7 +137,8 @@ class FuseTest : public ::testing::Test {
|
||||
// path, pretending there is an inode and avoid ENOENT when testing. If mode
|
||||
// is not given, it creates a regular file with mode 0600.
|
||||
void SetServerInodeLookup(const std::string& path,
|
||||
mode_t mode = S_IFREG | S_IRUSR | S_IWUSR);
|
||||
mode_t mode = S_IFREG | S_IRUSR | S_IWUSR,
|
||||
uint64_t size = 512);
|
||||
|
||||
// Called by the testing thread to ask the FUSE server for its next received
|
||||
// FUSE request. Be sure to use the corresponding struct of iovec to receive
|
||||
@@ -166,16 +167,16 @@ class FuseTest : public ::testing::Test {
|
||||
protected:
|
||||
TempPath mount_point_;
|
||||
|
||||
// Unmounts the mountpoint of the FUSE server.
|
||||
void UnmountFuse();
|
||||
|
||||
private:
|
||||
// Opens /dev/fuse and inherit the file descriptor for the FUSE server.
|
||||
void MountFuse();
|
||||
void MountFuse(const char* mountOpts = kMountOpts);
|
||||
|
||||
// Creates a socketpair for communication and forks FUSE server.
|
||||
void SetUpFuseServer();
|
||||
|
||||
// Unmounts the mountpoint of the FUSE server.
|
||||
void UnmountFuse();
|
||||
|
||||
private:
|
||||
// Sends a FuseTestCmd and gets a uint32_t data from the FUSE server.
|
||||
inline uint32_t GetServerData(uint32_t cmd);
|
||||
|
||||
|
||||
@@ -0,0 +1,390 @@
|
||||
// Copyright 2020 The gVisor Authors.
|
||||
//
|
||||
// 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.
|
||||
|
||||
#include <errno.h>
|
||||
#include <fcntl.h>
|
||||
#include <linux/fuse.h>
|
||||
#include <sys/stat.h>
|
||||
#include <sys/statfs.h>
|
||||
#include <sys/types.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include <string>
|
||||
#include <vector>
|
||||
|
||||
#include "gtest/gtest.h"
|
||||
#include "test/fuse/linux/fuse_base.h"
|
||||
#include "test/util/fuse_util.h"
|
||||
#include "test/util/test_util.h"
|
||||
|
||||
namespace gvisor {
|
||||
namespace testing {
|
||||
|
||||
namespace {
|
||||
|
||||
class ReadTest : public FuseTest {
|
||||
void SetUp() override {
|
||||
FuseTest::SetUp();
|
||||
test_file_path_ = JoinPath(mount_point_.path().c_str(), test_file_);
|
||||
}
|
||||
|
||||
// TearDown overrides the parent's function
|
||||
// to skip checking the unconsumed release request at the end.
|
||||
void TearDown() override { UnmountFuse(); }
|
||||
|
||||
protected:
|
||||
const std::string test_file_ = "test_file";
|
||||
const mode_t test_file_mode_ = S_IFREG | S_IRWXU | S_IRWXG | S_IRWXO;
|
||||
const uint64_t test_fh_ = 1;
|
||||
const uint32_t open_flag_ = O_RDWR;
|
||||
|
||||
std::string test_file_path_;
|
||||
|
||||
PosixErrorOr<FileDescriptor> OpenTestFile(const std::string &path,
|
||||
uint64_t size = 512) {
|
||||
SetServerInodeLookup(test_file_, test_file_mode_, size);
|
||||
|
||||
struct fuse_out_header out_header_open = {
|
||||
.len = sizeof(struct fuse_out_header) + sizeof(struct fuse_open_out),
|
||||
};
|
||||
struct fuse_open_out out_payload_open = {
|
||||
.fh = test_fh_,
|
||||
.open_flags = open_flag_,
|
||||
};
|
||||
auto iov_out_open = FuseGenerateIovecs(out_header_open, out_payload_open);
|
||||
SetServerResponse(FUSE_OPEN, iov_out_open);
|
||||
|
||||
auto res = Open(path.c_str(), open_flag_);
|
||||
if (res.ok()) {
|
||||
SkipServerActualRequest();
|
||||
}
|
||||
return res;
|
||||
}
|
||||
};
|
||||
|
||||
class ReadTestSmallMaxRead : public ReadTest {
|
||||
void SetUp() override {
|
||||
MountFuse(mountOpts);
|
||||
SetUpFuseServer();
|
||||
test_file_path_ = JoinPath(mount_point_.path().c_str(), test_file_);
|
||||
}
|
||||
|
||||
protected:
|
||||
constexpr static char mountOpts[] =
|
||||
"rootmode=755,user_id=0,group_id=0,max_read=4096";
|
||||
// 4096 is hard-coded as the max_read in mount options.
|
||||
const int size_fragment = 4096;
|
||||
};
|
||||
|
||||
TEST_F(ReadTest, ReadWhole) {
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_));
|
||||
|
||||
// Prepare for the read.
|
||||
const int n_read = 5;
|
||||
std::vector<char> data(n_read);
|
||||
RandomizeBuffer(data.data(), data.size());
|
||||
struct fuse_out_header out_header_read = {
|
||||
.len =
|
||||
static_cast<uint32_t>(sizeof(struct fuse_out_header) + data.size()),
|
||||
};
|
||||
auto iov_out_read = FuseGenerateIovecs(out_header_read, data);
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
|
||||
// Read the whole "file".
|
||||
std::vector<char> buf(n_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), n_read),
|
||||
SyscallSucceedsWithValue(n_read));
|
||||
|
||||
// Check the read request.
|
||||
struct fuse_in_header in_header_read;
|
||||
struct fuse_read_in in_payload_read;
|
||||
auto iov_in = FuseGenerateIovecs(in_header_read, in_payload_read);
|
||||
GetServerActualRequest(iov_in);
|
||||
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_header_read.len,
|
||||
sizeof(in_header_read) + sizeof(in_payload_read));
|
||||
EXPECT_EQ(in_header_read.opcode, FUSE_READ);
|
||||
EXPECT_EQ(in_payload_read.offset, 0);
|
||||
EXPECT_EQ(buf, data);
|
||||
}
|
||||
|
||||
TEST_F(ReadTest, ReadPartial) {
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_));
|
||||
|
||||
// Prepare for the read.
|
||||
const int n_data = 10;
|
||||
std::vector<char> data(n_data);
|
||||
RandomizeBuffer(data.data(), data.size());
|
||||
// Note: due to read ahead, current read implementation will treat any
|
||||
// response that is longer than requested as correct (i.e. not reach the EOF).
|
||||
// Therefore, the test below should make sure the size to read does not exceed
|
||||
// n_data.
|
||||
struct fuse_out_header out_header_read = {
|
||||
.len =
|
||||
static_cast<uint32_t>(sizeof(struct fuse_out_header) + data.size()),
|
||||
};
|
||||
auto iov_out_read = FuseGenerateIovecs(out_header_read, data);
|
||||
struct fuse_in_header in_header_read;
|
||||
struct fuse_read_in in_payload_read;
|
||||
auto iov_in = FuseGenerateIovecs(in_header_read, in_payload_read);
|
||||
|
||||
std::vector<char> buf(n_data);
|
||||
|
||||
// Read 1 bytes.
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), 1), SyscallSucceedsWithValue(1));
|
||||
|
||||
// Check the 1-byte read request.
|
||||
GetServerActualRequest(iov_in);
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_header_read.len,
|
||||
sizeof(in_header_read) + sizeof(in_payload_read));
|
||||
EXPECT_EQ(in_header_read.opcode, FUSE_READ);
|
||||
EXPECT_EQ(in_payload_read.offset, 0);
|
||||
|
||||
// Read 3 bytes.
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), 3), SyscallSucceedsWithValue(3));
|
||||
|
||||
// Check the 3-byte read request.
|
||||
GetServerActualRequest(iov_in);
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_payload_read.offset, 1);
|
||||
|
||||
// Read 5 bytes.
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), 5), SyscallSucceedsWithValue(5));
|
||||
|
||||
// Check the 5-byte read request.
|
||||
GetServerActualRequest(iov_in);
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_payload_read.offset, 4);
|
||||
}
|
||||
|
||||
TEST_F(ReadTest, PRead) {
|
||||
const int file_size = 512;
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_, file_size));
|
||||
|
||||
// Prepare for the read.
|
||||
const int n_read = 5;
|
||||
std::vector<char> data(n_read);
|
||||
RandomizeBuffer(data.data(), data.size());
|
||||
struct fuse_out_header out_header_read = {
|
||||
.len =
|
||||
static_cast<uint32_t>(sizeof(struct fuse_out_header) + data.size()),
|
||||
};
|
||||
auto iov_out_read = FuseGenerateIovecs(out_header_read, data);
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
|
||||
// Read some bytes.
|
||||
std::vector<char> buf(n_read);
|
||||
const int offset_read = file_size >> 1;
|
||||
EXPECT_THAT(pread(fd.get(), buf.data(), n_read, offset_read),
|
||||
SyscallSucceedsWithValue(n_read));
|
||||
|
||||
// Check the read request.
|
||||
struct fuse_in_header in_header_read;
|
||||
struct fuse_read_in in_payload_read;
|
||||
auto iov_in = FuseGenerateIovecs(in_header_read, in_payload_read);
|
||||
GetServerActualRequest(iov_in);
|
||||
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_header_read.len,
|
||||
sizeof(in_header_read) + sizeof(in_payload_read));
|
||||
EXPECT_EQ(in_header_read.opcode, FUSE_READ);
|
||||
EXPECT_EQ(in_payload_read.offset, offset_read);
|
||||
EXPECT_EQ(buf, data);
|
||||
}
|
||||
|
||||
TEST_F(ReadTest, ReadZero) {
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_));
|
||||
|
||||
// Issue the read.
|
||||
std::vector<char> buf;
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), 0), SyscallSucceedsWithValue(0));
|
||||
}
|
||||
|
||||
TEST_F(ReadTest, ReadShort) {
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_));
|
||||
|
||||
// Prepare for the short read.
|
||||
const int n_read = 5;
|
||||
std::vector<char> data(n_read >> 1);
|
||||
RandomizeBuffer(data.data(), data.size());
|
||||
struct fuse_out_header out_header_read = {
|
||||
.len =
|
||||
static_cast<uint32_t>(sizeof(struct fuse_out_header) + data.size()),
|
||||
};
|
||||
auto iov_out_read = FuseGenerateIovecs(out_header_read, data);
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
|
||||
// Read the whole "file".
|
||||
std::vector<char> buf(n_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), n_read),
|
||||
SyscallSucceedsWithValue(data.size()));
|
||||
|
||||
// Check the read request.
|
||||
struct fuse_in_header in_header_read;
|
||||
struct fuse_read_in in_payload_read;
|
||||
auto iov_in = FuseGenerateIovecs(in_header_read, in_payload_read);
|
||||
GetServerActualRequest(iov_in);
|
||||
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_header_read.len,
|
||||
sizeof(in_header_read) + sizeof(in_payload_read));
|
||||
EXPECT_EQ(in_header_read.opcode, FUSE_READ);
|
||||
EXPECT_EQ(in_payload_read.offset, 0);
|
||||
std::vector<char> short_buf(buf.begin(), buf.begin() + data.size());
|
||||
EXPECT_EQ(short_buf, data);
|
||||
}
|
||||
|
||||
TEST_F(ReadTest, ReadShortEOF) {
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_));
|
||||
|
||||
// Prepare for the short read.
|
||||
struct fuse_out_header out_header_read = {
|
||||
.len = static_cast<uint32_t>(sizeof(struct fuse_out_header)),
|
||||
};
|
||||
auto iov_out_read = FuseGenerateIovecs(out_header_read);
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
|
||||
// Read the whole "file".
|
||||
const int n_read = 10;
|
||||
std::vector<char> buf(n_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), n_read), SyscallSucceedsWithValue(0));
|
||||
|
||||
// Check the read request.
|
||||
struct fuse_in_header in_header_read;
|
||||
struct fuse_read_in in_payload_read;
|
||||
auto iov_in = FuseGenerateIovecs(in_header_read, in_payload_read);
|
||||
GetServerActualRequest(iov_in);
|
||||
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_header_read.len,
|
||||
sizeof(in_header_read) + sizeof(in_payload_read));
|
||||
EXPECT_EQ(in_header_read.opcode, FUSE_READ);
|
||||
EXPECT_EQ(in_payload_read.offset, 0);
|
||||
}
|
||||
|
||||
TEST_F(ReadTestSmallMaxRead, ReadSmallMaxRead) {
|
||||
const int n_fragment = 10;
|
||||
const int n_read = size_fragment * n_fragment;
|
||||
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_, n_read));
|
||||
|
||||
// Prepare for the read.
|
||||
std::vector<char> data(size_fragment);
|
||||
RandomizeBuffer(data.data(), data.size());
|
||||
struct fuse_out_header out_header_read = {
|
||||
.len =
|
||||
static_cast<uint32_t>(sizeof(struct fuse_out_header) + data.size()),
|
||||
};
|
||||
auto iov_out_read = FuseGenerateIovecs(out_header_read, data);
|
||||
|
||||
for (int i = 0; i < n_fragment; ++i) {
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
}
|
||||
|
||||
// Read the whole "file".
|
||||
std::vector<char> buf(n_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), n_read),
|
||||
SyscallSucceedsWithValue(n_read));
|
||||
|
||||
ASSERT_EQ(GetServerNumUnsentResponses(), 0);
|
||||
ASSERT_EQ(GetServerNumUnconsumedRequests(), n_fragment);
|
||||
|
||||
// Check each read segment.
|
||||
struct fuse_in_header in_header_read;
|
||||
struct fuse_read_in in_payload_read;
|
||||
auto iov_in = FuseGenerateIovecs(in_header_read, in_payload_read);
|
||||
|
||||
for (int i = 0; i < n_fragment; ++i) {
|
||||
GetServerActualRequest(iov_in);
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_header_read.len,
|
||||
sizeof(in_header_read) + sizeof(in_payload_read));
|
||||
EXPECT_EQ(in_header_read.opcode, FUSE_READ);
|
||||
EXPECT_EQ(in_payload_read.offset, i * size_fragment);
|
||||
EXPECT_EQ(in_payload_read.size, size_fragment);
|
||||
|
||||
auto it = buf.begin() + i * size_fragment;
|
||||
EXPECT_EQ(std::vector<char>(it, it + size_fragment), data);
|
||||
}
|
||||
}
|
||||
|
||||
TEST_F(ReadTestSmallMaxRead, ReadSmallMaxReadShort) {
|
||||
const int n_fragment = 10;
|
||||
const int n_read = size_fragment * n_fragment;
|
||||
|
||||
auto fd = ASSERT_NO_ERRNO_AND_VALUE(OpenTestFile(test_file_path_, n_read));
|
||||
|
||||
// Prepare for the read.
|
||||
std::vector<char> data(size_fragment);
|
||||
RandomizeBuffer(data.data(), data.size());
|
||||
struct fuse_out_header out_header_read = {
|
||||
.len =
|
||||
static_cast<uint32_t>(sizeof(struct fuse_out_header) + data.size()),
|
||||
};
|
||||
auto iov_out_read = FuseGenerateIovecs(out_header_read, data);
|
||||
|
||||
for (int i = 0; i < n_fragment - 1; ++i) {
|
||||
SetServerResponse(FUSE_READ, iov_out_read);
|
||||
}
|
||||
|
||||
// The last fragment is a short read.
|
||||
std::vector<char> half_data(data.begin(), data.begin() + (data.size() >> 1));
|
||||
struct fuse_out_header out_header_read_short = {
|
||||
.len = static_cast<uint32_t>(sizeof(struct fuse_out_header) +
|
||||
half_data.size()),
|
||||
};
|
||||
auto iov_out_read_short =
|
||||
FuseGenerateIovecs(out_header_read_short, half_data);
|
||||
SetServerResponse(FUSE_READ, iov_out_read_short);
|
||||
|
||||
// Read the whole "file".
|
||||
std::vector<char> buf(n_read);
|
||||
EXPECT_THAT(read(fd.get(), buf.data(), n_read),
|
||||
SyscallSucceedsWithValue(n_read - (data.size() >> 1)));
|
||||
|
||||
ASSERT_EQ(GetServerNumUnsentResponses(), 0);
|
||||
ASSERT_EQ(GetServerNumUnconsumedRequests(), n_fragment);
|
||||
|
||||
// Check each read segment.
|
||||
struct fuse_in_header in_header_read;
|
||||
struct fuse_read_in in_payload_read;
|
||||
auto iov_in = FuseGenerateIovecs(in_header_read, in_payload_read);
|
||||
|
||||
for (int i = 0; i < n_fragment; ++i) {
|
||||
GetServerActualRequest(iov_in);
|
||||
EXPECT_EQ(in_payload_read.fh, test_fh_);
|
||||
EXPECT_EQ(in_header_read.len,
|
||||
sizeof(in_header_read) + sizeof(in_payload_read));
|
||||
EXPECT_EQ(in_header_read.opcode, FUSE_READ);
|
||||
EXPECT_EQ(in_payload_read.offset, i * size_fragment);
|
||||
EXPECT_EQ(in_payload_read.size, size_fragment);
|
||||
|
||||
auto it = buf.begin() + i * size_fragment;
|
||||
if (i != n_fragment - 1) {
|
||||
EXPECT_EQ(std::vector<char>(it, it + data.size()), data);
|
||||
} else {
|
||||
EXPECT_EQ(std::vector<char>(it, it + half_data.size()), half_data);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
} // namespace
|
||||
|
||||
} // namespace testing
|
||||
} // namespace gvisor
|
||||
Reference in New Issue
Block a user