mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-05 23:26:38 +08:00
feat(upload): add programmatic Ingest service and file-upload skill
Introduce upload/ingest as the single domain entry for storing files, writing w_uploads records, and maintaining incremental stats. Refactor HTTP UploadFile and delete handlers to delegate to ingest, fix stats decrement ordering on Remove, and document usage in the file-upload skill.
This commit is contained in:
@@ -7,6 +7,7 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/cache"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/handler"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
uploadtask "github.com/Rain-kl/Wavelet/internal/apps/upload/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/util"
|
||||
@@ -28,6 +29,36 @@ var (
|
||||
ServeFileByID = filesrv.ServeFileByID
|
||||
)
|
||||
|
||||
// Programmatic ingest API
|
||||
var (
|
||||
Ingest = ingest.Ingest
|
||||
Remove = ingest.Remove
|
||||
RemoveOwned = ingest.RemoveOwned
|
||||
FindByHash = ingest.FindByHash
|
||||
)
|
||||
|
||||
// Ingest policy constants
|
||||
const (
|
||||
PolicyCreate = ingest.PolicyCreate
|
||||
PolicyDedupNewRecord = ingest.PolicyDedupNewRecord
|
||||
PolicyResolveExisting = ingest.PolicyResolveExisting
|
||||
)
|
||||
|
||||
type (
|
||||
// IngestRequest is the programmatic upload ingest payload.
|
||||
IngestRequest = ingest.Request
|
||||
// IngestResult reports ingest side effects.
|
||||
IngestResult = ingest.Result
|
||||
// IngestPolicy controls hash-collision behavior during ingest.
|
||||
IngestPolicy = ingest.Policy
|
||||
)
|
||||
|
||||
// Ingest errors
|
||||
var (
|
||||
ErrIngestForbidden = ingest.ErrForbidden
|
||||
ErrIngestStorageReadOnly = ingest.ErrStorageReadOnly
|
||||
)
|
||||
|
||||
// Cache management
|
||||
var (
|
||||
ResetAccessCaches = cache.ResetAccessCaches
|
||||
@@ -36,7 +67,9 @@ var (
|
||||
|
||||
// Stats
|
||||
var (
|
||||
ApplyUploadStatsAdd = uploadstats.ApplyUploadStatsAdd
|
||||
// Deprecated: use upload.Ingest or upload.Remove; stats are applied internally.
|
||||
ApplyUploadStatsAdd = uploadstats.ApplyUploadStatsAdd
|
||||
// Deprecated: use upload.Ingest or upload.Remove; stats are applied internally.
|
||||
ApplyUploadStatsRemove = uploadstats.ApplyUploadStatsRemove
|
||||
RebuildUploadStats = uploadstats.RebuildUploadStats
|
||||
)
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"strconv"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/oauth"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
@@ -235,7 +236,7 @@ func DeleteMyFile(c *gin.Context) {
|
||||
c.AbortWithStatus(http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
if err == errUploadForbidden {
|
||||
if err == ingest.ErrForbidden {
|
||||
c.AbortWithStatus(http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
@@ -289,7 +290,7 @@ func UpdateMyFile(c *gin.Context) {
|
||||
c.AbortWithStatus(http.StatusNotFound)
|
||||
return
|
||||
}
|
||||
if err == errUploadForbidden {
|
||||
if err == ingest.ErrForbidden {
|
||||
c.AbortWithStatus(http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -4,20 +4,13 @@
|
||||
package handler
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"sort"
|
||||
"strings"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
@@ -31,30 +24,11 @@ func listMyUploadFiles(ctx context.Context, userID uint64, filter repository.Upl
|
||||
}
|
||||
|
||||
func softDeleteUpload(ctx context.Context, uploadID uint64) (model.Upload, error) {
|
||||
upload, err := repository.GetActiveUploadByID(ctx, uploadID)
|
||||
if err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
if err := repository.SoftDeleteUpload(ctx, &upload); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
uploadstats.RecordUploadStatsRemove(ctx, &upload)
|
||||
return upload, nil
|
||||
return ingest.Remove(ctx, uploadID)
|
||||
}
|
||||
|
||||
func softDeleteOwnedUpload(ctx context.Context, userID, uploadID uint64) (model.Upload, error) {
|
||||
upload, err := repository.GetActiveUploadByID(ctx, uploadID)
|
||||
if err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
if upload.UserID != userID {
|
||||
return model.Upload{}, errUploadForbidden
|
||||
}
|
||||
if err := repository.SoftDeleteUpload(ctx, &upload); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
uploadstats.RecordUploadStatsRemove(ctx, &upload)
|
||||
return upload, nil
|
||||
return ingest.RemoveOwned(ctx, userID, uploadID)
|
||||
}
|
||||
|
||||
func listDistinctUploadTypes(ctx context.Context) ([]string, error) {
|
||||
@@ -77,7 +51,7 @@ func updateOwnedUpload(ctx context.Context, userID, uploadID uint64, input updat
|
||||
return model.Upload{}, err
|
||||
}
|
||||
if upload.UserID != userID {
|
||||
return model.Upload{}, errUploadForbidden
|
||||
return model.Upload{}, ingest.ErrForbidden
|
||||
}
|
||||
|
||||
updates := make(map[string]any)
|
||||
@@ -103,96 +77,10 @@ func listUploadsForBatchDownload(ctx context.Context, ids []uint64) ([]model.Upl
|
||||
return repository.ListUploadsByIDs(ctx, ids)
|
||||
}
|
||||
|
||||
type instantUploadInput struct {
|
||||
UserID uint64
|
||||
FileHash string
|
||||
Size int64
|
||||
MimeType string
|
||||
Extension string
|
||||
OrigName string
|
||||
UploadType string
|
||||
AccessMode int
|
||||
}
|
||||
|
||||
func createInstantUpload(ctx context.Context, existing model.Upload, input instantUploadInput) (model.Upload, error) {
|
||||
newUpload := model.Upload{
|
||||
ID: idgen.NextUint64ID(),
|
||||
UserID: input.UserID,
|
||||
FileName: input.OrigName,
|
||||
FilePath: existing.FilePath,
|
||||
FileSize: input.Size,
|
||||
MimeType: input.MimeType,
|
||||
Extension: input.Extension,
|
||||
Hash: input.FileHash,
|
||||
Type: input.UploadType,
|
||||
Status: model.UploadStatusUsed,
|
||||
AccessMode: input.AccessMode,
|
||||
Metadata: existing.Metadata,
|
||||
}
|
||||
if err := repository.CreateUpload(ctx, &newUpload); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
uploadstats.RecordUploadStatsAdd(ctx, &newUpload)
|
||||
logger.InfoF(ctx, "文件触发秒传成功! ID: %d, Path: %s", newUpload.ID, existing.FilePath)
|
||||
return newUpload, nil
|
||||
}
|
||||
|
||||
func findReusableUpload(ctx context.Context, hash string, size int64) (model.Upload, error) {
|
||||
return repository.FindReusableUploadByHash(ctx, hash, size)
|
||||
}
|
||||
|
||||
func saveNewUploadRecord(ctx context.Context, upload *model.Upload, filePath string) error {
|
||||
if err := repository.CreateUpload(ctx, upload); err != nil {
|
||||
_, backend, backendErr := storage.Active(ctx)
|
||||
if backendErr == nil {
|
||||
if deleteErr := backend.Delete(ctx, filePath); deleteErr != nil {
|
||||
logger.WarnF(ctx, "清理未写入数据库的上传对象失败: %v", deleteErr)
|
||||
}
|
||||
}
|
||||
return err
|
||||
}
|
||||
uploadstats.RecordUploadStatsAdd(ctx, upload)
|
||||
return nil
|
||||
}
|
||||
|
||||
func loadUploadStats(ctx context.Context) ([]model.UploadStat, error) {
|
||||
return repository.ListUploadStats(ctx)
|
||||
}
|
||||
|
||||
var errUploadForbidden = errors.New("upload forbidden")
|
||||
|
||||
func storeUploadObject(ctx context.Context, subPath string, size int64, mimeType string, buf *bytes.Buffer, meta *model.UploadMetadata) (string, error) {
|
||||
if uploadstorage.ReadOnly(ctx) {
|
||||
return "", errors.New(shared.ErrStorageReadOnly)
|
||||
}
|
||||
driver, backend, err := storage.Active(ctx)
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "初始化活动存储失败: %v", err)
|
||||
return "", errors.New(shared.ErrSaveFileFailed)
|
||||
}
|
||||
result, err := backend.Put(ctx, subPath, bytes.NewReader(buf.Bytes()), size, mimeType)
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "写入 %s 存储失败: %v", driver, err)
|
||||
return "", errors.New(shared.ErrSaveFileFailed)
|
||||
}
|
||||
meta.Bucket = result.Bucket
|
||||
return result.Key, nil
|
||||
}
|
||||
|
||||
func validateUploadAllowedExtension(ctx context.Context, ext string) string {
|
||||
sc, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUploadAllowedExtensions)
|
||||
if err != nil || sc.Value == "" {
|
||||
return ""
|
||||
}
|
||||
allowedExts := strings.Split(strings.ToLower(sc.Value), ",")
|
||||
for _, allowedExt := range allowedExts {
|
||||
if strings.TrimSpace(allowedExt) == ext {
|
||||
return ""
|
||||
}
|
||||
}
|
||||
return shared.ErrUnsupportedFormat
|
||||
}
|
||||
|
||||
func isRecordNotFound(err error) bool {
|
||||
return errors.Is(err, gorm.ErrRecordNotFound)
|
||||
}
|
||||
@@ -8,7 +8,6 @@ package handler
|
||||
import (
|
||||
"archive/zip"
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
@@ -21,16 +20,15 @@ import (
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/oauth"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/filesrv"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/ingest"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/util"
|
||||
"github.com/Rain-kl/Wavelet/internal/common"
|
||||
"github.com/Rain-kl/Wavelet/internal/common/response"
|
||||
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"github.com/gin-gonic/gin"
|
||||
@@ -91,11 +89,6 @@ func UploadFile(c *gin.Context) {
|
||||
ext = "bin"
|
||||
}
|
||||
|
||||
if errMsg := validateUploadAllowedExtension(ctx, ext); errMsg != "" {
|
||||
response.AbortBadRequest(c, errMsg)
|
||||
return
|
||||
}
|
||||
|
||||
hashWriter := sha256.New()
|
||||
var buf bytes.Buffer
|
||||
size, err := io.Copy(&buf, io.TeeReader(file, hashWriter))
|
||||
@@ -120,51 +113,43 @@ func UploadFile(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
handled, lookupErr := tryInstantUpload(ctx, c, currUser, fileHash, size, mimeType, ext, origName, accessMode)
|
||||
if handled {
|
||||
return
|
||||
}
|
||||
if lookupErr != nil && !errors.Is(lookupErr, gorm.ErrRecordNotFound) {
|
||||
response.AbortBadRequest(c, shared.ErrFileValidationFailed)
|
||||
return
|
||||
}
|
||||
|
||||
meta, errMsg := parseUploadMetadata(c, mimeType)
|
||||
if errMsg != "" {
|
||||
response.AbortBadRequest(c, errMsg)
|
||||
return
|
||||
}
|
||||
|
||||
id := idgen.NextUint64ID()
|
||||
subPath := fmt.Sprintf("uploads/%s/%d.%s", time.Now().Format("2006/01/02"), id, ext)
|
||||
|
||||
subPath, err = storeUploadObject(ctx, subPath, size, mimeType, &buf, &meta)
|
||||
if err != nil {
|
||||
response.AbortBadRequest(c, err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
newUpload := model.Upload{
|
||||
ID: id,
|
||||
result, err := ingest.Ingest(ctx, ingest.Request{
|
||||
UserID: currUser.ID,
|
||||
Reader: bytes.NewReader(buf.Bytes()),
|
||||
Size: size,
|
||||
FileName: origName,
|
||||
FilePath: subPath,
|
||||
FileSize: size,
|
||||
MimeType: mimeType,
|
||||
Extension: ext,
|
||||
Hash: fileHash,
|
||||
Type: uploadType,
|
||||
Status: model.UploadStatusUsed,
|
||||
AccessMode: accessMode,
|
||||
AccessMode: &accessMode,
|
||||
Metadata: meta,
|
||||
}
|
||||
|
||||
if err := saveNewUploadRecord(ctx, &newUpload, subPath); err != nil {
|
||||
Policy: ingest.PolicyDedupNewRecord,
|
||||
})
|
||||
if err != nil {
|
||||
if errors.Is(err, ingest.ErrStorageReadOnly) {
|
||||
response.AbortConflict(c, shared.ErrStorageReadOnly)
|
||||
return
|
||||
}
|
||||
if err.Error() == shared.ErrUnsupportedFormat {
|
||||
response.AbortBadRequest(c, shared.ErrUnsupportedFormat)
|
||||
return
|
||||
}
|
||||
if err.Error() == shared.ErrSaveFileFailed {
|
||||
response.AbortBadRequest(c, shared.ErrSaveFileFailed)
|
||||
return
|
||||
}
|
||||
response.AbortBadRequest(c, shared.ErrSaveUploadRecordFailed)
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, response.OK(newUpload))
|
||||
c.JSON(http.StatusOK, response.OK(result.Upload))
|
||||
}
|
||||
|
||||
// DownloadFile 通用单文件下载接口
|
||||
@@ -320,34 +305,6 @@ func resolveUploadAccessMode(c *gin.Context, uploadType string) (int, string) {
|
||||
return accessMode, ""
|
||||
}
|
||||
|
||||
func tryInstantUpload(ctx context.Context, c *gin.Context, currUser *model.User, fileHash string, size int64, mimeType, ext, origName string, accessMode int) (bool, error) {
|
||||
existing, err := findReusableUpload(ctx, fileHash, size)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if uploadstorage.ReadOnly(ctx) {
|
||||
response.AbortConflict(c, shared.ErrStorageReadOnly)
|
||||
return true, nil
|
||||
}
|
||||
|
||||
newUpload, err := createInstantUpload(ctx, existing, instantUploadInput{
|
||||
UserID: currUser.ID,
|
||||
FileHash: fileHash,
|
||||
Size: size,
|
||||
MimeType: mimeType,
|
||||
Extension: ext,
|
||||
OrigName: origName,
|
||||
UploadType: c.DefaultPostForm("type", "generic"),
|
||||
AccessMode: accessMode,
|
||||
})
|
||||
if err != nil {
|
||||
response.AbortBadRequest(c, shared.ErrSaveUploadRecordFailed)
|
||||
return true, err
|
||||
}
|
||||
c.JSON(http.StatusOK, response.OK(newUpload))
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func parseUploadMetadata(c *gin.Context, mimeType string) (model.UploadMetadata, string) {
|
||||
var meta model.UploadMetadata
|
||||
metadataStr := c.DefaultPostForm("metadata", "")
|
||||
|
||||
@@ -0,0 +1,16 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"errors"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
)
|
||||
|
||||
// ErrForbidden indicates the caller is not allowed to mutate the upload record.
|
||||
var ErrForbidden = errors.New("upload forbidden")
|
||||
|
||||
// ErrStorageReadOnly indicates the storage backend is in migration read-only mode.
|
||||
var ErrStorageReadOnly = errors.New(shared.ErrStorageReadOnly)
|
||||
@@ -0,0 +1,187 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload/shared"
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
uploadstorage "github.com/Rain-kl/Wavelet/internal/apps/upload/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func normalizeRequest(req *Request) {
|
||||
req.Extension = strings.ToLower(strings.TrimSpace(req.Extension))
|
||||
if req.Extension == "" {
|
||||
req.Extension = "bin"
|
||||
}
|
||||
if req.Type == "" {
|
||||
req.Type = "generic"
|
||||
}
|
||||
if req.Status == "" {
|
||||
req.Status = model.UploadStatusUsed
|
||||
}
|
||||
}
|
||||
|
||||
func resolveAccessMode(uploadType string, explicit *int) int {
|
||||
if explicit != nil {
|
||||
return *explicit
|
||||
}
|
||||
if uploadType == shared.DefaultPublicUploadType {
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func validateAllowedExtension(ctx context.Context, ext string) error {
|
||||
sc, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyUploadAllowedExtensions)
|
||||
if err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
if sc.Value == "" {
|
||||
return nil
|
||||
}
|
||||
|
||||
allowedExts := strings.Split(strings.ToLower(sc.Value), ",")
|
||||
for _, allowedExt := range allowedExts {
|
||||
if strings.TrimSpace(allowedExt) == ext {
|
||||
return nil
|
||||
}
|
||||
}
|
||||
return errors.New(shared.ErrUnsupportedFormat)
|
||||
}
|
||||
|
||||
func defaultObjectKey(id uint64, ext string) string {
|
||||
return fmt.Sprintf("uploads/%s/%d.%s", time.Now().Format("2006/01/02"), id, ext)
|
||||
}
|
||||
|
||||
func buildObjectKey(req Request, id uint64) string {
|
||||
if req.ObjectKeyFn != nil {
|
||||
return req.ObjectKeyFn(id, req.Extension)
|
||||
}
|
||||
return defaultObjectKey(id, req.Extension)
|
||||
}
|
||||
|
||||
func storeObject(ctx context.Context, objectKey string, reader io.Reader, size int64, mimeType string, meta *model.UploadMetadata) (string, error) {
|
||||
if uploadstorage.ReadOnly(ctx) {
|
||||
return "", ErrStorageReadOnly
|
||||
}
|
||||
|
||||
driver, backend, err := storage.Active(ctx)
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "初始化活动存储失败: %v", err)
|
||||
return "", errors.New(shared.ErrSaveFileFailed)
|
||||
}
|
||||
|
||||
result, err := backend.Put(ctx, objectKey, reader, size, mimeType)
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "写入 %s 存储失败: %v", driver, err)
|
||||
return "", errors.New(shared.ErrSaveFileFailed)
|
||||
}
|
||||
|
||||
meta.Bucket = result.Bucket
|
||||
return result.Key, nil
|
||||
}
|
||||
|
||||
func persistUploadRecord(ctx context.Context, upload *model.Upload, objectKey string) error {
|
||||
if err := repository.CreateUpload(ctx, upload); err != nil {
|
||||
_, backend, backendErr := storage.Active(ctx)
|
||||
if backendErr == nil {
|
||||
if deleteErr := backend.Delete(ctx, objectKey); deleteErr != nil {
|
||||
logger.WarnF(ctx, "清理未写入数据库的上传对象失败: %v", deleteErr)
|
||||
}
|
||||
}
|
||||
return err
|
||||
}
|
||||
uploadstats.RecordUploadStatsAdd(ctx, upload)
|
||||
return nil
|
||||
}
|
||||
|
||||
func createDedupRecord(ctx context.Context, existing model.Upload, req Request) (Result, error) {
|
||||
accessMode := resolveAccessMode(req.Type, req.AccessMode)
|
||||
newUpload := model.Upload{
|
||||
ID: idgen.NextUint64ID(),
|
||||
UserID: req.UserID,
|
||||
FileName: req.FileName,
|
||||
FilePath: existing.FilePath,
|
||||
FileSize: req.Size,
|
||||
MimeType: req.MimeType,
|
||||
Extension: req.Extension,
|
||||
Hash: req.Hash,
|
||||
Type: req.Type,
|
||||
Status: req.Status,
|
||||
AccessMode: accessMode,
|
||||
Metadata: existing.Metadata,
|
||||
}
|
||||
if err := persistUploadRecord(ctx, &newUpload, existing.FilePath); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
logger.InfoF(ctx, "文件触发秒传成功! ID: %d, Path: %s", newUpload.ID, existing.FilePath)
|
||||
return Result{
|
||||
Upload: newUpload,
|
||||
Created: true,
|
||||
Stored: false,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func uploadstorageReadOnly(ctx context.Context) bool {
|
||||
return uploadstorage.ReadOnly(ctx)
|
||||
}
|
||||
|
||||
func createNewUpload(ctx context.Context, req Request) (Result, error) {
|
||||
if uploadstorageReadOnly(ctx) {
|
||||
return Result{}, ErrStorageReadOnly
|
||||
}
|
||||
if !req.SkipExtensionCheck {
|
||||
if err := validateAllowedExtension(ctx, req.Extension); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
}
|
||||
|
||||
id := idgen.NextUint64ID()
|
||||
objectKey := buildObjectKey(req, id)
|
||||
storedKey, err := storeObject(ctx, objectKey, req.Reader, req.Size, req.MimeType, &req.Metadata)
|
||||
if err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
|
||||
accessMode := resolveAccessMode(req.Type, req.AccessMode)
|
||||
upload := model.Upload{
|
||||
ID: id,
|
||||
UserID: req.UserID,
|
||||
FileName: req.FileName,
|
||||
FilePath: storedKey,
|
||||
FileSize: req.Size,
|
||||
MimeType: req.MimeType,
|
||||
Extension: req.Extension,
|
||||
Hash: req.Hash,
|
||||
Type: req.Type,
|
||||
Status: req.Status,
|
||||
AccessMode: accessMode,
|
||||
Metadata: req.Metadata,
|
||||
}
|
||||
if err := persistUploadRecord(ctx, &upload, storedKey); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
|
||||
return Result{
|
||||
Upload: upload,
|
||||
Created: true,
|
||||
Stored: true,
|
||||
}, nil
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// Ingest stores or resolves an upload using the configured policy and side effects.
|
||||
func Ingest(ctx context.Context, req Request) (Result, error) {
|
||||
normalizeRequest(&req)
|
||||
if req.Hash == "" {
|
||||
return Result{}, errors.New("ingest hash is required")
|
||||
}
|
||||
if req.Reader == nil {
|
||||
return Result{}, errors.New("ingest reader is required")
|
||||
}
|
||||
if req.Size < 0 {
|
||||
return Result{}, errors.New("ingest size must be non-negative")
|
||||
}
|
||||
|
||||
switch req.Policy {
|
||||
case PolicyDedupNewRecord, PolicyResolveExisting:
|
||||
return ingestWithHashPolicy(ctx, req)
|
||||
case PolicyCreate:
|
||||
return createNewUpload(ctx, req)
|
||||
default:
|
||||
return Result{}, errors.New("unsupported ingest policy")
|
||||
}
|
||||
}
|
||||
|
||||
// FindByHash returns a reusable active upload with the same hash and size.
|
||||
func FindByHash(ctx context.Context, hash string, size int64) (model.Upload, error) {
|
||||
return repository.FindReusableUploadByHash(ctx, hash, size)
|
||||
}
|
||||
|
||||
func ingestWithHashPolicy(ctx context.Context, req Request) (Result, error) {
|
||||
existing, err := repository.FindReusableUploadByHash(ctx, req.Hash, req.Size)
|
||||
if err == nil {
|
||||
switch req.Policy {
|
||||
case PolicyResolveExisting:
|
||||
return Result{
|
||||
Upload: existing,
|
||||
Resolved: true,
|
||||
}, nil
|
||||
case PolicyDedupNewRecord:
|
||||
if uploadstorageReadOnly(ctx) {
|
||||
return Result{}, ErrStorageReadOnly
|
||||
}
|
||||
return createDedupRecord(ctx, existing, req)
|
||||
}
|
||||
}
|
||||
if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return Result{}, err
|
||||
}
|
||||
|
||||
return createNewUpload(ctx, req)
|
||||
}
|
||||
@@ -0,0 +1,283 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"io"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/storage"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
)
|
||||
|
||||
func TestIngestPolicyCreateIncrementsStats(t *testing.T) {
|
||||
_, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01")
|
||||
hash := sha256.Sum256(content)
|
||||
|
||||
restoreStorage, disableStorage := setupMockStorage(t, nil)
|
||||
defer restoreStorage()
|
||||
defer disableStorage()
|
||||
|
||||
result, err := Ingest(ctx, Request{
|
||||
UserID: 1001,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "mirror.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hex.EncodeToString(hash[:]),
|
||||
Type: "pixez_mirror",
|
||||
Policy: PolicyCreate,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Ingest(PolicyCreate) returned error: %v", err)
|
||||
}
|
||||
if !result.Created || !result.Stored || result.Resolved {
|
||||
t.Fatalf("Ingest(PolicyCreate) = %+v, want Created+Stored without Resolved", result)
|
||||
}
|
||||
|
||||
stats, err := loadTotalStats(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("loadTotalStats returned error: %v", err)
|
||||
}
|
||||
if stats.TotalCount != 1 || stats.TotalSize != int64(len(content)) {
|
||||
t.Fatalf("loadTotalStats() = count %d size %d, want count 1 size %d", stats.TotalCount, stats.TotalSize, len(content))
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngestPolicyResolveExistingSkipsStatsOnHit(t *testing.T) {
|
||||
dbConn, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01")
|
||||
hash := sha256.Sum256(content)
|
||||
hashStr := hex.EncodeToString(hash[:])
|
||||
|
||||
existing := model.Upload{
|
||||
ID: 88001,
|
||||
UserID: 42,
|
||||
FileName: "existing.png",
|
||||
FilePath: "uploads/existing.png",
|
||||
FileSize: int64(len(content)),
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "pixez_mirror",
|
||||
Status: model.UploadStatusUsed,
|
||||
CreatedAt: time.Now(),
|
||||
}
|
||||
if err := dbConn.Create(&existing).Error; err != nil {
|
||||
t.Fatalf("seed upload failed: %v", err)
|
||||
}
|
||||
|
||||
restoreStorage, disableStorage := setupMockStorage(t, nil)
|
||||
defer restoreStorage()
|
||||
defer disableStorage()
|
||||
|
||||
result, err := Ingest(ctx, Request{
|
||||
UserID: 1001,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "mirror.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "pixez_mirror",
|
||||
Policy: PolicyResolveExisting,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Ingest(PolicyResolveExisting) returned error: %v", err)
|
||||
}
|
||||
if !result.Resolved || result.Created || result.Stored {
|
||||
t.Fatalf("Ingest(PolicyResolveExisting) = %+v, want Resolved only", result)
|
||||
}
|
||||
if result.Upload.ID != existing.ID {
|
||||
t.Fatalf("Ingest(PolicyResolveExisting).Upload.ID = %d, want %d", result.Upload.ID, existing.ID)
|
||||
}
|
||||
|
||||
stats, err := loadTotalStats(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("loadTotalStats returned error: %v", err)
|
||||
}
|
||||
if stats.TotalCount != 0 || stats.TotalSize != 0 {
|
||||
t.Fatalf("loadTotalStats() = count %d size %d, want zero stats for resolved upload", stats.TotalCount, stats.TotalSize)
|
||||
}
|
||||
}
|
||||
|
||||
func TestIngestPolicyDedupNewRecordCreatesSecondRecord(t *testing.T) {
|
||||
dbConn, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01")
|
||||
hash := sha256.Sum256(content)
|
||||
hashStr := hex.EncodeToString(hash[:])
|
||||
putCount := 0
|
||||
|
||||
restoreStorage, disableStorage := setupMockStorage(t, &putCount)
|
||||
defer restoreStorage()
|
||||
defer disableStorage()
|
||||
|
||||
first, err := Ingest(ctx, Request{
|
||||
UserID: 1001,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "first.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "avatar",
|
||||
Policy: PolicyDedupNewRecord,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("first Ingest returned error: %v", err)
|
||||
}
|
||||
if putCount != 1 {
|
||||
t.Fatalf("putCount after first ingest = %d, want 1", putCount)
|
||||
}
|
||||
|
||||
second, err := Ingest(ctx, Request{
|
||||
UserID: 1002,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "second.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hashStr,
|
||||
Type: "avatar",
|
||||
Policy: PolicyDedupNewRecord,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("second Ingest returned error: %v", err)
|
||||
}
|
||||
if putCount != 1 {
|
||||
t.Fatalf("putCount after dedup ingest = %d, want 1", putCount)
|
||||
}
|
||||
if first.Upload.FilePath != second.Upload.FilePath {
|
||||
t.Fatalf("dedup file paths differ: %s vs %s", first.Upload.FilePath, second.Upload.FilePath)
|
||||
}
|
||||
if first.Upload.ID == second.Upload.ID {
|
||||
t.Fatal("dedup records should have unique IDs")
|
||||
}
|
||||
|
||||
var count int64
|
||||
if err := dbConn.Model(&model.Upload{}).Where("hash = ?", hashStr).Count(&count).Error; err != nil {
|
||||
t.Fatalf("count uploads failed: %v", err)
|
||||
}
|
||||
if count != 2 {
|
||||
t.Fatalf("upload count = %d, want 2", count)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRemoveDecrementsStats(t *testing.T) {
|
||||
_, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
defer cleanup()
|
||||
ctx := context.Background()
|
||||
|
||||
content := []byte("\x89PNG\r\n\x1a\n\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01")
|
||||
hash := sha256.Sum256(content)
|
||||
|
||||
restoreStorage, disableStorage := setupMockStorage(t, nil)
|
||||
defer restoreStorage()
|
||||
defer disableStorage()
|
||||
|
||||
result, err := Ingest(ctx, Request{
|
||||
UserID: 1001,
|
||||
Reader: bytes.NewReader(content),
|
||||
Size: int64(len(content)),
|
||||
FileName: "delete-me.png",
|
||||
MimeType: "image/png",
|
||||
Extension: "png",
|
||||
Hash: hex.EncodeToString(hash[:]),
|
||||
Type: "generic",
|
||||
Policy: PolicyCreate,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Ingest returned error: %v", err)
|
||||
}
|
||||
|
||||
if _, err := Remove(ctx, result.Upload.ID); err != nil {
|
||||
t.Fatalf("Remove(%d) returned error: %v", result.Upload.ID, err)
|
||||
}
|
||||
|
||||
stats, err := loadTotalStats(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("loadTotalStats returned error: %v", err)
|
||||
}
|
||||
if stats.TotalCount != 0 || stats.TotalSize != 0 {
|
||||
t.Fatalf("loadTotalStats() after remove = count %d size %d, want zero", stats.TotalCount, stats.TotalSize)
|
||||
}
|
||||
}
|
||||
|
||||
type totalStatsSnapshot struct {
|
||||
TotalCount int64
|
||||
TotalSize int64
|
||||
}
|
||||
|
||||
func loadTotalStats(ctx context.Context) (totalStatsSnapshot, error) {
|
||||
var rows []model.UploadStat
|
||||
if err := db.DB(ctx).Where("dimension = ?", model.UploadStatDimensionTotal).Find(&rows).Error; err != nil {
|
||||
return totalStatsSnapshot{}, err
|
||||
}
|
||||
if len(rows) == 0 {
|
||||
return totalStatsSnapshot{}, nil
|
||||
}
|
||||
return totalStatsSnapshot{
|
||||
TotalCount: rows[0].FileCount,
|
||||
TotalSize: rows[0].FileSize,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func setupMockStorage(t *testing.T, putCount *int) (restore func(), disable func()) {
|
||||
t.Helper()
|
||||
mockFiles := make(map[string][]byte)
|
||||
restore = storage.MockStorage(
|
||||
func(ctx context.Context, key string, body io.Reader, size int64, contentType string) error {
|
||||
data, err := io.ReadAll(body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
mockFiles[key] = data
|
||||
if putCount != nil {
|
||||
*putCount++
|
||||
}
|
||||
return nil
|
||||
},
|
||||
func(ctx context.Context, key string) (*storage.Object, error) {
|
||||
data, ok := mockFiles[key]
|
||||
if !ok {
|
||||
return nil, os.ErrNotExist
|
||||
}
|
||||
return &storage.Object{
|
||||
Body: io.NopCloser(bytes.NewReader(data)),
|
||||
ContentLength: int64(len(data)),
|
||||
ContentType: "application/octet-stream",
|
||||
}, nil
|
||||
},
|
||||
func(ctx context.Context, key string) error {
|
||||
delete(mockFiles, key)
|
||||
return nil
|
||||
},
|
||||
)
|
||||
storage.IsEnabledFunc = func() bool { return true }
|
||||
storage.ResetCache()
|
||||
disable = func() {
|
||||
storage.IsEnabledFunc = func() bool { return false }
|
||||
storage.ResetCache()
|
||||
}
|
||||
return restore, disable
|
||||
}
|
||||
@@ -0,0 +1,43 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
uploadstats "github.com/Rain-kl/Wavelet/internal/apps/upload/stats"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
)
|
||||
|
||||
// Remove soft-deletes an upload and decrements incremental stats.
|
||||
func Remove(ctx context.Context, uploadID uint64) (model.Upload, error) {
|
||||
upload, err := repository.GetActiveUploadByID(ctx, uploadID)
|
||||
if err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
uploadstats.RecordUploadStatsRemove(ctx, &upload)
|
||||
if err := repository.SoftDeleteUpload(ctx, &upload); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
upload.Status = model.UploadStatusDeleted
|
||||
return upload, nil
|
||||
}
|
||||
|
||||
// RemoveOwned soft-deletes an upload owned by userID and decrements incremental stats.
|
||||
func RemoveOwned(ctx context.Context, userID, uploadID uint64) (model.Upload, error) {
|
||||
upload, err := repository.GetActiveUploadByID(ctx, uploadID)
|
||||
if err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
if upload.UserID != userID {
|
||||
return model.Upload{}, ErrForbidden
|
||||
}
|
||||
uploadstats.RecordUploadStatsRemove(ctx, &upload)
|
||||
if err := repository.SoftDeleteUpload(ctx, &upload); err != nil {
|
||||
return model.Upload{}, err
|
||||
}
|
||||
upload.Status = model.UploadStatusDeleted
|
||||
return upload, nil
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package ingest provides the programmatic upload domain service for Wavelet.
|
||||
package ingest
|
||||
|
||||
import (
|
||||
"io"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
)
|
||||
|
||||
// Policy controls how ingest handles hash collisions and record creation.
|
||||
type Policy int
|
||||
|
||||
const (
|
||||
// PolicyCreate always stores a new object and creates a new upload record.
|
||||
PolicyCreate Policy = iota
|
||||
|
||||
// PolicyDedupNewRecord reuses an existing object path on hash match but creates a new record and stats delta.
|
||||
PolicyDedupNewRecord
|
||||
|
||||
// PolicyResolveExisting returns an existing upload on hash match without creating a record or stats delta.
|
||||
PolicyResolveExisting
|
||||
)
|
||||
|
||||
// ObjectKeyFn builds the storage object key for a new upload.
|
||||
type ObjectKeyFn func(id uint64, ext string) string
|
||||
|
||||
// Request describes a programmatic file ingest operation.
|
||||
type Request struct {
|
||||
UserID uint64
|
||||
Type string
|
||||
|
||||
AccessMode *int
|
||||
Status model.UploadStatus
|
||||
|
||||
Reader io.Reader
|
||||
Size int64
|
||||
FileName string
|
||||
MimeType string
|
||||
Extension string
|
||||
Hash string
|
||||
|
||||
Metadata model.UploadMetadata
|
||||
Policy Policy
|
||||
|
||||
ObjectKeyFn ObjectKeyFn
|
||||
|
||||
// SkipExtensionCheck bypasses the configured upload extension whitelist.
|
||||
SkipExtensionCheck bool
|
||||
}
|
||||
|
||||
// Result reports the outcome of an ingest operation.
|
||||
type Result struct {
|
||||
Upload model.Upload
|
||||
Created bool
|
||||
Stored bool
|
||||
Resolved bool
|
||||
}
|
||||
@@ -63,6 +63,7 @@ func GetActiveUploadByID(ctx context.Context, id uint64) (model.Upload, error) {
|
||||
}
|
||||
|
||||
// SoftDeleteUpload marks an upload as deleted.
|
||||
// External modules must use upload.Remove or upload.RemoveOwned; only internal/apps/upload may call this.
|
||||
func SoftDeleteUpload(ctx context.Context, upload *model.Upload) error {
|
||||
return db.DB(ctx).Model(upload).Update("status", model.UploadStatusDeleted).Error
|
||||
}
|
||||
@@ -97,6 +98,7 @@ func FindReusableUploadByHash(ctx context.Context, hash string, size int64) (mod
|
||||
}
|
||||
|
||||
// CreateUpload persists a new upload record.
|
||||
// External modules must use upload.Ingest; only internal/apps/upload may call this.
|
||||
func CreateUpload(ctx context.Context, upload *model.Upload) error {
|
||||
return db.DB(ctx).Create(upload).Error
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user