fix(client): use sync flush if open file as directIO

close:#1000464581

Signed-off-by: chihe <chihe@oppo.com>
(cherry picked from commit 34fb8cdeb2)
This commit is contained in:
chihe 2025-11-25 19:51:57 +08:00
parent cbbead5859
commit 7a618a44ad
11 changed files with 83 additions and 25 deletions

View File

@ -551,7 +551,8 @@ func (f *File) Write(ctx context.Context, req *fuse.WriteRequest, resp *fuse.Wri
var size int
if f.shouldAccessReplicaStorageClass() {
f.super.ec.GetStreamer(ino).SetParentInode(f.parentIno)
if size, err = f.super.ec.Write(ino, int(req.Offset), req.Data, flags, checkFunc, f.info.StorageClass, false); err == ParseError(syscall.ENOSPC) {
if size, err = f.super.ec.Write(ino, int(req.Offset), req.Data, flags, checkFunc, f.info.StorageClass,
false, waitForFlush); err == ParseError(syscall.ENOSPC) {
return
}
} else {

View File

@ -1057,7 +1057,7 @@ func (c *Client) write(f *File, offset int64, data []byte, flags int) (n int, er
}
return nil
}
n, err = c.ec.Write(f.ino, int(offset), data, flags, checkFunc, f.storageClass, false)
n, err = c.ec.Write(f.ino, int(offset), data, flags, checkFunc, f.storageClass, false, false)
} else {
n, err = f.fileWriter.Write(c.ctx(c.ID, f.ino), int(offset), data, flags)
}

View File

@ -1796,7 +1796,7 @@ func (c *client) write(f *file, offset int, data []byte, flags int) (n int, err
}
return nil
}
n, err = c.ec.Write(f.ino, offset, data, flags, checkFunc, f.storageClass, false)
n, err = c.ec.Write(f.ino, offset, data, flags, checkFunc, f.storageClass, false, false)
} else {
n, err = f.fileWriter.Write(c.ctx(c.id, f.ino), offset, data, flags)
}

View File

@ -31,7 +31,7 @@ type ExtentApi interface {
OpenStream(inode uint64, openForWrite, isCache bool, fullPath string) error
CloseStream(inode uint64) error
Read(inode uint64, data []byte, offset int, size int, storageClass uint32, isMigration bool) (read int, err error)
Write(inode uint64, offset int, data []byte, flags int, checkFunc func() error, storageClass uint32, isMigration bool) (write int, err error)
Write(inode uint64, offset int, data []byte, flags int, checkFunc func() error, storageClass uint32, isMigration, waitForFlush bool) (write int, err error)
Flush(inode uint64) error
Close() error
}
@ -103,7 +103,7 @@ func (t *TransitionMgr) migrate(e *proto.ScanDentry) (err error) {
return
}
if readN > 0 {
writeN, err = t.ecForW.Write(e.Inode, writeOffset, buf[:readN], 0, nil, proto.OpTypeToStorageType(e.Op), true)
writeN, err = t.ecForW.Write(e.Inode, writeOffset, buf[:readN], 0, nil, proto.OpTypeToStorageType(e.Op), true, false)
if err != nil {
err = fmt.Errorf("write dst file err(%v)", err)
log.LogWarnf("migrate: inode(%v), writeOffset(%v): %v", e.Inode, writeOffset, err)

View File

@ -52,7 +52,7 @@ func (m *MockExtentClient) Read(inode uint64, data []byte, offset int, size int,
return len(data), io.EOF
}
func (m *MockExtentClient) Write(inode uint64, offset int, data []byte, flags int, checkFunc func() error, storageClass uint32, isMigration bool) (write int, err error) {
func (m *MockExtentClient) Write(inode uint64, offset int, data []byte, flags int, checkFunc func() error, storageClass uint32, isMigration, waitForFlush bool) (write int, err error) {
m.data = data
return
}

View File

@ -1417,7 +1417,8 @@ func (v *Volume) streamWrite(inode uint64, reader io.Reader, h hash.Hash, storag
}
return nil
}
if writeN, err = v.ec.Write(inode, offset, buf[:readN], 0, checkFunc, storageClass, false); err != nil {
if writeN, err = v.ec.Write(inode, offset, buf[:readN], 0, checkFunc, storageClass,
false, false); err != nil {
log.LogErrorf("streamWrite: data write tmp file fail, inode(%v) offset(%v) err(%v)", inode, offset, err)
exporter.Warning(fmt.Sprintf("write data fail: volume(%v) inode(%v) offset(%v) size(%v) err(%v)",
v.name, inode, offset, readN, err))
@ -2874,7 +2875,8 @@ func (v *Volume) CopyFile(sv *Volume, sourcePath, targetPath, metaDirective stri
if proto.IsCold(v.volType) || proto.IsStorageClassBlobStore(tInodeInfo.StorageClass) {
writeN, err = ebsWriter.WriteWithoutPool(tctx, writeOffset, buf[:readN])
} else {
writeN, err = v.ec.Write(tInodeInfo.Inode, writeOffset, buf[:readN], 0, nil, tInodeInfo.StorageClass, false)
writeN, err = v.ec.Write(tInodeInfo.Inode, writeOffset, buf[:readN], 0, nil,
tInodeInfo.StorageClass, false, false)
}
if err != nil {
log.LogErrorf("CopyFile: write target path from volume (%v) path(%v) fail, volume(%v) path(%v) inode(%v) target offset(%v) err(%v)",

View File

@ -455,7 +455,7 @@ func MockCheckDataPartitionExistFalse(client *stream.ExtentClient, partitionID u
}
func MockWriteTrue(client *stream.ExtentClient, inode uint64, offset int, data []byte,
flags int, checkFunc func() error, storageClass uint32, isMigration bool,
flags int, checkFunc func() error, storageClass uint32, isMigration, waitForFlush bool,
) (write int, err error) {
return len(data), nil
}

View File

@ -30,6 +30,7 @@ type ExtentRequest struct {
Size int
Data []byte
ExtentKey *proto.ExtentKey
CreateNewEk bool
}
// String returns the string format of the extent request.
@ -531,7 +532,7 @@ func (cache *ExtentCache) PrepareWriteRequests(offset, size int, data []byte) []
log.LogDebugf("action[ExtentCache.PrepareWriteRequests] ek [%v], pivot[%v]", ek, pivot)
return false
})
createNewExtentKey := false
cache.root.AscendRange(lower, upper, func(i btree.Item) bool {
ek := i.(*proto.ExtentKey)
ekStart := int(ek.FileOffset)
@ -555,6 +556,10 @@ func (cache *ExtentCache) PrepareWriteRequests(offset, size int, data []byte) []
start = end
return false
} else {
// exact cover current extent: start == ekStart && end == ekEnd
if start == ekStart && end == ekEnd {
createNewExtentKey = true
}
return true
}
} else if start < ekEnd {
@ -580,6 +585,9 @@ func (cache *ExtentCache) PrepareWriteRequests(offset, size int, data []byte) []
if start < end {
// add hole (start, end)
req := NewExtentRequest(start, end-start, data[start-offset:end-offset], nil)
if createNewExtentKey {
req.CreateNewEk = true
}
requests = append(requests, req)
}

View File

@ -709,7 +709,8 @@ func (client *ExtentClient) SetFileSize(inode uint64, size int, sync bool) {
}
// Write writes the data.
func (client *ExtentClient) Write(inode uint64, offset int, data []byte, flags int, checkFunc func() error, storageClass uint32, isMigration bool) (write int, err error) {
func (client *ExtentClient) Write(inode uint64, offset int, data []byte, flags int, checkFunc func() error,
storageClass uint32, isMigration, waitForFlush bool) (write int, err error) {
prefix := fmt.Sprintf("Write{ino(%v)offset(%v)size(%v)}", inode, offset, len(data))
s := client.GetStreamer(inode)
if s == nil {
@ -727,7 +728,7 @@ func (client *ExtentClient) Write(inode uint64, offset int, data []byte, flags i
// TODO unhandled error
s.GetExtents(isMigration)
})
s.waitForFlush = waitForFlush
write, err = s.IssueWriteRequest(offset, data, flags, checkFunc, storageClass, isMigration)
if err != nil {
log.LogError(errors.Stack(err))

View File

@ -84,6 +84,7 @@ type Streamer struct {
writeProtectionLock sync.Mutex // protects write operation state
aheadReadBlockSize uint32
waitForFlush bool
}
type bcacheKey struct {
@ -149,8 +150,9 @@ func (s *Streamer) SetFullPath(fullPath string) {
// String returns the string format of the streamer.
func (s *Streamer) String() string {
return fmt.Sprintf("Streamer{ino(%v), fullPath(%v), refcnt(%v), isOpen(%v) openForWrite(%v), request(%v), eh(%v) addr(%p)}",
s.inode, s.fullPath, atomic.LoadInt32(&s.refcnt), s.isOpen, s.openForWrite, len(s.request), s.handler, s)
return fmt.Sprintf("Streamer{ino(%v), fullPath(%v), refcnt(%v), isOpen(%v) openForWrite(%v), request(%v), "+
"eh(%v) waitForFlush(%v) addr(%p)}",
s.inode, s.fullPath, atomic.LoadInt32(&s.refcnt), s.isOpen, s.openForWrite, len(s.request), s.handler, s.waitForFlush, s)
}
// TODO should we call it RefreshExtents instead?
@ -611,7 +613,7 @@ func (s *Streamer) completeAsyncFlush(req *AsyncFlushRequest) {
if nextReq.handler.id >= handler.id {
goto end
}
time.Sleep(100 * time.Millisecond)
time.Sleep(1 * time.Millisecond)
}
}
}
@ -692,14 +694,18 @@ func (s *Streamer) getActiveHandlerFlush(handlerID uint64) *AsyncFlushRequest {
// addPendingAsyncFlush adds a request to the pending map using handler.id as key
func (s *Streamer) addPendingAsyncFlush(handlerID uint64, req *AsyncFlushRequest) {
s.pendingAsyncFlushMap.Store(handlerID, req)
if log.EnableDebug() {
log.LogDebugf("addPendingAsyncFlush: streamer(%v) handler(%v) trace(%v)", s.inode, handlerID, string(debug.Stack()))
}
}
// removePendingAsyncFlush removes a request from the pending map
func (s *Streamer) removePendingAsyncFlush(handlerID uint64) {
s.pendingAsyncFlushMap.Delete(handlerID)
if log.EnableDebug() {
log.LogDebugf("removePendingAsyncFlush streamer(%v) handler(%v) trace(%v)", s.inode, handlerID, string(debug.Stack()))
}
}
// getPendingRequestsCount returns the number of pending requests
func (s *Streamer) getPendingRequests() []uint64 {

View File

@ -385,7 +385,7 @@ begin:
isChecked := false
// Must flush before doing overwrite
for _, req := range requests {
if req.ExtentKey == nil && offset >= fileSize {
if req.ExtentKey == nil && !req.CreateNewEk {
continue
}
err = s.flush(true, uuid.New().String())
@ -1004,7 +1004,7 @@ func (s *Streamer) doWriteAppendEx(data []byte, offset, size int, direct bool, r
return
}
func (s *Streamer) flush(wait bool, id string) (err error) {
func (s *Streamer) flushAsync(wait bool, id string) (err error) {
pending := make(map[*ExtentHandler]chan error)
asyncExtentHandler := make([]*ExtentHandler, 0)
for {
@ -1079,7 +1079,7 @@ func (s *Streamer) flush(wait bool, id string) (err error) {
}
}
if !progressed {
time.Sleep(time.Millisecond * 10)
time.Sleep(time.Millisecond * 1)
}
}
if firstErr != nil {
@ -1164,7 +1164,8 @@ func (s *Streamer) closeOpenHandler(wait bool) (err error) {
}(handler)
// Use async flush if async flush is enabled
if s.client.enableAsyncFlush {
// directIO use sync flush
if s.client.enableAsyncFlush && !s.waitForFlush {
log.LogDebugf("closeOpenHandler: using async flush for handler(%v) with inflight(%v) id(%v)",
handler, atomic.LoadInt32(&handler.inflight), id)
s.requestAsyncFlush(handler, cleanFunc)
@ -1340,3 +1341,42 @@ func (s *Streamer) tinySizeLimit() int {
func (s *Streamer) setError() {
atomic.StoreInt32(&s.status, StreamerError)
}
func (s *Streamer) flushSync() (err error) {
for {
element := s.dirtylist.Get()
if element == nil {
break
}
eh := element.Value.(*ExtentHandler)
log.LogDebugf("Streamer flush begin: eh(%v)", eh)
err = eh.flush()
if err != nil {
log.LogErrorf("Streamer flush failed: eh(%v)", eh)
return
}
eh.stream.dirtylist.Remove(element)
if eh.getStatus() == ExtentStatusOpen {
s.dirty = false
log.LogDebugf("Streamer flush handler open: eh(%v)", eh)
} else {
// TODO unhandled error
eh.cleanup()
log.LogDebugf("Streamer flush handler cleaned up: eh(%v)", eh)
}
// During LTP testing, consecutive write requests may be interspersed with direct I/O operations,
// and thus flushSync may flush asynchronous event handles (eh).
s.removePendingAsyncFlush(eh.id)
log.LogDebugf("Streamer flush end: eh(%v)", eh)
}
return
}
func (s *Streamer) flush(wait bool, id string) (err error) {
if s.client.enableAsyncFlush && !s.waitForFlush {
return s.flushAsync(wait, id)
} else {
return s.flushSync()
}
}