diff --git a/weed/command/filer_remote_gateway_buckets.go b/weed/command/filer_remote_gateway_buckets.go index 2bacf14b8..9dffd4cc2 100644 --- a/weed/command/filer_remote_gateway_buckets.go +++ b/weed/command/filer_remote_gateway_buckets.go @@ -64,6 +64,9 @@ func (option *RemoteGatewayOptions) followBucketUpdatesAndUploadToRemote(filerSo StartTsNs: lastOffsetTs.UnixNano(), StopTsNs: 0, EventErrorType: pb.RetryForeverOnError, + GetResumeTsNs: func() int64 { + return processor.processedTsWatermark.Load() + }, } return pb.FollowMetadata(pb.ServerAddress(*option.filerAddress), option.grpcDialOption, metadataFollowOption, processEventFnWithOffset) diff --git a/weed/command/filer_remote_sync_dir.go b/weed/command/filer_remote_sync_dir.go index 7f40a4287..688f95ef9 100644 --- a/weed/command/filer_remote_sync_dir.go +++ b/weed/command/filer_remote_sync_dir.go @@ -76,6 +76,9 @@ func followUpdatesAndUploadToRemote(option *RemoteSyncOptions, filerSource *sour StartTsNs: lastOffsetTs.UnixNano(), StopTsNs: 0, EventErrorType: pb.RetryForeverOnError, + GetResumeTsNs: func() int64 { + return processor.processedTsWatermark.Load() + }, } return pb.FollowMetadata(pb.ServerAddress(*option.filerAddress), option.grpcDialOption, metadataFollowOption, processEventFnWithOffset) diff --git a/weed/command/filer_sync.go b/weed/command/filer_sync.go index 64aebcc6f..1912f2f65 100644 --- a/weed/command/filer_sync.go +++ b/weed/command/filer_sync.go @@ -475,6 +475,9 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi StartTsNs: sourceFilerOffsetTsNs, StopTsNs: 0, EventErrorType: pb.RetryForeverOnError, + GetResumeTsNs: func() int64 { + return processor.processedTsWatermark.Load() + }, // While the source has only read activity it emits no metadata events, so // the watermark above never advances and sync_offset would look stuck. // The idle heartbeat moves the gauge to the source's current time once we diff --git a/weed/command/filer_sync_jobs.go b/weed/command/filer_sync_jobs.go index 732c1252d..9314d337b 100644 --- a/weed/command/filer_sync_jobs.go +++ b/weed/command/filer_sync_jobs.go @@ -16,6 +16,11 @@ import ( "github.com/seaweedfs/seaweedfs/weed/util" ) +// maxFailedSyncEvents bounds the failedTs ledger. A destination rejecting +// every event would otherwise add an entry per source event for the life of +// the processor. +var maxFailedSyncEvents = 1 << 16 + // tsMinHeap implements heap.Interface for int64 timestamps. type tsMinHeap []int64 @@ -60,6 +65,16 @@ type syncJobPaths struct { dataSize int64 } +// failedEventKey identifies an event for the failure ledger. A timestamp alone +// is not unique across events, so a success for one event must not clear an +// unresolved failure recorded for a different event at the same TsNs. +type failedEventKey struct { + tsNs int64 + path util.FullPath + newPath util.FullPath + kind jobKind +} + // syncStreamMetrics holds the metric children for one sync stream, curried // once so per-event updates skip the label lookup. type syncStreamMetrics struct { @@ -80,6 +95,7 @@ type MetadataProcessor struct { concurrencyLimit int fn pb.ProcessMetadataFunc processedTsWatermark atomic.Int64 + filteredTsNs int64 // Indexes for O(depth) conflict detection, replacing O(n) linear scan. // activeFilePaths counts active file jobs at each exact path. @@ -105,10 +121,15 @@ type MetadataProcessor struct { // used for O(log n) amortized watermark tracking. tsHeap tsMinHeap - // oldestFailedTsNs is the timestamp of the oldest event whose job returned - // an error, or 0 when none has. The watermark is never advanced to it or - // past it, so the persisted sync offset stays behind the failure and a - // restart replays the event instead of skipping it forever. + // failedTs records every event whose job returned an error and has not + // since completed, and oldestFailedTsNs caches its minimum (0 when empty). + // The watermark is never advanced to it or past it, so the persisted sync + // offset stays behind the failure and a restart replays the event instead + // of skipping it forever. Past maxFailedSyncEvents the set collapses to a + // sticky pin at the smallest failure seen: replay from the oldest failure + // still works, but individual recoveries no longer unpin until a restart. + failedTs map[failedEventKey]struct{} + failedSticky bool oldestFailedTsNs int64 // metrics is nil for callers that do not report per-event metrics. @@ -124,6 +145,7 @@ func NewMetadataProcessor(fn pb.ProcessMetadataFunc, concurrency int, offsetTsNs activeBarrierDirPaths: make(map[util.FullPath]int), activeNonBarrierDirPaths: make(map[util.FullPath]int), descendantCount: make(map[util.FullPath]int), + failedTs: make(map[failedEventKey]struct{}), } t.processedTsWatermark.Store(offsetTsNs) t.activeJobsCond = sync.NewCond(&t.activeJobsLock) @@ -281,6 +303,18 @@ func (t *MetadataProcessor) conflictsWith(resp *filer_pb.SubscribeMetadataRespon func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) { if filer_pb.IsEmpty(resp) { + // A filtered-progress marker means the source skipped everything below + // it for us; once all earlier work has finished, the watermark can move + // to it so idle stretches still advance the resume point. + t.activeJobsLock.Lock() + defer t.activeJobsLock.Unlock() + if resp.TsNs > t.filteredTsNs { + t.filteredTsNs = resp.TsNs + } + if len(t.activeJobs) == 0 && resp.TsNs > t.processedTsWatermark.Load() && + (t.oldestFailedTsNs == 0 || resp.TsNs < t.oldestFailedTsNs) { + t.processedTsWatermark.Store(resp.TsNs) + } return } @@ -327,12 +361,40 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) t.activeJobsLock.Lock() defer t.activeJobsLock.Unlock() + failedKey := failedEventKey{tsNs: resp.TsNs, path: jobPaths.path, newPath: jobPaths.newPath, kind: jobPaths.kind} if jobErr != nil { - if t.oldestFailedTsNs == 0 || resp.TsNs < t.oldestFailedTsNs { - t.oldestFailedTsNs = resp.TsNs - glog.Errorf("process %v: %v; holding sync offset at %v so this event is replayed on restart", resp, jobErr, time.Unix(0, resp.TsNs)) - } else { + if t.failedSticky { + if resp.TsNs < t.oldestFailedTsNs { + t.oldestFailedTsNs = resp.TsNs + } glog.Errorf("process %v: %v", resp, jobErr) + } else if _, recorded := t.failedTs[failedKey]; !recorded { + if len(t.failedTs) >= maxFailedSyncEvents { + t.failedSticky = true + t.failedTs = nil + if resp.TsNs < t.oldestFailedTsNs { + t.oldestFailedTsNs = resp.TsNs + } + glog.Warningf("process %v: %v; over %d unresolved failures, pinning sync offset at %v until restart", resp, jobErr, maxFailedSyncEvents, time.Unix(0, t.oldestFailedTsNs)) + } else { + t.failedTs[failedKey] = struct{}{} + if t.oldestFailedTsNs == 0 || resp.TsNs < t.oldestFailedTsNs { + t.oldestFailedTsNs = resp.TsNs + glog.Errorf("process %v: %v; holding sync offset at %v so this event is replayed on restart", resp, jobErr, time.Unix(0, resp.TsNs)) + } else { + glog.Errorf("process %v: %v", resp, jobErr) + } + } + } + } else if _, recorded := t.failedTs[failedKey]; recorded { + delete(t.failedTs, failedKey) + if resp.TsNs == t.oldestFailedTsNs { + t.oldestFailedTsNs = 0 + for k := range t.failedTs { + if t.oldestFailedTsNs == 0 || k.tsNs < t.oldestFailedTsNs { + t.oldestFailedTsNs = k.tsNs + } + } } } @@ -369,6 +431,10 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) t.processedTsWatermark.Store(resp.TsNs) } } + if len(t.activeJobs) == 0 && t.filteredTsNs > t.processedTsWatermark.Load() && + (t.oldestFailedTsNs == 0 || t.filteredTsNs < t.oldestFailedTsNs) { + t.processedTsWatermark.Store(t.filteredTsNs) + } t.activeJobsCond.Signal() }() } diff --git a/weed/command/filer_sync_jobs_test.go b/weed/command/filer_sync_jobs_test.go index 131bff00f..495f6f396 100644 --- a/weed/command/filer_sync_jobs_test.go +++ b/weed/command/filer_sync_jobs_test.go @@ -586,33 +586,6 @@ func BenchmarkConflictCheck(b *testing.B) { } } -// TestMetadataProcessorEmptyMarkerKeepsWatermarkStale: the MaxUnsyncedEvents -// marker (empty EventNotification, fresh timestamp) is dropped by AddSyncJob and -// does NOT advance processedTsWatermark, so offsetFunc keeps publishing the stale -// offset. This is why the client must not drive sync_offset off the watermark -// for these markers. -func TestMetadataProcessorEmptyMarkerKeepsWatermarkStale(t *testing.T) { - const staleOffset = int64(1_000_000_000) - freshTs := staleOffset + int64(time.Hour) // a "now"-ish source timestamp - - p := NewMetadataProcessor(func(*filer_pb.SubscribeMetadataResponse) error { return nil }, 4, staleOffset) - - marker := &filer_pb.SubscribeMetadataResponse{ - TsNs: freshTs, - EventNotification: &filer_pb.EventNotification{}, - } - if !filer_pb.IsEmpty(marker) { - t.Fatal("marker should be IsEmpty") - } - - p.AddSyncJob(marker) - - if got := p.processedTsWatermark.Load(); got != staleOffset { - t.Fatalf("empty marker advanced watermark to %d; want it to stay stale at %d", got, staleOffset) - } - t.Logf("marker carried fresh ts %d but watermark stayed stale at %d", freshTs, staleOffset) -} - // waitForJobsToDrain blocks until every job goroutine has finished bookkeeping. func waitForJobsToDrain(t *testing.T, p *MetadataProcessor) { t.Helper() @@ -742,3 +715,167 @@ func TestSyncStreamMetrics(t *testing.T) { } } } + +// TestFailedJobReplaySuccessClearsPin verifies that when the failed event is +// replayed (a reconnect resubscribing from the watermark) and succeeds this +// time, the failure pin clears and the watermark can move again. Without the +// clear, every later reconnect would replay the same backlog forever. +func TestFailedJobReplaySuccessClearsPin(t *testing.T) { + failed := true + fn := func(resp *filer_pb.SubscribeMetadataResponse) error { + if resp.TsNs == 200 && failed { + failed = false + return errors.New("AccessDenied: Access Denied") + } + return nil + } + p := NewMetadataProcessor(fn, 1, 0) + + p.AddSyncJob(makeResp("/dir", "a.txt", false, 100, true)) + p.AddSyncJob(makeResp("/dir", "b.txt", false, 200, true)) + p.AddSyncJob(makeResp("/dir", "c.txt", false, 300, true)) + waitForJobsToDrain(t, p) + if got := p.OldestFailedTsNs(); got != 200 { + t.Fatalf("oldest failed = %d, want 200", got) + } + if got := p.processedTsWatermark.Load(); got != 100 { + t.Fatalf("watermark = %d, want it held at 100 by the failure at 200", got) + } + + p.AddSyncJob(makeResp("/dir", "b.txt", false, 200, true)) + waitForJobsToDrain(t, p) + if got := p.OldestFailedTsNs(); got != 0 { + t.Fatalf("oldest failed = %d after a successful replay, want 0", got) + } + if got := p.processedTsWatermark.Load(); got != 200 { + t.Fatalf("watermark = %d after recovery, want 200", got) + } +} + +// TestFilteredMarkerAdvancesWatermark verifies that a filtered-progress marker +// (empty event with a timestamp) moves the watermark once all earlier work has +// finished, but never past an in-flight job or an unresolved failure. +func TestFilteredMarkerAdvancesWatermark(t *testing.T) { + marker := func(ts int64) *filer_pb.SubscribeMetadataResponse { + return &filer_pb.SubscribeMetadataResponse{TsNs: ts, EventNotification: &filer_pb.EventNotification{}} + } + + t.Run("idle", func(t *testing.T) { + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { return nil }, 100, 50) + p.AddSyncJob(marker(90)) + if got := p.processedTsWatermark.Load(); got != 90 { + t.Fatalf("watermark = %d after marker, want 90", got) + } + }) + + t.Run("behind in-flight job", func(t *testing.T) { + release := make(chan struct{}) + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + <-release + return nil + }, 100, 50) + p.AddSyncJob(makeResp("/dir", "f.txt", false, 60, true)) + p.AddSyncJob(marker(80)) + if got := p.processedTsWatermark.Load(); got != 50 { + t.Fatalf("watermark = %d with a job in flight, want 50", got) + } + close(release) + waitForJobsToDrain(t, p) + if got := p.processedTsWatermark.Load(); got != 80 { + t.Fatalf("watermark = %d after drain, want the retained marker at 80", got) + } + p.AddSyncJob(marker(90)) + if got := p.processedTsWatermark.Load(); got != 90 { + t.Fatalf("watermark = %d after drain and marker, want 90", got) + } + }) + + t.Run("behind a failure", func(t *testing.T) { + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + return errors.New("AccessDenied: Access Denied") + }, 100, 50) + p.AddSyncJob(makeResp("/dir", "f.txt", false, 100, true)) + waitForJobsToDrain(t, p) + p.AddSyncJob(marker(200)) + if got := p.processedTsWatermark.Load(); got != 50 { + t.Fatalf("watermark = %d past a failure pin, want 50", got) + } + p.AddSyncJob(marker(70)) + if got := p.processedTsWatermark.Load(); got != 70 { + t.Fatalf("watermark = %d behind the pin, want 70", got) + } + }) +} + +// TestFailedLedgerCapsAndStaysPinned verifies that a sustained run of distinct +// failures cannot grow failedTs without bound: past maxFailedSyncEvents the +// ledger collapses to a sticky pin at the oldest failure, so the watermark +// still replays from it while memory stays bounded. +func TestFailedLedgerCapsAndStaysPinned(t *testing.T) { + defer func(old int) { maxFailedSyncEvents = old }(maxFailedSyncEvents) + maxFailedSyncEvents = 4 + + fail := true + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + if fail { + return errors.New("AccessDenied: Access Denied") + } + return nil + }, 100, 0) + for i := int64(1); i <= 10; i++ { + p.AddSyncJob(makeResp("/dir", fmt.Sprintf("f%d.txt", i), false, i*100, true)) + } + waitForJobsToDrain(t, p) + + if !p.failedSticky { + t.Fatal("ledger did not collapse past the cap") + } + if got := p.OldestFailedTsNs(); got != 100 { + t.Fatalf("oldest failed = %d, want the pin at the oldest failure 100", got) + } + if got := p.processedTsWatermark.Load(); got != 0 { + t.Fatalf("watermark = %d, want it pinned at 0", got) + } + + fail = false + p.AddSyncJob(makeResp("/dir", "f1.txt", false, 100, true)) + waitForJobsToDrain(t, p) + if got := p.OldestFailedTsNs(); got != 100 { + t.Fatalf("oldest failed = %d after a collapsed replay, want the pin held at 100", got) + } + if got := p.processedTsWatermark.Load(); got != 0 { + t.Fatalf("watermark = %d after a collapsed replay, want it still pinned at 0", got) + } +} + +// TestFailedLedgerDistinguishesEventsAtSameTs verifies that a success for one +// event does not clear the pin recorded for a different event that happened to +// share its timestamp — the ledger keys on event identity, not just TsNs. +func TestFailedLedgerDistinguishesEventsAtSameTs(t *testing.T) { + fail := true + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + if fail && resp.EventNotification.NewEntry.GetName() == "bad.txt" { + return errors.New("AccessDenied: Access Denied") + } + return nil + }, 100, 0) + + p.AddSyncJob(makeResp("/dir", "bad.txt", false, 200, true)) + waitForJobsToDrain(t, p) + p.AddSyncJob(makeResp("/dir", "good.txt", false, 200, true)) + waitForJobsToDrain(t, p) + + if got := p.OldestFailedTsNs(); got != 200 { + t.Fatalf("oldest failed = %d, want the other event's pin held at 200", got) + } + if got := p.processedTsWatermark.Load(); got != 0 { + t.Fatalf("watermark = %d, want it still pinned at 0", got) + } + + fail = false + p.AddSyncJob(makeResp("/dir", "bad.txt", false, 200, true)) + waitForJobsToDrain(t, p) + if got := p.OldestFailedTsNs(); got != 0 { + t.Fatalf("oldest failed = %d after the failed event itself recovered, want 0", got) + } +} diff --git a/weed/pb/filer_pb_tail.go b/weed/pb/filer_pb_tail.go index 68538d7f4..c9fa54b83 100644 --- a/weed/pb/filer_pb_tail.go +++ b/weed/pb/filer_pb_tail.go @@ -44,6 +44,11 @@ type MetadataFollowOption struct { // a freshness signal only and does not advance StartTsNs, so the resume // checkpoint stays on the last real event. OnIdleHeartbeat func(tsNs int64) + // GetResumeTsNs, when non-nil, supplies the reconnect position instead of + // StartTsNs. It is read on every subscribe, so a callback can return the + // durably processed watermark while StartTsNs keeps tracking positions the + // stream has merely seen. + GetResumeTsNs func() int64 } type ProcessMetadataFunc func(resp *filer_pb.SubscribeMetadataResponse) error @@ -71,12 +76,16 @@ func makeSubscribeMetadataFunc(option *MetadataFollowOption, processEventFn Proc return func(client filer_pb.SeaweedFilerClient) error { ctx, cancel := context.WithCancel(context.Background()) defer cancel() + sinceNs := option.StartTsNs + if option.GetResumeTsNs != nil { + sinceNs = option.GetResumeTsNs() + } stream, err := client.SubscribeMetadata(ctx, &filer_pb.SubscribeMetadataRequest{ ClientName: option.ClientName, PathPrefix: option.PathPrefix, PathPrefixes: option.AdditionalPathPrefixes, Directories: option.DirectoriesToWatch, - SinceNs: option.StartTsNs, + SinceNs: sinceNs, Signature: option.SelfSignature, ClientId: option.ClientId, ClientEpoch: option.ClientEpoch, @@ -121,16 +130,31 @@ func makeSubscribeMetadataFunc(option *MetadataFollowOption, processEventFn Proc option.OnIdleHeartbeat(resp.TsNs) } // The marker advances the resume cursor past the filtered range; the - // heartbeat leaves StartTsNs put so a restart cannot outrun a straggler. + // heartbeat leaves it put so a restart cannot outrun a straggler. A + // consumer with a resume callback keeps its cursor in its processed + // watermark, so it sees the marker instead. if resp.EventNotification != nil && resp.TsNs > 0 { - option.StartTsNs = resp.TsNs + if option.GetResumeTsNs != nil { + if err := processEventFn(resp); err != nil { + handleErr(resp, err) + } + } else { + option.StartTsNs = resp.TsNs + } } return } if err := processEventFn(resp); err != nil { handleErr(resp, err) + // RetryForeverOnError only returns once the event was handled; + // other modes leave it failed and the cursor stays behind it. + if option.EventErrorType != RetryForeverOnError { + return + } + } + if option.GetResumeTsNs == nil { + option.StartTsNs = resp.TsNs } - option.StartTsNs = resp.TsNs } var pendingRefs []*filer_pb.LogFileChunkRef @@ -145,8 +169,15 @@ func makeSubscribeMetadataFunc(option *MetadataFollowOption, processEventFn Proc if len(pendingRefs) == 0 || option.LogFileReaderFn == nil { return nil } + readFromNs := option.StartTsNs + if option.GetResumeTsNs != nil { + // The resume cursor lives in the processed watermark, so a + // resubscribed ref replay must not be filtered by positions the + // previous stream had only seen. + readFromNs = sinceNs + } lastTs, readErr := ReadLogFileRefs(pendingRefs, option.LogFileReaderFn, - option.StartTsNs, option.StopTsNs, + readFromNs, option.StopTsNs, PathFilter{ PathPrefix: option.PathPrefix, AdditionalPathPrefixes: option.AdditionalPathPrefixes, @@ -156,7 +187,7 @@ func makeSubscribeMetadataFunc(option *MetadataFollowOption, processEventFn Proc if readErr != nil { return fmt.Errorf("%w: %w", ErrLogFileRead, readErr) } - if lastTs > 0 { + if lastTs > 0 && option.GetResumeTsNs == nil { option.StartTsNs = lastTs } pendingRefs = nil diff --git a/weed/pb/filer_pb_tail_test.go b/weed/pb/filer_pb_tail_test.go index 224d323bd..701a797c0 100644 --- a/weed/pb/filer_pb_tail_test.go +++ b/weed/pb/filer_pb_tail_test.go @@ -3,6 +3,8 @@ package pb import ( "context" "io" + "sync" + "sync/atomic" "testing" "time" @@ -60,8 +62,9 @@ func TestFilerSyncOffsetStaysFreshOnFilteredMarker(t *testing.T) { var timeline []gaugeWrite var heartbeatCalls, markerToProcessFn int - // AddSyncJob drops empty events and does not advance the watermark; the real - // processor is checked in command.TestMetadataProcessorEmptyMarkerKeepsWatermarkStale. + // This consumer sets no resume callback, so markers keep moving StartTsNs + // and never reach processEventFn; the callback path is checked in + // TestFilerSyncMarkerReachesCallbackConsumer. realProcessFn := func(resp *filer_pb.SubscribeMetadataResponse) error { if filer_pb.IsEmpty(resp) { markerToProcessFn++ @@ -192,3 +195,245 @@ func TestFilerSyncBatchedFreshnessSignalDoesNotCrash(t *testing.T) { t.Errorf("expected StartTsNs %d (marker), got %d (heartbeat must not advance the cursor)", markerTs, option.StartTsNs) } } + +// TestFilerSyncResumeFromProcessedWatermarkOnReconnect verifies that when GetResumeTsNs is +// configured, reconnection uses the processed watermark instead of skipping ahead to the +// latest received timestamp. +func TestFilerSyncResumeFromProcessedWatermarkOnReconnect(t *testing.T) { + const initialTs = int64(100) + const watermarkTs = int64(200) + const latestStreamTs = int64(500) + + var capturedSinceNs int64 + recordingClient := &recordingFilerClient{ + onSubscribe: func(req *filer_pb.SubscribeMetadataRequest) { + capturedSinceNs = req.SinceNs + }, + stream: &fakeSubscribeStream{ + responses: []*filer_pb.SubscribeMetadataResponse{ + { + Directory: "/watched", + TsNs: latestStreamTs, + EventNotification: &filer_pb.EventNotification{NewEntry: &filer_pb.Entry{Name: "file"}}, + }, + }, + }, + } + + option := &MetadataFollowOption{ + ClientName: "syncFrom_A_To_B", + StartTsNs: initialTs, + GetResumeTsNs: func() int64 { + return watermarkTs + }, + } + + processFn := func(resp *filer_pb.SubscribeMetadataResponse) error { + return nil + } + + fn := makeSubscribeMetadataFunc(option, processFn) + if err := fn(recordingClient); err != nil { + t.Fatalf("unexpected error: %v", err) + } + + if capturedSinceNs != watermarkTs { + t.Fatalf("expected subscribe SinceNs to be watermark %d, got %d", watermarkTs, capturedSinceNs) + } + // StartTsNs must not be mutated when GetResumeTsNs is set + if option.StartTsNs != initialTs { + t.Fatalf("expected option.StartTsNs to remain %d, got %d", initialTs, option.StartTsNs) + } +} + +// TestFilerSyncDoesNotAdvanceStartTsNsOnProcessError verifies that a failing synchronous +// processEventFn does not advance option.StartTsNs past the failed event. +func TestFilerSyncDoesNotAdvanceStartTsNsOnProcessError(t *testing.T) { + const initialTs = int64(100) + const failedTs = int64(200) + + option := &MetadataFollowOption{ + ClientName: "syncFrom_A_To_B", + StartTsNs: initialTs, + EventErrorType: TrivialOnError, + } + + stream := &fakeSubscribeStream{ + responses: []*filer_pb.SubscribeMetadataResponse{ + { + Directory: "/watched", + TsNs: failedTs, + EventNotification: &filer_pb.EventNotification{NewEntry: &filer_pb.Entry{Name: "bad"}}, + }, + }, + } + + processFn := func(resp *filer_pb.SubscribeMetadataResponse) error { + return io.ErrUnexpectedEOF + } + + fn := makeSubscribeMetadataFunc(option, processFn) + _ = fn(&fakeFilerClient{stream: stream}) + + if option.StartTsNs != initialTs { + t.Fatalf("expected StartTsNs to stay at %d on error, got %d", initialTs, option.StartTsNs) + } +} + +type recordingFilerClient struct { + filer_pb.SeaweedFilerClient + onSubscribe func(req *filer_pb.SubscribeMetadataRequest) + stream *fakeSubscribeStream +} + +func (c *recordingFilerClient) SubscribeMetadata(ctx context.Context, in *filer_pb.SubscribeMetadataRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[filer_pb.SubscribeMetadataResponse], error) { + if c.onSubscribe != nil { + c.onSubscribe(in) + } + return c.stream, nil +} + +// RetryForeverOnError resolves a failure inside handleErr, so the cursor must +// still move past the recovered event instead of replaying it on reconnect. +func TestFilerSyncRecoveredEventAdvancesCursor(t *testing.T) { + var calls int + processFn := func(resp *filer_pb.SubscribeMetadataResponse) error { + calls++ + if calls == 1 { + return io.ErrUnexpectedEOF + } + return nil + } + option := &MetadataFollowOption{ + ClientName: "syncFrom_A_To_B", + StartTsNs: 100, + EventErrorType: RetryForeverOnError, + } + stream := &fakeSubscribeStream{ + responses: []*filer_pb.SubscribeMetadataResponse{ + {Directory: "/watched", TsNs: 300, EventNotification: &filer_pb.EventNotification{ + NewEntry: &filer_pb.Entry{Name: "file"}, + }}, + }, + } + if err := makeSubscribeMetadataFunc(option, processFn)(&fakeFilerClient{stream: stream}); err != nil { + t.Fatalf("follow: %v", err) + } + if calls != 2 { + t.Fatalf("expected the event retried once (2 calls), got %d", calls) + } + if option.StartTsNs != 300 { + t.Fatalf("expected StartTsNs 300 after the retry succeeded, got %d", option.StartTsNs) + } +} + +// A consumer with a resume callback keeps its cursor in its processed +// watermark, so a filtered-progress marker is handed to processEventFn (which +// can count it as processed) instead of mutating StartTsNs. +func TestFilerSyncMarkerReachesCallbackConsumer(t *testing.T) { + var markers, events int + option := &MetadataFollowOption{ + ClientName: "syncFrom_A_To_B", + StartTsNs: 100, + GetResumeTsNs: func() int64 { + return 100 + }, + } + stream := &fakeSubscribeStream{ + responses: []*filer_pb.SubscribeMetadataResponse{ + {Directory: "/watched", TsNs: 300, EventNotification: &filer_pb.EventNotification{ + NewEntry: &filer_pb.Entry{Name: "file"}, + }}, + {TsNs: 500, EventNotification: &filer_pb.EventNotification{}}, + }, + } + fn := makeSubscribeMetadataFunc(option, func(resp *filer_pb.SubscribeMetadataResponse) error { + if filer_pb.IsEmpty(resp) { + markers++ + } else { + events++ + } + return nil + }) + if err := fn(&fakeFilerClient{stream: stream}); err != nil { + t.Fatalf("follow: %v", err) + } + if markers != 1 || events != 1 { + t.Fatalf("expected 1 marker and 1 event at processEventFn, got %d and %d", markers, events) + } + if option.StartTsNs != 100 { + t.Fatalf("callback consumer must not mutate StartTsNs, got %d", option.StartTsNs) + } +} + +func TestFilerSyncMarkerCallbackRetries(t *testing.T) { + var calls int + option := &MetadataFollowOption{ + StartTsNs: 100, + EventErrorType: RetryForeverOnError, + GetResumeTsNs: func() int64 { return 100 }, + } + stream := &fakeSubscribeStream{responses: []*filer_pb.SubscribeMetadataResponse{ + {TsNs: 500, EventNotification: &filer_pb.EventNotification{}}, + }} + fn := makeSubscribeMetadataFunc(option, func(resp *filer_pb.SubscribeMetadataResponse) error { + calls++ + if calls == 1 { + return io.ErrUnexpectedEOF + } + return nil + }) + if err := fn(&fakeFilerClient{stream: stream}); err != nil { + t.Fatalf("follow: %v", err) + } + if calls != 2 || option.StartTsNs != 100 { + t.Fatalf("calls = %d, cursor = %d; want 2 calls and unchanged cursor 100", calls, option.StartTsNs) + } +} + +// Each subscribe call re-reads the callback, so a reconnect after the consumer +// made progress resumes from the newer watermark. +func TestFilerSyncReconnectReadsWatermarkEachSubscribe(t *testing.T) { + var watermark atomic.Int64 + watermark.Store(100) + var sinceNs []int64 + var mu sync.Mutex + stream := &fakeSubscribeStream{ + responses: []*filer_pb.SubscribeMetadataResponse{ + {Directory: "/watched", TsNs: 300, EventNotification: &filer_pb.EventNotification{ + NewEntry: &filer_pb.Entry{Name: "file"}, + }}, + }, + } + client := &recordingFilerClient{ + onSubscribe: func(req *filer_pb.SubscribeMetadataRequest) { + mu.Lock() + sinceNs = append(sinceNs, req.SinceNs) + mu.Unlock() + }, + stream: stream, + } + option := &MetadataFollowOption{ + ClientName: "syncFrom_A_To_B", + StartTsNs: 100, + GetResumeTsNs: func() int64 { + return watermark.Load() + }, + } + fn := makeSubscribeMetadataFunc(option, func(resp *filer_pb.SubscribeMetadataResponse) error { + watermark.Store(resp.TsNs) + return nil + }) + if err := fn(client); err != nil { + t.Fatalf("first subscribe: %v", err) + } + client.stream = &fakeSubscribeStream{} + if err := fn(client); err != nil { + t.Fatalf("resubscribe: %v", err) + } + mu.Lock() + defer mu.Unlock() + if len(sinceNs) != 2 || sinceNs[0] != 100 || sinceNs[1] != 300 { + t.Fatalf("expected subscribes at 100 then 300, got %v", sinceNs) + } +} diff --git a/weed/util/retry.go b/weed/util/retry.go index 03b1fe5d7..fbed37a16 100644 --- a/weed/util/retry.go +++ b/weed/util/retry.go @@ -95,6 +95,9 @@ func IsTransientError(err error) bool { return true } if st, ok := ServerStatus(err); ok { + if st.Code() == codes.Canceled { + return strings.Contains(st.Message(), "the client connection is closing") + } return st.Code() == codes.Unavailable || st.Code() == codes.ResourceExhausted || IsTransientErrorMessage(st.Message()) } diff --git a/weed/util/retry_test.go b/weed/util/retry_test.go index f2782a533..9ec2ac17e 100644 --- a/weed/util/retry_test.go +++ b/weed/util/retry_test.go @@ -25,6 +25,9 @@ func TestIsTransientError(t *testing.T) { fmt.Errorf("send: %w", syscall.ETIMEDOUT), &net.DNSError{Err: "operation timed out", IsTimeout: true}, io.ErrUnexpectedEOF, + // transport teardown the peer reports as Canceled, not the caller's + // own context cancel + status.Error(codes.Canceled, "grpc: the client connection is closing"), } for _, err := range transient { if !IsTransientError(err) { @@ -37,6 +40,8 @@ func TestIsTransientError(t *testing.T) { errors.New("AccessDenied: Access Denied"), errors.New("NoSuchBucket: The specified bucket does not exist"), context.Canceled, + status.Error(codes.Canceled, context.Canceled.Error()), + fmt.Errorf("send: %w", status.Error(codes.Canceled, context.Canceled.Error())), fmt.Errorf("write: %w", context.DeadlineExceeded), } for _, err := range permanent {