mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-01 12:16:07 +00:00
Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
578fed7d30 | ||
|
|
6b701a94c4 | ||
|
|
18b51253da | ||
|
|
280f620b55 |
@@ -1000,10 +1000,15 @@ func (s3a *S3ApiServer) streamFromVolumeServers(w http.ResponseWriter, r *http.R
|
||||
entry = cachedEntry
|
||||
glog.V(1).Infof("streamFromVolumeServers: successfully cached remote object, got %d chunks", len(chunks))
|
||||
} else {
|
||||
// Caching failed - return error to client
|
||||
glog.Errorf("streamFromVolumeServers: failed to cache remote object for streaming")
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
return newStreamErrorWithResponse(fmt.Errorf("failed to cache remote object for streaming"))
|
||||
// Client disconnected: report cancellation, not 503.
|
||||
if ctxErr := r.Context().Err(); ctxErr != nil {
|
||||
return ctxErr
|
||||
}
|
||||
// Cache still filling: 503 + Retry-After so SDKs back off and retry.
|
||||
glog.V(1).Infof("streamFromVolumeServers: remote object %s/%s not cached yet, returning 503 for retry", bucket, object)
|
||||
w.Header().Set("Retry-After", "5")
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrServiceUnavailable)
|
||||
return newStreamErrorWithResponse(fmt.Errorf("remote object not cached yet"))
|
||||
}
|
||||
} else if totalSize > 0 && len(entry.Content) == 0 {
|
||||
// Not a remote entry but has size without content - this is a data integrity issue
|
||||
|
||||
+27
-9
@@ -3,6 +3,7 @@ package weed_server
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
@@ -39,6 +40,28 @@ var writePool = sync.Pool{New: func() interface{} {
|
||||
},
|
||||
}
|
||||
|
||||
// ErrCacheNotReady signals that a remote-only object's local cache is still
|
||||
// filling. Callers should map it to 503 + Retry-After so SDKs back off and retry.
|
||||
var ErrCacheNotReady = errors.New("remote object not cached yet")
|
||||
|
||||
// writePrepareWriteFnErr writes an HTTP response for an error from
|
||||
// prepareWriteFn, before any 2xx headers have been written. Client cancels are
|
||||
// silent; ErrCacheNotReady becomes 503 + Retry-After; everything else is 500.
|
||||
func writePrepareWriteFnErr(w http.ResponseWriter, err error) {
|
||||
w.Header().Del("Content-Length")
|
||||
switch {
|
||||
case errors.Is(err, context.Canceled):
|
||||
glog.V(3).Infof("ProcessRangeRequest: client disconnected: %v", err)
|
||||
case errors.Is(err, ErrCacheNotReady):
|
||||
glog.V(1).Infof("ProcessRangeRequest: cache not ready, returning 503: %v", err)
|
||||
w.Header().Set("Retry-After", "5")
|
||||
http.Error(w, err.Error(), http.StatusServiceUnavailable)
|
||||
default:
|
||||
glog.Errorf("ProcessRangeRequest: %v", err)
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
}
|
||||
}
|
||||
|
||||
func init() {
|
||||
serverStats = stats.NewServerStats()
|
||||
go serverStats.Start()
|
||||
@@ -290,9 +313,7 @@ func ProcessRangeRequest(r *http.Request, w http.ResponseWriter, totalSize int64
|
||||
w.Header().Set("Content-Length", strconv.FormatInt(totalSize, 10))
|
||||
writeFn, err := prepareWriteFn(0, totalSize)
|
||||
if err != nil {
|
||||
glog.Errorf("ProcessRangeRequest: %v", err)
|
||||
w.Header().Del("Content-Length")
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
writePrepareWriteFnErr(w, err)
|
||||
return fmt.Errorf("ProcessRangeRequest: %w", err)
|
||||
}
|
||||
if err = writeFn(bufferedWriter); err != nil {
|
||||
@@ -340,9 +361,7 @@ func ProcessRangeRequest(r *http.Request, w http.ResponseWriter, totalSize int64
|
||||
|
||||
writeFn, err := prepareWriteFn(ra.start, ra.length)
|
||||
if err != nil {
|
||||
glog.Errorf("ProcessRangeRequest range[0]: %+v err: %v", w.Header(), err)
|
||||
w.Header().Del("Content-Length")
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
writePrepareWriteFnErr(w, err)
|
||||
return fmt.Errorf("ProcessRangeRequest: %w", err)
|
||||
}
|
||||
w.WriteHeader(http.StatusPartialContent)
|
||||
@@ -365,9 +384,8 @@ func ProcessRangeRequest(r *http.Request, w http.ResponseWriter, totalSize int64
|
||||
}
|
||||
writeFn, err := prepareWriteFn(ra.start, ra.length)
|
||||
if err != nil {
|
||||
glog.Errorf("ProcessRangeRequest range[%d] err: %v", i, err)
|
||||
http.Error(w, "Internal Error", http.StatusInternalServerError)
|
||||
return fmt.Errorf("ProcessRangeRequest range[%d] err: %v", i, err)
|
||||
writePrepareWriteFnErr(w, err)
|
||||
return fmt.Errorf("ProcessRangeRequest range[%d]: %w", i, err)
|
||||
}
|
||||
writeFnByRange[i] = writeFn
|
||||
}
|
||||
|
||||
@@ -211,8 +211,13 @@ func (fs *FilerServer) GetOrHeadHandler(w http.ResponseWriter, r *http.Request)
|
||||
Name: name,
|
||||
}); err != nil {
|
||||
stats.FilerHandlerCounter.WithLabelValues(stats.ErrorReadCache).Inc()
|
||||
glog.ErrorfCtx(ctx, "CacheRemoteObjectToLocalCluster %s: %v", entry.FullPath, err)
|
||||
return nil, fmt.Errorf("cache %s: %v", entry.FullPath, err)
|
||||
// Client disconnected: surface ctx error so caller stays silent.
|
||||
if ctxErr := ctx.Err(); ctxErr != nil {
|
||||
return nil, ctxErr
|
||||
}
|
||||
// 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()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user