132 lines
5.4 KiB
Go
132 lines
5.4 KiB
Go
package system
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"time"
|
|
|
|
"kra/app/system/internal/biz"
|
|
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
)
|
|
|
|
type uploadSessionPO struct {
|
|
ID uint `gorm:"primaryKey"`
|
|
CreatedAt time.Time
|
|
UpdatedAt time.Time
|
|
DeletedAt gorm.DeletedAt `gorm:"index"`
|
|
UserID uint `gorm:"index"`
|
|
FileName string
|
|
FileHash string `gorm:"index"`
|
|
FileSize int64
|
|
ChunkSize int64
|
|
ChunkTotal int
|
|
Status string
|
|
StorageKey string
|
|
MediaID uint
|
|
}
|
|
|
|
func (uploadSessionPO) TableName() string { return "media_uploads" }
|
|
|
|
type uploadChunkPO struct {
|
|
ID uint `gorm:"primaryKey"`
|
|
CreatedAt time.Time
|
|
UpdatedAt time.Time
|
|
DeletedAt gorm.DeletedAt `gorm:"index"`
|
|
UploadID uint `gorm:"uniqueIndex:idx_upload_chunk"`
|
|
ChunkIndex int `gorm:"uniqueIndex:idx_upload_chunk"`
|
|
ChunkHash string
|
|
Size int64
|
|
}
|
|
|
|
func (uploadChunkPO) TableName() string { return "media_upload_chunks" }
|
|
|
|
func uploadFromPO(v uploadSessionPO) *biz.UploadSession {
|
|
return &biz.UploadSession{ID: v.ID, UserID: v.UserID, FileName: v.FileName, FileHash: v.FileHash, FileSize: v.FileSize, ChunkSize: v.ChunkSize, ChunkTotal: v.ChunkTotal, Status: v.Status, StorageKey: v.StorageKey, MediaID: v.MediaID}
|
|
}
|
|
func (r *mediaRepo) FindCompletedSession(ctx context.Context, userID uint, hash string) (*biz.UploadSession, error) {
|
|
var po uploadSessionPO
|
|
if err := r.data.DB().WithContext(ctx).Where("user_id = ? AND file_hash = ? AND status = ?", userID, hash, "completed").First(&po).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return nil, biz.ErrUploadSessionNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
return uploadFromPO(po), nil
|
|
}
|
|
func (r *mediaRepo) FindUploadingSession(ctx context.Context, userID uint, hash string) (*biz.UploadSession, error) {
|
|
var po uploadSessionPO
|
|
if err := r.data.DB().WithContext(ctx).Where("user_id = ? AND file_hash = ? AND status = ?", userID, hash, "uploading").First(&po).Error; err != nil {
|
|
if errors.Is(err, gorm.ErrRecordNotFound) {
|
|
return nil, biz.ErrUploadSessionNotFound
|
|
}
|
|
return nil, err
|
|
}
|
|
return uploadFromPO(po), nil
|
|
}
|
|
func (r *mediaRepo) CreateUploadSession(ctx context.Context, v *biz.UploadSession) error {
|
|
po := uploadSessionPO{UserID: v.UserID, FileName: v.FileName, FileHash: v.FileHash, FileSize: v.FileSize, ChunkSize: v.ChunkSize, ChunkTotal: v.ChunkTotal, Status: v.Status}
|
|
if err := r.data.DB().WithContext(ctx).Create(&po).Error; err != nil {
|
|
return err
|
|
}
|
|
v.ID = po.ID
|
|
return nil
|
|
}
|
|
func (r *mediaRepo) FindUploadSession(ctx context.Context, id uint) (*biz.UploadSession, error) {
|
|
var po uploadSessionPO
|
|
if err := r.data.DB().WithContext(ctx).First(&po, id).Error; err != nil {
|
|
return nil, err
|
|
}
|
|
return uploadFromPO(po), nil
|
|
}
|
|
func (r *mediaRepo) ClaimUploadSession(ctx context.Context, id uint) (bool, error) {
|
|
result := r.data.DB().WithContext(ctx).Model(&uploadSessionPO{}).Where("id = ? AND status = ?", id, "uploading").Update("status", "merging")
|
|
return result.RowsAffected == 1, result.Error
|
|
}
|
|
func (r *mediaRepo) FailUploadSession(ctx context.Context, id uint) error {
|
|
return r.data.DB().WithContext(ctx).Model(&uploadSessionPO{}).Where("id = ?", id).Update("status", "failed").Error
|
|
}
|
|
func (r *mediaRepo) CompleteUploadSession(ctx context.Context, id uint, key string, mediaID uint) error {
|
|
return r.data.DB().WithContext(ctx).Model(&uploadSessionPO{}).Where("id = ?", id).Updates(map[string]any{"status": "completed", "storage_key": key, "media_id": mediaID}).Error
|
|
}
|
|
func (r *mediaRepo) DeleteUploadSession(ctx context.Context, id uint) error {
|
|
// The compatible flow uses GORM's normal Delete here, retaining the soft-deleted session
|
|
// for audit/recovery rather than physically removing it.
|
|
return r.data.DB().WithContext(ctx).Delete(&uploadSessionPO{}, id).Error
|
|
}
|
|
func (r *mediaRepo) UpsertChunk(ctx context.Context, uploadID uint, v *biz.UploadChunk) error {
|
|
po := uploadChunkPO{UploadID: uploadID, ChunkIndex: v.Index, ChunkHash: v.Hash, Size: v.Size}
|
|
return r.data.DB().WithContext(ctx).Clauses(clause.OnConflict{
|
|
Columns: []clause.Column{{Name: "upload_id"}, {Name: "chunk_index"}},
|
|
DoUpdates: clause.AssignmentColumns([]string{"chunk_hash", "size", "updated_at", "deleted_at"}),
|
|
}).Create(&po).Error
|
|
}
|
|
func (r *mediaRepo) ListChunks(ctx context.Context, uploadID uint) ([]*biz.UploadChunk, error) {
|
|
var pos []uploadChunkPO
|
|
if err := r.data.DB().WithContext(ctx).Where("upload_id = ?", uploadID).Order("chunk_index").Find(&pos).Error; err != nil {
|
|
return nil, err
|
|
}
|
|
out := make([]*biz.UploadChunk, 0, len(pos))
|
|
for _, po := range pos {
|
|
out = append(out, &biz.UploadChunk{Index: po.ChunkIndex, Hash: po.ChunkHash, Size: po.Size})
|
|
}
|
|
return out, nil
|
|
}
|
|
func (r *mediaRepo) DeleteChunks(ctx context.Context, uploadID uint) error {
|
|
return r.data.DB().WithContext(ctx).Where("upload_id = ?", uploadID).Delete(&uploadChunkPO{}).Error
|
|
}
|
|
func (r *mediaRepo) StaleUploadSessionIDs(ctx context.Context, before time.Time) ([]uint, error) {
|
|
var ids []uint
|
|
err := r.data.DB().WithContext(ctx).Model(&uploadSessionPO{}).Where("status = ? AND updated_at < ?", "uploading", before).Pluck("id", &ids).Error
|
|
return ids, err
|
|
}
|
|
func (r *mediaRepo) DeleteUploadData(ctx context.Context, uploadID uint) error {
|
|
return r.data.DB().WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
|
if err := tx.Where("upload_id = ?", uploadID).Delete(&uploadChunkPO{}).Error; err != nil {
|
|
return err
|
|
}
|
|
return tx.Delete(&uploadSessionPO{}, uploadID).Error
|
|
})
|
|
}
|