mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-06 14:31:57 +02:00
* 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.
773 lines
32 KiB
Go
773 lines
32 KiB
Go
package command
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/filer"
|
|
"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/pb/remote_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/remote_storage"
|
|
"github.com/seaweedfs/seaweedfs/weed/replication/source"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
"google.golang.org/grpc"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/status"
|
|
"google.golang.org/protobuf/proto"
|
|
)
|
|
|
|
func followUpdatesAndUploadToRemote(option *RemoteSyncOptions, filerSource *source.FilerSource, mountedDir string) error {
|
|
|
|
// read filer remote storage mount mappings
|
|
_, _, remoteStorageMountLocation, remoteStorage, detectErr := filer.DetectMountInfo(option.grpcDialOption, pb.ServerAddress(*option.filerAddress), mountedDir)
|
|
if detectErr != nil {
|
|
return fmt.Errorf("read mount info: %w", detectErr)
|
|
}
|
|
|
|
eachEntryFunc, err := option.makeEventProcessor(remoteStorage, mountedDir, remoteStorageMountLocation, filerSource)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
lastOffsetTs := collectLastSyncOffset(option, option.grpcDialOption, pb.ServerAddress(*option.filerAddress), mountedDir, *option.timeAgo)
|
|
processor := NewMetadataProcessor(eachEntryFunc, 128, lastOffsetTs.UnixNano())
|
|
|
|
var lastLogTsNs = time.Now().UnixNano()
|
|
processEventFnWithOffset := pb.AddOffsetFunc(func(resp *filer_pb.SubscribeMetadataResponse) error {
|
|
processor.AddSyncJob(resp)
|
|
return nil
|
|
}, 3*time.Second, func(counter int64, lastTsNs int64) error {
|
|
offsetTsNs := processor.processedTsWatermark.Load()
|
|
if offsetTsNs == 0 {
|
|
return nil
|
|
}
|
|
// use processor.processedTsWatermark instead of the lastTsNs from the most recent job
|
|
now := time.Now().UnixNano()
|
|
glog.V(0).Infof("remote sync %s progressed to %v %0.2f/sec", *option.filerAddress, time.Unix(0, offsetTsNs), float64(counter)/(float64(now-lastLogTsNs)/1e9))
|
|
lastLogTsNs = now
|
|
return remote_storage.SetSyncOffset(option.grpcDialOption, pb.ServerAddress(*option.filerAddress), mountedDir, offsetTsNs)
|
|
})
|
|
|
|
option.clientEpoch++
|
|
|
|
prefix := mountedDir
|
|
if !strings.HasSuffix(prefix, "/") {
|
|
prefix = prefix + "/"
|
|
}
|
|
|
|
metadataFollowOption := &pb.MetadataFollowOption{
|
|
ClientName: "filer.remote.sync",
|
|
ClientId: option.clientId,
|
|
ClientEpoch: option.clientEpoch,
|
|
SelfSignature: 0,
|
|
PathPrefix: prefix,
|
|
AdditionalPathPrefixes: []string{filer.DirectoryEtcRemote},
|
|
DirectoriesToWatch: nil,
|
|
StartTsNs: lastOffsetTs.UnixNano(),
|
|
StopTsNs: 0,
|
|
EventErrorType: pb.RetryForeverOnError,
|
|
GetResumeTsNs: func() int64 {
|
|
return processor.processedTsWatermark.Load()
|
|
},
|
|
Resubscribe: processor.ResubscribeCh(),
|
|
}
|
|
|
|
return pb.FollowMetadata(pb.ServerAddress(*option.filerAddress), option.grpcDialOption, metadataFollowOption, processEventFnWithOffset)
|
|
}
|
|
|
|
func (option *RemoteSyncOptions) makeEventProcessor(remoteStorage *remote_pb.RemoteConf, mountedDir string, remoteStorageMountLocation *remote_pb.RemoteStorageLocation, filerSource *source.FilerSource) (pb.ProcessMetadataFunc, error) {
|
|
client, err := remote_storage.GetRemoteStorage(remoteStorage)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
handleEtcRemoteChanges := func(resp *filer_pb.SubscribeMetadataResponse) error {
|
|
message := resp.EventNotification
|
|
if metadataEventUpdatesDirectory(resp, filer.DirectoryEtcRemote) {
|
|
if message.NewEntry.Name == filer.REMOTE_STORAGE_MOUNT_FILE {
|
|
mappings, readErr := filer.UnmarshalRemoteStorageMappings(message.NewEntry.Content)
|
|
if readErr != nil {
|
|
return fmt.Errorf("unmarshal mappings: %w", readErr)
|
|
}
|
|
if remoteLoc, found := mappings.Mappings[mountedDir]; found {
|
|
if remoteStorageMountLocation.Bucket != remoteLoc.Bucket || remoteStorageMountLocation.Path != remoteLoc.Path {
|
|
glog.Fatalf("Unexpected mount changes %+v => %+v", remoteStorageMountLocation, remoteLoc)
|
|
}
|
|
} else {
|
|
glog.V(0).Infof("unmounted %s exiting ...", mountedDir)
|
|
os.Exit(0)
|
|
}
|
|
}
|
|
if message.NewEntry.Name == remoteStorage.Name+filer.REMOTE_STORAGE_CONF_SUFFIX {
|
|
conf := &remote_pb.RemoteConf{}
|
|
if err := proto.Unmarshal(message.NewEntry.Content, conf); err != nil {
|
|
return fmt.Errorf("unmarshal %s/%s: %v", filer.DirectoryEtcRemote, message.NewEntry.Name, err)
|
|
}
|
|
remoteStorage = conf
|
|
if newClient, err := remote_storage.GetRemoteStorage(remoteStorage); err == nil {
|
|
client = newClient
|
|
} else {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
if metadataEventRemovesFromDirectory(resp, filer.DirectoryEtcRemote) &&
|
|
message.OldEntry.Name == filer.REMOTE_STORAGE_MOUNT_FILE {
|
|
glog.V(0).Infof("unmounted %s exiting ...", mountedDir)
|
|
os.Exit(0)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
eachEntryFunc := func(resp *filer_pb.SubscribeMetadataResponse) error {
|
|
message := resp.EventNotification
|
|
sourceInEtcRemote, targetInEtcRemote := metadataEventDirectoryMembership(resp, filer.DirectoryEtcRemote)
|
|
if sourceInEtcRemote || targetInEtcRemote {
|
|
return handleEtcRemoteChanges(resp)
|
|
}
|
|
|
|
if filer_pb.IsEmpty(resp) {
|
|
return nil
|
|
}
|
|
if filer_pb.IsCreate(resp) {
|
|
if isMultipartUploadFile(message.NewParentPath, message.NewEntry.Name) {
|
|
return nil
|
|
}
|
|
// Propagate delete markers as deletions on the remote.
|
|
// Delete markers are zero-content version entries, so they
|
|
// would be filtered out by the HasData check below.
|
|
if isDeleteMarker(message.NewEntry) {
|
|
if newParent, newName, ok := rewriteVersionedSourcePath(message.NewParentPath, message.NewEntry.Name); ok {
|
|
dest := toRemoteStorageLocation(util.FullPath(mountedDir), util.NewFullPath(newParent, newName), remoteStorageMountLocation)
|
|
return syncDeleteMarker(client, option, message, dest)
|
|
}
|
|
return nil
|
|
}
|
|
if !filer.HasData(message.NewEntry) {
|
|
return nil
|
|
}
|
|
glog.V(2).Infof("create: %+v", resp)
|
|
if !shouldSendToRemote(message.NewEntry) {
|
|
glog.V(2).Infof("skipping creating: %+v", resp)
|
|
return nil
|
|
}
|
|
// Rewrite internal versioning paths to the original S3 key
|
|
// to prevent double-versioning when central also has versioning enabled
|
|
parentPath, entryName := message.NewParentPath, message.NewEntry.Name
|
|
isRewrittenVersion := false
|
|
if newParent, newName, ok := rewriteVersionedSourcePath(parentPath, entryName); ok {
|
|
glog.V(0).Infof("rewrite versioned path %s/%s -> %s/%s", parentPath, entryName, newParent, newName)
|
|
parentPath, entryName = newParent, newName
|
|
isRewrittenVersion = true
|
|
}
|
|
dest := toRemoteStorageLocation(util.FullPath(mountedDir), util.NewFullPath(parentPath, entryName), remoteStorageMountLocation)
|
|
if message.NewEntry.IsDirectory {
|
|
glog.V(0).Infof("mkdir %s", remote_storage.FormatLocation(dest))
|
|
return client.WriteDirectory(dest, remoteWriteEntry(message.NewEntry, *option.storageClass))
|
|
}
|
|
remoteEntry, writeErr := retriedWriteFile(client, filerSource, message.NewParentPath, remoteWriteEntry(message.NewEntry, *option.storageClass), dest)
|
|
if errors.Is(writeErr, errSuperseded) {
|
|
glog.Errorf("skipping %s: %v", remote_storage.FormatLocation(dest), writeErr)
|
|
return nil
|
|
}
|
|
if writeErr != nil {
|
|
return writeErr
|
|
}
|
|
// Skip updateLocalEntry for versioned rewrites: the logical
|
|
// object (e.g. file.xml) has no filer entry in versioned
|
|
// buckets, and stamping the internal v_* entry with a
|
|
// RemoteEntry for the logical key is semantically wrong.
|
|
// Replay is safe because S3 PutObject is idempotent.
|
|
if isRewrittenVersion {
|
|
return nil
|
|
}
|
|
return updateLocalEntry(option, message.NewParentPath, message.NewEntry, remoteEntry)
|
|
}
|
|
if filer_pb.IsDelete(resp) {
|
|
return processDeleteEvent(client, mountedDir, remoteStorageMountLocation, resp)
|
|
}
|
|
if message.OldEntry != nil && message.NewEntry != nil {
|
|
return processUpdateEvent(option, filerSource, *option.storageClass, client, mountedDir, remoteStorageMountLocation, resp)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
return eachEntryFunc, nil
|
|
}
|
|
|
|
// processDeleteEvent removes the remote object an entry mapped to. A remote
|
|
// object that is already absent counts as deleted: the entry was never
|
|
// uploaded (created and deleted faster than the sync ran, or a replay of an
|
|
// event the inline delete already handled), and GCS reports that case as
|
|
// ErrRemoteObjectNotFound where S3 and Azure answer an idempotent success.
|
|
// Returning the error would pin the sync offset on an event that has nothing
|
|
// left to do.
|
|
func processDeleteEvent(
|
|
client remote_storage.RemoteStorageClient,
|
|
mountedDir string,
|
|
remoteStorageMountLocation *remote_pb.RemoteStorageLocation,
|
|
resp *filer_pb.SubscribeMetadataResponse,
|
|
) error {
|
|
message := resp.EventNotification
|
|
// Skip deletion of internal version files; individual version
|
|
// deletes should not propagate to the remote object
|
|
if isVersionedPath(resp.Directory, message.OldEntry.Name, message.OldEntry.IsDirectory) {
|
|
glog.V(2).Infof("skipping delete of internal version path: %s/%s", resp.Directory, message.OldEntry.Name)
|
|
return nil
|
|
}
|
|
glog.V(2).Infof("delete: %+v", resp)
|
|
dest := toRemoteStorageLocation(util.FullPath(mountedDir), util.NewFullPath(resp.Directory, message.OldEntry.Name), remoteStorageMountLocation)
|
|
if message.OldEntry.IsDirectory {
|
|
glog.V(0).Infof("rmdir %s", remote_storage.FormatLocation(dest))
|
|
return client.RemoveDirectory(dest)
|
|
}
|
|
glog.V(0).Infof("delete %s", remote_storage.FormatLocation(dest))
|
|
return deleteRemoteFile(client, dest)
|
|
}
|
|
|
|
// deleteRemoteFile deletes the object and treats an already-absent object as
|
|
// deleted.
|
|
func deleteRemoteFile(client remote_storage.RemoteStorageClient, dest *remote_pb.RemoteStorageLocation) error {
|
|
if err := client.DeleteFile(dest); err != nil && !errors.Is(err, remote_storage.ErrRemoteObjectNotFound) {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func processUpdateEvent(
|
|
filerClient filer_pb.FilerClient,
|
|
filerSource filer_pb.FilerClient,
|
|
storageClass string,
|
|
client remote_storage.RemoteStorageClient,
|
|
mountedDir string,
|
|
remoteStorageMountLocation *remote_pb.RemoteStorageLocation,
|
|
resp *filer_pb.SubscribeMetadataResponse,
|
|
) error {
|
|
message := resp.EventNotification
|
|
if isMultipartUploadFile(message.NewParentPath, message.NewEntry.Name) {
|
|
return nil
|
|
}
|
|
if isVersionedPath(message.NewParentPath, message.NewEntry.Name, message.NewEntry.IsDirectory) {
|
|
glog.V(2).Infof("skipping update of internal version path: %s/%s", message.NewParentPath, message.NewEntry.Name)
|
|
return nil
|
|
}
|
|
oldDest := toRemoteStorageLocation(util.FullPath(mountedDir), util.NewFullPath(resp.Directory, message.OldEntry.Name), remoteStorageMountLocation)
|
|
dest := toRemoteStorageLocation(util.FullPath(mountedDir), util.NewFullPath(message.NewParentPath, message.NewEntry.Name), remoteStorageMountLocation)
|
|
if proto.Equal(oldDest, dest) && !shouldSendToRemote(message.NewEntry) {
|
|
glog.V(2).Infof("skipping updating: %+v", resp)
|
|
return nil
|
|
}
|
|
if message.NewEntry.IsDirectory {
|
|
return client.WriteDirectory(dest, remoteWriteEntry(message.NewEntry, storageClass))
|
|
}
|
|
if isMetadataOnlyUpdate(resp.Directory, message) {
|
|
remoteEntry, err := liveRemoteEntry(filerClient, message.NewParentPath, message.NewEntry)
|
|
if errors.Is(err, filer_pb.ErrNotFound) {
|
|
glog.V(2).Infof("skipping updating deleted entry: %+v", resp)
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if remoteEntry != nil {
|
|
glog.V(2).Infof("update meta: %+v", resp)
|
|
return client.UpdateFileMetadata(dest, message.OldEntry, remoteWriteEntry(message.NewEntry, storageClass))
|
|
}
|
|
glog.V(0).Infof("never replicated, uploading %s", remote_storage.FormatLocation(dest))
|
|
}
|
|
glog.V(2).Infof("update: %+v", resp)
|
|
if !proto.Equal(oldDest, dest) {
|
|
// A renamed entry that holds no local data now (the snapshot was
|
|
// remote-only, or remote.uncache ran since) has its content only as a
|
|
// remote object. Remote-only reads resolve by the entry's own path, so
|
|
// the rename completes only once the destination object holds it. The
|
|
// filer's current state decides, not the snapshot: an entry rewritten
|
|
// since the event has local data again, and the rewrite's bytes are
|
|
// what the destination must hold.
|
|
current, err := currentEntry(filerSource, message.NewParentPath, message.NewEntry.Name)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if isRemoteOnly(current) {
|
|
return completeRemoteOnlyRename(filerClient, client, message.NewParentPath, current, oldDest, dest, storageClass)
|
|
}
|
|
glog.V(0).Infof("delete %s", remote_storage.FormatLocation(oldDest))
|
|
if err := client.DeleteFile(oldDest); err != nil {
|
|
if isMultipartUploadFile(resp.Directory, message.OldEntry.Name) {
|
|
return nil
|
|
}
|
|
if !errors.Is(err, remote_storage.ErrRemoteObjectNotFound) {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
remoteEntry, writeErr := retriedWriteFile(client, filerSource, message.NewParentPath, remoteWriteEntry(message.NewEntry, storageClass), dest)
|
|
if errors.Is(writeErr, errSuperseded) {
|
|
glog.Errorf("skipping %s: %v", remote_storage.FormatLocation(dest), writeErr)
|
|
if !proto.Equal(oldDest, dest) {
|
|
return uploadCurrentEntry(filerClient, filerSource, client, message.NewParentPath, message.NewEntry.Name, oldDest, dest, storageClass)
|
|
}
|
|
return nil
|
|
}
|
|
if writeErr != nil {
|
|
return writeErr
|
|
}
|
|
return updateLocalEntry(filerClient, message.NewParentPath, message.NewEntry, remoteEntry)
|
|
}
|
|
|
|
// currentEntry returns what the filer holds at dir/name now, nil when the
|
|
// entry is gone.
|
|
func currentEntry(filerSource filer_pb.FilerClient, dir, name string) (*filer_pb.Entry, error) {
|
|
current, _, _, err := filer_pb.GetEntry(context.Background(), filerSource, util.NewFullPath(dir, name))
|
|
if errors.Is(err, filer_pb.ErrNotFound) {
|
|
return nil, nil
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return current, nil
|
|
}
|
|
|
|
// isRemoteOnly reports an entry whose content exists only as its remote
|
|
// object: no local data, a RemoteEntry with a size.
|
|
func isRemoteOnly(entry *filer_pb.Entry) bool {
|
|
return entry != nil && !filer.HasData(entry) && entry.IsInRemoteOnly()
|
|
}
|
|
|
|
// completeRemoteOnlyRename finishes a rename whose content exists only on the
|
|
// remote. When the destination object already holds what the entry's stamp
|
|
// describes (the rename ran before, its offset was not persisted, and
|
|
// remote.uncache followed), the old key goes. Otherwise the old object is
|
|
// copied over the destination, the entry is stamped with the copy, and the
|
|
// old key goes. With neither object present the content is lost: the event
|
|
// fails and holds the offset for recovery.
|
|
func completeRemoteOnlyRename(filerClient filer_pb.FilerClient, client remote_storage.RemoteStorageClient, dir string, current *filer_pb.Entry, oldDest, dest *remote_pb.RemoteStorageLocation, storageClass string) error {
|
|
if existing, err := client.StatFile(dest); err == nil {
|
|
if describes(current.RemoteEntry, existing) {
|
|
glog.V(0).Infof("%s already holds the renamed content", remote_storage.FormatLocation(dest))
|
|
return deleteRemoteFile(client, oldDest)
|
|
}
|
|
glog.V(0).Infof("%s holds an object the entry does not describe (size %d, want %d); replacing it from %s", remote_storage.FormatLocation(dest), existing.RemoteSize, current.RemoteEntry.RemoteSize, remote_storage.FormatLocation(oldDest))
|
|
} else if !errors.Is(err, remote_storage.ErrRemoteObjectNotFound) {
|
|
return err
|
|
}
|
|
stat, err := client.StatFile(oldDest)
|
|
if errors.Is(err, remote_storage.ErrRemoteObjectNotFound) {
|
|
return fmt.Errorf("%s: content is on neither %s nor %s", util.NewFullPath(dir, current.Name), remote_storage.FormatLocation(oldDest), remote_storage.FormatLocation(dest))
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
reader, err := openRemoteObject(client, oldDest, stat.RemoteSize)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer reader.Close()
|
|
glog.V(0).Infof("copy %s -> %s", remote_storage.FormatLocation(oldDest), remote_storage.FormatLocation(dest))
|
|
remoteEntry, err := client.WriteFile(dest, remoteWriteEntry(current, storageClass), reader)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := updateLocalEntry(filerClient, dir, current, remoteEntry); err != nil {
|
|
return err
|
|
}
|
|
return deleteRemoteFile(client, oldDest)
|
|
}
|
|
|
|
// describes reports whether the object a stat returned is the one the entry's
|
|
// stamp describes: same size, and the same ETag when both sides carry one (a
|
|
// copy written as one stream can legitimately carry a different ETag from a
|
|
// multipart original).
|
|
func describes(stamp, object *filer_pb.RemoteEntry) bool {
|
|
if stamp == nil || object == nil || stamp.RemoteSize != object.RemoteSize {
|
|
return false
|
|
}
|
|
return stamp.RemoteETag == "" || object.RemoteETag == "" || stamp.RemoteETag == object.RemoteETag
|
|
}
|
|
|
|
// openRemoteObject streams the object when the client can, and reads it whole
|
|
// otherwise.
|
|
func openRemoteObject(client remote_storage.RemoteStorageClient, loc *remote_pb.RemoteStorageLocation, size int64) (io.ReadCloser, error) {
|
|
if streamer, ok := client.(remote_storage.RemoteStorageStreamReader); ok {
|
|
return streamer.ReadFileAsStream(context.Background(), loc, 0, size)
|
|
}
|
|
data, err := client.ReadFile(loc, 0, size)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return io.NopCloser(bytes.NewReader(data)), nil
|
|
}
|
|
|
|
// uploadCurrentEntry uploads what the filer holds at dir/name now, when the
|
|
// event that superseded this rename will not. A rename whose snapshot is
|
|
// superseded has already deleted the old key; the rewrite behind it in the
|
|
// log uploads the destination itself unless shouldSendToRemote skips it on
|
|
// the inherited RemoteEntry, whose RemoteMtime can equal the rewrite's mtime
|
|
// within the same second. Only that case uploads here, so the content goes
|
|
// up once. A remote-only entry is finished the way completeRemoteOnlyRename
|
|
// finishes a rename: its stamp can already describe the destination (a sync
|
|
// plus remote.uncache in the meantime), and only then is it complete.
|
|
func uploadCurrentEntry(filerClient filer_pb.FilerClient, filerSource filer_pb.FilerClient, client remote_storage.RemoteStorageClient, dir, name string, oldDest, dest *remote_pb.RemoteStorageLocation, storageClass string) error {
|
|
current, _, _, err := filer_pb.GetEntry(context.Background(), filerSource, util.NewFullPath(dir, name))
|
|
if errors.Is(err, filer_pb.ErrNotFound) {
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if current.IsDirectory {
|
|
return nil
|
|
}
|
|
if isRemoteOnly(current) {
|
|
return completeRemoteOnlyRename(filerClient, client, dir, current, oldDest, dest, storageClass)
|
|
}
|
|
if shouldSendToRemote(current) {
|
|
glog.V(0).Infof("leaving %s to the rewrite that superseded the rename", remote_storage.FormatLocation(dest))
|
|
return nil
|
|
}
|
|
glog.V(0).Infof("uploading the current %s in place of the superseded rename", remote_storage.FormatLocation(dest))
|
|
remoteEntry, writeErr := retriedWriteFile(client, filerSource, dir, remoteWriteEntry(current, storageClass), dest)
|
|
if errors.Is(writeErr, errSuperseded) {
|
|
glog.Errorf("skipping %s: %v", remote_storage.FormatLocation(dest), writeErr)
|
|
return nil
|
|
}
|
|
if writeErr != nil {
|
|
return writeErr
|
|
}
|
|
return updateLocalEntry(filerClient, dir, current, remoteEntry)
|
|
}
|
|
|
|
// isSuperseded reports whether the filer has moved past the entry an event
|
|
// described: it is deleted, or it no longer references every chunk the event
|
|
// named. Those are the chunks the filer deletes when an entry is updated, so
|
|
// they are gone from the volume servers, or about to be, and no retry of the
|
|
// upload can succeed. Chunks are compared by file id, not content: a rewrite
|
|
// stores even identical bytes under new ids and drops the old ones. A replay
|
|
// from an earlier offset (-timeAgo) re-emits such events. Failing one holds the
|
|
// sync offset before it, so every restart of the subscription replays it into
|
|
// the same dead chunks, and progress on everything after it in the log is never
|
|
// persisted. Skipping is safe: the event that superseded this one follows in
|
|
// the log, and the delete removes the remote object or the rewrite uploads the
|
|
// current content. A lookup that fails for any other reason keeps the write
|
|
// failure, so the event is retried.
|
|
func isSuperseded(filerClient filer_pb.FilerClient, dir string, entry *filer_pb.Entry) bool {
|
|
current, _, _, err := filer_pb.GetEntry(context.Background(), filerClient, util.NewFullPath(dir, entry.Name))
|
|
if errors.Is(err, filer_pb.ErrNotFound) {
|
|
return true
|
|
}
|
|
if err != nil {
|
|
return false
|
|
}
|
|
if !filer.HasData(entry) && filer.HasData(current) {
|
|
// The event described an entry without data (remote-only, or empty);
|
|
// the filer has written to it since. Uploading the snapshot would put
|
|
// an empty object where the write belongs.
|
|
return true
|
|
}
|
|
if len(entry.Content) > 0 || len(current.Content) > 0 {
|
|
return !bytes.Equal(entry.Content, current.Content)
|
|
}
|
|
return len(filer.DoMinusChunks(entry.GetChunks(), current.GetChunks())) > 0
|
|
}
|
|
|
|
// errSuperseded marks an upload that failed for an entry the filer has since
|
|
// moved past (isSuperseded). The caller skips the event instead of failing it.
|
|
var errSuperseded = errors.New("deleted or rewritten since the event was logged")
|
|
|
|
// retriedWriteFile uploads the entry, retrying transient failures. The filer is
|
|
// asked whether the entry is superseded before the first attempt and after
|
|
// every failed one, and the upload stops at once when it is. Before: a backlog
|
|
// (a restart resuming from an old offset) carries every intermediate version
|
|
// of a hot file, and uploading each one in turn is wasted bandwidth that can
|
|
// trip the remote's per-object mutation rate limit; only the version the filer
|
|
// still holds is worth sending. After: a dead chunk reads as a transient
|
|
// "RequestError" from the SDK, and waiting out the backoff on it buys nothing.
|
|
func retriedWriteFile(client remote_storage.RemoteStorageClient, filerSource filer_pb.FilerClient, dir string, newEntry *filer_pb.Entry, dest *remote_pb.RemoteStorageLocation) (remoteEntry *filer_pb.RemoteEntry, err error) {
|
|
if isSuperseded(filerSource, dir, newEntry) {
|
|
return nil, fmt.Errorf("%s %w", util.NewFullPath(dir, newEntry.Name), errSuperseded)
|
|
}
|
|
err = util.RetryOnError("writeFile", func(err error) bool {
|
|
return !errors.Is(err, errSuperseded) && util.IsTransientError(err)
|
|
}, func() error {
|
|
reader := filer.NewFileReader(filerSource, newEntry)
|
|
glog.V(0).Infof("create %s", remote_storage.FormatLocation(dest))
|
|
var writeErr error
|
|
remoteEntry, writeErr = client.WriteFile(dest, newEntry, reader)
|
|
if writeErr != nil && isSuperseded(filerSource, dir, newEntry) {
|
|
return fmt.Errorf("%s %w: %w", util.NewFullPath(dir, newEntry.Name), errSuperseded, writeErr)
|
|
}
|
|
return writeErr
|
|
})
|
|
if err != nil && !errors.Is(err, errSuperseded) {
|
|
glog.Errorf("write to %s: %v", dest, err)
|
|
}
|
|
return
|
|
}
|
|
|
|
func collectLastSyncOffset(filerClient filer_pb.FilerClient, grpcDialOption grpc.DialOption, filerAddress pb.ServerAddress, mountedDir string, timeAgo time.Duration) time.Time {
|
|
// 1. specified by timeAgo
|
|
// 2. last offset timestamp for this directory
|
|
// 3. directory creation time
|
|
var lastOffsetTs time.Time
|
|
if timeAgo == 0 {
|
|
mountedDirEntry, _, _, err := filer_pb.GetEntry(context.Background(), filerClient, util.FullPath(mountedDir))
|
|
if err != nil {
|
|
glog.V(0).Infof("get mounted directory %s: %v", mountedDir, err)
|
|
return time.Now()
|
|
}
|
|
|
|
lastOffsetTsNs, err := remote_storage.GetSyncOffset(grpcDialOption, filerAddress, mountedDir)
|
|
if mountedDirEntry != nil {
|
|
if err == nil && mountedDirEntry.Attributes.Crtime < lastOffsetTsNs/1000000 {
|
|
lastOffsetTs = time.Unix(0, lastOffsetTsNs)
|
|
glog.V(0).Infof("resume from %v", lastOffsetTs)
|
|
} else {
|
|
lastOffsetTs = time.Unix(mountedDirEntry.Attributes.Crtime, 0)
|
|
}
|
|
} else {
|
|
lastOffsetTs = time.Now()
|
|
}
|
|
} 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
|
|
}
|
|
|
|
func toRemoteStorageLocation(mountDir, sourcePath util.FullPath, remoteMountLocation *remote_pb.RemoteStorageLocation) *remote_pb.RemoteStorageLocation {
|
|
source := string(sourcePath[len(mountDir):])
|
|
dest := util.FullPath(remoteMountLocation.Path).Child(source)
|
|
return &remote_pb.RemoteStorageLocation{
|
|
Name: remoteMountLocation.Name,
|
|
Bucket: remoteMountLocation.Bucket,
|
|
Path: string(dest),
|
|
}
|
|
}
|
|
|
|
// isMetadataOnlyUpdate reports whether an update leaves the entry at the same
|
|
// path with the same content, so the remote object needs at most its metadata
|
|
// rewritten -- provided it is already there, which liveRemoteEntry establishes.
|
|
func isMetadataOnlyUpdate(dir string, message *filer_pb.EventNotification) bool {
|
|
if dir != message.NewParentPath || message.OldEntry.Name != message.NewEntry.Name {
|
|
return false
|
|
}
|
|
return filer.IsSameData(message.OldEntry, message.NewEntry)
|
|
}
|
|
|
|
// liveRemoteEntry returns the RemoteEntry showing the entry's object is on the
|
|
// remote, or nil when it never got there. The event's own is not enough: a
|
|
// chmod right after a write is logged before the sync has uploaded the write
|
|
// and stamped the entry, and treating it as unreplicated would upload twice.
|
|
// Returns filer_pb.ErrNotFound when the entry has since been deleted.
|
|
func liveRemoteEntry(filerClient filer_pb.FilerClient, dir string, entry *filer_pb.Entry) (*filer_pb.RemoteEntry, error) {
|
|
if entry.RemoteEntry != nil {
|
|
return entry.RemoteEntry, nil
|
|
}
|
|
current, _, _, err := filer_pb.GetEntry(context.Background(), filerClient, util.NewFullPath(dir, entry.Name))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return current.RemoteEntry, nil
|
|
}
|
|
|
|
func shouldSendToRemote(entry *filer_pb.Entry) bool {
|
|
if entry.RemoteEntry == nil {
|
|
return true
|
|
}
|
|
if entry.RemoteEntry.RemoteMtime < entry.Attributes.Mtime {
|
|
return true
|
|
}
|
|
return false
|
|
}
|
|
|
|
// remoteWriteEntry returns the entry as remote storage should see it: the
|
|
// storage class attribute is dropped, or overridden by -storageClass. The
|
|
// event entry is left untouched so updateLocalEntry still compares the entry
|
|
// the filer stored.
|
|
func remoteWriteEntry(entry *filer_pb.Entry, storageClass string) *filer_pb.Entry {
|
|
clone := proto.Clone(entry).(*filer_pb.Entry)
|
|
if storageClass == "" {
|
|
delete(clone.Extended, s3_constants.AmzStorageClass)
|
|
} else {
|
|
if clone.Extended == nil {
|
|
clone.Extended = map[string][]byte{}
|
|
}
|
|
clone.Extended[s3_constants.AmzStorageClass] = []byte(storageClass)
|
|
}
|
|
return clone
|
|
}
|
|
|
|
// updateLocalEntry stamps the entry an event described with its RemoteEntry.
|
|
// The write carries IF_ENTRY_EQUAL over the event's entry: the filer deletes
|
|
// every stored chunk absent from an updated entry, so a snapshot older than
|
|
// the live entry (the file was rewritten while its upload was in flight, or
|
|
// the event is a replay) would delete the live chunks. A failed precondition
|
|
// means the filer moved past this event; the event that superseded it follows
|
|
// in the log and stamps the current entry, so the stale stamp is skipped the
|
|
// same way a superseded upload is.
|
|
func updateLocalEntry(filerClient filer_pb.FilerClient, dir string, entry *filer_pb.Entry, remoteEntry *filer_pb.RemoteEntry) error {
|
|
remoteEntry.LastLocalSyncTsNs = time.Now().UnixNano()
|
|
expected := proto.Clone(entry).(*filer_pb.Entry)
|
|
entry.RemoteEntry = remoteEntry
|
|
err := filerClient.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
|
|
_, err := client.UpdateEntry(context.Background(), &filer_pb.UpdateEntryRequest{
|
|
Directory: dir,
|
|
Entry: entry,
|
|
Condition: ifEntryEqual(expected),
|
|
})
|
|
return err
|
|
})
|
|
if isFailedPrecondition(err) {
|
|
glog.Errorf("skipping stale stamp of %s: %v", util.NewFullPath(dir, entry.Name), err)
|
|
return nil
|
|
}
|
|
if isEntryGone(err) {
|
|
glog.Errorf("skipping stamp of %s: deleted since the event was logged: %v", util.NewFullPath(dir, entry.Name), err)
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
// isEntryGone reports an UpdateEntry the filer refused because the entry no
|
|
// longer exists: a delete that followed the event superseded the stamp, and
|
|
// the delete's own event follows in the log. A filer with the typed answer
|
|
// returns codes.NotFound; an older filer returns a plain error of the form
|
|
// "not found <path>: <cause>", so the cause (the message's tail, never the
|
|
// path) is matched against filer_pb.ErrNotFound's text.
|
|
func isEntryGone(err error) bool {
|
|
if err == nil {
|
|
return false
|
|
}
|
|
if st, ok := status.FromError(err); ok && st.Code() == codes.NotFound {
|
|
return true
|
|
}
|
|
return strings.HasSuffix(strings.TrimSpace(err.Error()), filer_pb.ErrNotFound.Error())
|
|
}
|
|
|
|
// ifEntryEqual builds the precondition that the stored entry still equals the
|
|
// one the event described: chunk fids, inline content, and metadata alike.
|
|
func ifEntryEqual(entry *filer_pb.Entry) *filer_pb.WriteCondition {
|
|
return &filer_pb.WriteCondition{
|
|
Clauses: []*filer_pb.WriteCondition_Clause{{Kind: filer_pb.WriteCondition_IF_ENTRY_EQUAL, ExpectedEntry: entry}},
|
|
}
|
|
}
|
|
|
|
// isFailedPrecondition reports a write condition the filer refused, through
|
|
// any wrapping WithFilerClient added.
|
|
func isFailedPrecondition(err error) bool {
|
|
if err == nil {
|
|
return false
|
|
}
|
|
st, ok := status.FromError(err)
|
|
return ok && st.Code() == codes.FailedPrecondition
|
|
}
|
|
|
|
func isMultipartUploadFile(dir string, name string) bool {
|
|
return isMultipartUploadDir(dir) && strings.HasSuffix(name, ".part")
|
|
}
|
|
|
|
func isMultipartUploadDir(dir string) bool {
|
|
return strings.HasPrefix(dir, "/buckets/") &&
|
|
strings.Contains(dir, "/"+s3_constants.MultipartUploadsFolder+"/")
|
|
}
|
|
|
|
// isDeleteMarker returns true if the entry is an S3 delete marker
|
|
// (a zero-content version entry with ExtDeleteMarkerKey set to "true").
|
|
func isDeleteMarker(entry *filer_pb.Entry) bool {
|
|
if entry == nil || entry.Extended == nil {
|
|
return false
|
|
}
|
|
return string(entry.Extended[s3_constants.ExtDeleteMarkerKey]) == "true"
|
|
}
|
|
|
|
// syncDeleteMarker propagates a delete marker to the remote storage and
|
|
// persists a local sync marker so that replaying the same event is a no-op.
|
|
func syncDeleteMarker(
|
|
client remote_storage.RemoteStorageClient,
|
|
filerClient filer_pb.FilerClient,
|
|
message *filer_pb.EventNotification,
|
|
dest *remote_pb.RemoteStorageLocation,
|
|
) error {
|
|
glog.V(0).Infof("delete (marker) %s", remote_storage.FormatLocation(dest))
|
|
if err := deleteRemoteFile(client, dest); err != nil {
|
|
return err
|
|
}
|
|
return updateLocalEntry(filerClient, message.NewParentPath, message.NewEntry, &filer_pb.RemoteEntry{
|
|
StorageName: dest.Name,
|
|
RemoteMtime: message.NewEntry.Attributes.GetMtime(),
|
|
})
|
|
}
|
|
|
|
// isVersionedPath returns true if the dir/name refers to an internal
|
|
// versioning path (.versions directory or a version file inside it).
|
|
// These paths are SeaweedFS-internal and must not be synced to remote
|
|
// storage as-is, because the remote S3 endpoint may apply its own
|
|
// versioning, leading to double-versioned paths.
|
|
//
|
|
// For directories: matches only when the entry name ends with the
|
|
// VersionsFolder suffix (e.g. "file.xml.versions").
|
|
// For files: matches only when the parent directory ends with
|
|
// VersionsFolder and the file name has the "v_" prefix used by
|
|
// the internal version file naming convention.
|
|
func isVersionedPath(dir string, name string, isDir bool) bool {
|
|
if !strings.HasPrefix(dir, "/buckets/") {
|
|
return false
|
|
}
|
|
if isDir {
|
|
return strings.HasSuffix(name, s3_constants.VersionsFolder)
|
|
}
|
|
return strings.HasSuffix(dir, s3_constants.VersionsFolder) && strings.HasPrefix(name, "v_")
|
|
}
|
|
|
|
// rewriteVersionedSourcePath rewrites an internal versioning path to the
|
|
// original S3 object key. When a file is uploaded to a versioned bucket,
|
|
// SeaweedFS stores it internally as:
|
|
//
|
|
// /buckets/{bucket}/{key}.versions/v_{versionId}
|
|
//
|
|
// This function strips the ".versions/v_{versionId}" suffix and returns
|
|
// the original parent directory and object name, so the remote destination
|
|
// points to the logical S3 key rather than the internal version storage path.
|
|
//
|
|
// Returns (newDir, newName, true) if the path was rewritten, or
|
|
// (dir, name, false) if the path is not a versioned path.
|
|
func rewriteVersionedSourcePath(dir string, name string) (string, string, bool) {
|
|
if !strings.HasPrefix(dir, "/buckets/") {
|
|
return dir, name, false
|
|
}
|
|
if !strings.HasSuffix(dir, s3_constants.VersionsFolder) {
|
|
return dir, name, false
|
|
}
|
|
if !strings.HasPrefix(name, "v_") {
|
|
return dir, name, false
|
|
}
|
|
// dir = "/buckets/bucket/path/to/file.xml.versions"
|
|
// name = "v_abc123"
|
|
// Original object: dir without ".versions" suffix → "/buckets/bucket/path/to/file.xml"
|
|
originalObjectPath := dir[:len(dir)-len(s3_constants.VersionsFolder)]
|
|
lastSlash := strings.LastIndex(originalObjectPath, "/")
|
|
if lastSlash < 0 {
|
|
return dir, name, false
|
|
}
|
|
newDir := originalObjectPath[:lastSlash]
|
|
if lastSlash == 0 {
|
|
newDir = "/"
|
|
}
|
|
return newDir, originalObjectPath[lastSlash+1:], true
|
|
}
|