This commit is contained in:
leonrayang 2026-03-24 11:17:46 +08:00
parent 22c3cfcc96
commit cd74e4867c
3 changed files with 4 additions and 116 deletions

View File

@ -725,7 +725,7 @@ func (f *File) Flush(ctx context.Context, req *fuse.FlushRequest) (err error) {
if err != nil {
msg := fmt.Sprintf("Flush: ino(%v) err(%v)", f.info.Inode, err)
if isStreamerReleasingErr(err) {
log.LogWarnf("TRACE Flush benign err during release: ino(%v) err(%v)", f.info.Inode, err)
log.LogDebugf("TRACE Flush benign err during release: ino(%v) err(%v)", f.info.Inode, err)
} else {
f.super.handleError("Flush", msg)
log.LogErrorf("TRACE Flush err: ino(%v) err(%v)", f.info.Inode, err)

View File

@ -162,15 +162,6 @@ func (s *Streamer) String() string {
s.inode, s.fullPath, atomic.LoadInt32(&s.refcnt), s.isOpen, s.openForWrite, len(s.request), s.handler, s.waitForFlush, s)
}
func (s *Streamer) pendingAsyncFlushCount() int {
count := 0
s.pendingAsyncFlushMap.Range(func(_, _ interface{}) bool {
count++
return true
})
return count
}
// TODO should we call it RefreshExtents instead?
func (s *Streamer) GetExtents(isMigration bool) error {
if s.client.disableMetaCache || !s.needBCache {
@ -247,7 +238,7 @@ func (s *Streamer) prepareReadRequestsChecked(data []byte, offset, size, maxRetr
if unresolved {
// Give async write/append pipeline a short window to publish resolved extent keys.
time.Sleep(2 * time.Millisecond)
log.LogWarnf("streamer.read unresolved extentkey retry: ino(%v) offset(%v) size(%v) retry(%v/%v) reqs(%v)",
log.LogDebugf("streamer.read unresolved extentkey retry: ino(%v) offset(%v) size(%v) retry(%v/%v) reqs(%v)",
s.inode, offset, size, retry+1, maxRetry, current)
}
}
@ -282,21 +273,7 @@ func (s *Streamer) recoverReadHoleByFlush(requests []*ExtentRequest, data []byte
return
}
extentsSnapshot := func() interface{} {
if s.extents == nil {
return "<nil>"
}
return s.extents.List()
}
// TEMP FLUSH_TRACE: investigate write-read visibility window in LTP iogen01.
for retry := 0; holeInFileRange && retry < holeRecoverMaxRetry; retry++ {
log.LogWarnf("FLUSH_TRACE read_recover_pre: ino(%v) offset(%v) size(%v) retry(%v/%v) filesize(%v) dirty(%v) waitForFlush(%v) reqs(%v) extents(%v)",
s.inode, offset, size, retry+1, holeRecoverMaxRetry, updatedFileSize, s.dirty, s.waitForFlush, updated, extentsSnapshot())
log.LogWarnf("FLUSH_TRACE_PIPE read_recover_pre: ino(%v) offset(%v) size(%v) retry(%v/%v) pipe(%s)",
s.inode, offset, size, retry+1, holeRecoverMaxRetry, s.flushTracePipeSnapshot())
log.LogWarnf("streamer.read recover suspicious hole by flush: ino(%v) offset(%v) size(%v) retry(%v/%v) filesize(%v) reqs(%v)",
s.inode, offset, size, retry+1, holeRecoverMaxRetry, updatedFileSize, updated)
s.writeLock.Lock()
if err = s.IssueFlushRequest(); err != nil {
s.writeLock.Unlock()
@ -307,10 +284,6 @@ func (s *Streamer) recoverReadHoleByFlush(requests []*ExtentRequest, data []byte
updatedFileSize, _ = s.extents.Size()
holeBytes, holeInFileRange = calcHoleStats(updated, updatedFileSize)
log.LogWarnf("FLUSH_TRACE read_recover_post: ino(%v) offset(%v) size(%v) retry(%v/%v) filesize(%v) dirty(%v) waitForFlush(%v) holeInFileRange(%v) holeBytes(%v) reqs(%v) extents(%v)",
s.inode, offset, size, retry+1, holeRecoverMaxRetry, updatedFileSize, s.dirty, s.waitForFlush, holeInFileRange, holeBytes, updated, extentsSnapshot())
log.LogWarnf("FLUSH_TRACE_PIPE read_recover_post: ino(%v) offset(%v) size(%v) retry(%v/%v) pipe(%s)",
s.inode, offset, size, retry+1, holeRecoverMaxRetry, s.flushTracePipeSnapshot())
if holeInFileRange && retry+1 < holeRecoverMaxRetry {
// Give in-flight append/extent publication a short window before next prepare.
@ -319,8 +292,8 @@ func (s *Streamer) recoverReadHoleByFlush(requests []*ExtentRequest, data []byte
}
if holeInFileRange {
log.LogWarnf("streamer.read suspicious hole: ino(%v) offset(%v) size(%v) filesize(%v) holeBytes(%v) dirty(%v) waitForFlush(%v) reqs(%v)",
s.inode, offset, size, filesize, holeBytes, s.dirty, s.waitForFlush, requests)
log.LogWarnf("streamer.read hole remains after flush retries: ino(%v) offset(%v) size(%v) filesize(%v) holeBytes(%v) reqs(%v)",
s.inode, offset, size, filesize, holeBytes, requests)
}
return
}

View File

@ -50,48 +50,12 @@ const (
var errUnresolvedExtentKey = stderrs.New("unresolved extent key (partition id 0)")
func extentHandlerPipeSnapshot(eh *ExtentHandler) string {
if eh == nil {
return "<nil>"
}
return fmt.Sprintf("{id:%d off:%d size:%d status:%d inflight:%d pendingWrites:%d reqLen:%d replyLen:%d writeDataLen:%d key:%v lastKey:%v}",
eh.id,
eh.fileOffset,
eh.size,
eh.getStatus(),
atomic.LoadInt32(&eh.inflight),
atomic.LoadInt64(&eh.pendingWrites),
len(eh.request),
len(eh.reply),
len(eh.writeDataChan),
eh.key,
eh.lastKey)
}
func (s *Streamer) flushTracePipeSnapshot() string {
return fmt.Sprintf("{reqLen:%d dirty:%v dirtyListLen:%d asyncFlushChLen:%d pendingAsync:%d handler:%s}",
len(s.request),
s.dirty,
s.dirtylist.Len(),
len(s.asyncFlushCh),
s.pendingAsyncFlushCount(),
extentHandlerPipeSnapshot(s.handler))
}
// reconcileDirtyState keeps streamer dirty flag consistent with dirtylist.
// Some async close/flush paths may temporarily drift dirty/list state.
func (s *Streamer) reconcileDirtyState() {
s.dirty = s.dirtylist.Len() > 0
}
// logDirtyInvariant only logs potential state drift and does not change behavior.
func (s *Streamer) logDirtyInvariant(where string, wait bool, id string) {
if s.dirty && s.dirtylist.Len() == 0 {
log.LogWarnf("FLUSH_INVARIANT dirty_without_handlers: where(%v) ino(%v) wait(%v) id(%v) pipe(%s)",
where, s.inode, wait, id, s.flushTracePipeSnapshot())
}
}
const (
streamWriterFlushPeriod = 3
streamWriterIdleTimeoutPeriod = 10
@ -563,7 +527,6 @@ begin:
}
}
log.LogDebugf("Streamer write exit: ino(%v) filesize(%v) offset(%v) size(%v) done total(%v) err(%v)", s.inode, filesize, offset, size, total, err)
s.logDirtyInvariant("write_exit", false, "")
return
}
@ -1116,16 +1079,6 @@ func (s *Streamer) doWriteAppendEx(data []byte, offset, size int, direct bool, r
}
func (s *Streamer) flushAsync(wait bool, id string) (err error) {
extentsSnapshot := func() interface{} {
if s.extents == nil {
return "<nil>"
}
return s.extents.List()
}
log.LogWarnf("FLUSH_TRACE flush_async_enter: ino(%v) wait(%v) id(%v) dirty(%v) dirtylistLen(%v) handler(%v) extents(%v)",
s.inode, wait, id, s.dirty, s.dirtylist.Len(), s.handler, extentsSnapshot())
log.LogWarnf("FLUSH_TRACE_PIPE flush_async_enter: ino(%v) wait(%v) id(%v) pipe(%s)",
s.inode, wait, id, s.flushTracePipeSnapshot())
pending := make(map[*ExtentHandler]*AsyncFlushRequest)
asyncExtentHandler := make([]*ExtentHandler, 0)
for {
@ -1145,10 +1098,6 @@ func (s *Streamer) flushAsync(wait bool, id string) (err error) {
}
}
log.LogDebugf("Streamer(%v) flush begin: eh(%v) id(%v)", s.inode, eh, id)
log.LogWarnf("FLUSH_TRACE flush_async_handler: ino(%v) wait(%v) id(%v) ehID(%v) fileOffset(%v) size(%v) status(%v) inflight(%v) key(%v) lastKey(%v)",
s.inode, wait, id, eh.id, eh.fileOffset, eh.size, eh.getStatus(), atomic.LoadInt32(&eh.inflight), eh.key, eh.lastKey)
log.LogWarnf("FLUSH_TRACE_PIPE flush_async_handler: ino(%v) wait(%v) id(%v) ehPipe(%s) streamPipe(%s)",
s.inode, wait, id, extentHandlerPipeSnapshot(eh), s.flushTracePipeSnapshot())
// Use async flush for better performance if enabled
if s.client.enableAsyncFlush {
@ -1183,8 +1132,6 @@ func (s *Streamer) flushAsync(wait bool, id string) (err error) {
}
if !s.client.enableAsyncFlush {
s.reconcileDirtyState()
log.LogWarnf("FLUSH_TRACE flush_async_exit_syncpath: ino(%v) wait(%v) id(%v) err(%v) dirty(%v) dirtylistLen(%v) handler(%v) extents(%v)",
s.inode, wait, id, err, s.dirty, s.dirtylist.Len(), s.handler, extentsSnapshot())
return
}
log.LogDebugf("Streamer(%v) wait(%v) pending(%v) id(%v)", s.inode, wait, pending, id)
@ -1225,10 +1172,6 @@ func (s *Streamer) flushAsync(wait bool, id string) (err error) {
}
}
s.reconcileDirtyState()
log.LogWarnf("FLUSH_TRACE flush_async_exit: ino(%v) wait(%v) id(%v) err(%v) dirty(%v) dirtylistLen(%v) pendingLen(%v) handler(%v) extents(%v)",
s.inode, wait, id, err, s.dirty, s.dirtylist.Len(), len(pending), s.handler, extentsSnapshot())
log.LogWarnf("FLUSH_TRACE_PIPE flush_async_exit: ino(%v) wait(%v) id(%v) pipe(%s)",
s.inode, wait, id, s.flushTracePipeSnapshot())
return
}
@ -1322,7 +1265,6 @@ func (s *Streamer) closeOpenHandler(wait bool) (err error) {
log.LogDebugf("closeOpenHandler: streamer(%v) wait for wait (%v) id(%v)", s.inode, wait, id)
err = s.flush(wait, id)
log.LogDebugf("closeOpenHandler: streamer(%v) wait (%v) id(%v) end err(%v)", s.inode, wait, id, err)
s.logDirtyInvariant("close_open_handler_exit", wait, id)
if err != nil {
log.LogErrorf("closeOpenHandler: flush extent failed, err %s", err.Error())
return err
@ -1487,16 +1429,6 @@ func (s *Streamer) setError() {
}
func (s *Streamer) flushSync() (err error) {
extentsSnapshot := func() interface{} {
if s.extents == nil {
return "<nil>"
}
return s.extents.List()
}
log.LogWarnf("FLUSH_TRACE flush_sync_enter: ino(%v) dirty(%v) dirtylistLen(%v) handler(%v) extents(%v)",
s.inode, s.dirty, s.dirtylist.Len(), s.handler, extentsSnapshot())
log.LogWarnf("FLUSH_TRACE_PIPE flush_sync_enter: ino(%v) pipe(%s)",
s.inode, s.flushTracePipeSnapshot())
for {
element := s.dirtylist.Get()
if element == nil {
@ -1505,10 +1437,6 @@ func (s *Streamer) flushSync() (err error) {
eh := element.Value.(*ExtentHandler)
log.LogDebugf("Streamer flush begin: eh(%v)", eh)
log.LogWarnf("FLUSH_TRACE flush_sync_handler: ino(%v) ehID(%v) fileOffset(%v) size(%v) status(%v) inflight(%v) key(%v) lastKey(%v)",
s.inode, eh.id, eh.fileOffset, eh.size, eh.getStatus(), atomic.LoadInt32(&eh.inflight), eh.key, eh.lastKey)
log.LogWarnf("FLUSH_TRACE_PIPE flush_sync_handler: ino(%v) ehPipe(%s) streamPipe(%s)",
s.inode, extentHandlerPipeSnapshot(eh), s.flushTracePipeSnapshot())
err = eh.flush()
if err != nil {
log.LogErrorf("Streamer flush failed: eh(%v)", eh)
@ -1529,23 +1457,10 @@ func (s *Streamer) flushSync() (err error) {
log.LogDebugf("Streamer flush end: eh(%v)", eh)
}
s.reconcileDirtyState()
log.LogWarnf("FLUSH_TRACE flush_sync_exit: ino(%v) err(%v) dirty(%v) dirtylistLen(%v) handler(%v) extents(%v)",
s.inode, err, s.dirty, s.dirtylist.Len(), s.handler, extentsSnapshot())
log.LogWarnf("FLUSH_TRACE_PIPE flush_sync_exit: ino(%v) pipe(%s)",
s.inode, s.flushTracePipeSnapshot())
return
}
func (s *Streamer) flush(wait bool, id string) (err error) {
defer func() {
if wait {
s.logDirtyInvariant("flush_exit", wait, id)
}
}()
log.LogWarnf("FLUSH_TRACE flush_dispatch: ino(%v) wait(%v) id(%v) enableAsyncFlush(%v) waitForFlush(%v) dirty(%v) dirtylistLen(%v) handler(%v)",
s.inode, wait, id, s.client.enableAsyncFlush, s.waitForFlush, s.dirty, s.dirtylist.Len(), s.handler)
log.LogWarnf("FLUSH_TRACE_PIPE flush_dispatch: ino(%v) wait(%v) id(%v) pipe(%s)",
s.inode, wait, id, s.flushTracePipeSnapshot())
if s.client.enableAsyncFlush && !s.waitForFlush {
return s.flushAsync(wait, id)
} else {