mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-09-14 02:20:41 +02:00
withFilerClientFailover treated a filer's ErrNotFound like a transport failure: it kept the result, re-queried every other filer, and finally wrapped the answer as "all filers failed, last error: ... no entry is found in filer store". For workloads with many legitimate misses (e.g. GET object?versionId=X for a version that was deleted or expired), this turned each 404 into N filer round-trips and produced a misleading error string. A reachable filer that answers ErrNotFound has given an authoritative answer; failover exists to route around unreachable or unhealthy filers, not to look harder for an entry the store reports as absent. Return ErrNotFound directly instead of fanning out. Callers that need read-after-write retries already handle that at the S3 semantic layer (e.g. getLatestObjectVersion).
136 lines
4.3 KiB
Go
136 lines
4.3 KiB
Go
package s3api
|
|
|
|
import (
|
|
"encoding/base64"
|
|
"errors"
|
|
"fmt"
|
|
"net/http"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/glog"
|
|
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
|
"google.golang.org/grpc"
|
|
|
|
"github.com/seaweedfs/seaweedfs/weed/pb"
|
|
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
|
)
|
|
|
|
var _ = filer_pb.FilerClient(&S3ApiServer{})
|
|
|
|
func (s3a *S3ApiServer) WithFilerClient(streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error {
|
|
// Use filerClient for proper connection management and failover
|
|
if s3a.filerClient != nil {
|
|
return s3a.withFilerClientFailover(streamingMode, fn)
|
|
}
|
|
|
|
// Fallback to direct connection if filerClient not initialized
|
|
// This should only happen during initialization or testing
|
|
return pb.WithGrpcClient(streamingMode, s3a.randomClientId, func(grpcConnection *grpc.ClientConn) error {
|
|
client := filer_pb.NewSeaweedFilerClient(grpcConnection)
|
|
return fn(client)
|
|
}, s3a.getFilerAddress().ToGrpcAddress(), false, s3a.option.GrpcDialOption)
|
|
|
|
}
|
|
|
|
// withFilerClientFailover attempts to execute fn with automatic failover to other filers
|
|
func (s3a *S3ApiServer) withFilerClientFailover(streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error {
|
|
// Get current filer as starting point
|
|
currentFiler := s3a.filerClient.GetCurrentFiler()
|
|
|
|
// Try current filer first (fast path)
|
|
err := pb.WithGrpcClient(streamingMode, s3a.randomClientId, func(grpcConnection *grpc.ClientConn) error {
|
|
client := filer_pb.NewSeaweedFilerClient(grpcConnection)
|
|
return fn(client)
|
|
}, currentFiler.ToGrpcAddress(), false, s3a.option.GrpcDialOption)
|
|
|
|
if err == nil {
|
|
s3a.filerClient.RecordFilerSuccess(currentFiler)
|
|
return nil
|
|
}
|
|
|
|
// A reachable filer answering ErrNotFound is authoritative; failover is for
|
|
// unreachable/unhealthy filers, not for re-asking about an absent entry.
|
|
if errors.Is(err, filer_pb.ErrNotFound) {
|
|
return err
|
|
}
|
|
|
|
s3a.filerClient.RecordFilerFailure(currentFiler)
|
|
|
|
// Current filer failed - try all other filers with health-aware selection
|
|
filers := s3a.filerClient.GetAllFilers()
|
|
var lastErr error = err
|
|
|
|
for _, filer := range filers {
|
|
if filer == currentFiler {
|
|
continue // Already tried this one
|
|
}
|
|
|
|
// Skip filers known to be unhealthy (circuit breaker pattern)
|
|
if s3a.filerClient.ShouldSkipUnhealthyFiler(filer) {
|
|
glog.V(2).Infof("WithFilerClient: skipping unhealthy filer %s", filer)
|
|
continue
|
|
}
|
|
|
|
err = pb.WithGrpcClient(streamingMode, s3a.randomClientId, func(grpcConnection *grpc.ClientConn) error {
|
|
client := filer_pb.NewSeaweedFilerClient(grpcConnection)
|
|
return fn(client)
|
|
}, filer.ToGrpcAddress(), false, s3a.option.GrpcDialOption)
|
|
|
|
if err == nil {
|
|
// Success! Record success and update current filer for future requests
|
|
s3a.filerClient.RecordFilerSuccess(filer)
|
|
s3a.filerClient.SetCurrentFiler(filer)
|
|
glog.V(1).Infof("WithFilerClient: failover from %s to %s succeeded", currentFiler, filer)
|
|
return nil
|
|
}
|
|
|
|
// Authoritative not-found - stop failing over.
|
|
if errors.Is(err, filer_pb.ErrNotFound) {
|
|
return err
|
|
}
|
|
|
|
s3a.filerClient.RecordFilerFailure(filer)
|
|
glog.V(2).Infof("WithFilerClient: failover to %s failed: %v", filer, err)
|
|
lastErr = err
|
|
}
|
|
|
|
// All filers failed
|
|
return fmt.Errorf("all filers failed, last error: %w", lastErr)
|
|
}
|
|
|
|
func (s3a *S3ApiServer) AdjustedUrl(location *filer_pb.Location) string {
|
|
return location.Url
|
|
}
|
|
|
|
func (s3a *S3ApiServer) GetDataCenter() string {
|
|
return s3a.option.DataCenter
|
|
}
|
|
|
|
func writeSuccessResponseXML(w http.ResponseWriter, r *http.Request, response interface{}) {
|
|
s3err.WriteXMLResponse(w, r, http.StatusOK, response)
|
|
s3err.PostLog(r, http.StatusOK, s3err.ErrNone)
|
|
}
|
|
|
|
func writeSuccessResponseXMLBytes(w http.ResponseWriter, r *http.Request, response []byte) {
|
|
s3err.WriteResponse(w, r, http.StatusOK, response, s3err.MimeXML)
|
|
s3err.PostLog(r, http.StatusOK, s3err.ErrNone)
|
|
}
|
|
|
|
func writeSuccessResponseEmpty(w http.ResponseWriter, r *http.Request) {
|
|
s3err.WriteEmptyResponse(w, r, http.StatusOK)
|
|
}
|
|
|
|
func writeFailureResponse(w http.ResponseWriter, r *http.Request, errCode s3err.ErrorCode) {
|
|
s3err.WriteErrorResponse(w, r, errCode)
|
|
}
|
|
|
|
func validateContentMd5(h http.Header) ([]byte, error) {
|
|
md5B64, ok := h["Content-Md5"]
|
|
if ok {
|
|
if md5B64[0] == "" {
|
|
return nil, fmt.Errorf("Content-Md5 header set to empty value")
|
|
}
|
|
return base64.StdEncoding.DecodeString(md5B64[0])
|
|
}
|
|
return []byte{}, nil
|
|
}
|