diff --git a/weed/command/filer_remote_gateway_buckets.go b/weed/command/filer_remote_gateway_buckets.go index 9dffd4cc2..82f21b4e2 100644 --- a/weed/command/filer_remote_gateway_buckets.go +++ b/weed/command/filer_remote_gateway_buckets.go @@ -67,6 +67,7 @@ func (option *RemoteGatewayOptions) followBucketUpdatesAndUploadToRemote(filerSo GetResumeTsNs: func() int64 { return processor.processedTsWatermark.Load() }, + Resubscribe: processor.ResubscribeCh(), } 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 688f95ef9..0e02c3eaf 100644 --- a/weed/command/filer_remote_sync_dir.go +++ b/weed/command/filer_remote_sync_dir.go @@ -79,6 +79,7 @@ func followUpdatesAndUploadToRemote(option *RemoteSyncOptions, filerSource *sour GetResumeTsNs: func() int64 { return processor.processedTsWatermark.Load() }, + Resubscribe: processor.ResubscribeCh(), } return pb.FollowMetadata(pb.ServerAddress(*option.filerAddress), option.grpcDialOption, metadataFollowOption, processEventFnWithOffset) @@ -539,6 +540,11 @@ func collectLastSyncOffset(filerClient filer_pb.FilerClient, grpcDialOption grpc } } else { lastOffsetTs = time.Now().Add(-timeAgo) + if lastOffsetTsNs, err := remote_storage.GetSyncOffset(grpcDialOption, filerAddress, mountedDir); err == nil && lastOffsetTsNs > 0 { + if savedOffsetTs := time.Unix(0, lastOffsetTsNs); savedOffsetTs.Before(lastOffsetTs) { + lastOffsetTs = savedOffsetTs + } + } } return lastOffsetTs } diff --git a/weed/command/filer_sync.go b/weed/command/filer_sync.go index 1912f2f65..d987f5799 100644 --- a/weed/command/filer_sync.go +++ b/weed/command/filer_sync.go @@ -478,6 +478,7 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi GetResumeTsNs: func() int64 { return processor.processedTsWatermark.Load() }, + Resubscribe: processor.ResubscribeCh(), // 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 9314d337b..69554fd9a 100644 --- a/weed/command/filer_sync_jobs.go +++ b/weed/command/filer_sync_jobs.go @@ -89,7 +89,12 @@ type syncStreamMetrics struct { } type MetadataProcessor struct { - activeJobs map[int64]*syncJobPaths + // activeJobCount is the number of in-flight jobs and activeJobTs counts + // them per event timestamp. Several events can share a TsNs — batched + // writes log together — so a single slot per TsNs would drain early and + // let the resubscribe or the watermark outrun a sibling still running. + activeJobCount int + activeJobTs map[int64]int activeJobsLock sync.Mutex activeJobsCond *sync.Cond concurrencyLimit int @@ -132,6 +137,24 @@ type MetadataProcessor struct { failedSticky bool oldestFailedTsNs int64 + // resubscribeCh closes once a failure has stopped the processor and all + // in-flight jobs have drained, asking the metadata follower to drop the + // stream so the caller's reconnect replays the pinned events in order — + // and never races the replay against work still running in this abandoned + // processor. + resubscribeCh chan struct{} + resubscribeOnce sync.Once + + // stopped latches when a job failure pins the watermark. Admission then + // drops new events instead of queueing them into a processor that is about + // to be abandoned: every skipped event replays from the pinned watermark + // after the reconnect, and on a stream that never goes quiet the drain — + // and with it the replay — would otherwise never come. The cond broadcast + // releases a blocked AddSyncJob. A redelivery of an event still in + // failedTs is the one exception: it may run so its success shrinks the + // replay. + stopped bool + // metrics is nil for callers that do not report per-event metrics. metrics *syncStreamMetrics } @@ -139,13 +162,14 @@ type MetadataProcessor struct { func NewMetadataProcessor(fn pb.ProcessMetadataFunc, concurrency int, offsetTsNs int64) *MetadataProcessor { t := &MetadataProcessor{ fn: fn, - activeJobs: make(map[int64]*syncJobPaths), + activeJobTs: make(map[int64]int), concurrencyLimit: concurrency, activeFilePaths: make(map[util.FullPath]int), activeBarrierDirPaths: make(map[util.FullPath]int), activeNonBarrierDirPaths: make(map[util.FullPath]int), descendantCount: make(map[util.FullPath]int), failedTs: make(map[failedEventKey]struct{}), + resubscribeCh: make(chan struct{}), } t.processedTsWatermark.Store(offsetTsNs) t.activeJobsCond = sync.NewCond(&t.activeJobsLock) @@ -175,6 +199,13 @@ func (t *MetadataProcessor) OldestFailedTsNs() int64 { return t.oldestFailedTsNs } +// ResubscribeCh closes once a job failure has stopped the processor and its +// in-flight jobs have drained, signaling the metadata follower to drop the +// stream so a reconnect replays what the watermark still covers. +func (t *MetadataProcessor) ResubscribeCh() <-chan struct{} { + return t.resubscribeCh +} + // pathAncestors returns all proper ancestor directories of p. // For "/a/b/c", returns ["/a/b", "/a", "/"]. func pathAncestors(p util.FullPath) []util.FullPath { @@ -311,7 +342,7 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) if resp.TsNs > t.filteredTsNs { t.filteredTsNs = resp.TsNs } - if len(t.activeJobs) == 0 && resp.TsNs > t.processedTsWatermark.Load() && + if t.activeJobCount == 0 && resp.TsNs > t.processedTsWatermark.Load() && (t.oldestFailedTsNs == 0 || resp.TsNs < t.oldestFailedTsNs) { t.processedTsWatermark.Store(resp.TsNs) } @@ -320,24 +351,48 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) dataSize := eventDataSize(resp) - // counted before the admission wait: received-processed-failed is the - // number of events read off the stream but not yet done + t.activeJobsLock.Lock() + defer t.activeJobsLock.Unlock() + + p, newPath, kind := extractJobInfo(resp) + eventKey := failedEventKey{tsNs: resp.TsNs, path: p, newPath: newPath, kind: kind} + + for t.activeJobCount >= t.concurrencyLimit || t.conflictsWith(resp) { + // A stopped processor never queues: an event that cannot start drops + // and replays in order after the resubscribe. + if t.stopped { + return + } + t.activeJobsCond.Wait() + } + if t.stopped { + select { + case <-t.resubscribeCh: + // Already drained and signaled: nothing may start now, or it would + // race the replay the signal just asked for. + return + default: + } + // The one event still worth running is a failure still pinning the + // watermark, redelivered on this stream: its success clears the ledger + // entry and shrinks the replay. Everything else replays anyway. + if _, pinned := t.failedTs[eventKey]; !pinned { + return + } + } + + // counted once admitted: received-processed-failed is the number of events + // this processor read but has not finished. A dropped event is not counted + // here — the replay's own admission counts it. if t.metrics != nil { t.metrics.received.Inc() t.metrics.receivedBytes.Add(float64(dataSize)) } - t.activeJobsLock.Lock() - defer t.activeJobsLock.Unlock() - - for len(t.activeJobs) >= t.concurrencyLimit || t.conflictsWith(resp) { - t.activeJobsCond.Wait() - } - - p, newPath, kind := extractJobInfo(resp) jobPaths := &syncJobPaths{path: p, newPath: newPath, kind: kind, dataSize: dataSize} - t.activeJobs[resp.TsNs] = jobPaths + t.activeJobCount++ + t.activeJobTs[resp.TsNs]++ t.addPathToIndex(p, kind) if newPath != "" { t.addPathToIndex(newPath, kind) @@ -363,6 +418,13 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) failedKey := failedEventKey{tsNs: resp.TsNs, path: jobPaths.path, newPath: jobPaths.newPath, kind: jobPaths.kind} if jobErr != nil { + // Latch the stop the moment a failure lands: events behind the pin + // replay after the resubscribe anyway, and a stream that never goes + // quiet would otherwise keep the active count above zero so the + // failure stays pinned until an unrelated reconnect — the wait this + // whole mechanism exists to remove. + t.stopped = true + t.activeJobsCond.Broadcast() if t.failedSticky { if resp.TsNs < t.oldestFailedTsNs { t.oldestFailedTsNs = resp.TsNs @@ -398,7 +460,11 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) } } - delete(t.activeJobs, resp.TsNs) + t.activeJobCount-- + t.activeJobTs[resp.TsNs]-- + if t.activeJobTs[resp.TsNs] == 0 { + delete(t.activeJobTs, resp.TsNs) + } t.removePathFromIndex(jobPaths.path, jobPaths.kind) if jobPaths.newPath != "" { t.removePathFromIndex(jobPaths.newPath, jobPaths.kind) @@ -418,7 +484,7 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) // Lazy-clean stale entries from heap top (already-completed jobs). // Each entry is pushed once and popped once: O(log n) amortized. for t.tsHeap.Len() > 0 { - if _, active := t.activeJobs[t.tsHeap[0]]; active { + if t.activeJobTs[t.tsHeap[0]] > 0 { break } heap.Pop(&t.tsHeap) @@ -431,10 +497,15 @@ func (t *MetadataProcessor) AddSyncJob(resp *filer_pb.SubscribeMetadataResponse) t.processedTsWatermark.Store(resp.TsNs) } } - if len(t.activeJobs) == 0 && t.filteredTsNs > t.processedTsWatermark.Load() && + if t.activeJobCount == 0 && t.filteredTsNs > t.processedTsWatermark.Load() && (t.oldestFailedTsNs == 0 || t.filteredTsNs < t.oldestFailedTsNs) { t.processedTsWatermark.Store(t.filteredTsNs) } + // Signal once the stop has drained: even if a redelivery cleared the + // pin, events dropped while stopped still have to replay. + if t.stopped && t.activeJobCount == 0 { + t.resubscribeOnce.Do(func() { close(t.resubscribeCh) }) + } t.activeJobsCond.Signal() }() } diff --git a/weed/command/filer_sync_jobs_test.go b/weed/command/filer_sync_jobs_test.go index 495f6f396..968858638 100644 --- a/weed/command/filer_sync_jobs_test.go +++ b/weed/command/filer_sync_jobs_test.go @@ -94,8 +94,8 @@ func TestFileVsFileConflict(t *testing.T) { // Add a file job active := makeResp("/dir1", "file.txt", false, 1, true) - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) // Same file should conflict @@ -125,8 +125,8 @@ func TestFileUnderActiveDirConflict(t *testing.T) { // Add a directory job at /dir1 active := makeResp("/", "dir1", true, 1, true) - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) // File under /dir1 should conflict @@ -163,8 +163,8 @@ func TestDirWithActiveFileUnder(t *testing.T) { // Add file jobs under /dir1 f1 := makeResp("/dir1/sub", "file.txt", false, 1, true) - path, newPath, kind := extractJobInfo(f1) - p.activeJobs[f1.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(f1) + p.activeJobTs[f1.TsNs] = 1 p.addPathToIndex(path, kind) // Directory /dir1 should conflict (has active file under it) @@ -187,8 +187,8 @@ func TestDirVsDirConflict(t *testing.T) { // Add directory job at /a/b active := makeResp("/a", "b", true, 1, true) - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) // /a/b/c (descendant) should conflict @@ -224,8 +224,8 @@ func TestRenameConflict(t *testing.T) { // Add file job at /dir1/file.txt f1 := makeResp("/dir1", "file.txt", false, 1, true) - path, newPath, kind := extractJobInfo(f1) - p.activeJobs[f1.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(f1) + p.activeJobTs[f1.TsNs] = 1 p.addPathToIndex(path, kind) // Rename from /dir2/a.txt to /dir1/file.txt should conflict (newPath matches) @@ -255,7 +255,7 @@ func TestActiveRenameConflict(t *testing.T) { // Add active rename job: /dir1/old.txt -> /dir2/new.txt rename := makeRenameResp("/dir1", "old.txt", "/dir2", "new.txt", false, 1) path, newPath, kind := extractJobInfo(rename) - p.activeJobs[rename.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + p.activeJobTs[rename.TsNs] = 1 p.addPathToIndex(path, kind) if newPath != "" { p.addPathToIndex(newPath, kind) @@ -289,8 +289,8 @@ func TestRootDirConflict(t *testing.T) { // Note: a dir entry at "/" would be created as FullPath("/").Child("somedir") // But let's test what happens with an active dir at /some/path and check root active := makeResp("/some", "dir", true, 1, true) - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) // Root dir should conflict because active dir /some/dir is under / @@ -334,7 +334,7 @@ func TestWatermarkWithHeap(t *testing.T) { // Simulate adding jobs in order for _, ts := range []int64{10, 20, 30} { jobPath := util.FullPath("/file" + string(rune('0'+ts/10))) - p.activeJobs[ts] = &syncJobPaths{path: jobPath, kind: kindFile} + p.activeJobTs[ts] = 1 p.addPathToIndex(jobPath, kindFile) heap.Push(&p.tsHeap, ts) } @@ -344,11 +344,11 @@ func TestWatermarkWithHeap(t *testing.T) { } // Remove non-oldest (ts=20) — heap top should stay 10 - delete(p.activeJobs, 20) + delete(p.activeJobTs, 20) p.removePathFromIndex("/file2", kindFile) // Lazy clean: top is 10 which is still active, so no pop for p.tsHeap.Len() > 0 { - if _, active := p.activeJobs[p.tsHeap[0]]; active { + if p.activeJobTs[p.tsHeap[0]] > 0 { break } heap.Pop(&p.tsHeap) @@ -358,10 +358,10 @@ func TestWatermarkWithHeap(t *testing.T) { } // Remove oldest (ts=10) — lazy clean should find 30 - delete(p.activeJobs, 10) + delete(p.activeJobTs, 10) p.removePathFromIndex("/file1", kindFile) for p.tsHeap.Len() > 0 { - if _, active := p.activeJobs[p.tsHeap[0]]; active { + if p.activeJobTs[p.tsHeap[0]] > 0 { break } heap.Pop(&p.tsHeap) @@ -383,11 +383,11 @@ func TestNonBarrierDirUpdateDoesNotBlockDescendants(t *testing.T) { // Active non-barrier: attribute update on /dir1. active := makeDirUpdateResp("/", "dir1", 1) - path, newPath, kind := extractJobInfo(active) + path, _, kind := extractJobInfo(active) if kind != kindNonBarrierDir { t.Fatalf("expected kindNonBarrierDir for dir attribute update, got %v", kind) } - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) // File under /dir1 should NOT conflict with the attribute update. @@ -407,11 +407,11 @@ func TestNonBarrierDirUpdateDoesNotBlockDescendants(t *testing.T) { p := NewMetadataProcessor(noop, 100, 0) active := makeResp("/", "dir1", true, 1, true) // create - path, newPath, kind := extractJobInfo(active) + path, _, kind := extractJobInfo(active) if kind != kindBarrierDir { t.Fatalf("expected kindBarrierDir for dir create, got %v", kind) } - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) under := makeResp("/dir1", "file.txt", false, 2, true) @@ -425,8 +425,8 @@ func TestNonBarrierDirUpdateDoesNotBlockDescendants(t *testing.T) { // Active file under /dir1. f := makeResp("/dir1", "file.txt", false, 1, true) - path, newPath, kind := extractJobInfo(f) - p.activeJobs[f.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(f) + p.activeJobTs[f.TsNs] = 1 p.addPathToIndex(path, kind) // Incoming barrier delete on /dir1 should still wait for the @@ -442,8 +442,8 @@ func TestNonBarrierDirUpdateDoesNotBlockDescendants(t *testing.T) { // Active non-barrier dir update at /a/b. upd := makeDirUpdateResp("/a", "b", 1) - path, newPath, kind := extractJobInfo(upd) - p.activeJobs[upd.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(upd) + p.activeJobTs[upd.TsNs] = 1 p.addPathToIndex(path, kind) // A barrier delete on /a (the ancestor) should wait for it. @@ -464,8 +464,8 @@ func TestSamePathBarrierSerialization(t *testing.T) { t.Run("barrier dir at p blocks same-path file", func(t *testing.T) { p := NewMetadataProcessor(noop, 100, 0) active := makeResp("/", "dir1", true, 1, true) // dir create - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) file := makeResp("/", "dir1", false, 2, true) @@ -477,8 +477,8 @@ func TestSamePathBarrierSerialization(t *testing.T) { t.Run("barrier dir at p blocks another same-path barrier dir", func(t *testing.T) { p := NewMetadataProcessor(noop, 100, 0) active := makeResp("/", "dir1", true, 1, true) // dir create - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) del := makeResp("/", "dir1", true, 2, false) // dir delete, same path @@ -490,8 +490,8 @@ func TestSamePathBarrierSerialization(t *testing.T) { t.Run("barrier dir at p blocks non-barrier update at same path", func(t *testing.T) { p := NewMetadataProcessor(noop, 100, 0) active := makeResp("/", "dir1", true, 1, true) // dir create - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) upd := makeDirUpdateResp("/", "dir1", 2) @@ -503,8 +503,8 @@ func TestSamePathBarrierSerialization(t *testing.T) { t.Run("file at p blocks same-path barrier dir", func(t *testing.T) { p := NewMetadataProcessor(noop, 100, 0) active := makeResp("/", "thing", false, 1, true) // file create at /thing - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) // Barrier dir at /thing (e.g. a file→dir promotion) must wait. @@ -520,11 +520,11 @@ func TestSamePathBarrierSerialization(t *testing.T) { // delete/rename/create on /dir1. p := NewMetadataProcessor(noop, 100, 0) active := makeDirUpdateResp("/", "dir1", 1) - path, newPath, kind := extractJobInfo(active) + path, _, kind := extractJobInfo(active) if kind != kindNonBarrierDir { t.Fatalf("expected kindNonBarrierDir, got %v", kind) } - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) del := makeResp("/", "dir1", true, 2, false) // dir delete @@ -541,8 +541,8 @@ func TestSamePathBarrierSerialization(t *testing.T) { t.Run("non-barrier update at p does NOT block same-path non-barrier update", func(t *testing.T) { p := NewMetadataProcessor(noop, 100, 0) active := makeDirUpdateResp("/", "dir1", 1) - path, newPath, kind := extractJobInfo(active) - p.activeJobs[active.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(active) + p.activeJobTs[active.TsNs] = 1 p.addPathToIndex(path, kind) // Concurrent attribute bumps are allowed: last writer wins. @@ -569,8 +569,8 @@ func BenchmarkConflictCheck(b *testing.B) { dir := fmt.Sprintf("/dir%d/sub%d", i/100, i%100) name := fmt.Sprintf("file%d.txt", i) resp := makeResp(dir, name, false, int64(i+1), true) - path, newPath, kind := extractJobInfo(resp) - p.activeJobs[resp.TsNs] = &syncJobPaths{path: path, newPath: newPath, kind: kind} + path, _, kind := extractJobInfo(resp) + p.activeJobTs[resp.TsNs] = 1 p.addPathToIndex(path, kind) } @@ -592,7 +592,7 @@ func waitForJobsToDrain(t *testing.T, p *MetadataProcessor) { deadline := time.Now().Add(10 * time.Second) for time.Now().Before(deadline) { p.activeJobsLock.Lock() - remaining := len(p.activeJobs) + remaining := p.activeJobCount p.activeJobsLock.Unlock() if remaining == 0 { return @@ -624,35 +624,54 @@ func TestFailedJobHoldsWatermark(t *testing.T) { t.Fatalf("watermark = %d after a successful job, want 100", got) } - p.AddSyncJob(makeResp("/dir", "b.txt", false, failedTsNs, true)) - waitForJobsToDrain(t, p) - if got := p.processedTsWatermark.Load(); got != 100 { - t.Fatalf("watermark = %d after a failed job, want it held at 100", got) + // a later event still in flight when the failure lands finishes fine, but + // the offset stays behind the failure + release := make(chan struct{}) + slowFn := func(resp *filer_pb.SubscribeMetadataResponse) error { + if resp.TsNs == 300 { + <-release + } + return fn(resp) } - - // later events keep flowing, but the offset stays behind the failure - p.AddSyncJob(makeResp("/dir", "c.txt", false, 300, true)) - waitForJobsToDrain(t, p) - if got := p.processedTsWatermark.Load(); got != 100 { + p2 := NewMetadataProcessor(slowFn, 10, 0) + // admit the slow job first so it is in flight when the failure lands + p2.AddSyncJob(makeResp("/dir", "c.txt", false, 300, true)) + p2.AddSyncJob(makeResp("/dir", "a.txt", false, 100, true)) + p2.AddSyncJob(makeResp("/dir", "b.txt", false, failedTsNs, true)) + deadline := time.Now().Add(10 * time.Second) + for p2.OldestFailedTsNs() == 0 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + if p2.OldestFailedTsNs() != failedTsNs { + t.Fatalf("oldest failed = %d, want the pin at %d", p2.OldestFailedTsNs(), failedTsNs) + } + close(release) + waitForJobsToDrain(t, p2) + if got := p2.processedTsWatermark.Load(); got != 100 { t.Fatalf("watermark = %d after a later success, want it held at 100", got) } } // TestFailedJobHoldsWatermarkAtOldestFailure verifies that the watermark is -// pinned by the oldest failure, not the most recent one. +// pinned by the oldest failure, not the most recent one. Once the processor +// drains it stops accepting events for the resubscribe, so both failures have +// to be in the same drained batch. func TestFailedJobHoldsWatermarkAtOldestFailure(t *testing.T) { + release := make(chan struct{}) fn := func(resp *filer_pb.SubscribeMetadataResponse) error { + <-release if resp.TsNs == 200 || resp.TsNs == 400 { return errors.New("AccessDenied: Access Denied") } return nil } - p := NewMetadataProcessor(fn, 1, 0) + p := NewMetadataProcessor(fn, 10, 0) for _, ts := range []int64{100, 200, 300, 400, 500} { p.AddSyncJob(makeResp("/dir", fmt.Sprintf("f%d.txt", ts), false, ts, true)) - waitForJobsToDrain(t, p) } + close(release) + waitForJobsToDrain(t, p) if got := p.processedTsWatermark.Load(); got != 100 { t.Fatalf("watermark = %d, want it held at 100 by the failure at 200", got) @@ -665,14 +684,18 @@ func TestFailedJobHoldsWatermarkAtOldestFailure(t *testing.T) { // update's new 60-byte chunk but not its shared one, and nothing for the // delete despite its chunk. func TestSyncStreamMetrics(t *testing.T) { + release := make(chan struct{}) fn := func(resp *filer_pb.SubscribeMetadataResponse) error { + <-release if resp.TsNs == 2 { return errors.New("AccessDenied: Access Denied") } return nil } p := NewMetadataProcessor(fn, 100, 0) - p.SetMetrics("srcFiler", "dstFiler", "TestSyncStreamMetrics", "/") + // the counters are process-global: a unique client name keeps repeated + // runs (-count>1) from accumulating into each other + p.SetMetrics("srcFiler", "dstFiler", fmt.Sprintf("TestSyncStreamMetrics-%d", time.Now().UnixNano()), "/") create := makeResp("/dir1", "a.txt", false, 1, true) create.EventNotification.NewEntry.Chunks = []*filer_pb.FileChunk{{FileId: "1,a0", Size: 100}} @@ -694,6 +717,7 @@ func TestSyncStreamMetrics(t *testing.T) { for _, resp := range []*filer_pb.SubscribeMetadataResponse{create, failing, update, del} { p.AddSyncJob(resp) } + close(release) waitForJobsToDrain(t, p) for _, tc := range []struct { @@ -717,38 +741,67 @@ 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. +// redelivered while the processor is still alive and succeeds this time, the +// failure pin clears and the watermark can move again. It is the one event a +// stopped processor still runs. The resubscribe still signals once the jobs +// drain: anything dropped after the stop has to replay too. func TestFailedJobReplaySuccessClearsPin(t *testing.T) { failed := true + release := make(chan struct{}) fn := func(resp *filer_pb.SubscribeMetadataResponse) error { + if resp.TsNs == 300 { + <-release + } if resp.TsNs == 200 && failed { failed = false return errors.New("AccessDenied: Access Denied") } return nil } - p := NewMetadataProcessor(fn, 1, 0) + p := NewMetadataProcessor(fn, 10, 0) - p.AddSyncJob(makeResp("/dir", "a.txt", false, 100, true)) - p.AddSyncJob(makeResp("/dir", "b.txt", false, 200, true)) + // the slow job is admitted first so the processor has in-flight work when + // the failure lands, keeping the drain — and the resubscribe — open p.AddSyncJob(makeResp("/dir", "c.txt", false, 300, true)) - waitForJobsToDrain(t, p) + p.AddSyncJob(makeResp("/dir", "b.txt", false, 200, true)) + + deadline := time.Now().Add(10 * time.Second) + for p.OldestFailedTsNs() != 200 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } 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) + + // once stopped, a new event drops instead of queueing into a processor + // that is about to be abandoned — it replays after the resubscribe + p.AddSyncJob(makeResp("/dir", "d.txt", false, 400, true)) + p.activeJobsLock.Lock() + dropped := p.activeJobTs[400] == 0 + p.activeJobsLock.Unlock() + if !dropped { + t.Fatal("new event admitted after the failure stopped the processor") } + // the redelivery is the exception: it runs and its success clears the pin p.AddSyncJob(makeResp("/dir", "b.txt", false, 200, true)) - waitForJobsToDrain(t, p) + deadline = time.Now().Add(10 * time.Second) + for p.OldestFailedTsNs() != 0 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } 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) + + close(release) + waitForJobsToDrain(t, p) + select { + case <-p.ResubscribeCh(): + case <-time.After(time.Second): + t.Fatal("resubscribe never signaled after the stopped processor drained") + } + if got := p.processedTsWatermark.Load(); got != 300 { + t.Fatalf("watermark = %d after recovery, want 300", got) } } @@ -816,7 +869,9 @@ func TestFailedLedgerCapsAndStaysPinned(t *testing.T) { maxFailedSyncEvents = 4 fail := true + release := make(chan struct{}) p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + <-release if fail { return errors.New("AccessDenied: Access Denied") } @@ -825,6 +880,7 @@ func TestFailedLedgerCapsAndStaysPinned(t *testing.T) { for i := int64(1); i <= 10; i++ { p.AddSyncJob(makeResp("/dir", fmt.Sprintf("f%d.txt", i), false, i*100, true)) } + close(release) waitForJobsToDrain(t, p) if !p.failedSticky { @@ -850,21 +906,36 @@ func TestFailedLedgerCapsAndStaysPinned(t *testing.T) { // 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. +// share its timestamp — the ledger keys on event identity, not just TsNs. A +// slow job holds the drain open so the redelivery still lands on this +// processor generation. func TestFailedLedgerDistinguishesEventsAtSameTs(t *testing.T) { fail := true + release := make(chan struct{}) + hold := make(chan struct{}) p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { - if fail && resp.EventNotification.NewEntry.GetName() == "bad.txt" { + if resp.EventNotification.NewEntry.GetName() == "slow.txt" { + <-hold + } else { + <-release + } + if resp.EventNotification.NewEntry.GetName() == "bad.txt" && fail { return errors.New("AccessDenied: Access Denied") } return nil }, 100, 0) + // both same-ts events admit before either resolves, so the success lands + // while the failure is already pinned + p.AddSyncJob(makeResp("/dir", "slow.txt", false, 900, true)) p.AddSyncJob(makeResp("/dir", "bad.txt", false, 200, true)) - waitForJobsToDrain(t, p) p.AddSyncJob(makeResp("/dir", "good.txt", false, 200, true)) - waitForJobsToDrain(t, p) + close(release) + deadline := time.Now().Add(10 * time.Second) + for p.OldestFailedTsNs() != 200 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } if got := p.OldestFailedTsNs(); got != 200 { t.Fatalf("oldest failed = %d, want the other event's pin held at 200", got) } @@ -872,10 +943,133 @@ func TestFailedLedgerDistinguishesEventsAtSameTs(t *testing.T) { t.Fatalf("watermark = %d, want it still pinned at 0", got) } + // redelivering the failed event itself is what clears the pin — and it is + // the one event a stopped processor still admits fail = false p.AddSyncJob(makeResp("/dir", "bad.txt", false, 200, true)) + deadline = time.Now().Add(10 * time.Second) + for p.OldestFailedTsNs() != 0 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + close(hold) waitForJobsToDrain(t, p) if got := p.OldestFailedTsNs(); got != 0 { t.Fatalf("oldest failed = %d after the failed event itself recovered, want 0", got) } + if got := p.processedTsWatermark.Load(); got != 900 { + t.Fatalf("watermark = %d after the pin cleared and the rest drained, want 900", got) + } +} + +// TestFailedJobSignalsResubscribe verifies that a job exhausting its retries +// closes ResubscribeCh so the follower drops the stream and the reconnect +// replays the pinned event — the path that used to wait for a restart. +func TestFailedJobSignalsResubscribe(t *testing.T) { + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + return errors.New("AccessDenied: Access Denied") + }, 100, 0) + + p.AddSyncJob(makeResp("/dir", "a.txt", false, 100, true)) + waitForJobsToDrain(t, p) + + select { + case <-p.ResubscribeCh(): + case <-time.After(time.Second): + t.Fatal("resubscribe channel never closed after the failure pinned the watermark") + } + if got := p.OldestFailedTsNs(); got != 100 { + t.Fatalf("oldest failed = %d, want 100", got) + } +} + +// TestResubscribeWaitsForInFlightJobs verifies the signal stays open while +// jobs admitted before the failure are still running — replaying behind them +// could restore older state over their writes — and closes once they drain. +// An event arriving after the stop drops instead of keeping the drain open, +// so a busy stream cannot starve the replay. +func TestResubscribeWaitsForInFlightJobs(t *testing.T) { + release := make(chan struct{}) + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + if resp.TsNs == 100 { + return errors.New("AccessDenied: Access Denied") + } + <-release + return nil + }, 100, 0) + + // the slow job admits first so it is in flight when the failure lands + p.AddSyncJob(makeResp("/dir", "b.txt", false, 200, true)) + p.AddSyncJob(makeResp("/dir", "a.txt", false, 100, true)) + + deadline := time.Now().Add(10 * time.Second) + for p.OldestFailedTsNs() == 0 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + if p.OldestFailedTsNs() != 100 { + t.Fatalf("oldest failed = %d, want the pin at 100", p.OldestFailedTsNs()) + } + select { + case <-p.ResubscribeCh(): + t.Fatal("resubscribe signaled while an in-flight job could still race the replay") + case <-time.After(50 * time.Millisecond): + } + + // the processor stopped on the failure, so this event drops — it replays + // after the resubscribe — instead of starving the drain + p.AddSyncJob(makeResp("/dir", "c.txt", false, 300, true)) + p.activeJobsLock.Lock() + dropped := p.activeJobTs[300] == 0 + p.activeJobsLock.Unlock() + if !dropped { + t.Fatal("event admitted after the processor stopped") + } + + close(release) + waitForJobsToDrain(t, p) + select { + case <-p.ResubscribeCh(): + case <-time.After(time.Second): + t.Fatal("resubscribe channel never closed after the in-flight jobs drained") + } +} + +// TestResubscribeWaitsForSameTsSibling guards the per-timestamp job +// counting: events in one batch can share a TsNs, and the failed job must +// not free the bookkeeping of a sibling still running at that timestamp. +// Otherwise its completion could report the processor drained and the +// resubscribe would replay over the sibling's writes. +func TestResubscribeWaitsForSameTsSibling(t *testing.T) { + release := make(chan struct{}) + p := NewMetadataProcessor(func(resp *filer_pb.SubscribeMetadataResponse) error { + if resp.EventNotification.NewEntry.GetName() == "bad.txt" { + return errors.New("AccessDenied: Access Denied") + } + <-release + return nil + }, 100, 0) + + // the slow job admits first so it is in flight when the same-ts failure lands + p.AddSyncJob(makeResp("/dir", "slow.txt", false, 200, true)) + p.AddSyncJob(makeResp("/dir", "bad.txt", false, 200, true)) + + deadline := time.Now().Add(10 * time.Second) + for p.OldestFailedTsNs() == 0 && time.Now().Before(deadline) { + time.Sleep(time.Millisecond) + } + if p.OldestFailedTsNs() != 200 { + t.Fatalf("oldest failed = %d, want the pin at 200", p.OldestFailedTsNs()) + } + select { + case <-p.ResubscribeCh(): + t.Fatal("resubscribe signaled while a same-ts job was still running") + case <-time.After(50 * time.Millisecond): + } + + close(release) + waitForJobsToDrain(t, p) + select { + case <-p.ResubscribeCh(): + case <-time.After(time.Second): + t.Fatal("resubscribe channel never closed after the same-ts jobs drained") + } } diff --git a/weed/pb/filer_pb_tail.go b/weed/pb/filer_pb_tail.go index c9fa54b83..edbd68c01 100644 --- a/weed/pb/filer_pb_tail.go +++ b/weed/pb/filer_pb_tail.go @@ -2,6 +2,7 @@ package pb import ( "context" + "errors" "fmt" "io" "time" @@ -49,8 +50,18 @@ type MetadataFollowOption struct { // durably processed watermark while StartTsNs keeps tracking positions the // stream has merely seen. GetResumeTsNs func() int64 + // Resubscribe, when closed, drops the stream with ErrResubscribe so the + // caller's reconnect loop resubscribes from GetResumeTsNs and replays the + // events still pinning the processed watermark. Target-side job failures + // never surface on this stream, so without it a pinned event waits for an + // unrelated source-stream reconnect (or a restart) to be replayed. + Resubscribe <-chan struct{} } +// ErrResubscribe ends a metadata follow when the consumer asked for the +// stream to drop so the reconnect replays what the resume watermark pins. +var ErrResubscribe = errors.New("resubscribe to replay events behind the failed offset") + type ProcessMetadataFunc func(resp *filer_pb.SubscribeMetadataResponse) error func FollowMetadata(filerAddress ServerAddress, grpcDialOption grpc.DialOption, option *MetadataFollowOption, processEventFn ProcessMetadataFunc) error { @@ -98,6 +109,16 @@ func makeSubscribeMetadataFunc(option *MetadataFollowOption, processEventFn Proc return fmt.Errorf("subscribe: %w", err) } + if option.Resubscribe != nil { + go func() { + select { + case <-option.Resubscribe: + cancel() + case <-ctx.Done(): + } + }() + } + handleErr := func(resp *filer_pb.SubscribeMetadataResponse, err error) { switch option.EventErrorType { case TrivialOnError: @@ -105,12 +126,21 @@ func makeSubscribeMetadataFunc(option *MetadataFollowOption, processEventFn Proc case FatalOnError: glog.Fatalf("process %v: %v", resp, err) case RetryForeverOnError: - util.RetryUntil("followMetaUpdates", func() error { - return processEventFn(resp) - }, func(err error) bool { - glog.Errorf("process %v: %v", resp, err) - return true - }) + waitTime := time.Second + for ctx.Err() == nil { + if err := processEventFn(resp); err == nil { + break + } else { + glog.Errorf("process %v: %v", resp, err) + } + select { + case <-ctx.Done(): + case <-time.After(waitTime): + } + if waitTime < util.RetryWaitTime { + waitTime += waitTime / 2 + } + } case DontLogError: // pass default: @@ -200,6 +230,11 @@ func makeSubscribeMetadataFunc(option *MetadataFollowOption, processEventFn Proc return drainPendingRefs() } if listenErr != nil { + select { + case <-option.Resubscribe: + return ErrResubscribe + default: + } return listenErr } diff --git a/weed/pb/filer_pb_tail_test.go b/weed/pb/filer_pb_tail_test.go index 701a797c0..63909f4d4 100644 --- a/weed/pb/filer_pb_tail_test.go +++ b/weed/pb/filer_pb_tail_test.go @@ -2,6 +2,7 @@ package pb import ( "context" + "errors" "io" "sync" "sync/atomic" @@ -437,3 +438,52 @@ func TestFilerSyncReconnectReadsWatermarkEachSubscribe(t *testing.T) { t.Fatalf("expected subscribes at 100 then 300, got %v", sinceNs) } } + +// blockingFilerClient returns a stream whose Recv parks until the subscribe +// context ends, standing in for a quiet source stream during a sink outage. +type blockingFilerClient struct { + filer_pb.SeaweedFilerClient +} + +func (c *blockingFilerClient) SubscribeMetadata(ctx context.Context, in *filer_pb.SubscribeMetadataRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[filer_pb.SubscribeMetadataResponse], error) { + return &blockingSubscribeStream{ctx: ctx}, nil +} + +type blockingSubscribeStream struct { + grpc.ClientStream + ctx context.Context +} + +func (s *blockingSubscribeStream) Recv() (*filer_pb.SubscribeMetadataResponse, error) { + <-s.ctx.Done() + return nil, s.ctx.Err() +} + +// A consumer that pins an event asks for the stream to drop so the reconnect +// replays it; the follower must end with ErrResubscribe, not sit on Recv. +func TestFilerSyncResubscribeSignalEndsStream(t *testing.T) { + resubscribe := make(chan struct{}) + option := &MetadataFollowOption{ + ClientName: "syncFrom_A_To_B", + EventErrorType: DontLogError, + Resubscribe: resubscribe, + } + fn := makeSubscribeMetadataFunc(option, func(resp *filer_pb.SubscribeMetadataResponse) error { + return nil + }) + done := make(chan error, 1) + go func() { + done <- fn(&blockingFilerClient{}) + }() + + close(resubscribe) + + select { + case err := <-done: + if !errors.Is(err, ErrResubscribe) { + t.Fatalf("expected ErrResubscribe, got %v", err) + } + case <-time.After(5 * time.Second): + t.Fatal("follower did not end after the resubscribe signal") + } +} diff --git a/weed/replication/sink/filersink/filer_sink.go b/weed/replication/sink/filersink/filer_sink.go index 0f7a42e10..868cb01fc 100644 --- a/weed/replication/sink/filersink/filer_sink.go +++ b/weed/replication/sink/filersink/filer_sink.go @@ -172,7 +172,7 @@ func (fs *FilerSink) DeleteEntry(key string, isDirectory, deleteIncludeChunks bo err := filer_pb.Remove(context.Background(), fs, dir, name, deleteIncludeChunks, true, true, true, signatures) if err != nil { glog.V(0).Infof("delete entry %s: %v", key, err) - return fmt.Errorf("delete entry %s: %v", key, err) + return fmt.Errorf("delete entry %s: %w", key, err) } return nil } @@ -292,7 +292,7 @@ func (fs *FilerSink) CreateEntry(key string, entry *filer_pb.Entry, signatures [ glog.V(3).Infof("create: %v", request) if err := filer_pb.CreateEntry(context.Background(), client, request); err != nil { glog.V(0).Infof("create entry %s: %v", key, err) - return fmt.Errorf("create entry %s: %v", key, err) + return fmt.Errorf("create entry %s: %w", key, err) } return nil @@ -325,7 +325,7 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent }) if err != nil { - return false, fmt.Errorf("lookup %s: %v", key, err) + return false, fmt.Errorf("lookup %s: %w", key, err) } glog.V(4).Infof("oldEntry %+v, newEntry %+v, existingEntry: %+v", oldEntry, newEntry, existingEntry) @@ -360,7 +360,7 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent // source-side chunks resolve via source filer; sink volume IDs may collide. deletedChunks, newChunks, err := compareChunks(context.Background(), filer.LookupFn(fs.filerSource), oldEntry, newEntry) if err != nil { - return true, fmt.Errorf("replicate %s compare chunks error: %v", key, err) + return true, fmt.Errorf("replicate %s compare chunks error: %w", key, err) } // delete the chunks that are deleted from the source @@ -397,7 +397,7 @@ func (fs *FilerSink) UpdateEntry(key string, oldEntry *filer_pb.Entry, newParent } if _, err := client.UpdateEntry(context.Background(), request); err != nil { - return fmt.Errorf("update existingEntry %s: %v", key, err) + return fmt.Errorf("update existingEntry %s: %w", key, err) } return nil diff --git a/weed/util/retry_test.go b/weed/util/retry_test.go index 9ec2ac17e..64db256c5 100644 --- a/weed/util/retry_test.go +++ b/weed/util/retry_test.go @@ -26,8 +26,9 @@ func TestIsTransientError(t *testing.T) { &net.DNSError{Err: "operation timed out", IsTimeout: true}, io.ErrUnexpectedEOF, // transport teardown the peer reports as Canceled, not the caller's - // own context cancel + // own context cancel; also reachable through a caller's %w wrap status.Error(codes.Canceled, "grpc: the client connection is closing"), + fmt.Errorf("create entry /x: %w", status.Error(codes.Canceled, "grpc: the client connection is closing")), } for _, err := range transient { if !IsTransientError(err) {