diff --git a/weed/pb/filer_pb_direct_read.go b/weed/pb/filer_pb_direct_read.go index 294b78c7a..2a1a78c32 100644 --- a/weed/pb/filer_pb_direct_read.go +++ b/weed/pb/filer_pb_direct_read.go @@ -144,6 +144,23 @@ func readMultiFilersMerged( streams := make([]filerStream, len(filerOrder)) var wg sync.WaitGroup + // A genuine (non chunk-not-found) read error must fail the whole replay: the + // caller advances its cursor only on success, so a swallowed error leaves a + // permanent gap. stop aborts the other readers; fatalErr is read on exit. + stop := make(chan struct{}) + var stopOnce sync.Once + closeStop := func() { stopOnce.Do(func() { close(stop) }) } + var fatalMu sync.Mutex + var fatalErr error + setFatal := func(e error) { + fatalMu.Lock() + if fatalErr == nil { + fatalErr = e + } + fatalMu.Unlock() + closeStop() + } + for i, filerId := range filerOrder { entryCh := make(chan *filer_pb.LogEntry, 512) streams[i] = filerStream{filerId: filerId, entryCh: entryCh} @@ -152,10 +169,20 @@ func readMultiFilersMerged( go func(refs []*filer_pb.LogFileChunkRef, ch chan *filer_pb.LogEntry) { defer wg.Done() defer close(ch) - readFilerFilesToChannel(refs, newReader, startTsNs, stopTsNs, ch) + readFilerFilesToChannel(refs, newReader, startTsNs, stopTsNs, ch, stop, setFatal) }(perFiler[filerId], entryCh) } + // Stop readers, drain channels so none block on a send, then wait for exit. + drainAndWait := func() { + closeStop() + for i := range streams { + for range streams[i].entryCh { + } + } + wg.Wait() + } + // Seed the min-heap with the first entry from each filer pq := &logEntryHeap{} heap.Init(pq) @@ -167,15 +194,23 @@ func readMultiFilersMerged( // Merge loop for pq.Len() > 0 { + // stop is closed only by setFatal here, so a closed stop means a reader + // aborted; lock-free bail on the hot path. + select { + case <-stop: + drainAndWait() + fatalMu.Lock() + fe := fatalErr + fatalMu.Unlock() + return lastTsNs, fe + default: + } + item := heap.Pop(pq).(*logEntryHeapItem) lastTsNs, err = processOneLogEntry(item.entry, filter, processEventFn) if err != nil { - for i := range streams { - for range streams[i].entryCh { - } - } - wg.Wait() + drainAndWait() return } @@ -184,7 +219,13 @@ func readMultiFilersMerged( } } - wg.Wait() + drainAndWait() + fatalMu.Lock() + fe := fatalErr + fatalMu.Unlock() + if fe != nil { + return lastTsNs, fe + } return } @@ -193,6 +234,8 @@ func readFilerFilesToChannel( newReader LogFileReaderFn, startTsNs, stopTsNs int64, ch chan *filer_pb.LogEntry, + stop <-chan struct{}, + setFatal func(error), ) { type prefetchResult struct { entries []*filer_pb.LogEntry @@ -214,7 +257,12 @@ func readFilerFilesToChannel( } for i, ref := range refs { - result := <-pendingCh + var result prefetchResult + select { + case result = <-pendingCh: + case <-stop: + return + } if i+1 < len(refs) { pendingCh = startPrefetch(refs[i+1]) @@ -223,14 +271,19 @@ func readFilerFilesToChannel( if result.err != nil { if isChunkNotFound(result.err) { glog.V(0).Infof("skip log file filer=%s ts=%d: %v", ref.FilerId, ref.FileTsNs, result.err) - } else { - glog.Errorf("read log file filer=%s ts=%d: %v", ref.FilerId, ref.FileTsNs, result.err) + continue } - continue + glog.Errorf("read log file filer=%s ts=%d: %v", ref.FilerId, ref.FileTsNs, result.err) + setFatal(fmt.Errorf("read log file filer=%s ts=%d: %w", ref.FilerId, ref.FileTsNs, result.err)) + return } for _, entry := range result.entries { - ch <- entry + select { + case ch <- entry: + case <-stop: + return + } } } } diff --git a/weed/pb/filer_pb_direct_read_test.go b/weed/pb/filer_pb_direct_read_test.go index 68f71a306..19a8e20d8 100644 --- a/weed/pb/filer_pb_direct_read_test.go +++ b/weed/pb/filer_pb_direct_read_test.go @@ -278,3 +278,54 @@ func TestDirectReadVsServerSideThroughput(t *testing.T) { t.Logf("Speedup: %.1fx (parallel + prefetch + no gRPC vs server-side sequential)", directRate/serverRate) } } + +// failingReaderFn returns err for the given file key, delegating otherwise. +func failingReaderFn(base LogFileReaderFn, failKey string, err error) LogFileReaderFn { + return func(chunks []*filer_pb.FileChunk) (io.ReadCloser, error) { + if len(chunks) > 0 && chunks[0].FileId == failKey { + return nil, err + } + return base(chunks) + } +} + +// A real (non not-found) read error must fail the whole replay, not silently +// drop the file and advance the cursor. +func TestReadLogFileRefsMultiFilerGenuineErrorAborts(t *testing.T) { + files := newTestLogFiles(3, 2, 10, 0) + failKey := files.refs[2].Chunks[0].FileId // filer01's first file + readerFn := failingReaderFn(files.readerFn(), failKey, fmt.Errorf("failed to locate %s", failKey)) + + var count int64 + _, err := ReadLogFileRefs(files.refs, readerFn, 0, 0, + PathFilter{PathPrefix: "/"}, + func(resp *filer_pb.SubscribeMetadataResponse) error { + atomic.AddInt64(&count, 1) + return nil + }) + if err == nil { + t.Fatalf("expected error from genuine read failure, got nil (delivered=%d)", count) + } +} + +// A chunk-not-found error skips only that file (volume gone), not the replay. +func TestReadLogFileRefsMultiFilerNotFoundSkips(t *testing.T) { + files := newTestLogFiles(3, 2, 10, 0) + skipKey := files.refs[2].Chunks[0].FileId // filer01's first file + readerFn := failingReaderFn(files.readerFn(), skipKey, fmt.Errorf("volume not found: %s", skipKey)) + + var count int64 + _, err := ReadLogFileRefs(files.refs, readerFn, 0, 0, + PathFilter{PathPrefix: "/"}, + func(resp *filer_pb.SubscribeMetadataResponse) error { + atomic.AddInt64(&count, 1) + return nil + }) + if err != nil { + t.Fatalf("chunk-not-found should be skipped, got error: %v", err) + } + expected := int64(files.totalEvents() - 10) // one skipped file's events + if count != expected { + t.Fatalf("expected %d events after skipping one file, got %d", expected, count) + } +}