filer: resume metadata subscriber from processed watermark on reconnect (#11574)

* filer: resume metadata subscriber from processed watermark on reconnect

* filer: take the reconnect position from GetResumeTsNs verbatim

The callback is the subscriber's durable resume point; falling back to
StartTsNs when it returns zero can resume from a cursor the log-chunk
reader advanced past still-pending work.

* filer: advance the stream cursor once a retried event recovers

RetryForeverOnError resolves the failure inside handleErr, so returning
without moving StartTsNs replays work the event already did when the
stream reconnects before the next one arrives.

* filer: let filtered-progress markers move the processed watermark

A marker means the source examined everything up to its timestamp and
skipped what did not match the subscription. With a resume callback the
marker now reaches the consumer, and AddSyncJob advances the watermark
to it once every earlier job finished and no failure pins the offset.
Idle filtered stretches no longer rescan on every reconnect, while the
guards keep the watermark behind pending or failed work.

* filer: unpin the watermark once a failed event completes

oldestFailedTsNs was only ever set, so a failure that a replay later
fixed still held the resume offset, and every reconnect re-read the
same backlog. Track outstanding failures in a set and recompute the
pin when the failed event's job finally succeeds.

* filer.remote.gateway: resume bucket sync from the processed watermark

The bucket-sync subscriber runs the same MetadataProcessor queue as
filer.remote.sync; give it the same GetResumeTsNs callback so a
reconnect resumes from durably processed work, not the last seen event.

* util: treat a peer-sent gRPC Canceled as transient

A peer tearing down its end of the transport reports codes.Canceled
("the client connection is closing"), which IsTransientError used to
reject: the sync job then failed on the first try and held the offset
until a restart. Caller's own cancels are still excluded up front by
errors.Is(err, context.Canceled), so only teardown-style statuses take
the new branch.

* fix: preserve filtered progress and distinguish caller cancellation

* filer: bound the failed-event ledger past a persistent outage

A destination rejecting every event grew failedTs by one entry per source
event for the life of the processor. Past maxFailedSyncEvents the set now
collapses to a sticky pin at the smallest failure seen, so the watermark
still replays from the oldest failure while memory stays bounded; a
restart re-derives the exact set.

Also keep a resume-callback consumer's chunk-ref replay filter at the
subscribe-time position instead of option.StartTsNs, so a resubscribe does
not filter out events whose async processing is still pending.

* filer: key the failed-event ledger by event, not just timestamp

A success for one event cleared the pin recorded for a different event
that shared its TsNs, letting the watermark pass an unresolved failure.
The ledger now keys on the event's path identity, so recovery unblocks
only the event that actually failed.

---------

Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com>
Co-authored-by: Chris Lu <chris.lu@gmail.com>
This commit is contained in:
authored and GitHub committed 2026-10-03 20:57:50 +08:00
1 parent 8c1ebbee32
commit 562afa8ec9
9 files changed
+539 -43

No files matched your search

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