mirror of
https://github.com/DullJZ/s3-balance.git
synced 2026-09-05 15:56:39 +08:00
Fix the deleted key issue
This commit is contained in:
+163
-170
@@ -54,20 +54,20 @@ func NewS3Handler(
|
|||||||
func (h *S3Handler) RegisterS3Routes(router *mux.Router) {
|
func (h *S3Handler) RegisterS3Routes(router *mux.Router) {
|
||||||
// Service operations
|
// Service operations
|
||||||
router.HandleFunc("/", h.handleListBuckets).Methods("GET")
|
router.HandleFunc("/", h.handleListBuckets).Methods("GET")
|
||||||
|
|
||||||
// Bucket operations
|
// Bucket operations
|
||||||
router.HandleFunc("/{bucket}", h.handleBucketOperations).Methods("GET", "HEAD", "PUT", "DELETE")
|
router.HandleFunc("/{bucket}", h.handleBucketOperations).Methods("GET", "HEAD", "PUT", "DELETE")
|
||||||
|
|
||||||
// Object operations
|
// Object operations
|
||||||
router.HandleFunc("/{bucket}/{key:.*}", h.handleObjectOperations).Methods("GET", "HEAD", "PUT", "DELETE")
|
router.HandleFunc("/{bucket}/{key:.*}", h.handleObjectOperations).Methods("GET", "HEAD", "PUT", "DELETE")
|
||||||
|
|
||||||
// Multipart upload operations
|
// Multipart upload operations
|
||||||
router.HandleFunc("/{bucket}/{key:.*}", h.handleMultipartUpload).Methods("POST").Queries("uploads", "")
|
router.HandleFunc("/{bucket}/{key:.*}", h.handleMultipartUpload).Methods("POST").Queries("uploads", "")
|
||||||
router.HandleFunc("/{bucket}/{key:.*}", h.handleListMultipartUploads).Methods("GET").Queries("uploads", "")
|
router.HandleFunc("/{bucket}/{key:.*}", h.handleListMultipartUploads).Methods("GET").Queries("uploads", "")
|
||||||
router.HandleFunc("/{bucket}/{key:.*}", h.handleListMultipartParts).Methods("GET").Queries("uploadId", "")
|
router.HandleFunc("/{bucket}/{key:.*}", h.handleListMultipartParts).Methods("GET").Queries("uploadId", "")
|
||||||
router.HandleFunc("/{bucket}/{key:.*}", h.handleCompleteMultipartUpload).Methods("POST").Queries("uploadId", "")
|
router.HandleFunc("/{bucket}/{key:.*}", h.handleCompleteMultipartUpload).Methods("POST").Queries("uploadId", "")
|
||||||
router.HandleFunc("/{bucket}/{key:.*}", h.handleAbortMultipartUpload).Methods("DELETE").Queries("uploadId", "")
|
router.HandleFunc("/{bucket}/{key:.*}", h.handleAbortMultipartUpload).Methods("DELETE").Queries("uploadId", "")
|
||||||
|
|
||||||
// 添加认证中间件
|
// 添加认证中间件
|
||||||
router.Use(h.s3AuthMiddleware)
|
router.Use(h.s3AuthMiddleware)
|
||||||
}
|
}
|
||||||
@@ -120,54 +120,54 @@ type CommonPrefix struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type InitiateMultipartUploadResult struct {
|
type InitiateMultipartUploadResult struct {
|
||||||
XMLName xml.Name `xml:"InitiateMultipartUploadResult"`
|
XMLName xml.Name `xml:"InitiateMultipartUploadResult"`
|
||||||
Xmlns string `xml:"xmlns,attr"`
|
Xmlns string `xml:"xmlns,attr"`
|
||||||
Bucket string `xml:"Bucket"`
|
Bucket string `xml:"Bucket"`
|
||||||
Key string `xml:"Key"`
|
Key string `xml:"Key"`
|
||||||
UploadID string `xml:"UploadId"`
|
UploadID string `xml:"UploadId"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type ListMultipartUploadsResult struct {
|
type ListMultipartUploadsResult struct {
|
||||||
XMLName xml.Name `xml:"ListMultipartUploadsResult"`
|
XMLName xml.Name `xml:"ListMultipartUploadsResult"`
|
||||||
Xmlns string `xml:"xmlns,attr"`
|
Xmlns string `xml:"xmlns,attr"`
|
||||||
Bucket string `xml:"Bucket"`
|
Bucket string `xml:"Bucket"`
|
||||||
KeyMarker string `xml:"KeyMarker"`
|
KeyMarker string `xml:"KeyMarker"`
|
||||||
UploadIdMarker string `xml:"UploadIdMarker"`
|
UploadIdMarker string `xml:"UploadIdMarker"`
|
||||||
NextKeyMarker string `xml:"NextKeyMarker"`
|
NextKeyMarker string `xml:"NextKeyMarker"`
|
||||||
NextUploadIdMarker string `xml:"NextUploadIdMarker"`
|
NextUploadIdMarker string `xml:"NextUploadIdMarker"`
|
||||||
MaxUploads int `xml:"MaxUploads"`
|
MaxUploads int `xml:"MaxUploads"`
|
||||||
IsTruncated bool `xml:"IsTruncated"`
|
IsTruncated bool `xml:"IsTruncated"`
|
||||||
Uploads []Upload `xml:"Upload"`
|
Uploads []Upload `xml:"Upload"`
|
||||||
CommonPrefixes []CommonPrefix `xml:"CommonPrefixes,omitempty"`
|
CommonPrefixes []CommonPrefix `xml:"CommonPrefixes,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type Upload struct {
|
type Upload struct {
|
||||||
Key string `xml:"Key"`
|
Key string `xml:"Key"`
|
||||||
UploadID string `xml:"UploadId"`
|
UploadID string `xml:"UploadId"`
|
||||||
Initiator Owner `xml:"Initiator"`
|
Initiator Owner `xml:"Initiator"`
|
||||||
Owner Owner `xml:"Owner"`
|
Owner Owner `xml:"Owner"`
|
||||||
StorageClass string `xml:"StorageClass"`
|
StorageClass string `xml:"StorageClass"`
|
||||||
Initiated time.Time `xml:"Initiated"`
|
Initiated time.Time `xml:"Initiated"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type ListPartsResult struct {
|
type ListPartsResult struct {
|
||||||
XMLName xml.Name `xml:"ListPartsResult"`
|
XMLName xml.Name `xml:"ListPartsResult"`
|
||||||
Xmlns string `xml:"xmlns,attr"`
|
Xmlns string `xml:"xmlns,attr"`
|
||||||
Bucket string `xml:"Bucket"`
|
Bucket string `xml:"Bucket"`
|
||||||
Key string `xml:"Key"`
|
Key string `xml:"Key"`
|
||||||
UploadID string `xml:"UploadId"`
|
UploadID string `xml:"UploadId"`
|
||||||
PartNumberMarker int `xml:"PartNumberMarker"`
|
PartNumberMarker int `xml:"PartNumberMarker"`
|
||||||
NextPartNumberMarker int `xml:"NextPartNumberMarker"`
|
NextPartNumberMarker int `xml:"NextPartNumberMarker"`
|
||||||
MaxParts int `xml:"MaxParts"`
|
MaxParts int `xml:"MaxParts"`
|
||||||
IsTruncated bool `xml:"IsTruncated"`
|
IsTruncated bool `xml:"IsTruncated"`
|
||||||
Parts []Part `xml:"Part"`
|
Parts []Part `xml:"Part"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type Part struct {
|
type Part struct {
|
||||||
PartNumber int `xml:"PartNumber"`
|
PartNumber int `xml:"PartNumber"`
|
||||||
LastModified time.Time `xml:"LastModified"`
|
LastModified time.Time `xml:"LastModified"`
|
||||||
ETag string `xml:"ETag"`
|
ETag string `xml:"ETag"`
|
||||||
Size int64 `xml:"Size"`
|
Size int64 `xml:"Size"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type CompleteMultipartUpload struct {
|
type CompleteMultipartUpload struct {
|
||||||
@@ -202,7 +202,7 @@ func (h *S3Handler) s3AuthMiddleware(next http.Handler) http.Handler {
|
|||||||
// 允许匿名访问(用于测试)
|
// 允许匿名访问(用于测试)
|
||||||
// 在生产环境中应该要求认证
|
// 在生产环境中应该要求认证
|
||||||
}
|
}
|
||||||
|
|
||||||
next.ServeHTTP(w, r)
|
next.ServeHTTP(w, r)
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -210,7 +210,7 @@ func (h *S3Handler) s3AuthMiddleware(next http.Handler) http.Handler {
|
|||||||
// handleListBuckets 处理列出所有存储桶请求
|
// handleListBuckets 处理列出所有存储桶请求
|
||||||
func (h *S3Handler) handleListBuckets(w http.ResponseWriter, r *http.Request) {
|
func (h *S3Handler) handleListBuckets(w http.ResponseWriter, r *http.Request) {
|
||||||
buckets := h.bucketManager.GetAllBuckets()
|
buckets := h.bucketManager.GetAllBuckets()
|
||||||
|
|
||||||
result := ListBucketsResult{
|
result := ListBucketsResult{
|
||||||
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
||||||
Owner: Owner{
|
Owner: Owner{
|
||||||
@@ -221,7 +221,7 @@ func (h *S3Handler) handleListBuckets(w http.ResponseWriter, r *http.Request) {
|
|||||||
Bucket: make([]BucketInfo, 0, len(buckets)),
|
Bucket: make([]BucketInfo, 0, len(buckets)),
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
for _, b := range buckets {
|
for _, b := range buckets {
|
||||||
// 只显示启用的虚拟存储桶,对客户端隐藏底层真实存储桶
|
// 只显示启用的虚拟存储桶,对客户端隐藏底层真实存储桶
|
||||||
if b.IsAvailable() && b.Config.Enabled && b.Config.Virtual {
|
if b.IsAvailable() && b.Config.Enabled && b.Config.Virtual {
|
||||||
@@ -231,7 +231,7 @@ func (h *S3Handler) handleListBuckets(w http.ResponseWriter, r *http.Request) {
|
|||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
h.sendXMLResponse(w, http.StatusOK, result)
|
h.sendXMLResponse(w, http.StatusOK, result)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -239,7 +239,7 @@ func (h *S3Handler) handleListBuckets(w http.ResponseWriter, r *http.Request) {
|
|||||||
func (h *S3Handler) handleBucketOperations(w http.ResponseWriter, r *http.Request) {
|
func (h *S3Handler) handleBucketOperations(w http.ResponseWriter, r *http.Request) {
|
||||||
vars := mux.Vars(r)
|
vars := mux.Vars(r)
|
||||||
bucketName := vars["bucket"]
|
bucketName := vars["bucket"]
|
||||||
|
|
||||||
switch r.Method {
|
switch r.Method {
|
||||||
case "GET":
|
case "GET":
|
||||||
h.handleListObjects(w, r, bucketName)
|
h.handleListObjects(w, r, bucketName)
|
||||||
@@ -260,13 +260,13 @@ func (h *S3Handler) handleListObjects(w http.ResponseWriter, r *http.Request, bu
|
|||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 如果是虚拟存储桶,列出虚拟存储桶中的对象
|
// 如果是虚拟存储桶,列出虚拟存储桶中的对象
|
||||||
if bucket.IsVirtual() {
|
if bucket.IsVirtual() {
|
||||||
h.handleListObjectsForVirtualBucket(w, r, bucketName)
|
h.handleListObjectsForVirtualBucket(w, r, bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 如果不是虚拟存储桶,拒绝客户端访问真实存储桶
|
// 如果不是虚拟存储桶,拒绝客户端访问真实存储桶
|
||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
}
|
}
|
||||||
@@ -278,21 +278,21 @@ func (h *S3Handler) handleListObjectsForVirtualBucket(w http.ResponseWriter, r *
|
|||||||
marker := r.URL.Query().Get("marker")
|
marker := r.URL.Query().Get("marker")
|
||||||
maxKeysStr := r.URL.Query().Get("max-keys")
|
maxKeysStr := r.URL.Query().Get("max-keys")
|
||||||
// delimiter := r.URL.Query().Get("delimiter") // 暂时不支持delimiter
|
// delimiter := r.URL.Query().Get("delimiter") // 暂时不支持delimiter
|
||||||
|
|
||||||
maxKeys := 1000
|
maxKeys := 1000
|
||||||
if maxKeysStr != "" {
|
if maxKeysStr != "" {
|
||||||
if mk, err := strconv.Atoi(maxKeysStr); err == nil {
|
if mk, err := strconv.Atoi(maxKeysStr); err == nil {
|
||||||
maxKeys = mk
|
maxKeys = mk
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 从存储服务获取虚拟存储桶中的对象
|
// 从存储服务获取虚拟存储桶中的对象
|
||||||
objects, err := h.storage.GetVirtualBucketObjects(bucketName)
|
objects, err := h.storage.GetVirtualBucketObjects(bucketName)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
h.sendS3Error(w, "InternalError", "Failed to list virtual bucket objects", bucketName)
|
h.sendS3Error(w, "InternalError", "Failed to list virtual bucket objects", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
result := ListBucketResult{
|
result := ListBucketResult{
|
||||||
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
||||||
Name: bucketName,
|
Name: bucketName,
|
||||||
@@ -302,19 +302,19 @@ func (h *S3Handler) handleListObjectsForVirtualBucket(w http.ResponseWriter, r *
|
|||||||
IsTruncated: false,
|
IsTruncated: false,
|
||||||
Contents: make([]ObjectInfo, 0, len(objects)),
|
Contents: make([]ObjectInfo, 0, len(objects)),
|
||||||
}
|
}
|
||||||
|
|
||||||
// 过滤对象并转换为S3格式
|
// 过滤对象并转换为S3格式
|
||||||
for _, obj := range objects {
|
for _, obj := range objects {
|
||||||
// 前缀过滤
|
// 前缀过滤
|
||||||
if prefix != "" && !strings.HasPrefix(obj.Key, prefix) {
|
if prefix != "" && !strings.HasPrefix(obj.Key, prefix) {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
// Marker过滤
|
// Marker过滤
|
||||||
if marker != "" && obj.Key <= marker {
|
if marker != "" && obj.Key <= marker {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
result.Contents = append(result.Contents, ObjectInfo{
|
result.Contents = append(result.Contents, ObjectInfo{
|
||||||
Key: obj.Key,
|
Key: obj.Key,
|
||||||
LastModified: obj.UpdatedAt,
|
LastModified: obj.UpdatedAt,
|
||||||
@@ -322,13 +322,13 @@ func (h *S3Handler) handleListObjectsForVirtualBucket(w http.ResponseWriter, r *
|
|||||||
Size: obj.Size,
|
Size: obj.Size,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
// 如果超过了最大数量,设置截断标志
|
// 如果超过了最大数量,设置截断标志
|
||||||
if len(result.Contents) > maxKeys {
|
if len(result.Contents) > maxKeys {
|
||||||
result.Contents = result.Contents[:maxKeys]
|
result.Contents = result.Contents[:maxKeys]
|
||||||
result.IsTruncated = true
|
result.IsTruncated = true
|
||||||
}
|
}
|
||||||
|
|
||||||
h.sendXMLResponse(w, http.StatusOK, result)
|
h.sendXMLResponse(w, http.StatusOK, result)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -339,13 +339,13 @@ func (h *S3Handler) handleHeadBucket(w http.ResponseWriter, r *http.Request, buc
|
|||||||
w.WriteHeader(http.StatusNotFound)
|
w.WriteHeader(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 虚拟存储桶也应该返回成功状态
|
// 虚拟存储桶也应该返回成功状态
|
||||||
if bucket.IsVirtual() {
|
if bucket.IsVirtual() {
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 如果不是虚拟存储桶,拒绝客户端访问真实存储桶
|
// 如果不是虚拟存储桶,拒绝客户端访问真实存储桶
|
||||||
w.WriteHeader(http.StatusNotFound)
|
w.WriteHeader(http.StatusNotFound)
|
||||||
}
|
}
|
||||||
@@ -365,7 +365,7 @@ func (h *S3Handler) handleCreateBucket(w http.ResponseWriter, r *http.Request, b
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 检查是否为虚拟存储桶
|
// 检查是否为虚拟存储桶
|
||||||
if requestedBucket, exists := h.bucketManager.GetBucket(bucketName); exists && requestedBucket.IsVirtual() {
|
if requestedBucket, exists := h.bucketManager.GetBucket(bucketName); exists && requestedBucket.IsVirtual() {
|
||||||
// 虚拟存储桶需要选择一个真实存储桶进行映射
|
// 虚拟存储桶需要选择一个真实存储桶进行映射
|
||||||
@@ -374,21 +374,21 @@ func (h *S3Handler) handleCreateBucket(w http.ResponseWriter, r *http.Request, b
|
|||||||
h.sendS3Error(w, "InternalError", "No real buckets available for virtual bucket mapping", bucketName)
|
h.sendS3Error(w, "InternalError", "No real buckets available for virtual bucket mapping", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 简化:选择第一个可用的真实存储桶
|
// 简化:选择第一个可用的真实存储桶
|
||||||
// 实际应用中可能需要更复杂的策略
|
// 实际应用中可能需要更复杂的策略
|
||||||
targetBucket := realBuckets[0]
|
targetBucket := realBuckets[0]
|
||||||
|
|
||||||
// 创建虚拟存储桶到真实存储桶的映射
|
// 创建虚拟存储桶到真实存储桶的映射
|
||||||
if err := h.storage.CreateVirtualBucketMapping(bucketName, "", targetBucket.Config.Name); err != nil {
|
if err := h.storage.CreateVirtualBucketMapping(bucketName, "", targetBucket.Config.Name); err != nil {
|
||||||
h.sendS3Error(w, "InternalError", "Failed to create virtual bucket mapping", bucketName)
|
h.sendS3Error(w, "InternalError", "Failed to create virtual bucket mapping", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 在负载均衡场景下,不真正创建bucket,只返回成功
|
// 在负载均衡场景下,不真正创建bucket,只返回成功
|
||||||
// 实际的bucket应该在配置中预先定义
|
// 实际的bucket应该在配置中预先定义
|
||||||
w.Header().Set("Location", "/" + bucketName)
|
w.Header().Set("Location", "/"+bucketName)
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -401,7 +401,7 @@ func (h *S3Handler) handleDeleteBucket(w http.ResponseWriter, r *http.Request, b
|
|||||||
w.WriteHeader(http.StatusNoContent)
|
w.WriteHeader(http.StatusNoContent)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 虚拟存储桶需要删除映射关系
|
// 虚拟存储桶需要删除映射关系
|
||||||
if bucket.IsVirtual() {
|
if bucket.IsVirtual() {
|
||||||
// 删除虚拟存储桶映射
|
// 删除虚拟存储桶映射
|
||||||
@@ -410,7 +410,7 @@ func (h *S3Handler) handleDeleteBucket(w http.ResponseWriter, r *http.Request, b
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 在负载均衡场景下,不真正删除真实bucket
|
// 在负载均衡场景下,不真正删除真实bucket
|
||||||
w.WriteHeader(http.StatusNoContent)
|
w.WriteHeader(http.StatusNoContent)
|
||||||
}
|
}
|
||||||
@@ -420,7 +420,7 @@ func (h *S3Handler) handleObjectOperations(w http.ResponseWriter, r *http.Reques
|
|||||||
vars := mux.Vars(r)
|
vars := mux.Vars(r)
|
||||||
bucketName := vars["bucket"]
|
bucketName := vars["bucket"]
|
||||||
key := vars["key"]
|
key := vars["key"]
|
||||||
|
|
||||||
switch r.Method {
|
switch r.Method {
|
||||||
case "GET":
|
case "GET":
|
||||||
h.handleGetObject(w, r, bucketName, key)
|
h.handleGetObject(w, r, bucketName, key)
|
||||||
@@ -441,7 +441,7 @@ func (h *S3Handler) handleGetObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 如果是虚拟存储桶,需要通过映射查找真实存储桶
|
// 如果是虚拟存储桶,需要通过映射查找真实存储桶
|
||||||
var err error
|
var err error
|
||||||
var bucket1 *bucket.BucketInfo
|
var bucket1 *bucket.BucketInfo
|
||||||
@@ -453,7 +453,7 @@ func (h *S3Handler) handleGetObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
h.sendS3Error(w, "NoSuchKey", "The specified key does not exist", key)
|
h.sendS3Error(w, "NoSuchKey", "The specified key does not exist", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 获取映射到的真实存储桶
|
// 获取映射到的真实存储桶
|
||||||
bucket1, ok = h.bucketManager.GetBucket(mapping.RealBucketName)
|
bucket1, ok = h.bucketManager.GetBucket(mapping.RealBucketName)
|
||||||
if !ok {
|
if !ok {
|
||||||
@@ -461,7 +461,7 @@ func (h *S3Handler) handleGetObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// 生成预签名下载URL
|
// 生成预签名下载URL
|
||||||
downloadInfo, err := h.presigner.GenerateDownloadURL(
|
downloadInfo, err := h.presigner.GenerateDownloadURL(
|
||||||
context.Background(),
|
context.Background(),
|
||||||
@@ -472,7 +472,7 @@ func (h *S3Handler) handleGetObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
h.sendS3Error(w, "InternalError", "Failed to generate download URL", key)
|
h.sendS3Error(w, "InternalError", "Failed to generate download URL", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 默认使用预签名重定向模式,只有明确指定时才使用代理模式
|
// 默认使用预签名重定向模式,只有明确指定时才使用代理模式
|
||||||
if r.URL.Query().Get("proxy") == "true" {
|
if r.URL.Query().Get("proxy") == "true" {
|
||||||
// 代理模式:服务器下载内容并返回给客户端
|
// 代理模式:服务器下载内容并返回给客户端
|
||||||
@@ -482,12 +482,12 @@ func (h *S3Handler) handleGetObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
// 复制响应头
|
// 复制响应头
|
||||||
for k, v := range resp.Header {
|
for k, v := range resp.Header {
|
||||||
w.Header()[k] = v
|
w.Header()[k] = v
|
||||||
}
|
}
|
||||||
|
|
||||||
// 复制响应体
|
// 复制响应体
|
||||||
io.Copy(w, resp.Body)
|
io.Copy(w, resp.Body)
|
||||||
} else {
|
} else {
|
||||||
@@ -504,7 +504,7 @@ func (h *S3Handler) handleHeadObject(w http.ResponseWriter, r *http.Request, buc
|
|||||||
w.WriteHeader(http.StatusNotFound)
|
w.WriteHeader(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 如果是虚拟存储桶,需要通过映射查找真实存储桶
|
// 如果是虚拟存储桶,需要通过映射查找真实存储桶
|
||||||
if requestedBucket.IsVirtual() {
|
if requestedBucket.IsVirtual() {
|
||||||
// 获取虚拟存储桶映射
|
// 获取虚拟存储桶映射
|
||||||
@@ -513,20 +513,20 @@ func (h *S3Handler) handleHeadObject(w http.ResponseWriter, r *http.Request, buc
|
|||||||
w.WriteHeader(http.StatusNotFound)
|
w.WriteHeader(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
_ = mapping // 使用mapping变量,避免编译错误
|
_ = mapping // 使用mapping变量,避免编译错误
|
||||||
|
|
||||||
// 查找对象信息(在映射的真实存储桶中)
|
// 查找对象信息(在映射的真实存储桶中)
|
||||||
obj, err := h.storage.GetObjectInfo(key)
|
obj, err := h.storage.GetObjectInfo(key)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
w.WriteHeader(http.StatusNotFound)
|
w.WriteHeader(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
h.setObjectHeaders(w, obj)
|
h.setObjectHeaders(w, obj)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 真实存储桶的直接处理
|
// 真实存储桶的直接处理
|
||||||
// 从存储中获取对象信息
|
// 从存储中获取对象信息
|
||||||
obj, err := h.storage.GetObjectInfo(key)
|
obj, err := h.storage.GetObjectInfo(key)
|
||||||
@@ -534,7 +534,7 @@ func (h *S3Handler) handleHeadObject(w http.ResponseWriter, r *http.Request, buc
|
|||||||
w.WriteHeader(http.StatusNotFound)
|
w.WriteHeader(http.StatusNotFound)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
h.setObjectHeaders(w, obj)
|
h.setObjectHeaders(w, obj)
|
||||||
w.WriteHeader(http.StatusOK)
|
w.WriteHeader(http.StatusOK)
|
||||||
}
|
}
|
||||||
@@ -551,7 +551,7 @@ func (h *S3Handler) setObjectHeaders(w http.ResponseWriter, obj *storage.Object)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// handlePutObject 上传对象(默认使用预签名URL重定向)
|
// handlePutObject 上传对象
|
||||||
func (h *S3Handler) handlePutObject(w http.ResponseWriter, r *http.Request, bucketName string, key string) {
|
func (h *S3Handler) handlePutObject(w http.ResponseWriter, r *http.Request, bucketName string, key string) {
|
||||||
// 检查请求的存储桶是否为虚拟存储桶
|
// 检查请求的存储桶是否为虚拟存储桶
|
||||||
requestedBucket, ok := h.bucketManager.GetBucket(bucketName)
|
requestedBucket, ok := h.bucketManager.GetBucket(bucketName)
|
||||||
@@ -559,17 +559,17 @@ func (h *S3Handler) handlePutObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 获取内容长度
|
// 获取内容长度
|
||||||
contentLength := r.ContentLength
|
contentLength := r.ContentLength
|
||||||
if contentLength < 0 {
|
if contentLength < 0 {
|
||||||
h.sendS3Error(w, "MissingContentLength", "Content-Length header is required", key)
|
h.sendS3Error(w, "MissingContentLength", "Content-Length header is required", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
var targetBucket *bucket.BucketInfo
|
var targetBucket *bucket.BucketInfo
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
// 如果是虚拟存储桶,需要选择真实存储桶并创建映射
|
// 如果是虚拟存储桶,需要选择真实存储桶并创建映射
|
||||||
if requestedBucket.IsVirtual() {
|
if requestedBucket.IsVirtual() {
|
||||||
// 获取虚拟存储桶文件映射,如果不存在则创建
|
// 获取虚拟存储桶文件映射,如果不存在则创建
|
||||||
@@ -581,7 +581,7 @@ func (h *S3Handler) handlePutObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
h.sendS3Error(w, "InsufficientStorage", "No bucket has enough space", key)
|
h.sendS3Error(w, "InsufficientStorage", "No bucket has enough space", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 创建虚拟存储桶文件级映射
|
// 创建虚拟存储桶文件级映射
|
||||||
if err := h.storage.CreateVirtualBucketMapping(bucketName, key, targetBucket.Config.Name); err != nil {
|
if err := h.storage.CreateVirtualBucketMapping(bucketName, key, targetBucket.Config.Name); err != nil {
|
||||||
h.sendS3Error(w, "InternalError", "Failed to create virtual bucket file mapping", key)
|
h.sendS3Error(w, "InternalError", "Failed to create virtual bucket file mapping", key)
|
||||||
@@ -600,7 +600,7 @@ func (h *S3Handler) handlePutObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 生成预签名上传URL
|
// 生成预签名上传URL
|
||||||
uploadInfo, err := h.presigner.GenerateUploadURL(
|
uploadInfo, err := h.presigner.GenerateUploadURL(
|
||||||
context.Background(),
|
context.Background(),
|
||||||
@@ -613,52 +613,45 @@ func (h *S3Handler) handlePutObject(w http.ResponseWriter, r *http.Request, buck
|
|||||||
h.sendS3Error(w, "InternalError", "Failed to generate upload URL", key)
|
h.sendS3Error(w, "InternalError", "Failed to generate upload URL", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 默认使用预签名重定向模式,只有明确指定时才使用代理模式
|
// 只使用反向代理上传到真实预签名URL,不再返回307重定向
|
||||||
if r.URL.Query().Get("proxy") == "true" {
|
// 创建新的请求
|
||||||
// 代理模式:读取请求体并上传到预签名URL
|
req, err := http.NewRequest(uploadInfo.Method, uploadInfo.URL, r.Body)
|
||||||
// 创建新的请求
|
if err != nil {
|
||||||
req, err := http.NewRequest(uploadInfo.Method, uploadInfo.URL, r.Body)
|
h.sendS3Error(w, "InternalError", "Failed to create upload request", key)
|
||||||
if err != nil {
|
return
|
||||||
h.sendS3Error(w, "InternalError", "Failed to create upload request", key)
|
}
|
||||||
return
|
|
||||||
}
|
// 设置必要的头
|
||||||
|
req.ContentLength = contentLength
|
||||||
// 设置必要的头
|
if ct := r.Header.Get("Content-Type"); ct != "" {
|
||||||
req.ContentLength = contentLength
|
req.Header.Set("Content-Type", ct)
|
||||||
if ct := r.Header.Get("Content-Type"); ct != "" {
|
}
|
||||||
req.Header.Set("Content-Type", ct)
|
|
||||||
}
|
// 添加预签名URL所需的额外头
|
||||||
|
for k, v := range uploadInfo.Headers {
|
||||||
// 添加预签名URL所需的额外头
|
req.Header.Set(k, v)
|
||||||
for k, v := range uploadInfo.Headers {
|
}
|
||||||
req.Header.Set(k, v)
|
|
||||||
}
|
// 执行上传
|
||||||
|
client := &http.Client{Timeout: 30 * time.Minute}
|
||||||
// 执行上传
|
resp, err := client.Do(req)
|
||||||
client := &http.Client{Timeout: 30 * time.Minute}
|
if err != nil {
|
||||||
resp, err := client.Do(req)
|
h.sendS3Error(w, "InternalError", "Failed to upload object", key)
|
||||||
if err != nil {
|
return
|
||||||
h.sendS3Error(w, "InternalError", "Failed to upload object", key)
|
}
|
||||||
return
|
defer resp.Body.Close()
|
||||||
}
|
|
||||||
defer resp.Body.Close()
|
if resp.StatusCode >= 200 && resp.StatusCode < 300 {
|
||||||
|
// 记录对象元数据
|
||||||
if resp.StatusCode >= 200 && resp.StatusCode < 300 {
|
h.storage.RecordObject(key, targetBucket.Config.Name, contentLength, nil)
|
||||||
// 记录对象元数据
|
targetBucket.UpdateUsedSize(contentLength)
|
||||||
h.storage.RecordObject(key, targetBucket.Config.Name, contentLength, nil)
|
|
||||||
targetBucket.UpdateUsedSize(contentLength)
|
// 返回成功响应
|
||||||
|
w.Header().Set("ETag", fmt.Sprintf("\"%x\"", time.Now().UnixNano()))
|
||||||
// 返回成功响应
|
w.WriteHeader(http.StatusOK)
|
||||||
w.Header().Set("ETag", fmt.Sprintf("\"%x\"", time.Now().UnixNano()))
|
|
||||||
w.WriteHeader(http.StatusOK)
|
|
||||||
} else {
|
|
||||||
h.sendS3Error(w, "InternalError", "Upload failed", key)
|
|
||||||
}
|
|
||||||
} else {
|
} else {
|
||||||
// 重定向模式:返回307临时重定向让客户端直接上传(默认)
|
h.sendS3Error(w, "InternalError", "Upload failed", key)
|
||||||
w.Header().Set("Location", uploadInfo.URL)
|
|
||||||
w.WriteHeader(http.StatusTemporaryRedirect)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -670,10 +663,10 @@ func (h *S3Handler) handleDeleteObject(w http.ResponseWriter, r *http.Request, b
|
|||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
var bucket *bucket.BucketInfo
|
var bucket *bucket.BucketInfo
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
if requestedBucket.IsVirtual() {
|
if requestedBucket.IsVirtual() {
|
||||||
// 获取虚拟存储桶文件映射
|
// 获取虚拟存储桶文件映射
|
||||||
mapping, err := h.storage.GetVirtualBucketMapping(bucketName, key)
|
mapping, err := h.storage.GetVirtualBucketMapping(bucketName, key)
|
||||||
@@ -682,7 +675,7 @@ func (h *S3Handler) handleDeleteObject(w http.ResponseWriter, r *http.Request, b
|
|||||||
w.WriteHeader(http.StatusNoContent)
|
w.WriteHeader(http.StatusNoContent)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 获取映射到的真实存储桶
|
// 获取映射到的真实存储桶
|
||||||
bucket, ok = h.bucketManager.GetBucket(mapping.RealBucketName)
|
bucket, ok = h.bucketManager.GetBucket(mapping.RealBucketName)
|
||||||
if !ok {
|
if !ok {
|
||||||
@@ -694,7 +687,7 @@ func (h *S3Handler) handleDeleteObject(w http.ResponseWriter, r *http.Request, b
|
|||||||
w.WriteHeader(http.StatusNoContent)
|
w.WriteHeader(http.StatusNoContent)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 生成预签名删除URL
|
// 生成预签名删除URL
|
||||||
deleteInfo, err := h.presigner.GenerateDeleteURL(
|
deleteInfo, err := h.presigner.GenerateDeleteURL(
|
||||||
context.Background(),
|
context.Background(),
|
||||||
@@ -705,7 +698,7 @@ func (h *S3Handler) handleDeleteObject(w http.ResponseWriter, r *http.Request, b
|
|||||||
h.sendS3Error(w, "InternalError", "Failed to generate delete URL", key)
|
h.sendS3Error(w, "InternalError", "Failed to generate delete URL", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 执行删除
|
// 执行删除
|
||||||
req, _ := http.NewRequest("DELETE", deleteInfo.URL, nil)
|
req, _ := http.NewRequest("DELETE", deleteInfo.URL, nil)
|
||||||
client := &http.Client{Timeout: 30 * time.Second}
|
client := &http.Client{Timeout: 30 * time.Second}
|
||||||
@@ -715,15 +708,15 @@ func (h *S3Handler) handleDeleteObject(w http.ResponseWriter, r *http.Request, b
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
// 从数据库中删除对象记录
|
// 从数据库中删除对象记录
|
||||||
h.storage.DeleteObject(key)
|
h.storage.DeleteObject(key)
|
||||||
|
|
||||||
// 如果是虚拟存储桶,还需要删除文件级别映射
|
// 如果是虚拟存储桶,还需要删除文件级别映射
|
||||||
if requestedBucket.IsVirtual() {
|
if requestedBucket.IsVirtual() {
|
||||||
h.storage.DeleteVirtualBucketFileMapping(bucketName, key)
|
h.storage.DeleteVirtualBucketFileMapping(bucketName, key)
|
||||||
}
|
}
|
||||||
|
|
||||||
// S3规范要求删除操作总是返回204
|
// S3规范要求删除操作总是返回204
|
||||||
w.WriteHeader(http.StatusNoContent)
|
w.WriteHeader(http.StatusNoContent)
|
||||||
}
|
}
|
||||||
@@ -732,14 +725,14 @@ func (h *S3Handler) handleDeleteObject(w http.ResponseWriter, r *http.Request, b
|
|||||||
func (h *S3Handler) handleMultipartUpload(w http.ResponseWriter, r *http.Request) {
|
func (h *S3Handler) handleMultipartUpload(w http.ResponseWriter, r *http.Request) {
|
||||||
vars := mux.Vars(r)
|
vars := mux.Vars(r)
|
||||||
key := vars["key"]
|
key := vars["key"]
|
||||||
|
|
||||||
// 选择目标存储桶
|
// 选择目标存储桶
|
||||||
targetBucket, err := h.balancer.SelectBucket(key, 0) // 分片上传时不检查空间
|
targetBucket, err := h.balancer.SelectBucket(key, 0) // 分片上传时不检查空间
|
||||||
if err != nil {
|
if err != nil {
|
||||||
h.sendS3Error(w, "InternalError", "Failed to select bucket for upload", key)
|
h.sendS3Error(w, "InternalError", "Failed to select bucket for upload", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 初始化分片上传
|
// 初始化分片上传
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
createResp, err := targetBucket.Client.CreateMultipartUpload(ctx, &s3.CreateMultipartUploadInput{
|
createResp, err := targetBucket.Client.CreateMultipartUpload(ctx, &s3.CreateMultipartUploadInput{
|
||||||
@@ -750,14 +743,14 @@ func (h *S3Handler) handleMultipartUpload(w http.ResponseWriter, r *http.Request
|
|||||||
h.sendS3Error(w, "InternalError", "Failed to initiate multipart upload", key)
|
h.sendS3Error(w, "InternalError", "Failed to initiate multipart upload", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
result := InitiateMultipartUploadResult{
|
result := InitiateMultipartUploadResult{
|
||||||
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
||||||
Bucket: targetBucket.Config.Name,
|
Bucket: targetBucket.Config.Name,
|
||||||
Key: key,
|
Key: key,
|
||||||
UploadID: *createResp.UploadId,
|
UploadID: *createResp.UploadId,
|
||||||
}
|
}
|
||||||
|
|
||||||
h.sendXMLResponse(w, http.StatusOK, result)
|
h.sendXMLResponse(w, http.StatusOK, result)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -766,13 +759,13 @@ func (h *S3Handler) handleListMultipartUploads(w http.ResponseWriter, r *http.Re
|
|||||||
vars := mux.Vars(r)
|
vars := mux.Vars(r)
|
||||||
bucketName := vars["bucket"]
|
bucketName := vars["bucket"]
|
||||||
key := vars["key"]
|
key := vars["key"]
|
||||||
|
|
||||||
// 检查bucket是否存在
|
// 检查bucket是否存在
|
||||||
if _, ok := h.bucketManager.GetBucket(bucketName); !ok {
|
if _, ok := h.bucketManager.GetBucket(bucketName); !ok {
|
||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 简化实现:返回空列表
|
// 简化实现:返回空列表
|
||||||
result := ListMultipartUploadsResult{
|
result := ListMultipartUploadsResult{
|
||||||
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
||||||
@@ -782,7 +775,7 @@ func (h *S3Handler) handleListMultipartUploads(w http.ResponseWriter, r *http.Re
|
|||||||
IsTruncated: false,
|
IsTruncated: false,
|
||||||
Uploads: make([]Upload, 0),
|
Uploads: make([]Upload, 0),
|
||||||
}
|
}
|
||||||
|
|
||||||
h.sendXMLResponse(w, http.StatusOK, result)
|
h.sendXMLResponse(w, http.StatusOK, result)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -792,25 +785,25 @@ func (h *S3Handler) handleListMultipartParts(w http.ResponseWriter, r *http.Requ
|
|||||||
bucketName := vars["bucket"]
|
bucketName := vars["bucket"]
|
||||||
key := vars["key"]
|
key := vars["key"]
|
||||||
uploadID := r.URL.Query().Get("uploadId")
|
uploadID := r.URL.Query().Get("uploadId")
|
||||||
|
|
||||||
// 检查bucket是否存在
|
// 检查bucket是否存在
|
||||||
if _, ok := h.bucketManager.GetBucket(bucketName); !ok {
|
if _, ok := h.bucketManager.GetBucket(bucketName); !ok {
|
||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 简化实现:返回空列表
|
// 简化实现:返回空列表
|
||||||
result := ListPartsResult{
|
result := ListPartsResult{
|
||||||
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
||||||
Bucket: bucketName,
|
Bucket: bucketName,
|
||||||
Key: key,
|
Key: key,
|
||||||
UploadID: uploadID,
|
UploadID: uploadID,
|
||||||
PartNumberMarker: 0,
|
PartNumberMarker: 0,
|
||||||
MaxParts: 1000,
|
MaxParts: 1000,
|
||||||
IsTruncated: false,
|
IsTruncated: false,
|
||||||
Parts: make([]Part, 0),
|
Parts: make([]Part, 0),
|
||||||
}
|
}
|
||||||
|
|
||||||
h.sendXMLResponse(w, http.StatusOK, result)
|
h.sendXMLResponse(w, http.StatusOK, result)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -820,19 +813,19 @@ func (h *S3Handler) handleCompleteMultipartUpload(w http.ResponseWriter, r *http
|
|||||||
bucketName := vars["bucket"]
|
bucketName := vars["bucket"]
|
||||||
key := vars["key"]
|
key := vars["key"]
|
||||||
uploadID := r.URL.Query().Get("uploadId")
|
uploadID := r.URL.Query().Get("uploadId")
|
||||||
|
|
||||||
// 查找对象所在的实际存储桶(简化实现,使用配置的bucket)
|
// 查找对象所在的实际存储桶(简化实现,使用配置的bucket)
|
||||||
bucket, ok := h.bucketManager.GetBucket(bucketName)
|
bucket, ok := h.bucketManager.GetBucket(bucketName)
|
||||||
if !ok {
|
if !ok {
|
||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 解析请求体以获取分片列表
|
// 解析请求体以获取分片列表
|
||||||
var completeReq CompleteMultipartUpload
|
var completeReq CompleteMultipartUpload
|
||||||
body, _ := io.ReadAll(r.Body)
|
body, _ := io.ReadAll(r.Body)
|
||||||
xml.Unmarshal(body, &completeReq)
|
xml.Unmarshal(body, &completeReq)
|
||||||
|
|
||||||
// 完成分片上传
|
// 完成分片上传
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
var parts []types.CompletedPart
|
var parts []types.CompletedPart
|
||||||
@@ -842,7 +835,7 @@ func (h *S3Handler) handleCompleteMultipartUpload(w http.ResponseWriter, r *http
|
|||||||
PartNumber: aws.Int32(int32(part.PartNumber)),
|
PartNumber: aws.Int32(int32(part.PartNumber)),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
completeResp, err := bucket.Client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{
|
completeResp, err := bucket.Client.CompleteMultipartUpload(ctx, &s3.CompleteMultipartUploadInput{
|
||||||
Bucket: aws.String(bucket.Config.Name),
|
Bucket: aws.String(bucket.Config.Name),
|
||||||
Key: aws.String(key),
|
Key: aws.String(key),
|
||||||
@@ -855,18 +848,18 @@ func (h *S3Handler) handleCompleteMultipartUpload(w http.ResponseWriter, r *http
|
|||||||
h.sendS3Error(w, "InternalError", "Failed to complete multipart upload", key)
|
h.sendS3Error(w, "InternalError", "Failed to complete multipart upload", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
result := CompleteMultipartUploadResult{
|
result := CompleteMultipartUploadResult{
|
||||||
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
Xmlns: "http://s3.amazonaws.com/doc/2006-03-01/",
|
||||||
Location: "/" + bucket.Config.Name + "/" + key,
|
Location: "/" + bucket.Config.Name + "/" + key,
|
||||||
Bucket: bucket.Config.Name,
|
Bucket: bucket.Config.Name,
|
||||||
Key: key,
|
Key: key,
|
||||||
ETag: *completeResp.ETag,
|
ETag: *completeResp.ETag,
|
||||||
}
|
}
|
||||||
|
|
||||||
// 记录对象元数据(简化:假设总大小)
|
// 记录对象元数据(简化:假设总大小)
|
||||||
h.storage.RecordObject(key, bucket.Config.Name, 0, nil)
|
h.storage.RecordObject(key, bucket.Config.Name, 0, nil)
|
||||||
|
|
||||||
h.sendXMLResponse(w, http.StatusOK, result)
|
h.sendXMLResponse(w, http.StatusOK, result)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -876,14 +869,14 @@ func (h *S3Handler) handleAbortMultipartUpload(w http.ResponseWriter, r *http.Re
|
|||||||
bucketName := vars["bucket"]
|
bucketName := vars["bucket"]
|
||||||
key := vars["key"]
|
key := vars["key"]
|
||||||
uploadID := r.URL.Query().Get("uploadId")
|
uploadID := r.URL.Query().Get("uploadId")
|
||||||
|
|
||||||
// 查找对象所在的实际存储桶(简化实现,使用配置的bucket)
|
// 查找对象所在的实际存储桶(简化实现,使用配置的bucket)
|
||||||
bucket, ok := h.bucketManager.GetBucket(bucketName)
|
bucket, ok := h.bucketManager.GetBucket(bucketName)
|
||||||
if !ok {
|
if !ok {
|
||||||
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
h.sendS3Error(w, "NoSuchBucket", "The specified bucket does not exist", bucketName)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 中止分片上传
|
// 中止分片上传
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
_, err := bucket.Client.AbortMultipartUpload(ctx, &s3.AbortMultipartUploadInput{
|
_, err := bucket.Client.AbortMultipartUpload(ctx, &s3.AbortMultipartUploadInput{
|
||||||
@@ -895,7 +888,7 @@ func (h *S3Handler) handleAbortMultipartUpload(w http.ResponseWriter, r *http.Re
|
|||||||
h.sendS3Error(w, "InternalError", "Failed to abort multipart upload", key)
|
h.sendS3Error(w, "InternalError", "Failed to abort multipart upload", key)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
w.WriteHeader(http.StatusNoContent)
|
w.WriteHeader(http.StatusNoContent)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -903,13 +896,13 @@ func (h *S3Handler) handleAbortMultipartUpload(w http.ResponseWriter, r *http.Re
|
|||||||
func (h *S3Handler) sendXMLResponse(w http.ResponseWriter, statusCode int, data interface{}) {
|
func (h *S3Handler) sendXMLResponse(w http.ResponseWriter, statusCode int, data interface{}) {
|
||||||
w.Header().Set("Content-Type", "application/xml")
|
w.Header().Set("Content-Type", "application/xml")
|
||||||
w.WriteHeader(statusCode)
|
w.WriteHeader(statusCode)
|
||||||
|
|
||||||
encoder := xml.NewEncoder(w)
|
encoder := xml.NewEncoder(w)
|
||||||
encoder.Indent("", " ")
|
encoder.Indent("", " ")
|
||||||
|
|
||||||
// 写入XML声明
|
// 写入XML声明
|
||||||
w.Write([]byte(xml.Header))
|
w.Write([]byte(xml.Header))
|
||||||
|
|
||||||
if err := encoder.Encode(data); err != nil {
|
if err := encoder.Encode(data); err != nil {
|
||||||
// 如果编码失败,记录错误
|
// 如果编码失败,记录错误
|
||||||
http.Error(w, "Internal Server Error", http.StatusInternalServerError)
|
http.Error(w, "Internal Server Error", http.StatusInternalServerError)
|
||||||
@@ -924,7 +917,7 @@ func (h *S3Handler) sendS3Error(w http.ResponseWriter, code string, message stri
|
|||||||
Resource: resource,
|
Resource: resource,
|
||||||
RequestID: fmt.Sprintf("%d", time.Now().UnixNano()),
|
RequestID: fmt.Sprintf("%d", time.Now().UnixNano()),
|
||||||
}
|
}
|
||||||
|
|
||||||
statusCode := http.StatusBadRequest
|
statusCode := http.StatusBadRequest
|
||||||
switch code {
|
switch code {
|
||||||
case "NoSuchBucket", "NoSuchKey":
|
case "NoSuchBucket", "NoSuchKey":
|
||||||
@@ -938,7 +931,7 @@ func (h *S3Handler) sendS3Error(w http.ResponseWriter, code string, message stri
|
|||||||
case "InsufficientStorage":
|
case "InsufficientStorage":
|
||||||
statusCode = http.StatusInsufficientStorage
|
statusCode = http.StatusInsufficientStorage
|
||||||
}
|
}
|
||||||
|
|
||||||
h.sendXMLResponse(w, statusCode, errorResp)
|
h.sendXMLResponse(w, statusCode, errorResp)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -946,14 +939,14 @@ func (h *S3Handler) sendS3Error(w http.ResponseWriter, code string, message stri
|
|||||||
func parseS3Path(requestPath string) (bucket string, key string) {
|
func parseS3Path(requestPath string) (bucket string, key string) {
|
||||||
requestPath = strings.TrimPrefix(requestPath, "/")
|
requestPath = strings.TrimPrefix(requestPath, "/")
|
||||||
parts := strings.SplitN(requestPath, "/", 2)
|
parts := strings.SplitN(requestPath, "/", 2)
|
||||||
|
|
||||||
if len(parts) > 0 {
|
if len(parts) > 0 {
|
||||||
bucket = parts[0]
|
bucket = parts[0]
|
||||||
}
|
}
|
||||||
if len(parts) > 1 {
|
if len(parts) > 1 {
|
||||||
key = parts[1]
|
key = parts[1]
|
||||||
}
|
}
|
||||||
|
|
||||||
return bucket, key
|
return bucket, key
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -11,7 +11,7 @@ import (
|
|||||||
// Object 对象信息模型
|
// Object 对象信息模型
|
||||||
type Object struct {
|
type Object struct {
|
||||||
ID uint `gorm:"primaryKey" json:"id"`
|
ID uint `gorm:"primaryKey" json:"id"`
|
||||||
Key string `gorm:"uniqueIndex;size:512;not null" json:"key"`
|
Key string `gorm:"size:512;not null" json:"key"`
|
||||||
BucketName string `gorm:"index;size:255;not null" json:"bucket_name"`
|
BucketName string `gorm:"index;size:255;not null" json:"bucket_name"`
|
||||||
Size int64 `gorm:"not null;default:0" json:"size"`
|
Size int64 `gorm:"not null;default:0" json:"size"`
|
||||||
Metadata JSON `gorm:"type:json" json:"metadata,omitempty"`
|
Metadata JSON `gorm:"type:json" json:"metadata,omitempty"`
|
||||||
|
|||||||
@@ -21,6 +21,15 @@ func NewService(db *gorm.DB) *Service {
|
|||||||
|
|
||||||
// RecordObject 记录对象信息
|
// RecordObject 记录对象信息
|
||||||
func (s *Service) RecordObject(key, bucketName string, size int64, metadata map[string]string) error {
|
func (s *Service) RecordObject(key, bucketName string, size int64, metadata map[string]string) error {
|
||||||
|
// 首先检查是否存在已删除的同名对象
|
||||||
|
var deletedObj Object
|
||||||
|
if err := s.db.Unscoped().Where("`key` = ?", key).Where("`deleted_at` IS NOT NULL").First(&deletedObj).Error; err == nil {
|
||||||
|
// 存在已删除的同名对象,永久删除它
|
||||||
|
if err := s.db.Unscoped().Delete(&deletedObj).Error; err != nil {
|
||||||
|
return fmt.Errorf("failed to permanently delete soft-deleted object: %w", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
obj := &Object{
|
obj := &Object{
|
||||||
Key: key,
|
Key: key,
|
||||||
BucketName: bucketName,
|
BucketName: bucketName,
|
||||||
@@ -37,7 +46,7 @@ func (s *Service) RecordObject(key, bucketName string, size int64, metadata map[
|
|||||||
}
|
}
|
||||||
|
|
||||||
// 使用 Upsert(更新或插入)
|
// 使用 Upsert(更新或插入)
|
||||||
result := s.db.Where("`key` = ?", key).FirstOrCreate(&obj)
|
result := s.db.Where("`key` = ?", key).Where("`deleted_at` IS NULL").FirstOrCreate(&obj)
|
||||||
if result.Error != nil {
|
if result.Error != nil {
|
||||||
return fmt.Errorf("failed to record object: %w", result.Error)
|
return fmt.Errorf("failed to record object: %w", result.Error)
|
||||||
}
|
}
|
||||||
@@ -197,7 +206,10 @@ func (s *Service) updateBucketStats(bucketName string) error {
|
|||||||
|
|
||||||
// 获取或创建统计记录
|
// 获取或创建统计记录
|
||||||
result := s.db.Where("bucket_name = ?", bucketName).FirstOrCreate(&stats, BucketStats{
|
result := s.db.Where("bucket_name = ?", bucketName).FirstOrCreate(&stats, BucketStats{
|
||||||
BucketName: bucketName,
|
BucketName: bucketName,
|
||||||
|
LastCheckedAt: time.Now(),
|
||||||
|
CreatedAt: time.Now(),
|
||||||
|
UpdatedAt: time.Now(),
|
||||||
})
|
})
|
||||||
if result.Error != nil {
|
if result.Error != nil {
|
||||||
return fmt.Errorf("failed to get bucket stats: %w", result.Error)
|
return fmt.Errorf("failed to get bucket stats: %w", result.Error)
|
||||||
|
|||||||
Reference in New Issue
Block a user