mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-01 20:26:27 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
98fcda66f2 | ||
|
|
0f5c6a8ae0 |
@@ -457,6 +457,15 @@ func (s3a *S3ApiServer) completeMultipartUpload(r *http.Request, input *s3.Compl
|
||||
}
|
||||
} else if vErr == nil && versioningState == s3_constants.VersioningSuspended {
|
||||
// For suspended versioning, add "null" version ID metadata and return "null" version ID
|
||||
// If conditional headers are present, acquire a distributed lock to atomically
|
||||
// re-check conditions and create the entry, preventing TOCTOU races.
|
||||
condLock, condErr := s3a.lockAndRecheckConditionalHeaders(r, *input.Bucket, *input.Key)
|
||||
if condErr != s3err.ErrNone {
|
||||
return nil, condErr
|
||||
}
|
||||
if condLock != nil {
|
||||
defer condLock.StopShortLivedLock()
|
||||
}
|
||||
err = s3a.mkFile(dirName, entryName, finalParts, func(entry *filer_pb.Entry) {
|
||||
if entry.Extended == nil {
|
||||
entry.Extended = make(map[string][]byte)
|
||||
@@ -512,6 +521,15 @@ func (s3a *S3ApiServer) completeMultipartUpload(r *http.Request, input *s3.Compl
|
||||
}
|
||||
} else {
|
||||
// For non-versioned buckets, create main object file
|
||||
// If conditional headers are present, acquire a distributed lock to atomically
|
||||
// re-check conditions and create the entry, preventing TOCTOU races.
|
||||
condLock, condErr := s3a.lockAndRecheckConditionalHeaders(r, *input.Bucket, *input.Key)
|
||||
if condErr != s3err.ErrNone {
|
||||
return nil, condErr
|
||||
}
|
||||
if condLock != nil {
|
||||
defer condLock.StopShortLivedLock()
|
||||
}
|
||||
err = s3a.mkFile(dirName, entryName, finalParts, func(entry *filer_pb.Entry) {
|
||||
if entry.Extended == nil {
|
||||
entry.Extended = make(map[string][]byte)
|
||||
|
||||
@@ -17,6 +17,7 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/pquerna/cachecontrol/cacheobject"
|
||||
"github.com/seaweedfs/seaweedfs/weed/cluster"
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
@@ -697,6 +698,19 @@ func (s3a *S3ApiServer) putToFiler(r *http.Request, filePath string, dataReader
|
||||
// Step 4: Save metadata to filer via gRPC
|
||||
// Use context.Background() to ensure metadata save completes even if HTTP request is cancelled
|
||||
// This matches the chunk upload behavior and prevents orphaned chunks
|
||||
//
|
||||
// If conditional headers are present, acquire a distributed lock to atomically
|
||||
// re-check conditions and create the entry. This prevents TOCTOU races where
|
||||
// concurrent requests all pass the early check before any write completes.
|
||||
condBucket, condObject := s3_constants.GetBucketAndObject(r)
|
||||
condLock, condErr := s3a.lockAndRecheckConditionalHeaders(r, condBucket, condObject)
|
||||
if condErr != s3err.ErrNone {
|
||||
s3a.deleteOrphanedChunks(chunkResult.FileChunks)
|
||||
return "", condErr, SSEResponseMetadata{}
|
||||
}
|
||||
if condLock != nil {
|
||||
defer condLock.StopShortLivedLock()
|
||||
}
|
||||
glog.V(3).Infof("putToFiler: About to create entry - dir=%s, name=%s, chunks=%d, extended keys=%d",
|
||||
path.Dir(filePath), path.Base(filePath), len(entry.Chunks), len(entry.Extended))
|
||||
createErr := s3a.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
|
||||
@@ -1822,6 +1836,43 @@ func (s3a *S3ApiServer) checkConditionalHeaders(r *http.Request, bucket, object
|
||||
return s3a.validateConditionalHeaders(r, headers, entry, bucket, object)
|
||||
}
|
||||
|
||||
// lockAndRecheckConditionalHeaders acquires a distributed lock and re-checks conditional
|
||||
// headers atomically before a write operation. This prevents TOCTOU races where concurrent
|
||||
// requests could all pass the early (unlocked) check before any write completes.
|
||||
//
|
||||
// Returns (nil, ErrNone) if no conditional headers are present (no lock needed).
|
||||
// Returns (lock, ErrNone) if the re-check passes — the caller MUST defer lock.StopShortLivedLock()
|
||||
// to hold the lock across the subsequent write and release it after.
|
||||
// Returns (nil, errCode) if the re-check fails — no lock to release.
|
||||
func (s3a *S3ApiServer) lockAndRecheckConditionalHeaders(r *http.Request, bucket, object string) (*cluster.LiveLock, s3err.ErrorCode) {
|
||||
headers, errCode := parseConditionalHeaders(r)
|
||||
if errCode != s3err.ErrNone {
|
||||
return nil, errCode
|
||||
}
|
||||
if !headers.isSet {
|
||||
return nil, s3err.ErrNone
|
||||
}
|
||||
|
||||
lockKey := s3a.toFilerPath(bucket, object)
|
||||
lock := s3a.lockClient.NewShortLivedLock(lockKey, fmt.Sprintf("s3-cond-%d", s3a.randomClientId))
|
||||
|
||||
entry, err := s3a.resolveObjectEntry(bucket, object)
|
||||
if err != nil {
|
||||
if errors.Is(err, filer_pb.ErrNotFound) {
|
||||
entry = nil
|
||||
} else {
|
||||
lock.StopShortLivedLock()
|
||||
glog.Errorf("lockAndRecheckConditionalHeaders: error resolving object entry for %s/%s: %v", bucket, object, err)
|
||||
return nil, s3err.ErrInternalError
|
||||
}
|
||||
}
|
||||
if errCode = s3a.validateConditionalHeaders(r, headers, entry, bucket, object); errCode != s3err.ErrNone {
|
||||
lock.StopShortLivedLock()
|
||||
return nil, errCode
|
||||
}
|
||||
return lock, s3err.ErrNone
|
||||
}
|
||||
|
||||
// validateConditionalHeadersForReads checks conditional headers for read operations against the provided entry
|
||||
func (s3a *S3ApiServer) validateConditionalHeadersForReads(r *http.Request, headers conditionalHeaders, entry *filer_pb.Entry, bucket, object string) ConditionalHeaderResult {
|
||||
if !headers.isSet {
|
||||
|
||||
@@ -82,6 +82,7 @@ type S3ApiServer struct {
|
||||
embeddedIam *EmbeddedIamApi // Embedded IAM API server (when enabled)
|
||||
stsHandlers *STSHandlers // STS HTTP handlers for AssumeRoleWithWebIdentity
|
||||
cipher bool // encrypt data on volume servers
|
||||
lockClient *cluster.LockClient
|
||||
}
|
||||
|
||||
func NewS3ApiServer(router *mux.Router, option *S3ApiServerOption) (s3ApiServer *S3ApiServer, err error) {
|
||||
@@ -180,6 +181,7 @@ func NewS3ApiServerWithStore(router *mux.Router, option *S3ApiServerOption, expl
|
||||
policyEngine: policyEngine, // Initialize bucket policy engine
|
||||
inFlightDataLimitCond: sync.NewCond(new(sync.Mutex)),
|
||||
cipher: option.Cipher,
|
||||
lockClient: cluster.NewLockClient(option.GrpcDialOption, option.Filers[0]),
|
||||
}
|
||||
|
||||
// Set s3a reference in circuit breaker for upload limiting
|
||||
|
||||
Reference in New Issue
Block a user