mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-20 13:30:46 +02:00
A versioned write's only contended mutation is the .versions directory's latest pointer; the version file itself goes to a unique <object>/.versions/<versionId> path. Add a FinalizeVersionedWrite filer op that, under one exclusive lock on the object key, evaluates the precondition against the current latest, stamps the previous latest noncurrent (before the pointer flip so the lifecycle router observes it), then merges the latest pointer / cached metadata into the .versions entry. Key names are passed in, so the filer carries no S3 semantics. Routing and the lock are keyed on the object (objectWriteOwner + lock_key), the same key normal and suspended writes use, so all writes to one object resolve the same owner and serialize on the same lock regardless of versioning state — a versioned and a non-versioned write to the same object can't race on different owners during a versioning-state change. The gateway routes a versioned PutObject's finalize to that owner and tells putToFiler the version path is unique, so it skips the object write lock and the gateway precondition (the op does both atomically). When the owner is unknown or the condition can't reduce to one primitive, it stays on the lock path; on op error it returns InternalError. Versioned COPY, delete markers, suspended versioning, and multipart completion still use the lock and adopt the same op as follow-ups.
681 lines
24 KiB
Go
681 lines
24 KiB
Go
package weed_server
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/filer"
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/operation"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
|
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/status"
|
|
)
|
|
|
|
func (fs *FilerServer) LookupDirectoryEntry(ctx context.Context, req *filer_pb.LookupDirectoryEntryRequest) (*filer_pb.LookupDirectoryEntryResponse, error) {
|
|
|
|
glog.V(4).InfofCtx(ctx, "LookupDirectoryEntry %s", filepath.Join(req.Directory, req.Name))
|
|
|
|
entry, err := fs.filer.FindEntry(ctx, util.JoinPath(req.Directory, req.Name))
|
|
if err == filer_pb.ErrNotFound {
|
|
return &filer_pb.LookupDirectoryEntryResponse{}, err
|
|
}
|
|
if err != nil {
|
|
glog.V(3).InfofCtx(ctx, "LookupDirectoryEntry %s: %+v, ", filepath.Join(req.Directory, req.Name), err)
|
|
return nil, err
|
|
}
|
|
|
|
return &filer_pb.LookupDirectoryEntryResponse{
|
|
Entry: entry.ToProtoEntry(),
|
|
}, nil
|
|
}
|
|
|
|
func (fs *FilerServer) ListEntries(req *filer_pb.ListEntriesRequest, stream filer_pb.SeaweedFiler_ListEntriesServer) (err error) {
|
|
|
|
glog.V(4).Infof("ListEntries %v", req)
|
|
|
|
limit := int(req.Limit)
|
|
if limit == 0 {
|
|
limit = fs.option.DirListingLimit
|
|
}
|
|
|
|
paginationLimit := filer.PaginationSize
|
|
if limit < paginationLimit {
|
|
paginationLimit = limit
|
|
}
|
|
|
|
lastFileName := req.StartFromFileName
|
|
includeLastFile := req.InclusiveStartFrom
|
|
snapshotTsNs := req.SnapshotTsNs
|
|
if snapshotTsNs == 0 {
|
|
snapshotTsNs = time.Now().UnixNano()
|
|
}
|
|
sentSnapshot := false
|
|
var listErr error
|
|
for limit > 0 {
|
|
var hasEntries bool
|
|
lastFileName, listErr = fs.filer.StreamListDirectoryEntries(stream.Context(), util.FullPath(req.Directory), lastFileName, includeLastFile, int64(paginationLimit), req.Prefix, "", "", func(entry *filer.Entry) (bool, error) {
|
|
hasEntries = true
|
|
resp := &filer_pb.ListEntriesResponse{
|
|
Entry: entry.ToProtoEntry(),
|
|
}
|
|
if !sentSnapshot {
|
|
resp.SnapshotTsNs = snapshotTsNs
|
|
sentSnapshot = true
|
|
}
|
|
if err = stream.Send(resp); err != nil {
|
|
return false, err
|
|
}
|
|
|
|
limit--
|
|
if limit == 0 {
|
|
return false, nil
|
|
}
|
|
return true, nil
|
|
})
|
|
|
|
if listErr != nil {
|
|
return listErr
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if !hasEntries {
|
|
break
|
|
}
|
|
|
|
includeLastFile = false
|
|
|
|
}
|
|
|
|
// For empty directories we intentionally do NOT send a snapshot-only
|
|
// response (Entry == nil). Many consumers (Java FilerClient, S3 listing,
|
|
// etc.) treat any received response as an entry. The Go client-side
|
|
// DoSeaweedListWithSnapshot generates a client-side cutoff when the
|
|
// server sends no snapshot, so snapshot consistency is preserved
|
|
// without a server-side send.
|
|
|
|
return nil
|
|
}
|
|
|
|
func (fs *FilerServer) LookupVolume(ctx context.Context, req *filer_pb.LookupVolumeRequest) (*filer_pb.LookupVolumeResponse, error) {
|
|
|
|
resp := &filer_pb.LookupVolumeResponse{
|
|
LocationsMap: make(map[string]*filer_pb.Locations),
|
|
}
|
|
|
|
// Use master client's lookup with fallback - it handles cache and master query
|
|
vidLocations, err := fs.filer.MasterClient.LookupVolumeIdsWithFallback(ctx, req.VolumeIds)
|
|
|
|
// Convert wdclient.Location to filer_pb.Location
|
|
// Return partial results even if there was an error
|
|
for vidString, locations := range vidLocations {
|
|
resp.LocationsMap[vidString] = &filer_pb.Locations{
|
|
Locations: wdclientLocationsToPb(locations),
|
|
}
|
|
}
|
|
|
|
return resp, err
|
|
}
|
|
|
|
func wdclientLocationsToPb(locations []wdclient.Location) []*filer_pb.Location {
|
|
locs := make([]*filer_pb.Location, 0, len(locations))
|
|
for _, loc := range locations {
|
|
locs = append(locs, &filer_pb.Location{
|
|
Url: loc.Url,
|
|
PublicUrl: loc.PublicUrl,
|
|
GrpcPort: uint32(loc.GrpcPort),
|
|
DataCenter: loc.DataCenter,
|
|
})
|
|
}
|
|
return locs
|
|
}
|
|
|
|
func (fs *FilerServer) lookupFileId(ctx context.Context, fileId string) (targetUrls []string, err error) {
|
|
fid, err := needle.ParseFileIdFromString(fileId)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
locations, found := fs.filer.MasterClient.GetLocations(uint32(fid.VolumeId))
|
|
if !found || len(locations) == 0 {
|
|
return nil, fmt.Errorf("not found volume %d in %s", fid.VolumeId, fileId)
|
|
}
|
|
for _, loc := range locations {
|
|
targetUrls = append(targetUrls, fmt.Sprintf("http://%s/%s", loc.Url, fileId))
|
|
}
|
|
return
|
|
}
|
|
|
|
func (fs *FilerServer) CreateEntry(ctx context.Context, req *filer_pb.CreateEntryRequest) (resp *filer_pb.CreateEntryResponse, err error) {
|
|
|
|
glog.V(4).InfofCtx(ctx, "CreateEntry %v/%v", req.Directory, req.Entry.Name)
|
|
if len(req.Entry.HardLinkId) > 0 {
|
|
glog.V(4).InfofCtx(ctx, "CreateEntry %s/%s with HardLinkId %x counter=%d", req.Directory, req.Entry.Name, req.Entry.HardLinkId, req.Entry.HardLinkCounter)
|
|
}
|
|
|
|
resp = &filer_pb.CreateEntryResponse{}
|
|
|
|
chunks, garbage, err2 := fs.cleanupChunks(ctx, util.Join(req.Directory, req.Entry.Name), nil, req.Entry)
|
|
if err2 != nil {
|
|
return &filer_pb.CreateEntryResponse{}, fmt.Errorf("CreateEntry cleanupChunks %s %s: %v", req.Directory, req.Entry.Name, err2)
|
|
}
|
|
|
|
so, err := fs.detectStorageOption(ctx, string(util.NewFullPath(req.Directory, req.Entry.Name)), "", "", 0, "", "", "", "")
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
newEntry := filer.FromPbEntry(req.Directory, req.Entry)
|
|
newEntry.Chunks = chunks
|
|
// Don't apply TTL to remote entries - they're managed by remote storage
|
|
if newEntry.Remote == nil {
|
|
if newEntry.TtlSec == 0 {
|
|
newEntry.TtlSec = so.TtlSeconds
|
|
}
|
|
} else {
|
|
newEntry.TtlSec = 0
|
|
}
|
|
|
|
// Serialize concurrent mutations to the same path on this filer so the
|
|
// read (existence/condition) and the write are atomic. Callers route a
|
|
// key's writes to this owner filer, making this local lock sufficient.
|
|
fullpath := util.NewFullPath(req.Directory, req.Entry.Name)
|
|
pathLock := fs.entryLockTable.AcquireLock("CreateEntry", fullpath, util.ExclusiveLock)
|
|
defer fs.entryLockTable.ReleaseLock(fullpath, pathLock)
|
|
|
|
// Evaluate the optional precondition against the current entry while the
|
|
// path lock is held, so the check and the write are atomic on this filer.
|
|
if req.Condition != nil && req.Condition.Kind != filer_pb.WriteCondition_NONE {
|
|
current, findErr := fs.filer.FindEntry(ctx, fullpath)
|
|
if findErr != nil && findErr != filer_pb.ErrNotFound {
|
|
return &filer_pb.CreateEntryResponse{}, fmt.Errorf("CreateEntry condition check %s: %w", fullpath, findErr)
|
|
}
|
|
if findErr == filer_pb.ErrNotFound {
|
|
current = nil
|
|
}
|
|
if !writeConditionSatisfied(req.Condition, current) {
|
|
glog.V(3).InfofCtx(ctx, "CreateEntry %s: precondition %v failed", fullpath, req.Condition.Kind)
|
|
return &filer_pb.CreateEntryResponse{
|
|
Error: "precondition failed",
|
|
ErrorCode: filer_pb.FilerError_PRECONDITION_FAILED,
|
|
}, nil
|
|
}
|
|
}
|
|
|
|
ctx, eventSink := filer.WithMetadataEventSink(ctx)
|
|
createErr := fs.filer.CreateEntry(ctx, newEntry, req.OExcl, req.IsFromOtherCluster, req.Signatures, req.SkipCheckParentDirectory, so.MaxFileNameLength)
|
|
|
|
if createErr == nil {
|
|
fs.filer.DeleteChunksNotRecursive(garbage)
|
|
resp.MetadataEvent = eventSink.Last()
|
|
} else {
|
|
glog.V(3).InfofCtx(ctx, "CreateEntry %s: %v", filepath.Join(req.Directory, req.Entry.Name), createErr)
|
|
resp.Error = createErr.Error()
|
|
switch {
|
|
case errors.Is(createErr, filer_pb.ErrEntryNameTooLong):
|
|
resp.ErrorCode = filer_pb.FilerError_ENTRY_NAME_TOO_LONG
|
|
case errors.Is(createErr, filer_pb.ErrParentIsFile):
|
|
resp.ErrorCode = filer_pb.FilerError_PARENT_IS_FILE
|
|
case errors.Is(createErr, filer_pb.ErrExistingIsDirectory):
|
|
resp.ErrorCode = filer_pb.FilerError_EXISTING_IS_DIRECTORY
|
|
case errors.Is(createErr, filer_pb.ErrExistingIsFile):
|
|
resp.ErrorCode = filer_pb.FilerError_EXISTING_IS_FILE
|
|
case errors.Is(createErr, filer_pb.ErrEntryAlreadyExists):
|
|
resp.ErrorCode = filer_pb.FilerError_ENTRY_ALREADY_EXISTS
|
|
}
|
|
}
|
|
|
|
return
|
|
}
|
|
|
|
// writeConditionSatisfied reports whether the precondition holds against the
|
|
// current entry (nil if absent). The caller reduces request semantics to one
|
|
// primitive; each kind is a single comparison evaluated under the path lock.
|
|
func writeConditionSatisfied(cond *filer_pb.WriteCondition, current *filer.Entry) bool {
|
|
exists := current != nil
|
|
switch cond.Kind {
|
|
case filer_pb.WriteCondition_IF_NOT_EXISTS:
|
|
return !exists
|
|
case filer_pb.WriteCondition_IF_EXISTS:
|
|
return exists
|
|
case filer_pb.WriteCondition_IF_ETAG_MATCH:
|
|
return exists && storedEntryETag(current) == normalizeETag(cond.Etag)
|
|
case filer_pb.WriteCondition_IF_ETAG_NOT_MATCH:
|
|
return !exists || storedEntryETag(current) != normalizeETag(cond.Etag)
|
|
case filer_pb.WriteCondition_IF_UNMODIFIED_SINCE:
|
|
return !exists || current.Attr.Mtime.Unix() <= cond.UnixTime
|
|
case filer_pb.WriteCondition_IF_MODIFIED_SINCE:
|
|
return !exists || current.Attr.Mtime.Unix() > cond.UnixTime
|
|
default:
|
|
return true
|
|
}
|
|
}
|
|
|
|
// storedEntryETag mirrors the S3 gateway's ETag precedence (the stored
|
|
// Seaweed ETag extended attribute, then the chunk/Md5 fallback) so conditional
|
|
// comparisons match what the gateway computes, without coupling the filer to
|
|
// S3 request handling.
|
|
func storedEntryETag(entry *filer.Entry) string {
|
|
if v, ok := entry.Extended[s3_constants.ExtETagKey]; ok && len(v) > 0 {
|
|
return normalizeETag(string(v))
|
|
}
|
|
return normalizeETag(filer.ETagEntry(entry))
|
|
}
|
|
|
|
func normalizeETag(etag string) string {
|
|
return strings.Trim(etag, `"`)
|
|
}
|
|
|
|
// versionedConditionSatisfied evaluates a precondition against the current
|
|
// latest version, derived from the .versions directory entry's cached metadata
|
|
// (whether a latest exists and isn't a delete marker, and its ETag). It mirrors
|
|
// writeConditionSatisfied but reads the cached pointer fields instead of a plain
|
|
// object entry.
|
|
func versionedConditionSatisfied(req *filer_pb.FinalizeVersionedWriteRequest, versionsEntry *filer.Entry) bool {
|
|
exists := false
|
|
etag := ""
|
|
if versionsEntry != nil && versionsEntry.Extended != nil {
|
|
latestFile := string(versionsEntry.Extended[req.PriorLatestKey])
|
|
isDeleteMarker := req.LatestDeleteMarkerKey != "" && string(versionsEntry.Extended[req.LatestDeleteMarkerKey]) == "true"
|
|
exists = latestFile != "" && !isDeleteMarker
|
|
if req.LatestEtagKey != "" {
|
|
etag = normalizeETag(string(versionsEntry.Extended[req.LatestEtagKey]))
|
|
}
|
|
}
|
|
want := normalizeETag(req.Condition.Etag)
|
|
switch req.Condition.Kind {
|
|
case filer_pb.WriteCondition_IF_NOT_EXISTS:
|
|
return !exists
|
|
case filer_pb.WriteCondition_IF_EXISTS:
|
|
return exists
|
|
case filer_pb.WriteCondition_IF_ETAG_MATCH:
|
|
return exists && etag == want
|
|
case filer_pb.WriteCondition_IF_ETAG_NOT_MATCH:
|
|
return !exists || etag != want
|
|
default:
|
|
return true
|
|
}
|
|
}
|
|
|
|
// FinalizeVersionedWrite performs an S3 versioned-write finalize atomically on
|
|
// the owner of the object's .versions directory: it creates the new version
|
|
// file, stamps the previous latest as noncurrent (before the pointer flip so the
|
|
// lifecycle router observes it), then merges the latest-pointer / cached
|
|
// metadata into the .versions directory entry. A single exclusive lock on the
|
|
// .versions directory serializes the whole sequence against concurrent versioned
|
|
// writes to the same object, replacing the distributed lock. Key names come from
|
|
// the caller so the filer carries no S3 versioning semantics.
|
|
func (fs *FilerServer) FinalizeVersionedWrite(ctx context.Context, req *filer_pb.FinalizeVersionedWriteRequest) (*filer_pb.FinalizeVersionedWriteResponse, error) {
|
|
resp := &filer_pb.FinalizeVersionedWriteResponse{}
|
|
versionsDir := util.FullPath(req.VersionsDir)
|
|
|
|
// Serialize on the object key — the same lock all of this object's writes
|
|
// (normal, suspended, versioned) share — not the .versions directory, so a
|
|
// versioned write and a non-versioned write to the same object can't run on
|
|
// different owners during a versioning-state change. Fall back to the
|
|
// .versions directory for older callers that don't set lock_key.
|
|
lockKey := util.FullPath(req.LockKey)
|
|
if lockKey == "" {
|
|
lockKey = versionsDir
|
|
}
|
|
lock := fs.entryLockTable.AcquireLock("FinalizeVersionedWrite", lockKey, util.ExclusiveLock)
|
|
defer fs.entryLockTable.ReleaseLock(lockKey, lock)
|
|
|
|
// The caller already created the new version file under this .versions
|
|
// directory, so it exists; read the directory entry that holds the pointer.
|
|
versionsEntry, err := fs.filer.FindEntry(ctx, versionsDir)
|
|
if err != nil {
|
|
resp.Error = err.Error()
|
|
return resp, nil
|
|
}
|
|
|
|
// 1. Evaluate the optional precondition against the current latest version
|
|
// (the .versions directory's cached pointer metadata) before the flip, so the
|
|
// check and the flip are atomic under this lock.
|
|
if req.Condition != nil && req.Condition.Kind != filer_pb.WriteCondition_NONE {
|
|
if !versionedConditionSatisfied(req, versionsEntry) {
|
|
glog.V(3).InfofCtx(ctx, "FinalizeVersionedWrite %s: precondition %v failed", versionsDir, req.Condition.Kind)
|
|
resp.Error = "precondition failed"
|
|
resp.ErrorCode = filer_pb.FilerError_PRECONDITION_FAILED
|
|
return resp, nil
|
|
}
|
|
}
|
|
|
|
// 2. Stamp the previously-latest version as noncurrent BEFORE the pointer
|
|
// flip (so the lifecycle router observes it). Best-effort.
|
|
newFileName := string(req.SetExtended[req.PriorLatestKey])
|
|
if req.NoncurrentSinceNs > 0 && req.PriorLatestKey != "" && req.NoncurrentSinceKey != "" {
|
|
prior := string(versionsEntry.Extended[req.PriorLatestKey])
|
|
if prior != "" && prior != newFileName {
|
|
demotedPath := util.NewFullPath(req.VersionsDir, prior)
|
|
if demoted, derr := fs.filer.FindEntry(ctx, demotedPath); derr == nil {
|
|
if demoted.Extended == nil {
|
|
demoted.Extended = make(map[string][]byte)
|
|
}
|
|
demoted.Extended[req.NoncurrentSinceKey] = []byte(strconv.FormatInt(req.NoncurrentSinceNs, 10))
|
|
if uerr := fs.filer.UpdateEntry(ctx, demoted, demoted); uerr != nil {
|
|
glog.Warningf("FinalizeVersionedWrite: demote stamp %s: %v", demotedPath, uerr)
|
|
}
|
|
} else {
|
|
glog.V(2).InfofCtx(ctx, "FinalizeVersionedWrite: demote target %s not found: %v", demotedPath, derr)
|
|
}
|
|
}
|
|
}
|
|
|
|
// 3. Flip the latest pointer: merge the requested Extended keys into the
|
|
// .versions directory entry and write it back (emits the meta-log event the
|
|
// lifecycle router consumes).
|
|
if versionsEntry.Extended == nil {
|
|
versionsEntry.Extended = make(map[string][]byte)
|
|
}
|
|
for _, k := range req.DeleteExtended {
|
|
delete(versionsEntry.Extended, k)
|
|
}
|
|
for k, v := range req.SetExtended {
|
|
versionsEntry.Extended[k] = v
|
|
}
|
|
if err := fs.filer.CreateEntry(ctx, versionsEntry, false, req.IsFromOtherCluster, req.Signatures, true, fs.filer.MaxFilenameLength); err != nil {
|
|
resp.Error = err.Error()
|
|
return resp, nil
|
|
}
|
|
|
|
return resp, nil
|
|
}
|
|
|
|
func (fs *FilerServer) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntryRequest) (*filer_pb.UpdateEntryResponse, error) {
|
|
|
|
glog.V(4).InfofCtx(ctx, "UpdateEntry %v", req)
|
|
if len(req.Entry.HardLinkId) > 0 {
|
|
glog.V(4).InfofCtx(ctx, "UpdateEntry %s/%s with HardLinkId %x counter=%d", req.Directory, req.Entry.Name, req.Entry.HardLinkId, req.Entry.HardLinkCounter)
|
|
}
|
|
|
|
fullpath := util.Join(req.Directory, req.Entry.Name)
|
|
entry, err := fs.filer.FindEntry(ctx, util.FullPath(fullpath))
|
|
if err != nil {
|
|
return &filer_pb.UpdateEntryResponse{}, fmt.Errorf("not found %s: %v", fullpath, err)
|
|
}
|
|
if err := validateUpdateEntryPreconditions(entry, req.ExpectedExtended); err != nil {
|
|
return &filer_pb.UpdateEntryResponse{}, err
|
|
}
|
|
|
|
chunks, garbage, err2 := fs.cleanupChunks(ctx, fullpath, entry, req.Entry)
|
|
if err2 != nil {
|
|
return &filer_pb.UpdateEntryResponse{}, fmt.Errorf("UpdateEntry cleanupChunks %s: %v", fullpath, err2)
|
|
}
|
|
|
|
newEntry := filer.FromPbEntry(req.Directory, req.Entry)
|
|
newEntry.Chunks = chunks
|
|
|
|
// Don't apply TTL to remote entries - they're managed by remote storage
|
|
if newEntry.Remote != nil {
|
|
newEntry.TtlSec = 0
|
|
}
|
|
|
|
if filer.EqualEntry(entry, newEntry) {
|
|
return &filer_pb.UpdateEntryResponse{}, err
|
|
}
|
|
|
|
ctx, eventSink := filer.WithMetadataEventSink(ctx)
|
|
resp := &filer_pb.UpdateEntryResponse{}
|
|
if err = fs.filer.UpdateEntry(ctx, entry, newEntry); err == nil {
|
|
fs.filer.DeleteChunksNotRecursive(garbage)
|
|
|
|
fs.filer.NotifyUpdateEvent(ctx, entry, newEntry, true, req.IsFromOtherCluster, req.Signatures)
|
|
resp.MetadataEvent = eventSink.Last()
|
|
|
|
} else {
|
|
glog.V(3).InfofCtx(ctx, "UpdateEntry %s: %v", filepath.Join(req.Directory, req.Entry.Name), err)
|
|
}
|
|
|
|
return resp, err
|
|
}
|
|
|
|
func validateUpdateEntryPreconditions(entry *filer.Entry, expectedExtended map[string][]byte) error {
|
|
if len(expectedExtended) == 0 {
|
|
return nil
|
|
}
|
|
|
|
for key, expectedValue := range expectedExtended {
|
|
var actualValue []byte
|
|
var ok bool
|
|
if entry != nil {
|
|
actualValue, ok = entry.Extended[key]
|
|
}
|
|
if ok {
|
|
if !bytes.Equal(actualValue, expectedValue) {
|
|
return status.Errorf(codes.FailedPrecondition, "extended attribute %q changed", key)
|
|
}
|
|
continue
|
|
}
|
|
if len(expectedValue) > 0 {
|
|
return status.Errorf(codes.FailedPrecondition, "extended attribute %q changed", key)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (fs *FilerServer) cleanupChunks(ctx context.Context, fullpath string, existingEntry *filer.Entry, newEntry *filer_pb.Entry) (chunks, garbage []*filer_pb.FileChunk, err error) {
|
|
|
|
// remove old chunks if not included in the new ones
|
|
if existingEntry != nil {
|
|
garbage, err = filer.MinusChunks(ctx, fs.lookupFileId, existingEntry.GetChunks(), newEntry.GetChunks())
|
|
if err != nil {
|
|
return newEntry.GetChunks(), nil, fmt.Errorf("MinusChunks: %w", err)
|
|
}
|
|
}
|
|
|
|
// files with manifest chunks are usually large and append only, skip calculating covered chunks
|
|
manifestChunks, nonManifestChunks := filer.SeparateManifestChunks(newEntry.GetChunks())
|
|
|
|
chunks, coveredChunks := filer.CompactFileChunks(ctx, fs.lookupFileId, nonManifestChunks)
|
|
garbage = append(garbage, coveredChunks...)
|
|
|
|
if newEntry.Attributes != nil {
|
|
so, _ := fs.detectStorageOption(ctx, fullpath,
|
|
"",
|
|
"",
|
|
newEntry.Attributes.TtlSec,
|
|
"",
|
|
"",
|
|
"",
|
|
"",
|
|
) // ignore readonly error for capacity needed to manifestize
|
|
chunks, err = filer.MaybeManifestize(fs.saveAsChunk(ctx, so), chunks)
|
|
if err != nil {
|
|
// not good, but should be ok
|
|
glog.V(0).InfofCtx(ctx, "MaybeManifestize: %v", err)
|
|
}
|
|
}
|
|
|
|
chunks = append(manifestChunks, chunks...)
|
|
|
|
return
|
|
}
|
|
|
|
func (fs *FilerServer) AppendToEntry(ctx context.Context, req *filer_pb.AppendToEntryRequest) (*filer_pb.AppendToEntryResponse, error) {
|
|
|
|
glog.V(4).InfofCtx(ctx, "AppendToEntry %v", req)
|
|
fullpath := util.NewFullPath(req.Directory, req.EntryName)
|
|
|
|
// Serialize the read-modify-write against concurrent mutations to the same
|
|
// path on this filer. The append must route to this entry's owner filer for
|
|
// this local lock to be authoritative.
|
|
pathLock := fs.entryLockTable.AcquireLock("AppendToEntry", fullpath, util.ExclusiveLock)
|
|
defer fs.entryLockTable.ReleaseLock(fullpath, pathLock)
|
|
|
|
var offset int64 = 0
|
|
entry, err := fs.filer.FindEntry(ctx, fullpath)
|
|
if err == filer_pb.ErrNotFound {
|
|
entry = &filer.Entry{
|
|
FullPath: fullpath,
|
|
Attr: filer.Attr{
|
|
Crtime: time.Now(),
|
|
Mtime: time.Now(),
|
|
Mode: os.FileMode(0644),
|
|
Uid: OS_UID,
|
|
Gid: OS_GID,
|
|
},
|
|
}
|
|
} else {
|
|
offset = int64(filer.TotalSize(entry.GetChunks()))
|
|
}
|
|
|
|
for _, chunk := range req.Chunks {
|
|
chunk.Offset = offset
|
|
offset += int64(chunk.Size)
|
|
}
|
|
|
|
entry.Chunks = append(entry.GetChunks(), req.Chunks...)
|
|
so, err := fs.detectStorageOption(ctx, string(fullpath), "", "", entry.TtlSec, "", "", "", "")
|
|
if err != nil {
|
|
glog.WarningfCtx(ctx, "detectStorageOption: %v", err)
|
|
return &filer_pb.AppendToEntryResponse{}, err
|
|
}
|
|
entry.Chunks, err = filer.MaybeManifestize(fs.saveAsChunk(ctx, so), entry.GetChunks())
|
|
if err != nil {
|
|
// not good, but should be ok
|
|
glog.V(0).InfofCtx(ctx, "MaybeManifestize: %v", err)
|
|
}
|
|
|
|
err = fs.filer.CreateEntry(context.Background(), entry, false, false, nil, false, fs.filer.MaxFilenameLength)
|
|
|
|
return &filer_pb.AppendToEntryResponse{}, err
|
|
}
|
|
|
|
func (fs *FilerServer) DeleteEntry(ctx context.Context, req *filer_pb.DeleteEntryRequest) (resp *filer_pb.DeleteEntryResponse, err error) {
|
|
|
|
glog.V(4).InfofCtx(ctx, "DeleteEntry %v", req)
|
|
|
|
fullpath := util.JoinPath(req.Directory, req.Name)
|
|
resp = &filer_pb.DeleteEntryResponse{}
|
|
|
|
// A single-entry delete serializes against concurrent mutations to the same
|
|
// path so a conditional delete's check and the removal are atomic. Recursive
|
|
// deletes span many paths and keep the prior unlocked behavior.
|
|
if !req.IsRecursive {
|
|
pathLock := fs.entryLockTable.AcquireLock("DeleteEntry", fullpath, util.ExclusiveLock)
|
|
defer fs.entryLockTable.ReleaseLock(fullpath, pathLock)
|
|
|
|
if req.Condition != nil && req.Condition.Kind != filer_pb.WriteCondition_NONE {
|
|
current, findErr := fs.filer.FindEntry(ctx, fullpath)
|
|
if findErr != nil && findErr != filer_pb.ErrNotFound {
|
|
return &filer_pb.DeleteEntryResponse{}, fmt.Errorf("DeleteEntry condition check %s: %w", fullpath, findErr)
|
|
}
|
|
if findErr == filer_pb.ErrNotFound {
|
|
current = nil
|
|
}
|
|
if !writeConditionSatisfied(req.Condition, current) {
|
|
glog.V(3).InfofCtx(ctx, "DeleteEntry %s: precondition %v failed", fullpath, req.Condition.Kind)
|
|
resp.Error = "precondition failed"
|
|
resp.ErrorCode = filer_pb.FilerError_PRECONDITION_FAILED
|
|
return resp, nil
|
|
}
|
|
}
|
|
}
|
|
|
|
ctx, eventSink := filer.WithMetadataEventSink(ctx)
|
|
err = fs.filer.DeleteEntryMetaAndData(ctx, fullpath, req.IsRecursive, req.IgnoreRecursiveError, req.IsDeleteData, req.IsFromOtherCluster, req.Signatures, req.IfNotModifiedAfter)
|
|
if err != nil && err != filer_pb.ErrNotFound {
|
|
resp.Error = err.Error()
|
|
} else {
|
|
resp.MetadataEvent = eventSink.Last()
|
|
}
|
|
return resp, nil
|
|
}
|
|
|
|
func (fs *FilerServer) AssignVolume(ctx context.Context, req *filer_pb.AssignVolumeRequest) (resp *filer_pb.AssignVolumeResponse, err error) {
|
|
|
|
so, err := fs.resolveAssignStorageOption(ctx, req)
|
|
if err != nil {
|
|
glog.V(3).InfofCtx(ctx, "AssignVolume: %v", err)
|
|
return &filer_pb.AssignVolumeResponse{Error: fmt.Sprintf("assign volume: %v", err)}, nil
|
|
}
|
|
|
|
assignRequest, altRequest := so.ToAssignRequests(int(req.Count))
|
|
assignRequest.ExpectedDataSize = req.ExpectedDataSize
|
|
if altRequest != nil {
|
|
altRequest.ExpectedDataSize = req.ExpectedDataSize
|
|
}
|
|
|
|
assignResult, err := operation.Assign(ctx, fs.filer.GetMaster, fs.grpcDialOption, assignRequest, altRequest)
|
|
if err != nil {
|
|
glog.V(3).InfofCtx(ctx, "AssignVolume: %v", err)
|
|
return &filer_pb.AssignVolumeResponse{Error: fmt.Sprintf("assign volume: %v", err)}, nil
|
|
}
|
|
if assignResult.Error != "" {
|
|
glog.V(3).InfofCtx(ctx, "AssignVolume error: %v", assignResult.Error)
|
|
return &filer_pb.AssignVolumeResponse{Error: fmt.Sprintf("assign volume result: %v", assignResult.Error)}, nil
|
|
}
|
|
|
|
return &filer_pb.AssignVolumeResponse{
|
|
FileId: assignResult.Fid,
|
|
Count: int32(assignResult.Count),
|
|
Location: &filer_pb.Location{
|
|
Url: assignResult.Url,
|
|
PublicUrl: assignResult.PublicUrl,
|
|
GrpcPort: uint32(assignResult.GrpcPort),
|
|
},
|
|
Auth: string(assignResult.Auth),
|
|
Collection: so.Collection,
|
|
Replication: so.Replication,
|
|
}, nil
|
|
}
|
|
|
|
func (fs *FilerServer) resolveAssignStorageOption(ctx context.Context, req *filer_pb.AssignVolumeRequest) (*operation.StorageOption, error) {
|
|
so, err := fs.detectStorageOption(ctx, req.Path, req.Collection, req.Replication, req.TtlSec, req.DiskType, req.DataCenter, req.Rack, req.DataNode)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
// Mirror the HTTP write path: only apply the filer's default disk when the
|
|
// matched locationPrefix rule did not already select one.
|
|
if so.DiskType == "" {
|
|
so.DiskType = fs.option.DiskType
|
|
}
|
|
|
|
return so, nil
|
|
}
|
|
|
|
func (fs *FilerServer) CollectionList(ctx context.Context, req *filer_pb.CollectionListRequest) (resp *filer_pb.CollectionListResponse, err error) {
|
|
|
|
glog.V(4).InfofCtx(ctx, "CollectionList %v", req)
|
|
resp = &filer_pb.CollectionListResponse{}
|
|
|
|
err = fs.filer.MasterClient.WithClient(false, func(client master_pb.SeaweedClient) error {
|
|
masterResp, err := client.CollectionList(context.Background(), &master_pb.CollectionListRequest{
|
|
IncludeNormalVolumes: req.IncludeNormalVolumes,
|
|
IncludeEcVolumes: req.IncludeEcVolumes,
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, c := range masterResp.Collections {
|
|
resp.Collections = append(resp.Collections, &filer_pb.Collection{Name: c.Name})
|
|
}
|
|
return nil
|
|
})
|
|
|
|
return
|
|
}
|
|
|
|
func (fs *FilerServer) DeleteCollection(ctx context.Context, req *filer_pb.DeleteCollectionRequest) (resp *filer_pb.DeleteCollectionResponse, err error) {
|
|
|
|
glog.V(4).InfofCtx(ctx, "DeleteCollection %v", req)
|
|
|
|
err = fs.filer.DoDeleteCollection(req.GetCollection())
|
|
|
|
return &filer_pb.DeleteCollectionResponse{}, err
|
|
}
|