mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-01 12:16:07 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4e421400e4 | ||
|
|
56c51e7f50 |
@@ -222,6 +222,29 @@ message CreateEntryRequest {
|
||||
bool is_from_other_cluster = 4;
|
||||
repeated int32 signatures = 5;
|
||||
bool skip_check_parent_directory = 6;
|
||||
// Optional precondition evaluated against the current entry atomically with
|
||||
// the write, under the filer's per-path lock. The caller must route the
|
||||
// key's writes to this entry's owner filer for the check to be authoritative.
|
||||
WriteCondition condition = 7;
|
||||
}
|
||||
|
||||
// WriteCondition is a single precondition the filer evaluates against the
|
||||
// existing entry before writing. The client reduces request semantics (e.g.
|
||||
// RFC 7232 precedence) to one primitive; the filer just compares. A failed
|
||||
// condition returns FilerError PRECONDITION_FAILED.
|
||||
message WriteCondition {
|
||||
enum Kind {
|
||||
NONE = 0; // unconditional
|
||||
IF_NOT_EXISTS = 1; // fail if the entry exists (If-None-Match: *)
|
||||
IF_EXISTS = 2; // fail if the entry is absent (If-Match: *)
|
||||
IF_ETAG_MATCH = 3; // fail if absent or etag != expected (If-Match: <etag>)
|
||||
IF_ETAG_NOT_MATCH = 4; // fail if present and etag == expected (If-None-Match: <etag>)
|
||||
IF_UNMODIFIED_SINCE = 5; // fail if present and mtime > unix_time
|
||||
IF_MODIFIED_SINCE = 6; // fail if present and mtime <= unix_time
|
||||
}
|
||||
Kind kind = 1;
|
||||
string etag = 2; // expected ETag for IF_ETAG_* kinds
|
||||
int64 unix_time = 3; // bound (unix seconds) for IF_*_SINCE kinds
|
||||
}
|
||||
|
||||
// Structured error codes for filer entry operations.
|
||||
@@ -233,6 +256,7 @@ enum FilerError {
|
||||
EXISTING_IS_DIRECTORY = 3; // cannot overwrite directory with file
|
||||
EXISTING_IS_FILE = 4; // cannot overwrite file with directory
|
||||
ENTRY_ALREADY_EXISTS = 5; // O_EXCL and entry already exists
|
||||
PRECONDITION_FAILED = 6; // WriteCondition not satisfied
|
||||
}
|
||||
|
||||
message CreateEntryResponse {
|
||||
@@ -421,6 +445,7 @@ message SubscribeMetadataRequest {
|
||||
repeated string directories = 10; // exact directory to watch
|
||||
bool client_supports_batching = 11; // client can unpack SubscribeMetadataResponse.events
|
||||
bool client_supports_metadata_chunks = 12; // client can read log file chunks from volume servers
|
||||
bool client_supports_idle_heartbeat = 13; // server may send empty responses carrying the current time while the client is caught up
|
||||
}
|
||||
message SubscribeMetadataResponse {
|
||||
string directory = 1;
|
||||
|
||||
@@ -222,6 +222,29 @@ message CreateEntryRequest {
|
||||
bool is_from_other_cluster = 4;
|
||||
repeated int32 signatures = 5;
|
||||
bool skip_check_parent_directory = 6;
|
||||
// Optional precondition evaluated against the current entry atomically with
|
||||
// the write, under the filer's per-path lock. The caller must route the
|
||||
// key's writes to this entry's owner filer for the check to be authoritative.
|
||||
WriteCondition condition = 7;
|
||||
}
|
||||
|
||||
// WriteCondition is a single precondition the filer evaluates against the
|
||||
// existing entry before writing. The client reduces request semantics (e.g.
|
||||
// RFC 7232 precedence) to one primitive; the filer just compares. A failed
|
||||
// condition returns FilerError PRECONDITION_FAILED.
|
||||
message WriteCondition {
|
||||
enum Kind {
|
||||
NONE = 0; // unconditional
|
||||
IF_NOT_EXISTS = 1; // fail if the entry exists (If-None-Match: *)
|
||||
IF_EXISTS = 2; // fail if the entry is absent (If-Match: *)
|
||||
IF_ETAG_MATCH = 3; // fail if absent or etag != expected (If-Match: <etag>)
|
||||
IF_ETAG_NOT_MATCH = 4; // fail if present and etag == expected (If-None-Match: <etag>)
|
||||
IF_UNMODIFIED_SINCE = 5; // fail if present and mtime > unix_time
|
||||
IF_MODIFIED_SINCE = 6; // fail if present and mtime <= unix_time
|
||||
}
|
||||
Kind kind = 1;
|
||||
string etag = 2; // expected ETag for IF_ETAG_* kinds
|
||||
int64 unix_time = 3; // bound (unix seconds) for IF_*_SINCE kinds
|
||||
}
|
||||
|
||||
// Structured error codes for filer entry operations.
|
||||
@@ -233,6 +256,7 @@ enum FilerError {
|
||||
EXISTING_IS_DIRECTORY = 3; // cannot overwrite directory with file
|
||||
EXISTING_IS_FILE = 4; // cannot overwrite file with directory
|
||||
ENTRY_ALREADY_EXISTS = 5; // O_EXCL and entry already exists
|
||||
PRECONDITION_FAILED = 6; // WriteCondition not satisfied
|
||||
}
|
||||
|
||||
message CreateEntryResponse {
|
||||
|
||||
+563
-406
File diff suppressed because it is too large
Load Diff
@@ -7,15 +7,15 @@ import (
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/cluster"
|
||||
|
||||
"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"
|
||||
@@ -186,6 +186,32 @@ func (fs *FilerServer) CreateEntry(ctx context.Context, req *filer_pb.CreateEntr
|
||||
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)
|
||||
|
||||
@@ -212,6 +238,44 @@ func (fs *FilerServer) CreateEntry(ctx context.Context, req *filer_pb.CreateEntr
|
||||
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, `"`)
|
||||
}
|
||||
|
||||
func (fs *FilerServer) UpdateEntry(ctx context.Context, req *filer_pb.UpdateEntryRequest) (*filer_pb.UpdateEntryResponse, error) {
|
||||
|
||||
glog.V(4).InfofCtx(ctx, "UpdateEntry %v", req)
|
||||
@@ -328,9 +392,11 @@ func (fs *FilerServer) AppendToEntry(ctx context.Context, req *filer_pb.AppendTo
|
||||
glog.V(4).InfofCtx(ctx, "AppendToEntry %v", req)
|
||||
fullpath := util.NewFullPath(req.Directory, req.EntryName)
|
||||
|
||||
lockClient := cluster.NewLockClient(fs.grpcDialOption, fs.option.Host)
|
||||
lock := lockClient.NewShortLivedLock(string(fullpath), string(fs.option.Host))
|
||||
defer lock.StopShortLivedLock()
|
||||
// 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)
|
||||
|
||||
@@ -0,0 +1,108 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
func entryWithETag(etag string, mtime time.Time) *filer.Entry {
|
||||
return &filer.Entry{
|
||||
FullPath: "/test/obj",
|
||||
Attr: filer.Attr{Mtime: mtime},
|
||||
Extended: map[string][]byte{s3_constants.ExtETagKey: []byte(etag)},
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriteConditionSatisfied(t *testing.T) {
|
||||
base := time.Unix(1700000000, 0)
|
||||
present := entryWithETag("abc", base)
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
cond *filer_pb.WriteCondition
|
||||
cur *filer.Entry
|
||||
want bool
|
||||
}{
|
||||
{"none-absent", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_NONE}, nil, true},
|
||||
{"ifnotexists-absent", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_NOT_EXISTS}, nil, true},
|
||||
{"ifnotexists-present", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_NOT_EXISTS}, present, false},
|
||||
{"ifexists-absent", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_EXISTS}, nil, false},
|
||||
{"ifexists-present", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_EXISTS}, present, true},
|
||||
{"etagmatch-hit", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_ETAG_MATCH, Etag: `"abc"`}, present, true},
|
||||
{"etagmatch-miss", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_ETAG_MATCH, Etag: `"zzz"`}, present, false},
|
||||
{"etagmatch-absent", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_ETAG_MATCH, Etag: `"abc"`}, nil, false},
|
||||
{"etagnotmatch-hit", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_ETAG_NOT_MATCH, Etag: `"abc"`}, present, false},
|
||||
{"etagnotmatch-miss", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_ETAG_NOT_MATCH, Etag: `"zzz"`}, present, true},
|
||||
{"etagnotmatch-absent", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_ETAG_NOT_MATCH, Etag: `"abc"`}, nil, true},
|
||||
{"unmodsince-ok", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_UNMODIFIED_SINCE, UnixTime: base.Unix()}, present, true},
|
||||
{"unmodsince-fail", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_UNMODIFIED_SINCE, UnixTime: base.Unix() - 1}, present, false},
|
||||
{"modsince-ok", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_MODIFIED_SINCE, UnixTime: base.Unix() - 1}, present, true},
|
||||
{"modsince-fail", &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_MODIFIED_SINCE, UnixTime: base.Unix()}, present, false},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
if got := writeConditionSatisfied(tc.cond, tc.cur); got != tc.want {
|
||||
t.Errorf("%s: got %v want %v", tc.name, got, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// storedEntryETag prefers the stored Seaweed ETag attribute and falls back to
|
||||
// the Md5-derived ETag, matching the S3 gateway.
|
||||
func TestStoredEntryETag(t *testing.T) {
|
||||
withExt := entryWithETag("explicit", time.Unix(0, 0))
|
||||
if got := storedEntryETag(withExt); got != "explicit" {
|
||||
t.Errorf("extended etag: got %q", got)
|
||||
}
|
||||
md5Only := &filer.Entry{Attr: filer.Attr{Md5: []byte{0xab, 0xcd}}}
|
||||
if got := storedEntryETag(md5Only); got != "abcd" {
|
||||
t.Errorf("md5 fallback: got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
// The CreateEntry handler enforces the precondition atomically: a matching
|
||||
// If-Match overwrites, a non-matching one returns PRECONDITION_FAILED.
|
||||
func TestCreateEntryConditionEnforced(t *testing.T) {
|
||||
store := newRenameTestStore()
|
||||
store.entries["/test/obj"] = &filer.Entry{
|
||||
FullPath: "/test/obj",
|
||||
Attr: filer.Attr{Inode: 1, Mtime: time.Unix(1700000000, 0)},
|
||||
Extended: map[string][]byte{s3_constants.ExtETagKey: []byte("abc")},
|
||||
}
|
||||
f := newRenameTestFiler(store)
|
||||
f.DirBucketsPath = "/buckets"
|
||||
fs := &FilerServer{filer: f, option: &FilerOption{}, entryLockTable: util.NewLockTable[util.FullPath]()}
|
||||
|
||||
req := func(etag string) *filer_pb.CreateEntryRequest {
|
||||
return &filer_pb.CreateEntryRequest{
|
||||
Directory: "/test",
|
||||
SkipCheckParentDirectory: true,
|
||||
Entry: &filer_pb.Entry{
|
||||
Name: "obj",
|
||||
Attributes: &filer_pb.FuseAttributes{Mtime: 1700000001, FileMode: 0644, Inode: 2},
|
||||
},
|
||||
Condition: &filer_pb.WriteCondition{Kind: filer_pb.WriteCondition_IF_ETAG_MATCH, Etag: etag},
|
||||
}
|
||||
}
|
||||
|
||||
resp, err := fs.CreateEntry(context.Background(), req(`"zzz"`))
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected err: %v", err)
|
||||
}
|
||||
if resp.ErrorCode != filer_pb.FilerError_PRECONDITION_FAILED {
|
||||
t.Fatalf("mismatched etag: want PRECONDITION_FAILED, got %v (%q)", resp.ErrorCode, resp.Error)
|
||||
}
|
||||
|
||||
resp, err = fs.CreateEntry(context.Background(), req(`"abc"`))
|
||||
if err != nil {
|
||||
t.Fatalf("unexpected err: %v", err)
|
||||
}
|
||||
if resp.Error != "" {
|
||||
t.Fatalf("matching etag should overwrite, got error %q", resp.Error)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
// Concurrent OExcl creates for the same path must yield exactly one winner. The
|
||||
// filer's CreateEntry is a FindEntry-then-Insert; without the per-path lock both
|
||||
// racers observe "not found" and both insert. The exclusive entry lock makes the
|
||||
// check-then-act atomic so the losers see ErrEntryAlreadyExists.
|
||||
func TestCreateEntryOExclSerialized(t *testing.T) {
|
||||
store := newRenameTestStore()
|
||||
store.findDelay = 5 * time.Millisecond
|
||||
f := newRenameTestFiler(store)
|
||||
f.DirBucketsPath = "/buckets"
|
||||
|
||||
fs := &FilerServer{
|
||||
filer: f,
|
||||
option: &FilerOption{},
|
||||
entryLockTable: util.NewLockTable[util.FullPath](),
|
||||
}
|
||||
|
||||
const racers = 8
|
||||
var success int32
|
||||
var wg sync.WaitGroup
|
||||
start := make(chan struct{})
|
||||
for i := 0; i < racers; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
<-start
|
||||
resp, err := fs.CreateEntry(context.Background(), &filer_pb.CreateEntryRequest{
|
||||
Directory: "/test",
|
||||
OExcl: true,
|
||||
SkipCheckParentDirectory: true,
|
||||
Entry: &filer_pb.Entry{
|
||||
Name: "obj",
|
||||
Attributes: &filer_pb.FuseAttributes{Mtime: 1700000000, FileMode: 0644, Inode: 1},
|
||||
},
|
||||
})
|
||||
if err == nil && resp.Error == "" {
|
||||
atomic.AddInt32(&success, 1)
|
||||
}
|
||||
}()
|
||||
}
|
||||
close(start)
|
||||
wg.Wait()
|
||||
|
||||
if success != 1 {
|
||||
t.Fatalf("expected exactly 1 OExcl winner, got %d", success)
|
||||
}
|
||||
}
|
||||
@@ -29,6 +29,7 @@ type renameTestStore struct {
|
||||
findCalls map[string]int
|
||||
commitErr error
|
||||
deleteErr error
|
||||
findDelay time.Duration // optional: widen check-then-act windows in tests
|
||||
}
|
||||
|
||||
func newRenameTestStore() *renameTestStore {
|
||||
@@ -69,6 +70,9 @@ func (s *renameTestStore) UpdateEntry(_ context.Context, entry *filer.Entry) err
|
||||
}
|
||||
|
||||
func (s *renameTestStore) FindEntry(_ context.Context, p util.FullPath) (*filer.Entry, error) {
|
||||
if s.findDelay > 0 {
|
||||
time.Sleep(s.findDelay)
|
||||
}
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
s.findCalls[string(p)]++
|
||||
|
||||
@@ -124,6 +124,13 @@ type FilerServer struct {
|
||||
// mountPeerRegistry backs the MountRegister / MountList RPCs for peer
|
||||
// chunk sharing (tier 1). Always populated.
|
||||
mountPeerRegistry *filer.MountPeerRegistry
|
||||
|
||||
// entryLockTable serializes mutations to the same entry path on this filer.
|
||||
// It is the local serialization point for read-modify-write operations
|
||||
// (conditional create, append) once writers for a key are routed to this
|
||||
// node, replacing the distributed lock for that purpose. Idle keys are
|
||||
// evicted automatically, so the table stays bounded.
|
||||
entryLockTable *util.LockTable[util.FullPath]
|
||||
}
|
||||
|
||||
func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption) (fs *FilerServer, err error) {
|
||||
@@ -162,6 +169,7 @@ func NewFilerServer(defaultMux, readonlyMux *http.ServeMux, option *FilerOption)
|
||||
inFlightDataLimitCond: sync.NewCond(new(sync.Mutex)),
|
||||
recentCopyRequests: make(map[string]recentCopyRequest),
|
||||
CredentialManager: option.CredentialManager,
|
||||
entryLockTable: util.NewLockTable[util.FullPath](),
|
||||
}
|
||||
fs.mountPeerRegistry = filer.NewMountPeerRegistry()
|
||||
go fs.runMountPeerRegistrySweeper()
|
||||
|
||||
Reference in New Issue
Block a user