diff --git a/weed/filer/meta_aggregator.go b/weed/filer/meta_aggregator.go index 1635bb91a..6c1e0d1d1 100644 --- a/weed/filer/meta_aggregator.go +++ b/weed/filer/meta_aggregator.go @@ -19,6 +19,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/glog" "github.com/seaweedfs/seaweedfs/weed/pb" "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/stats" "github.com/seaweedfs/seaweedfs/weed/util/log_buffer" ) @@ -187,10 +188,7 @@ func (ma *MetaAggregator) doSubscribeToOneFiler(f *Filer, self pb.ServerAddress, var counter int64 var synced bool maybeReplicateMetadataChange = func(event *filer_pb.SubscribeMetadataResponse) { - if err := Replay(f.Store, event); err != nil { - glog.Errorf("failed to reply metadata change from %v: %v", peer, err) - return - } + replicateMetadataChange(f.Store, peer, event) counter++ if lastPersistTime.Add(time.Minute).Before(time.Now()) { if err := ma.updateOffset(f, peer, peerSignature, event.TsNs); err == nil { @@ -317,6 +315,25 @@ func (ma *MetaAggregator) doSubscribeToOneFiler(f *Filer, self pb.ServerAddress, return lastTsNs, err } +// replicateMetadataChange retries transient Replay failures with bounded +// backoff. A failure that outlives the retry budget is counted and logged, +// then skipped: blocking on an event that can never replay would stall every +// later event from this peer, which is worse than one entry staying stale. +func replicateMetadataChange(store FilerStore, peer pb.ServerAddress, event *filer_pb.SubscribeMetadataResponse) { + err := util.Retry("replicate metadata change from "+string(peer), func() error { + return Replay(store, event) + }) + if err == nil { + return + } + stats.FilerMetaAggregatorReplayFailures.WithLabelValues(string(peer)).Inc() + name := event.GetEventNotification().GetNewEntry().GetName() + if name == "" { + name = event.GetEventNotification().GetOldEntry().GetName() + } + glog.Errorf("giving up replicating metadata change from %s for %s/%s (ts=%d): %v", peer, event.Directory, name, event.TsNs, err) +} + // traversePeerMetadata does a full BFS traversal of a peer filer's metadata // and inserts all entries into the local store. This is used when a filer // connects to a peer for the first time and needs to bootstrap pre-existing data. diff --git a/weed/filer/meta_aggregator_replay_test.go b/weed/filer/meta_aggregator_replay_test.go new file mode 100644 index 000000000..43990b292 --- /dev/null +++ b/weed/filer/meta_aggregator_replay_test.go @@ -0,0 +1,111 @@ +package filer + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + dto "github.com/prometheus/client_model/go" + + "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb" + "github.com/seaweedfs/seaweedfs/weed/stats" + "github.com/seaweedfs/seaweedfs/weed/util" +) + +// flakyReplayStore fails InsertEntry for its first failUntil calls, then +// delegates to the embedded stubFilerStore. +type flakyReplayStore struct { + *stubFilerStore + countMu sync.Mutex + failUntil int + err error + attempts int +} + +func (s *flakyReplayStore) InsertEntry(ctx context.Context, entry *Entry) error { + s.countMu.Lock() + s.attempts++ + fail := s.attempts <= s.failUntil + s.countMu.Unlock() + if fail { + return s.err + } + return s.stubFilerStore.InsertEntry(ctx, entry) +} + +func (s *flakyReplayStore) attemptCount() int { + s.countMu.Lock() + defer s.countMu.Unlock() + return s.attempts +} + +func replayFailureCount(t *testing.T, peer pb.ServerAddress) float64 { + t.Helper() + var m dto.Metric + if err := stats.FilerMetaAggregatorReplayFailures.WithLabelValues(string(peer)).Write(&m); err != nil { + t.Fatalf("read counter: %v", err) + } + return m.GetCounter().GetValue() +} + +func quotaChangeEvent(bucket string, quota int64) *filer_pb.SubscribeMetadataResponse { + return &filer_pb.SubscribeMetadataResponse{ + Directory: "/buckets", + EventNotification: &filer_pb.EventNotification{ + NewEntry: &filer_pb.Entry{Name: bucket, IsDirectory: true, Quota: quota}, + }, + TsNs: time.Now().UnixNano(), + } +} + +// The old body logged a Replay error and returned, so one transient failure +// permanently diverged the entry from the peer. +func TestReplicateMetadataChangeRetriesTransientFailure(t *testing.T) { + // int64: untyped, the value defaults to int and overflows a 32-bit build + const wantQuota int64 = 131072 << 20 + + store := &flakyReplayStore{ + stubFilerStore: newStubFilerStore(), + failUntil: 1, + err: errors.New("i/o timeout talking to store"), + } + peer := pb.ServerAddress("peer-transient:1") + + replicateMetadataChange(store, peer, quotaChangeEvent("my-bucket", wantQuota)) + + inserted, err := store.FindEntry(context.Background(), util.NewFullPath("/buckets", "my-bucket")) + if err != nil { + t.Fatal("expected the entry to be inserted once the transient failure is retried past") + } + if inserted.Quota != wantQuota { + t.Fatalf("quota = %d, want %d", inserted.Quota, wantQuota) + } + if got := store.attemptCount(); got < 2 { + t.Fatalf("attempts = %d, want at least 2", got) + } +} + +// An event that can never replay must fail fast instead of blocking the +// subscribe stream, and must be counted so the divergence is not silent. +func TestReplicateMetadataChangeGivesUpLoudlyOnPermanentFailure(t *testing.T) { + store := &flakyReplayStore{ + stubFilerStore: newStubFilerStore(), + failUntil: 1 << 30, + err: errors.New("entry checksum mismatch"), + } + peer := pb.ServerAddress("peer-permanent:1") + + before := replayFailureCount(t, peer) + + replicateMetadataChange(store, peer, quotaChangeEvent("poison-bucket", 65536<<20)) + + if got := store.attemptCount(); got != 1 { + t.Fatalf("attempts = %d, want exactly 1: a non-transient error must fail fast, not retry", got) + } + if got := replayFailureCount(t, peer); got != before+1 { + t.Fatalf("FilerMetaAggregatorReplayFailures[%s] = %v, want %v", peer, got, before+1) + } +} diff --git a/weed/filer/meta_replay.go b/weed/filer/meta_replay.go index 51c4e6987..58542d305 100644 --- a/weed/filer/meta_replay.go +++ b/weed/filer/meta_replay.go @@ -8,6 +8,8 @@ import ( "github.com/seaweedfs/seaweedfs/weed/util" ) +// Replay applies the delete and the insert as separate, non-transactional +// store calls, so retrying it after a partial failure is not atomic. func Replay(filerStore FilerStore, resp *filer_pb.SubscribeMetadataResponse) error { message := resp.EventNotification var oldPath util.FullPath diff --git a/weed/stats/metrics.go b/weed/stats/metrics.go index b6161709e..a5c0646af 100644 --- a/weed/stats/metrics.go +++ b/weed/stats/metrics.go @@ -250,6 +250,14 @@ var ( Help: "Number of metadata subscribers currently parked waiting to read past a gap in the metadata log.", }, []string{"scope"}) + FilerMetaAggregatorReplayFailures = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Namespace: Namespace, + Subsystem: subsystemFiler, + Name: "meta_aggregator_replay_failures", + Help: "Number of peer metadata events skipped after replay retries were exhausted, leaving that entry diverged from the peer.", + }, []string{"peer"}) + // Sampled only on first creation, so counts track distinct objects. FilerObjectSizeBytesHistogram = prometheus.NewHistogram( prometheus.HistogramOpts{ @@ -899,6 +907,7 @@ func init() { Gather.MustRegister(FilerServerLastSendTsOfSubscribeGauge) Gather.MustRegister(FilerSubscribeGapStalledGauge) Gather.MustRegister(FilerSubscribeUnprovenGapCrossings) + Gather.MustRegister(FilerMetaAggregatorReplayFailures) Gather.MustRegister(FilerObjectSizeBytesHistogram) Gather.MustRegister(collectors.NewGoCollector()) Gather.MustRegister(collectors.NewProcessCollector(collectors.ProcessCollectorOpts{}))