mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-05 22:12:04 +02:00
filer.sync: resubscribe the metadata stream when a failure pins the offset (#11581)
* filer sink: keep the gRPC status inside wrapped errors
%v stringifies the status, so a peer teardown reported as Canceled ("the
client connection is closing") reached IsTransientError as plain text and
matched nothing: the sync job failed on the first attempt and pinned the
offset. %w keeps the status reachable, so the retry runs on a fresh
connection once the target is back.
* pb: let a consumer drop the metadata stream to force a resubscribe
A MetadataProcessor job that exhausts its retries pins the processed
watermark so the event replays on the next subscribe — but nothing on the
source stream notices a target-side failure, so the replay waited for an
unrelated reconnect or a restart. The new Resubscribe channel cancels the
stream's context; the Recv loop answers it with ErrResubscribe so the
caller's retry loop resubscribes from GetResumeTsNs and replays the pinned
events in order.
* pb: stop the event retry loop once the stream context is done
RetryUntil ignores context, so a subscriber parked on a failing offset
write would keep retrying past a resubscribe signal until the sink came
back. Stop retrying when the stream is being dropped so the resubscribe
takes effect promptly.
* filer.sync: signal resubscribe when a job failure pins the offset
A job that exhausts its in-job retries leaves the event pinned behind oldestFailedTsNs, replayable only on a reconnect. Closing resubscribeCh on the first recorded failure lets the metadata follower drop the stream so the reconnect replays the pinned events instead of waiting for a process restart (#11572).
* filer.sync: wire the resubscribe signal into the follow options
filer.sync, filer.remote.sync, and the remote gateway bucket sync all run their subscription inside an outer retry loop, so ErrResubscribe resurfaces as a resubscribe from the persisted watermark.
* filer.sync: wait for in-flight jobs before signaling resubscribe
* remote sync: never resume past the saved offset when -timeAgo is set
* filer.sync: drop events that arrive after the drain signals resubscribe
* pb: interrupt the event retry backoff when the stream context ends
* filer.sync: stop admitting once a failure pins, and count jobs per timestamp
A pinned watermark only released once the processor went fully quiet, so a busy stream could starve the resubscribe — the failed event would wait for an unrelated reconnect anyway, the wait this mechanism exists to remove. The processor now latches stopped when a job fails: admission drops new events (they replay from the pinned watermark after the reconnect), a broadcast releases blocked waiters, and the resubscribe signals as soon as the jobs already in flight drain. A redelivery of an event still in the failure ledger may still run so its success shrinks the replay, but nothing starts once the signal has fired, or it would race the replay it asked for.
Dropped events no longer inflate the received counters — an event counts only once admitted, and the replay's own admission counts it.
While here: activeJobs keyed by TsNs collapsed events sharing a timestamp, so one completion could empty the map while a same-ts sibling was still running — letting the drain gate and the watermark outrun it. Jobs are now counted per timestamp, and the drain and lazy heap cleanup go through the counts.
This commit is contained in:
1 parent
10b0f2b8ad
commit
f4ef37e752
9 files changed
+458
-99
No files matched your search
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
}()
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in new issue
Block a user