diff --git a/weed/filer/filer_notify.go b/weed/filer/filer_notify.go index 311874875..c66235afc 100644 --- a/weed/filer/filer_notify.go +++ b/weed/filer/filer_notify.go @@ -200,7 +200,27 @@ func isChunkNotFoundError(err error) bool { httpNotFoundPattern.MatchString(errMsg) } -func (f *Filer) ReadPersistedLogBuffer(startPosition log_buffer.MessagePosition, stopTsNs int64, eachLogEntryFn log_buffer.EachLogEntryFuncType) (lastTsNs int64, isDone bool, err error) { +// persistedLogReplayLimit caps concurrent legacy replays; each holds a chunk +// reader per source filer, so a reconnect storm of pre-offload clients would +// otherwise pin many GB. Metadata-chunks clients take sendLogFileRefs and never +// reach this path. +const persistedLogReplayLimit = 16 + +var persistedLogReplaySem = make(chan struct{}, persistedLogReplayLimit) + +func (f *Filer) ReadPersistedLogBuffer(ctx context.Context, startPosition log_buffer.MessagePosition, stopTsNs int64, eachLogEntryFn log_buffer.EachLogEntryFuncType) (lastTsNs int64, isDone bool, err error) { + + // Cap concurrent replays; bail if the stream is already gone so cancelled + // clients do not park on the semaphore. + if err := ctx.Err(); err != nil { + return 0, false, err + } + select { + case persistedLogReplaySem <- struct{}{}: + defer func() { <-persistedLogReplaySem }() + case <-ctx.Done(): + return 0, false, ctx.Err() + } visitor, visitErr := f.collectPersistedLogBuffer(startPosition, stopTsNs) if visitErr != nil { diff --git a/weed/server/filer_grpc_server_sub_meta.go b/weed/server/filer_grpc_server_sub_meta.go index 6bec4e2c7..0b5a1c3f8 100644 --- a/weed/server/filer_grpc_server_sub_meta.go +++ b/weed/server/filer_grpc_server_sub_meta.go @@ -215,7 +215,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest, if req.ClientSupportsMetadataChunks { processedTsNs, isDone, readPersistedLogErr = fs.sendLogFileRefs(ctx, stream, lastReadTime, req.UntilNs) } else { - processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(lastReadTime, req.UntilNs, eachLogEntryFn) + processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, eachLogEntryFn) } if readPersistedLogErr != nil { return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr) @@ -367,7 +367,7 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq if req.ClientSupportsMetadataChunks { processedTsNs, isDone, readPersistedLogErr = fs.sendLogFileRefs(ctx, stream, lastReadTime, req.UntilNs) } else { - processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(lastReadTime, req.UntilNs, eachLogEntryFn) + processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, eachLogEntryFn) } if readPersistedLogErr != nil { glog.V(0).Infof("read on disk %v local subscribe %s from %+v: %v", clientName, req.PathPrefix, lastReadTime, readPersistedLogErr)