fix(log_buffer): re-check buffer before bailing with ResumeFromDiskError (#9804)

ReadFromBuffer and HasData() take the read lock separately, so a write
that lands between them can make a subscriber which just read a
momentarily empty buffer return ResumeFromDiskError even though the data
is now servable from memory. Re-read under a fresh lock and only bail
when the position is genuinely behind the in-memory window (flushed to
disk); otherwise loop back and read it.
This commit is contained in:
Chris Lu
2026-06-02 21:37:15 -07:00
committed by GitHub
parent 24159fbff9
commit f711868fb6
+22 -2
View File
@@ -147,7 +147,17 @@ func (logBuffer *LogBuffer) LoopProcessLogData(readerName string, startPosition
glog.V(4).Infof("%s: Caught up to disk head, backing off to %s poll", readerName, caughtUpDiskPollInterval)
}
} else if logBuffer.HasData() {
return lastReadPosition, isDone, ResumeFromDiskError
// HasData() and ReadFromBuffer lock separately, so a racing write can make
// HasData() see data the empty-buffer read missed. Re-read; only bail if the
// position is genuinely behind the in-memory window (flushed to disk).
reBuf, _, reErr := logBuffer.ReadFromBuffer(lastReadPosition)
if reErr == ResumeFromDiskError {
return lastReadPosition, isDone, ResumeFromDiskError
}
if reBuf != nil {
logBuffer.ReleaseMemory(reBuf)
}
continue
}
// CRITICAL: Check if client is still connected
@@ -362,7 +372,17 @@ func (logBuffer *LogBuffer) LoopProcessLogDataWithOffset(readerName string, star
glog.V(4).Infof("%s: Caught up to disk head, backing off to %s poll", readerName, caughtUpDiskPollInterval)
}
} else if logBuffer.HasData() {
return lastReadPosition, isDone, ResumeFromDiskError
// HasData() and ReadFromBuffer lock separately, so a racing write can make
// HasData() see data the empty-buffer read missed. Re-read; only bail if the
// position is genuinely behind the in-memory window (flushed to disk).
reBuf, _, reErr := logBuffer.ReadFromBuffer(lastReadPosition)
if reErr == ResumeFromDiskError {
return lastReadPosition, isDone, ResumeFromDiskError
}
if reBuf != nil {
logBuffer.ReleaseMemory(reBuf)
}
continue
}
// CRITICAL: Check if client is still connected after disk read