diff --git a/weed/filer/meta_aggregator.go b/weed/filer/meta_aggregator.go index 2a922f438..0e2435243 100644 --- a/weed/filer/meta_aggregator.go +++ b/weed/filer/meta_aggregator.go @@ -2,6 +2,7 @@ package filer import ( "context" + "errors" "fmt" "io" "strings" @@ -139,16 +140,37 @@ func (ma *MetaAggregator) doSubscribeToOneFiler(f *Filer, self pb.ServerAddress, if peerSignature != f.Signature { if prevTsNs, err := ma.readOffset(f, peer, peerSignature); err == nil { lastTsNs = prevTsNs - defer func(prevTsNs int64) { - if lastTsNs != prevTsNs && lastTsNs != lastPersistTime.UnixNano() { - if err := ma.updateOffset(f, peer, peerSignature, lastTsNs); err == nil { - glog.V(0).Infof("last sync time with %s at %v (%d)", peer, time.Unix(0, lastTsNs), lastTsNs) - } else { - glog.Errorf("failed to save last sync time with %s at %v (%d)", peer, time.Unix(0, lastTsNs), lastTsNs) - } - } - }(prevTsNs) + } else if errors.Is(err, ErrKvNotFound) { + // No stored offset — this is the first time connecting to this peer. + // Traverse the peer's full metadata tree so we get pre-existing data. + // Record time before traversal and subtract a safety margin to + // account for clock skew between this filer and the peer. Any + // duplicate events replayed during the overlap are harmless since + // Replay does upserts. We use wall-clock time (same domain as the + // metadata stream TsNs) rather than entry Mtime which is a + // different concept and can be set to arbitrary values. + preTraverseTime := time.Now() + glog.V(0).Infof("no previous offset for peer %s, starting full metadata sync", peer) + if traverseErr := ma.traversePeerMetadata(f, peer); traverseErr != nil { + return lastTsNs, fmt.Errorf("initial metadata sync from %s: %v", peer, traverseErr) + } + lastTsNs = preTraverseTime.Add(-time.Minute).UnixNano() + if err := ma.updateOffset(f, peer, peerSignature, lastTsNs); err != nil { + return lastTsNs, fmt.Errorf("save bootstrap offset for peer %s: %w", peer, err) + } + glog.V(0).Infof("completed full metadata sync from peer %s, will stream changes from %v", peer, time.Unix(0, lastTsNs)) + } else { + return lastTsNs, fmt.Errorf("read offset for peer %s: %w", peer, err) } + defer func(prevTsNs int64) { + if lastTsNs != prevTsNs && lastTsNs != lastPersistTime.UnixNano() { + if err := ma.updateOffset(f, peer, peerSignature, lastTsNs); err == nil { + glog.V(0).Infof("last sync time with %s at %v (%d)", peer, time.Unix(0, lastTsNs), lastTsNs) + } else { + glog.Errorf("failed to save last sync time with %s at %v (%d)", peer, time.Unix(0, lastTsNs), lastTsNs) + } + } + }(lastTsNs) glog.V(0).Infof("follow peer: %v, last %v (%d)", peer, time.Unix(0, lastTsNs), lastTsNs) var counter int64 @@ -279,6 +301,59 @@ func (ma *MetaAggregator) doSubscribeToOneFiler(f *Filer, self pb.ServerAddress, return lastTsNs, 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. +func (ma *MetaAggregator) traversePeerMetadata(f *Filer, peer pb.ServerAddress) error { + return pb.WithFilerClient(true, 0, peer, ma.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + stream, err := client.TraverseBfsMetadata(ctx, &filer_pb.TraverseBfsMetadataRequest{ + Directory: "/", + ExcludedPrefixes: []string{SystemLogDir}, + }) + if err != nil { + return fmt.Errorf("traverse bfs metadata: %w", err) + } + var count int64 + for { + resp, recvErr := stream.Recv() + if recvErr == io.EOF { + break + } + if recvErr != nil { + return fmt.Errorf("traverse bfs metadata recv: %w", recvErr) + } + if resp.Entry == nil { + continue + } + fullpath := util.Join(resp.Directory, resp.Entry.Name) + entry := FromPbEntry(resp.Directory, resp.Entry) + if insertErr := f.Store.InsertEntry(context.Background(), entry); insertErr != nil { + // Entry may already exist (root dir, or partial previous bootstrap). + existing, findErr := f.Store.FindEntry(context.Background(), entry.FullPath) + if findErr != nil { + return fmt.Errorf("insert entry %s: %w", fullpath, insertErr) + } + // Only overwrite if the peer's entry is newer. + if entry.Attr.Mtime.After(existing.Attr.Mtime) { + if updateErr := f.Store.UpdateEntry(context.Background(), entry); updateErr != nil { + return fmt.Errorf("update entry %s: %w", fullpath, updateErr) + } + } else { + glog.V(1).Infof("skip older peer entry %s (peer mtime %v <= local mtime %v)", fullpath, entry.Attr.Mtime, existing.Attr.Mtime) + } + } + count++ + if count%10000 == 0 { + glog.V(0).Infof("synced %d entries from peer %s", count, peer) + } + } + glog.V(0).Infof("synced %d entries total from peer %s", count, peer) + return nil + }) +} + func (ma *MetaAggregator) readFilerStoreSignature(peer pb.ServerAddress) (sig int32, err error) { err = pb.WithFilerClient(false, 0, peer, ma.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error { resp, err := client.GetFilerConfiguration(context.Background(), &filer_pb.GetFilerConfigurationRequest{}) @@ -308,7 +383,7 @@ func (ma *MetaAggregator) readOffset(f *Filer, peer pb.ServerAddress, peerSignat value, err := f.Store.KvGet(context.Background(), key) if err != nil { - return 0, fmt.Errorf("readOffset %s : %v", peer, err) + return 0, fmt.Errorf("readOffset %s : %w", peer, err) } lastTsNs = int64(util.BytesToUint64(value))