diff --git a/weed/filer/meta_aggregator.go b/weed/filer/meta_aggregator.go index 6095815db..f993d68dc 100644 --- a/weed/filer/meta_aggregator.go +++ b/weed/filer/meta_aggregator.go @@ -61,6 +61,7 @@ type MetaAggregator struct { lowFlushWatermarkTsNs int64 deliveryAdvanced chan struct{} flushAdvanced chan struct{} + remotePeerArrived chan struct{} peerWatermarksLock sync.Mutex } @@ -95,6 +96,9 @@ func (ma *MetaAggregator) OnPeerUpdate(update *master_pb.ClusterNodeUpdate, star } stopChan := make(chan struct{}) ma.peerChans[address] = stopChan + if address != ma.self { + ma.noteRemotePeerArrivalLocked() + } // Account for the peer before its stream signals; keep prior values // on reconnect. ma.initPeerWatermark(address) @@ -272,7 +276,10 @@ func lowWatermarkOf(watermarks map[pb.ServerAddress]int64) int64 { func (ma *MetaAggregator) HasRemotePeers() bool { ma.peerChansLock.Lock() defer ma.peerChansLock.Unlock() + return ma.hasRemotePeersLocked() +} +func (ma *MetaAggregator) hasRemotePeersLocked() bool { for address := range ma.peerChans { if address != ma.self { return true @@ -281,6 +288,30 @@ func (ma *MetaAggregator) HasRemotePeers() bool { return false } +// RemotePeerArrivedChan returns a channel closed when the aggregator learns +// its first remote peer, or nil if one is already known. +func (ma *MetaAggregator) RemotePeerArrivedChan() <-chan struct{} { + ma.peerChansLock.Lock() + defer ma.peerChansLock.Unlock() + if ma.hasRemotePeersLocked() { + return nil + } + return ma.remotePeerArrivedChanLocked() +} + +// remotePeerArrivedChanLocked arms and returns the arrival channel. Caller must hold peerChansLock. +func (ma *MetaAggregator) remotePeerArrivedChanLocked() <-chan struct{} { + if ma.remotePeerArrived == nil { + ma.remotePeerArrived = make(chan struct{}) + } + return ma.remotePeerArrived +} + +// noteRemotePeerArrivalLocked closes the arrival channel when a remote peer is added. Caller must hold peerChansLock. +func (ma *MetaAggregator) noteRemotePeerArrivalLocked() { + ma.remotePeerArrived = closeWatermarkChan(ma.remotePeerArrived) +} + // HasPeer reports whether address is currently a tracked filer peer (or this // filer's own address). Callers use this to gate operations on known cluster // members. diff --git a/weed/filer/meta_aggregator_testing.go b/weed/filer/meta_aggregator_testing.go index b7443a28a..1c805f584 100644 --- a/weed/filer/meta_aggregator_testing.go +++ b/weed/filer/meta_aggregator_testing.go @@ -10,6 +10,9 @@ import ( func (ma *MetaAggregator) TrackPeerForTesting(peer pb.ServerAddress) { ma.peerChansLock.Lock() ma.peerChans[peer] = make(chan struct{}) + if peer != ma.self { + ma.noteRemotePeerArrivalLocked() + } ma.peerChansLock.Unlock() ma.initPeerWatermark(peer) } diff --git a/weed/server/filer_grpc_server_sub_meta.go b/weed/server/filer_grpc_server_sub_meta.go index b4c801065..858f64f4f 100644 --- a/weed/server/filer_grpc_server_sub_meta.go +++ b/weed/server/filer_grpc_server_sub_meta.go @@ -41,6 +41,13 @@ var ( metadataGapSettledHorizon = 2 * filer.LogFlushInterval ) +// errAggregationUpgrade ends a delegated local stream when remote peers +// appear, so the client reconnects to the aggregated stream. It is an +// error, not a clean end: RetryUntil-driven followers treat a clean end as +// "following finished" and stop reconnecting. It wraps StopReadingError so +// LoopProcessLogData does not log it. +var errAggregationUpgrade = fmt.Errorf("remote filer peers discovered after subscription started; reconnect for aggregated metadata: %w", log_buffer.StopReadingError) + const ( // MaxUnsyncedEvents send empty notification with timestamp when certain amount of events have been filtered MaxUnsyncedEvents = 1e3 @@ -72,6 +79,13 @@ type metadataStreamSender interface { Send(*filer_pb.SubscribeMetadataResponse) error } +// metadataLocalStream is the subset of the local-subscribe gRPC server stream +// the local loop uses, so the aggregated stream type can delegate to it. +type metadataLocalStream interface { + Send(*filer_pb.SubscribeMetadataResponse) error + Context() context.Context +} + const ( // batchBehindThreshold: when an event's timestamp is older than this // relative to wall clock, the sender switches to batch mode for throughput. @@ -403,17 +417,17 @@ func (r *gapStallReporter) close() { // corresponds to the event being waited for. The park is where a stalled // subscriber spends all its time, so every exit the read loop relies on has to // be checked here too. -func (fs *FilerServer) parkOnGap(ctx context.Context, req *filer_pb.SubscribeMetadataRequest, gapStall *gapStallReporter, evictedTsNs func() int64, cursor log_buffer.MessagePosition, notifyChan <-chan struct{}, reason string) (skipToTsNs int64, skip bool, done bool) { +func (fs *FilerServer) parkOnGap(ctx context.Context, req *filer_pb.SubscribeMetadataRequest, gapStall *gapStallReporter, evictedTsNs func() int64, cursor log_buffer.MessagePosition, notifyChan <-chan struct{}, reason string, upgradeOnRemotePeer <-chan struct{}) (skipToTsNs int64, skip bool, done bool, upgrade bool) { // Done exits run before park(): a finished stream was never parked, and // marking it so leaves a false "still behind" trace. A cursor at UntilNs is // finished - the bound is inclusive, cursors are exclusive, and // LoopProcessLogData (the only place UntilNs ends a stream) is unreachable // from a park. if req.UntilNs != 0 && cursor.Time.UnixNano() >= req.UntilNs { - return 0, false, true + return 0, false, true, false } if !fs.hasClient(req.ClientId, req.ClientEpoch) { - return 0, false, true + return 0, false, true, false } gapStall.park(cursor.Time, reason) if gapStall.stalledFor() >= maxGapStall { @@ -421,7 +435,7 @@ func (fs *FilerServer) parkOnGap(ctx context.Context, req *filer_pb.SubscribeMet // strictly after it, so the recorded loss is exactly (cursor, skipTo]. if evicted := evictedTsNs(); evicted > cursor.Time.UnixNano() { gapStall.gaveUp(cursor.Time, evicted, reason) - return evicted, true, false + return evicted, true, false, false } // Nothing was withheld past the cursor - nothing to skip, nothing being // lost; keep waiting on a fresh stall cycle. @@ -445,13 +459,15 @@ func (fs *FilerServer) parkOnGap(ctx context.Context, req *filer_pb.SubscribeMet continue } case <-ctx.Done(): - return 0, false, true + return 0, false, true, false + case <-upgradeOnRemotePeer: + return 0, false, false, true case <-retry: } if !fs.hasClient(req.ClientId, req.ClientEpoch) { - return 0, false, true + return 0, false, true, false } - return 0, false, false + return 0, false, false, false } } @@ -497,15 +513,16 @@ func resolveGapResume(currentTsNs, currentOffset, earliestMemTsNs, flushedTsNs, // the two subscribe loops; everything else about them must stay identical, and // this PR's history shows they drift when edited separately. type gapPass struct { - fs *FilerServer - req *filer_pb.SubscribeMetadataRequest - gapStall *gapStallReporter - earliest func() time.Time - evicted func() int64 // gap-proof watermark; aggregated uses the received-ts space - flushed func() int64 // what the last disk read proved covered: flushed AND inside its listing - gapChan <-chan struct{} - dataChan <-chan struct{} - gapReason func(earliest time.Time, evictedTsNs int64) string + fs *FilerServer + req *filer_pb.SubscribeMetadataRequest + gapStall *gapStallReporter + earliest func() time.Time + evicted func() int64 // gap-proof watermark; aggregated uses the received-ts space + flushed func() int64 // what the last disk read proved covered: flushed AND inside its listing + gapChan <-chan struct{} + dataChan <-chan struct{} + gapReason func(earliest time.Time, evictedTsNs int64) string + upgradeOnRemotePeer <-chan struct{} } type gapOutcome int @@ -514,6 +531,7 @@ const ( gapProceed gapOutcome = iota // read memory gapContinue // restart the pass gapDone // the stream is over + gapUpgrade // remote peer arrived; end for reconnect ) // resolve is the gap decision both loops run between the disk pass and the @@ -565,10 +583,13 @@ func (p *gapPass) resolve(ctx context.Context, cursor *log_buffer.MessagePositio } func (p *gapPass) park(ctx context.Context, cursor *log_buffer.MessagePosition, latch *error, notifyChan <-chan struct{}, reason string) gapOutcome { - skipTo, skip, done := p.fs.parkOnGap(ctx, p.req, p.gapStall, p.evicted, *cursor, notifyChan, reason) + skipTo, skip, done, upgrade := p.fs.parkOnGap(ctx, p.req, p.gapStall, p.evicted, *cursor, notifyChan, reason, p.upgradeOnRemotePeer) if done { return gapDone } + if upgrade { + return gapUpgrade + } if skip { *cursor = log_buffer.NewMessagePosition(skipTo, gapResumeCursorOffset) *latch = nil @@ -577,8 +598,17 @@ func (p *gapPass) park(ctx context.Context, cursor *log_buffer.MessagePosition, } func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest, stream filer_pb.SeaweedFiler_SubscribeMetadataServer) error { - if fs.filer.MetaAggregator == nil || !fs.filer.MetaAggregator.HasRemotePeers() { - return fs.SubscribeLocalMetadata(req, stream) + // A filer that has not learned remote peers yet serves the local log and + // upgrades when the first one appears. RemotePeerArrivedChan takes the + // arrival channel under the same lock as the peer check, so a peer + // learned in between returns nil and the stream goes straight to the + // aggregated path. + if fs.filer.MetaAggregator != nil { + if arrival := fs.filer.MetaAggregator.RemotePeerArrivedChan(); arrival != nil { + return fs.subscribeLocalMetadata(req, stream, arrival) + } + } else { + return fs.subscribeLocalMetadata(req, stream, nil) } ctx := stream.Context() @@ -744,7 +774,7 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest, diskPassProvenTsNs = refsStopTsNs } if refsStopTsNs > lastReadTime.Time.UnixNano() { - processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs) + processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, refsStopTsNs, sentRefs, nil) } else { processedTsNs, isDone, readPersistedLogErr = 0, false, nil } @@ -914,6 +944,14 @@ func (fs *FilerServer) SubscribeMetadata(req *filer_pb.SubscribeMetadataRequest, } func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataRequest, stream filer_pb.SeaweedFiler_SubscribeLocalMetadataServer) error { + return fs.subscribeLocalMetadata(req, stream, nil) +} + +// subscribeLocalMetadata serves the filer's own log to the stream. Peer +// aggregation streams pass upgradeOnRemotePeer == nil; the SubscribeMetadata +// delegation passes the aggregator's arrival channel so the stream ends +// when a remote peer appears and the client reconnects to the aggregated path. +func (fs *FilerServer) subscribeLocalMetadata(req *filer_pb.SubscribeMetadataRequest, stream metadataLocalStream, upgradeOnRemotePeer <-chan struct{}) error { ctx := stream.Context() peerAddress := findClientAddress(ctx, 0) @@ -962,6 +1000,13 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq var lastFlushReportNs int64 baseEachLogEntryFn := eachLogEntryFn(req, sender, eachEventNotificationFn, &unsyncedEvents) eachLogEntryFn := func(logEntry *filer_pb.LogEntry) (bool, error) { + if upgradeOnRemotePeer != nil { + select { + case <-upgradeOnRemotePeer: + return false, errAggregationUpgrade + default: + } + } lastSeenTsNs = logEntry.TsNs return baseEachLogEntryFn(logEntry) } @@ -974,16 +1019,19 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq var lastDiskReadTsNs int64 = -1 // Track the last read position we used for disk read sentRefs := make(map[string]sentRefState) + var upgradedToAggregation bool + localBuffer := fs.filer.LocalMetaLogBuffer gaps := &gapPass{ - fs: fs, - req: req, - gapStall: gapStall, - earliest: localBuffer.GetEarliestTime, - evicted: localBuffer.GetLastEvictedTsNs, // local disk carries the ring's own timestamps - flushed: func() int64 { return lastCheckedFlushTsNs }, - gapChan: localFlushChan, - dataChan: localFlushChan, + fs: fs, + req: req, + gapStall: gapStall, + earliest: localBuffer.GetEarliestTime, + evicted: localBuffer.GetLastEvictedTsNs, // local disk carries the ring's own timestamps + flushed: func() int64 { return lastCheckedFlushTsNs }, + gapChan: localFlushChan, + dataChan: localFlushChan, + upgradeOnRemotePeer: upgradeOnRemotePeer, gapReason: func(earliest time.Time, evictedTsNs int64) string { return fmt.Sprintf("gap is not flushed yet (earliest memory %v, flushed through %v)", earliest, time.Unix(0, lastCheckedFlushTsNs)) @@ -991,6 +1039,15 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq } for { + if upgradeOnRemotePeer != nil { + select { + case <-upgradeOnRemotePeer: + glog.V(0).Infof("remote peer discovered after local subscribe %s started: ending stream so the client reconnects to the aggregated stream", clientName) + return errAggregationUpgrade + default: + } + } + // Check if new data has been flushed to disk since last check, or if read position advanced currentFlushTsNs := fs.filer.LocalMetaLogBuffer.GetLastFlushTsNs() currentReadTsNs := lastReadTime.Time.UnixNano() @@ -1005,11 +1062,14 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq lastDiskReadTsNs = currentReadTsNs glog.V(4).Infof("read on disk %v local subscribe %s from %+v (lastFlushed: %v)", clientName, req.PathPrefix, lastReadTime, time.Unix(0, currentFlushTsNs)) if req.ClientSupportsMetadataChunks { - processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, req.UntilNs, sentRefs) + processedTsNs, isDone, readPersistedLogErr = fs.chunkDiskPass(ctx, sender, lastReadTime, req.UntilNs, sentRefs, upgradeOnRemotePeer) } else { processedTsNs, isDone, readPersistedLogErr = fs.filer.ReadPersistedLogBuffer(ctx, lastReadTime, req.UntilNs, eachLogEntryFn) } if readPersistedLogErr != nil { + if errors.Is(readPersistedLogErr, errAggregationUpgrade) { + return errAggregationUpgrade + } glog.V(0).Infof("read on disk %v local subscribe %s from %+v: %v", clientName, req.PathPrefix, lastReadTime, readPersistedLogErr) return fmt.Errorf("reading from persisted logs: %w", readPersistedLogErr) } @@ -1043,6 +1103,8 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq switch gaps.resolve(ctx, &lastReadTime, &readInMemoryLogErr, diskAdvanced) { case gapDone: return nil + case gapUpgrade: + return errAggregationUpgrade case gapContinue: continue } @@ -1055,6 +1117,14 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq return false default: } + if upgradeOnRemotePeer != nil { + select { + case <-upgradeOnRemotePeer: + upgradedToAggregation = true + return false + default: + } + } if !fs.hasClient(req.ClientId, req.ClientEpoch) { return false } @@ -1062,7 +1132,14 @@ func (fs *FilerServer) SubscribeLocalMetadata(req *filer_pb.SubscribeMetadataReq lastFlushReportNs = fs.maybeSendFlushReport(req, sender, lastFlushReportNs) return true }, eachLogEntryFn) + if upgradedToAggregation { + glog.V(0).Infof("remote peer discovered after local subscribe %s started: ending stream so the client reconnects to the aggregated stream", clientName) + return errAggregationUpgrade + } if readInMemoryLogErr != nil { + if errors.Is(readInMemoryLogErr, errAggregationUpgrade) { + return errAggregationUpgrade + } if errors.Is(readInMemoryLogErr, log_buffer.ResumeFromDiskError) { // Fell behind the ring: back to the disk pass (it re-runs when // the flush or the cursor moved), and from there to the gap @@ -1223,7 +1300,7 @@ func (fs *FilerServer) maybeSendIdleHeartbeat(req *filer_pb.SubscribeMetadataReq // empty-notification marker: both chunk consumers buffer refs until a non-ref // message, so an idle source would otherwise strand the backlog in the // client's pending list until the next mutation. -func (fs *FilerServer) chunkDiskPass(ctx context.Context, sender metadataStreamSender, startPos log_buffer.MessagePosition, untilNs int64, sent map[string]sentRefState) (processedTsNs int64, isDone bool, err error) { +func (fs *FilerServer) chunkDiskPass(ctx context.Context, sender metadataStreamSender, startPos log_buffer.MessagePosition, untilNs int64, sent map[string]sentRefState, upgradeOnRemotePeer <-chan struct{}) (processedTsNs int64, isDone bool, err error) { collected, _, err := fs.filer.CollectLogFileRefs(ctx, startPos, untilNs) if err != nil { return 0, false, err @@ -1232,9 +1309,16 @@ func (fs *FilerServer) chunkDiskPass(ctx context.Context, sender metadataStreamS if len(refs) == 0 { return startPos.Time.UnixNano(), false, nil } - if err := fs.sendRefsBatched(sender, refs); err != nil { + if err := fs.sendRefsBatched(sender, refs, upgradeOnRemotePeer); err != nil { return 0, false, err } + if upgradeOnRemotePeer != nil { + select { + case <-upgradeOnRemotePeer: + return 0, false, errAggregationUpgrade + default: + } + } // Shipped content end, read from the shipped chunks alone - a fresh // listing here could see a concurrent append and move the cursor past @@ -1285,9 +1369,16 @@ func (fs *FilerServer) chunkDiskPass(ctx context.Context, sender metadataStreamS // sendRefsBatched sends refs through the pipelined sender, which keeps them // out of Events batches; gRPC allows one sending goroutine per stream and the // sender's goroutine is it. -func (fs *FilerServer) sendRefsBatched(sender metadataStreamSender, refs []*filer_pb.LogFileChunkRef) error { +func (fs *FilerServer) sendRefsBatched(sender metadataStreamSender, refs []*filer_pb.LogFileChunkRef, upgradeOnRemotePeer <-chan struct{}) error { const maxRefsPerMessage = 64 for i := 0; i < len(refs); i += maxRefsPerMessage { + if upgradeOnRemotePeer != nil { + select { + case <-upgradeOnRemotePeer: + return errAggregationUpgrade + default: + } + } end := i + maxRefsPerMessage if end > len(refs) { end = len(refs) diff --git a/weed/server/filer_grpc_server_sub_meta_gap_test.go b/weed/server/filer_grpc_server_sub_meta_gap_test.go index 6505b0828..962871974 100644 --- a/weed/server/filer_grpc_server_sub_meta_gap_test.go +++ b/weed/server/filer_grpc_server_sub_meta_gap_test.go @@ -285,7 +285,7 @@ func TestParkOnGapExits(t *testing.T) { gapStall := newStall() defer gapStall.resumed() start := time.Now() - _, skip, done := fs.parkOnGap(context.Background(), req, gapStall, noEviction, cursor, nil, "test") + _, skip, done, _ := fs.parkOnGap(context.Background(), req, gapStall, noEviction, cursor, nil, "test", nil) if skip || done { t.Fatalf("skip=%v done=%v, want a plain retry", skip, done) } @@ -299,7 +299,7 @@ func TestParkOnGapExits(t *testing.T) { // old one keeps scanning the filer store until its TCP connection dies. gapStall := newStall() superseded := &filer_pb.SubscribeMetadataRequest{ClientId: 7, ClientEpoch: 2} - if _, _, done := fs.parkOnGap(context.Background(), superseded, gapStall, noEviction, cursor, nil, "test"); !done { + if _, _, done, _ := fs.parkOnGap(context.Background(), superseded, gapStall, noEviction, cursor, nil, "test", nil); !done { t.Fatal("want done") } if !gapStall.since.IsZero() { @@ -310,7 +310,7 @@ func TestParkOnGapExits(t *testing.T) { t.Run("a bounded subscription past its window ends without parking", func(t *testing.T) { gapStall := newStall() bounded := &filer_pb.SubscribeMetadataRequest{ClientId: 7, ClientEpoch: 3, UntilNs: cursorTs - 1} - if _, _, done := fs.parkOnGap(context.Background(), bounded, gapStall, noEviction, cursor, nil, "test"); !done { + if _, _, done, _ := fs.parkOnGap(context.Background(), bounded, gapStall, noEviction, cursor, nil, "test", nil); !done { t.Fatal("want done") } if !gapStall.since.IsZero() { @@ -324,7 +324,7 @@ func TestParkOnGapExits(t *testing.T) { // everything <= UntilNs delivered. Parking would make fs.verify hang. gapStall := newStall() bounded := &filer_pb.SubscribeMetadataRequest{ClientId: 7, ClientEpoch: 3, UntilNs: cursorTs} - if _, _, done := fs.parkOnGap(context.Background(), bounded, gapStall, noEviction, cursor, nil, "test"); !done { + if _, _, done, _ := fs.parkOnGap(context.Background(), bounded, gapStall, noEviction, cursor, nil, "test", nil); !done { t.Fatal("want done") } if !gapStall.since.IsZero() { @@ -338,7 +338,7 @@ func TestParkOnGapExits(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) cancel() start := time.Now() - if _, _, done := fs.parkOnGap(ctx, req, gapStall, noEviction, cursor, nil, "test"); !done { + if _, _, done, _ := fs.parkOnGap(ctx, req, gapStall, noEviction, cursor, nil, "test", nil); !done { t.Fatal("want done") } if elapsed := time.Since(start); elapsed >= unflushedGapRetryInterval { @@ -354,7 +354,7 @@ func TestParkOnGapExits(t *testing.T) { closed := make(chan struct{}) close(closed) start := time.Now() - _, skip, done := fs.parkOnGap(context.Background(), req, gapStall, noEviction, cursor, closed, "test") + _, skip, done, _ := fs.parkOnGap(context.Background(), req, gapStall, noEviction, cursor, closed, "test", nil) if skip || done { t.Fatalf("skip=%v done=%v, want a plain retry", skip, done) } @@ -369,7 +369,7 @@ func TestParkOnGapExits(t *testing.T) { notify := make(chan struct{}, 1) notify <- struct{}{} start := time.Now() - _, skip, done := fs.parkOnGap(context.Background(), req, gapStall, noEviction, cursor, notify, "test") + _, skip, done, _ := fs.parkOnGap(context.Background(), req, gapStall, noEviction, cursor, notify, "test", nil) if skip || done { t.Fatalf("skip=%v done=%v, want a plain retry", skip, done) } @@ -507,7 +507,7 @@ func TestParkOnGapStallOutcomes(t *testing.T) { gapStall.since = time.Now().Add(-maxGapStall) before := crossings() - skipTo, skip, done := fs.parkOnGap(context.Background(), req, gapStall, lb.GetLastEvictedTsNs, cursor, nil, "test") + skipTo, skip, done, _ := fs.parkOnGap(context.Background(), req, gapStall, lb.GetLastEvictedTsNs, cursor, nil, "test", nil) if done || !skip { t.Fatalf("skip=%v done=%v, want a forced skip", skip, done) } @@ -530,7 +530,7 @@ func TestParkOnGapStallOutcomes(t *testing.T) { defer gapStall.close() // release the gauge this test's park holds before := crossings() - _, skip, done := fs.parkOnGap(context.Background(), req, gapStall, lb.GetLastEvictedTsNs, cursor, nil, "test") + _, skip, done, _ := fs.parkOnGap(context.Background(), req, gapStall, lb.GetLastEvictedTsNs, cursor, nil, "test", nil) if skip || done { t.Fatalf("skip=%v done=%v, want neither: nothing is being lost", skip, done) } diff --git a/weed/server/filer_subscribe_upgrade_test.go b/weed/server/filer_subscribe_upgrade_test.go new file mode 100644 index 000000000..1d39b693a --- /dev/null +++ b/weed/server/filer_subscribe_upgrade_test.go @@ -0,0 +1,116 @@ +package weed_server + +// Regression tests for the aggregated-subscribe entrypoint's peer-discovery +// window. SubscribeMetadata delegates to the local loop when the aggregator +// knows no remote peers yet; these tests pin the upgrade: the local stream +// ends when a remote peer appears, so the client reconnects into the +// aggregated stream. + +import ( + "testing" + "time" + + "github.com/seaweedfs/seaweedfs/weed/filer" +) + +// startAggregatorWithoutPeers gives the harness an aggregator that tracks +// only self, the state of a filer whose gRPC server is up but whose master +// connection has not announced the other filers yet. +func (h *subscribeHarness) startAggregatorWithoutPeers() *filer.MetaAggregator { + ma := filer.NewMetaAggregator(h.f, testSelfAddress, nil) + ma.TrackPeerForTesting(testSelfAddress) + h.f.MetaAggregator = ma + h.t.Cleanup(ma.MetaLogBuffer.ShutdownLogBuffer) + return ma +} + +// TestSubscribeMetadataLocalStreamEndsWhenRemotePeerAppears pins the upgrade +// path: a stream that started local-only must end when the first remote peer +// appears, so the client reconnects into the aggregated stream. +func TestSubscribeMetadataLocalStreamEndsWhenRemotePeerAppears(t *testing.T) { + h := newSubscribeHarness(t) + + ma := h.startAggregatorWithoutPeers() + if ma.HasRemotePeers() { + t.Fatal("test setup: aggregator must start without remote peers") + } + + localTs := time.Now().UnixNano() + h.append(localTs) + + // The subscriber connects inside the discovery window. + r := h.subscribeAggregated(0) + waitForEvents(t, r, []int64{localTs}, 3*time.Second) + + // The master announces a second filer, as OnPeerUpdate does in production. + ma.TrackPeerForTesting(testPeerAddress) + + // The local stream must end so the client reconnects to the aggregated path. + select { + case err := <-r.done: + if err == nil { + t.Fatal("local stream ended cleanly; the upgrade must surface as an error so RetryUntil-driven followers reconnect") + } + case <-time.After(5 * time.Second): + t.Fatal("local-only stream stayed alive after a remote peer appeared; the subscriber is pinned to a partial cluster view") + } +} + +// TestSubscribeMetadataAggregatedPathAfterPeerArrival pins the other half: +// after the upgrade, a re-subscribe from the last delivered offset takes the +// aggregated path and receives the peer's events. +func TestSubscribeMetadataAggregatedPathAfterPeerArrival(t *testing.T) { + h := newSubscribeHarness(t) + + ma := h.startAggregatorWithoutPeers() + if ma.HasRemotePeers() { + t.Fatal("test setup: aggregator must start without remote peers") + } + + localTs := time.Now().UnixNano() + h.append(localTs) + r := h.subscribeAggregated(0) + waitForEvents(t, r, []int64{localTs}, 3*time.Second) + + ma.TrackPeerForTesting(testPeerAddress) + select { + case <-r.done: + case <-time.After(5 * time.Second): + t.Fatal("local-only stream stayed alive after a remote peer appeared") + } + + // The client reconnects from the last delivered offset; the aggregated + // path now serves the peer's events. + peerTs := time.Now().UnixNano() + h.appendAggregated(peerTs) + reportPeers(ma, peerTs) + r2 := h.subscribeAggregated(localTs) + waitForEvents(t, r2, []int64{peerTs}, 3*time.Second) +} + +// TestSubscribeMetadataLocalStreamPersistsWithoutPeers pins the boundary: on +// a standalone filer the local stream keeps serving indefinitely. +func TestSubscribeMetadataLocalStreamPersistsWithoutPeers(t *testing.T) { + h := newSubscribeHarness(t) + + ma := h.startAggregatorWithoutPeers() + if ma.HasRemotePeers() { + t.Fatal("test setup: aggregator must start without remote peers") + } + + localTs := time.Now().UnixNano() + h.append(localTs) + r := h.subscribeAggregated(0) + waitForEvents(t, r, []int64{localTs}, 3*time.Second) + + // No peer ever appears; the stream must still be alive and still deliver. + moreTs := time.Now().UnixNano() + h.append(moreTs) + waitForEventsAtLeastOnce(t, r, []int64{localTs, moreTs}, 3*time.Second) + + select { + case err := <-r.done: + t.Fatalf("local stream ended (%v) on a standalone filer; it must keep serving", err) + case <-time.After(200 * time.Millisecond): + } +}