mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-19 04:50:54 +02:00
* filer: claim the base fid when minting a volume read token GenJwtForVolumeServer stamped the fid verbatim, but the volume server strips a trailing _N delta suffix before comparing the claim, so a token minted for a batch-assigned fid like 3,01637037d6_1 was checked against 3,01637037d6 and never matched. Reading such a chunk through the filer returned 401 wherever jwt.signing.read.key was configured. Strip the suffix before minting, via a helper shared with the proxyChunkId validation that was already doing the same thing inline. * filer: don't mint a volume write token for an anonymous proxy caller The ?proxyChunkId= branch dispatches and returns before the JWT gate, so whatever credential the proxy attaches is reachable without authentication. It attached a token from maybeGetVolumeReadJwtAuthorizationToken, which fell back to the write signing key when jwt.signing.read.key was unset -- the configuration scaffold/security.toml recommends for a filer, since read JWTs are only supported in a master+volume setup. An anonymous DELETE /?proxyChunkId=<fid> therefore arrived at the volume server holding a write-key token scoped to that fid, and the volume server honored it. Sign read tokens with the read key only. The fallback bought nothing on a read anyway: a volume server enforces read JWTs solely when that same key is set, so when the fallback fired the read was unchecked regardless. Mint only for reads. Writers proxied through the filer carry their own volume JWT from AssignVolume, forwarded with the rest of the caller's headers, so weed mount -filerProxy uploads are unaffected. Moving the dispatch below the JWT gate instead would have broken them, since that token is signed with jwt.signing rather than jwt.filer_signing. On a read with nothing to mint, drop the caller's Authorization rather than relaying it: there it is a filer credential, and forwarding it would hand a volume server a token it never used to see. * filer: keep proxied writes out of the read concurrency semaphore The semaphore is named and documented for reads -- it exists so replication bursts can't open hundreds of connections to one volume server -- but it was applied to every proxied method. A write queued behind sixteen in-flight reads can wait past the 10s default expiry of the AssignVolume token it carries, and the volume server then answers 401. shouldReassignUpload treats a 4xx as final, so the uploader replays the same expired token instead of re-assigning and the write fails up to the caller. This only became reachable once the filer stopped re-minting a fresh token after the wait.
275 lines
9.1 KiB
Go
275 lines
9.1 KiB
Go
package weed_server
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"math"
|
|
"mime"
|
|
"net/http"
|
|
"path/filepath"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/filer"
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
|
"github.com/seaweedfs/seaweedfs/weed/security"
|
|
"github.com/seaweedfs/seaweedfs/weed/stats"
|
|
"github.com/seaweedfs/seaweedfs/weed/util"
|
|
)
|
|
|
|
// Validates the preconditions. Returns true if GET/HEAD operation should not proceed.
|
|
// Preconditions supported are:
|
|
//
|
|
// If-Modified-Since
|
|
// If-Unmodified-Since
|
|
// If-Match
|
|
// If-None-Match
|
|
func checkPreconditions(w http.ResponseWriter, r *http.Request, entry *filer.Entry) bool {
|
|
|
|
etag := filer.ETagEntry(entry)
|
|
/// When more than one conditional request header field is present in a
|
|
/// request, the order in which the fields are evaluated becomes
|
|
/// important. In practice, the fields defined in this document are
|
|
/// consistently implemented in a single, logical order, since "lost
|
|
/// update" preconditions have more strict requirements than cache
|
|
/// validation, a validated cache is more efficient than a partial
|
|
/// response, and entity tags are presumed to be more accurate than date
|
|
/// validators. https://tools.ietf.org/html/rfc7232#section-5
|
|
if entry.Attr.Mtime.IsZero() {
|
|
return false
|
|
}
|
|
w.Header().Set("Last-Modified", entry.Attr.Mtime.UTC().Format(http.TimeFormat))
|
|
|
|
ifMatchETagHeader := r.Header.Get("If-Match")
|
|
ifUnmodifiedSinceHeader := r.Header.Get("If-Unmodified-Since")
|
|
if ifMatchETagHeader != "" {
|
|
if util.CanonicalizeETag(etag) != util.CanonicalizeETag(ifMatchETagHeader) {
|
|
w.WriteHeader(http.StatusPreconditionFailed)
|
|
return true
|
|
}
|
|
} else if ifUnmodifiedSinceHeader != "" {
|
|
if t, parseError := time.Parse(http.TimeFormat, ifUnmodifiedSinceHeader); parseError == nil {
|
|
if t.Before(entry.Attr.Mtime) {
|
|
w.WriteHeader(http.StatusPreconditionFailed)
|
|
return true
|
|
}
|
|
}
|
|
}
|
|
|
|
ifNoneMatchETagHeader := r.Header.Get("If-None-Match")
|
|
ifModifiedSinceHeader := r.Header.Get("If-Modified-Since")
|
|
if ifNoneMatchETagHeader != "" {
|
|
if util.CanonicalizeETag(etag) == util.CanonicalizeETag(ifNoneMatchETagHeader) {
|
|
SetEtag(w, etag)
|
|
w.WriteHeader(http.StatusNotModified)
|
|
return true
|
|
}
|
|
} else if ifModifiedSinceHeader != "" {
|
|
if t, parseError := time.Parse(http.TimeFormat, ifModifiedSinceHeader); parseError == nil {
|
|
if !t.Before(entry.Attr.Mtime) {
|
|
SetEtag(w, etag)
|
|
w.WriteHeader(http.StatusNotModified)
|
|
return true
|
|
}
|
|
}
|
|
}
|
|
|
|
return false
|
|
}
|
|
|
|
func (fs *FilerServer) GetOrHeadHandler(w http.ResponseWriter, r *http.Request) {
|
|
ctx := r.Context()
|
|
path := r.URL.Path
|
|
isForDirectory := strings.HasSuffix(path, "/")
|
|
if isForDirectory && len(path) > 1 {
|
|
path = path[:len(path)-1]
|
|
}
|
|
|
|
entry, err := fs.filer.FindEntry(ctx, util.FullPath(path))
|
|
if err != nil {
|
|
if path == "/" {
|
|
fs.listDirectoryHandler(w, r)
|
|
return
|
|
}
|
|
if err == filer_pb.ErrNotFound {
|
|
glog.V(2).InfofCtx(ctx, "Not found %s: %v", path, err)
|
|
stats.FilerHandlerCounter.WithLabelValues(stats.ErrorReadNotFound).Inc()
|
|
w.WriteHeader(http.StatusNotFound)
|
|
} else {
|
|
glog.ErrorfCtx(ctx, "Internal %s: %v", path, err)
|
|
stats.FilerHandlerCounter.WithLabelValues(stats.ErrorReadInternal).Inc()
|
|
w.WriteHeader(http.StatusInternalServerError)
|
|
}
|
|
return
|
|
}
|
|
|
|
query := r.URL.Query()
|
|
|
|
if entry.IsDirectory() {
|
|
if fs.option.DisableDirListing {
|
|
w.WriteHeader(http.StatusForbidden)
|
|
return
|
|
}
|
|
if query.Get("metadata") == "true" {
|
|
writeJsonQuiet(w, r, http.StatusOK, entry)
|
|
return
|
|
}
|
|
// listDirectoryHandler checks ExposeDirectoryData internally
|
|
fs.listDirectoryHandler(w, r)
|
|
return
|
|
}
|
|
|
|
if query.Get("metadata") == "true" {
|
|
if query.Get("resolveManifest") == "true" {
|
|
if entry.Chunks, _, err = filer.ResolveChunkManifest(
|
|
ctx,
|
|
fs.filer.MasterClient.GetLookupFileIdFunction(),
|
|
entry.GetChunks(), 0, math.MaxInt64); err != nil {
|
|
err = fmt.Errorf("failed to resolve chunk manifest, err: %s", err.Error())
|
|
writeJsonError(w, r, http.StatusInternalServerError, err)
|
|
return
|
|
}
|
|
}
|
|
writeJsonQuiet(w, r, http.StatusOK, entry)
|
|
return
|
|
}
|
|
|
|
if checkPreconditions(w, r, entry) {
|
|
return
|
|
}
|
|
|
|
// Generate ETag for response
|
|
etag := filer.ETagEntry(entry)
|
|
w.Header().Set("Accept-Ranges", "bytes")
|
|
|
|
// mime type
|
|
mimeType := entry.Attr.Mime
|
|
if mimeType == "" {
|
|
if ext := filepath.Ext(entry.Name()); ext != "" {
|
|
mimeType = mime.TypeByExtension(ext)
|
|
}
|
|
}
|
|
if mimeType != "" {
|
|
w.Header().Set("Content-Type", mimeType)
|
|
} else {
|
|
w.Header().Set("Content-Type", "application/octet-stream")
|
|
}
|
|
|
|
// print out the header from extended properties
|
|
// Filter out xattr-* (filesystem extended attributes) and internal SeaweedFS headers
|
|
for k, v := range entry.Extended {
|
|
if !strings.HasPrefix(k, "xattr-") && !s3_constants.IsSeaweedFSInternalHeader(k) {
|
|
w.Header().Set(k, string(v))
|
|
}
|
|
}
|
|
|
|
//Seaweed custom header are not visible to Vue or javascript
|
|
seaweedHeaders := []string{}
|
|
for header := range w.Header() {
|
|
if strings.HasPrefix(header, "Seaweed-") {
|
|
seaweedHeaders = append(seaweedHeaders, header)
|
|
}
|
|
}
|
|
seaweedHeaders = append(seaweedHeaders, "Content-Disposition")
|
|
w.Header().Set("Access-Control-Expose-Headers", strings.Join(seaweedHeaders, ","))
|
|
|
|
SetEtag(w, etag)
|
|
|
|
filename := entry.Name()
|
|
AdjustPassthroughHeaders(w, r, filename)
|
|
|
|
// For range processing, use the original content size, not the encrypted size
|
|
// entry.Size() returns max(chunk_sizes, file_size) where chunk_sizes include encryption overhead
|
|
// For SSE objects, we need the original unencrypted size for proper range validation
|
|
totalSize := int64(entry.FileSize)
|
|
|
|
if r.Method == http.MethodHead {
|
|
w.Header().Set("Content-Length", strconv.FormatInt(totalSize, 10))
|
|
return
|
|
}
|
|
|
|
if entry.Remote != nil && entry.Remote.RemoteSize > 0 {
|
|
// inline content is served locally without chunks
|
|
hit := !entry.IsInRemoteOnly() || len(entry.Content) > 0
|
|
stats.RecordRemoteCacheRead(stats.RemoteCacheSourceFiler, fs.filer.DetectBucket(entry.FullPath), hit)
|
|
}
|
|
|
|
ProcessRangeRequest(r, w, totalSize, mimeType, func(offset int64, size int64) (filer.DoStreamContent, error) {
|
|
if offset+size <= int64(len(entry.Content)) {
|
|
return func(writer io.Writer) error {
|
|
_, err := writer.Write(entry.Content[offset : offset+size])
|
|
if err != nil {
|
|
stats.FilerHandlerCounter.WithLabelValues(stats.ErrorWriteEntry).Inc()
|
|
glog.ErrorfCtx(ctx, "failed to write entry content: %v", err)
|
|
}
|
|
return err
|
|
}, nil
|
|
}
|
|
chunks := entry.GetChunks()
|
|
if entry.IsInRemoteOnly() {
|
|
dir, name := entry.FullPath.DirAndName()
|
|
if resp, err := fs.CacheRemoteObjectToLocalCluster(ctx, &filer_pb.CacheRemoteObjectToLocalClusterRequest{
|
|
Directory: dir,
|
|
Name: name,
|
|
}); err != nil {
|
|
stats.FilerHandlerCounter.WithLabelValues(stats.ErrorReadCache).Inc()
|
|
// Client disconnected: surface ctx error so caller stays silent.
|
|
if ctxErr := ctx.Err(); ctxErr != nil {
|
|
return nil, ctxErr
|
|
}
|
|
// Entry vanished mid-cache: forward NotFound so caller maps to 404,
|
|
// not the 503 retry-loop.
|
|
if errors.Is(err, filer_pb.ErrNotFound) {
|
|
return nil, err
|
|
}
|
|
// Cache still filling: tag with sentinel so caller maps to 503 + Retry-After.
|
|
glog.WarningfCtx(ctx, "CacheRemoteObjectToLocalCluster %s: %v", entry.FullPath, err)
|
|
return nil, fmt.Errorf("cache %s: %w", entry.FullPath, ErrCacheNotReady)
|
|
} else {
|
|
chunks = resp.Entry.GetChunks()
|
|
}
|
|
}
|
|
|
|
// Use a detached context for streaming so client disconnects/cancellations don't abort volume server operations,
|
|
// while preserving request-scoped values like tracing IDs.
|
|
// Matches S3 API behavior. Request context (ctx) is used for metadata operations above.
|
|
streamCtx, streamCancel := context.WithCancel(context.WithoutCancel(ctx))
|
|
|
|
streamFn, err := filer.PrepareStreamContentWithPrefetch(streamCtx, fs.filer.MasterClient, fs.maybeGetVolumeReadJwtAuthorizationToken, chunks, offset, size, fs.option.DownloadMaxBytesPs, 4)
|
|
if err != nil {
|
|
streamCancel()
|
|
stats.FilerHandlerCounter.WithLabelValues(stats.ErrorReadStream).Inc()
|
|
glog.ErrorfCtx(ctx, "failed to prepare stream content %s: %v", r.URL, err)
|
|
return nil, err
|
|
}
|
|
return func(writer io.Writer) error {
|
|
defer streamCancel()
|
|
err := streamFn(writer)
|
|
if err != nil {
|
|
stats.FilerHandlerCounter.WithLabelValues(stats.ErrorReadStream).Inc()
|
|
glog.ErrorfCtx(ctx, "failed to stream content %s: %v", r.URL, err)
|
|
}
|
|
return err
|
|
}, nil
|
|
})
|
|
}
|
|
|
|
func (fs *FilerServer) maybeGetVolumeReadJwtAuthorizationToken(fileId string) string {
|
|
// Only ever sign with the read key. A volume server enforces read JWTs
|
|
// solely when jwt.signing.read.key is set, so falling back to the write key
|
|
// buys no access on a read -- it only hands out a token that would authorize
|
|
// a write.
|
|
key := fs.volumeGuard.ReadSigningKey()
|
|
if len(key) == 0 {
|
|
return ""
|
|
}
|
|
// Claim the base fid: the volume server strips a _N delta suffix before
|
|
// comparing, so a token claiming the suffixed form never matches.
|
|
return string(security.GenJwtForVolumeServer(key, fs.volumeGuard.ReadExpiresAfterSec(), baseFileId(fileId)))
|
|
}
|