From 7a618a44ad8d23130c529faa0ddb4f193b892013 Mon Sep 17 00:00:00 2001 From: chihe Date: Tue, 25 Nov 2025 19:51:57 +0800 Subject: [PATCH] fix(client): use sync flush if open file as directIO close:#1000464581 Signed-off-by: chihe (cherry picked from commit 34fb8cdeb227c880acca8d6040c4de7569c3e2c4) --- client/fs/file.go | 3 +- client/gosdk/cfs_client.go | 2 +- client/libsdk/libsdk.go | 2 +- lcnode/lc_transition.go | 4 +-- lcnode/lc_transition_test.go | 2 +- objectnode/fs_volume.go | 6 ++-- sdk/data/blobstore/reader_test.go | 2 +- sdk/data/stream/extent_cache.go | 18 ++++++++---- sdk/data/stream/extent_client.go | 5 ++-- sdk/data/stream/stream_reader.go | 16 +++++++---- sdk/data/stream/stream_writer.go | 48 ++++++++++++++++++++++++++++--- 11 files changed, 83 insertions(+), 25 deletions(-) diff --git a/client/fs/file.go b/client/fs/file.go index 5e4cf5e01..7ed692280 100644 --- a/client/fs/file.go +++ b/client/fs/file.go @@ -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 { diff --git a/client/gosdk/cfs_client.go b/client/gosdk/cfs_client.go index 38b418e90..678e2ecd6 100644 --- a/client/gosdk/cfs_client.go +++ b/client/gosdk/cfs_client.go @@ -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) } diff --git a/client/libsdk/libsdk.go b/client/libsdk/libsdk.go index 3795a8746..04e54d25d 100644 --- a/client/libsdk/libsdk.go +++ b/client/libsdk/libsdk.go @@ -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) } diff --git a/lcnode/lc_transition.go b/lcnode/lc_transition.go index e3f0f7311..88dcde8b5 100644 --- a/lcnode/lc_transition.go +++ b/lcnode/lc_transition.go @@ -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) diff --git a/lcnode/lc_transition_test.go b/lcnode/lc_transition_test.go index bf253d195..3a29f593c 100644 --- a/lcnode/lc_transition_test.go +++ b/lcnode/lc_transition_test.go @@ -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 } diff --git a/objectnode/fs_volume.go b/objectnode/fs_volume.go index 1bff367dc..84f95a419 100644 --- a/objectnode/fs_volume.go +++ b/objectnode/fs_volume.go @@ -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)", diff --git a/sdk/data/blobstore/reader_test.go b/sdk/data/blobstore/reader_test.go index 6fb2875b7..81bf3e54f 100644 --- a/sdk/data/blobstore/reader_test.go +++ b/sdk/data/blobstore/reader_test.go @@ -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 } diff --git a/sdk/data/stream/extent_cache.go b/sdk/data/stream/extent_cache.go index 8047024c3..049760db2 100644 --- a/sdk/data/stream/extent_cache.go +++ b/sdk/data/stream/extent_cache.go @@ -26,10 +26,11 @@ import ( // ExtentRequest defines the struct for the request of read or write an extent. type ExtentRequest struct { - FileOffset int - Size int - Data []byte - ExtentKey *proto.ExtentKey + FileOffset int + 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) } diff --git a/sdk/data/stream/extent_client.go b/sdk/data/stream/extent_client.go index d849d8953..d4f694315 100644 --- a/sdk/data/stream/extent_client.go +++ b/sdk/data/stream/extent_client.go @@ -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)) diff --git a/sdk/data/stream/stream_reader.go b/sdk/data/stream/stream_reader.go index 97254ad27..f376995a0 100644 --- a/sdk/data/stream/stream_reader.go +++ b/sdk/data/stream/stream_reader.go @@ -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,13 +694,17 @@ 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) - log.LogDebugf("addPendingAsyncFlush: streamer(%v) handler(%v) trace(%v)", s.inode, handlerID, string(debug.Stack())) + 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) - log.LogDebugf("removePendingAsyncFlush streamer(%v) handler(%v) trace(%v)", s.inode, handlerID, string(debug.Stack())) + 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 diff --git a/sdk/data/stream/stream_writer.go b/sdk/data/stream/stream_writer.go index 791c61275..489cd8fa0 100644 --- a/sdk/data/stream/stream_writer.go +++ b/sdk/data/stream/stream_writer.go @@ -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() + } +}