feat(uploads): add native resumable upload support
Implement a native chunked resumable upload API and frontend integration to support reliable large file uploads. Changes include: - Added a 3-step resumable upload API flow (create session, upload chunks, complete session). - Introduced configuration options for chunk size, retention hours, and toggling the feature. - Updated the frontend to utilize resumable uploads with progress tracking. - Configured temporary chunk storage under `data/tmp/uploads` with automatic cleanup. - Documented the API flow and configuration in the README.
This commit is contained in:
454
backend/libs/services/resumable.go
Normal file
454
backend/libs/services/resumable.go
Normal file
@@ -0,0 +1,454 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"go.etcd.io/bbolt"
|
||||
)
|
||||
|
||||
var resumableUploadsBucket = []byte("resumable_uploads")
|
||||
|
||||
const (
|
||||
ResumableStatusUploading = "uploading"
|
||||
ResumableStatusCompleted = "completed"
|
||||
ResumableStatusCancelled = "cancelled"
|
||||
)
|
||||
|
||||
type ResumableFileInput struct {
|
||||
Name string `json:"name"`
|
||||
Size int64 `json:"size"`
|
||||
ContentType string `json:"contentType"`
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
}
|
||||
|
||||
type ResumableSession struct {
|
||||
ID string `json:"id"`
|
||||
Options UploadOptions `json:"options"`
|
||||
Files []ResumableFile `json:"files"`
|
||||
ChunkSize int64 `json:"chunkSize"`
|
||||
Status string `json:"status"`
|
||||
BoxID string `json:"boxId,omitempty"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
ExpiresAt time.Time `json:"expiresAt"`
|
||||
}
|
||||
|
||||
type ResumableFile struct {
|
||||
ID string `json:"id"`
|
||||
Name string `json:"name"`
|
||||
Size int64 `json:"size"`
|
||||
ContentType string `json:"contentType"`
|
||||
Fingerprint string `json:"fingerprint,omitempty"`
|
||||
ChunkCount int `json:"chunkCount"`
|
||||
UploadedChunks []int `json:"uploadedChunks"`
|
||||
}
|
||||
|
||||
func (s *UploadService) ensureResumableBucket() error {
|
||||
return s.db.Update(func(tx *bbolt.Tx) error {
|
||||
_, err := tx.CreateBucketIfNotExists(resumableUploadsBucket)
|
||||
return err
|
||||
})
|
||||
}
|
||||
|
||||
func (s *UploadService) CreateResumableSession(files []ResumableFileInput, opts UploadOptions, chunkSize int64, retention time.Duration) (ResumableSession, error) {
|
||||
if len(files) == 0 {
|
||||
return ResumableSession{}, fmt.Errorf("no files were uploaded")
|
||||
}
|
||||
if chunkSize <= 0 {
|
||||
return ResumableSession{}, fmt.Errorf("chunk size must be positive")
|
||||
}
|
||||
if retention <= 0 {
|
||||
return ResumableSession{}, fmt.Errorf("retention must be positive")
|
||||
}
|
||||
if strings.TrimSpace(opts.Password) != "" {
|
||||
opts.PasswordSalt, opts.PasswordHash = hashPassword(opts.Password)
|
||||
opts.Password = ""
|
||||
}
|
||||
sessionFiles, err := s.resumableFilesFromInput(files, opts, chunkSize, nil)
|
||||
if err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
session := ResumableSession{
|
||||
ID: randomID(12),
|
||||
Options: opts,
|
||||
Files: sessionFiles,
|
||||
ChunkSize: chunkSize,
|
||||
Status: ResumableStatusUploading,
|
||||
CreatedAt: now,
|
||||
UpdatedAt: now,
|
||||
ExpiresAt: now.Add(retention),
|
||||
}
|
||||
if err := s.saveResumableSession(session); err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
return session, nil
|
||||
}
|
||||
|
||||
func (s *UploadService) AddResumableFiles(sessionID string, files []ResumableFileInput) (ResumableSession, error) {
|
||||
if len(files) == 0 {
|
||||
return s.GetResumableSession(sessionID)
|
||||
}
|
||||
session, err := s.GetResumableSession(sessionID)
|
||||
if err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
if err := resumableSessionWritable(session); err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
existing := make(map[string]bool)
|
||||
for _, file := range session.Files {
|
||||
existing[resumableFileKey(file.Name, file.Size, file.Fingerprint)] = true
|
||||
}
|
||||
newFiles, err := s.resumableFilesFromInput(files, session.Options, session.ChunkSize, existing)
|
||||
if err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
if len(newFiles) == 0 {
|
||||
return session, nil
|
||||
}
|
||||
session.Files = append(session.Files, newFiles...)
|
||||
session.UpdatedAt = time.Now().UTC()
|
||||
if err := s.saveResumableSession(session); err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
return session, nil
|
||||
}
|
||||
|
||||
func (s *UploadService) GetResumableSession(id string) (ResumableSession, error) {
|
||||
var session ResumableSession
|
||||
err := s.db.View(func(tx *bbolt.Tx) error {
|
||||
bucket := tx.Bucket(resumableUploadsBucket)
|
||||
if bucket == nil {
|
||||
return os.ErrNotExist
|
||||
}
|
||||
data := bucket.Get([]byte(id))
|
||||
if data == nil {
|
||||
return os.ErrNotExist
|
||||
}
|
||||
return json.Unmarshal(data, &session)
|
||||
})
|
||||
if err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
return session, nil
|
||||
}
|
||||
|
||||
func (s *UploadService) PutResumableChunk(ctx context.Context, sessionID, fileID string, index int, body io.Reader) (ResumableSession, error) {
|
||||
session, err := s.GetResumableSession(sessionID)
|
||||
if err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
if err := resumableSessionWritable(session); err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
fileIndex := -1
|
||||
for i, file := range session.Files {
|
||||
if file.ID == fileID {
|
||||
fileIndex = i
|
||||
break
|
||||
}
|
||||
}
|
||||
if fileIndex < 0 {
|
||||
return ResumableSession{}, os.ErrNotExist
|
||||
}
|
||||
file := session.Files[fileIndex]
|
||||
if index < 0 || index >= file.ChunkCount {
|
||||
return ResumableSession{}, fmt.Errorf("chunk index is invalid")
|
||||
}
|
||||
expectedSize := expectedChunkSize(file.Size, session.ChunkSize, index)
|
||||
chunkDir := s.resumableFileDir(session.ID, file.ID)
|
||||
if err := os.MkdirAll(chunkDir, 0o755); err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
chunkPath := s.resumableChunkPath(session.ID, file.ID, index)
|
||||
tempPath := chunkPath + ".tmp"
|
||||
target, err := os.OpenFile(tempPath, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600)
|
||||
if err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
written, copyErr := io.Copy(target, io.LimitReader(body, expectedSize+1))
|
||||
closeErr := target.Close()
|
||||
if copyErr != nil {
|
||||
_ = os.Remove(tempPath)
|
||||
return ResumableSession{}, copyErr
|
||||
}
|
||||
if closeErr != nil {
|
||||
_ = os.Remove(tempPath)
|
||||
return ResumableSession{}, closeErr
|
||||
}
|
||||
if written != expectedSize {
|
||||
_ = os.Remove(tempPath)
|
||||
return ResumableSession{}, fmt.Errorf("chunk size mismatch")
|
||||
}
|
||||
if err := os.Rename(tempPath, chunkPath); err != nil {
|
||||
_ = os.Remove(tempPath)
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
session.Files[fileIndex].UploadedChunks = addChunkIndex(session.Files[fileIndex].UploadedChunks, index)
|
||||
session.UpdatedAt = time.Now().UTC()
|
||||
if err := s.saveResumableSession(session); err != nil {
|
||||
return ResumableSession{}, err
|
||||
}
|
||||
return session, nil
|
||||
}
|
||||
|
||||
func (s *UploadService) CompleteResumableSession(ctx context.Context, sessionID string) (UploadResult, ResumableSession, error) {
|
||||
session, err := s.GetResumableSession(sessionID)
|
||||
if err != nil {
|
||||
return UploadResult{}, ResumableSession{}, err
|
||||
}
|
||||
if err := resumableSessionWritable(session); err != nil {
|
||||
return UploadResult{}, ResumableSession{}, err
|
||||
}
|
||||
staged, cleanup, err := s.assembleResumableFiles(ctx, session)
|
||||
if err != nil {
|
||||
return UploadResult{}, ResumableSession{}, err
|
||||
}
|
||||
defer cleanup()
|
||||
|
||||
result, err := s.CreateBoxFromIncoming(staged, session.Options)
|
||||
if err != nil {
|
||||
return UploadResult{}, ResumableSession{}, err
|
||||
}
|
||||
if err := os.RemoveAll(s.resumableSessionDir(session.ID)); err != nil {
|
||||
return UploadResult{}, ResumableSession{}, err
|
||||
}
|
||||
session.Status = ResumableStatusCompleted
|
||||
session.BoxID = result.BoxID
|
||||
session.UpdatedAt = time.Now().UTC()
|
||||
if err := s.saveResumableSession(session); err != nil {
|
||||
return UploadResult{}, ResumableSession{}, err
|
||||
}
|
||||
return result, session, nil
|
||||
}
|
||||
|
||||
func (s *UploadService) CancelResumableSession(sessionID string) error {
|
||||
session, err := s.GetResumableSession(sessionID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
session.Status = ResumableStatusCancelled
|
||||
session.UpdatedAt = time.Now().UTC()
|
||||
if err := s.saveResumableSession(session); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.RemoveAll(s.resumableSessionDir(session.ID))
|
||||
}
|
||||
|
||||
func (s *UploadService) CleanupExpiredResumableSessions(now time.Time) (int, error) {
|
||||
candidates := make([]ResumableSession, 0)
|
||||
err := s.db.View(func(tx *bbolt.Tx) error {
|
||||
bucket := tx.Bucket(resumableUploadsBucket)
|
||||
if bucket == nil {
|
||||
return nil
|
||||
}
|
||||
return bucket.ForEach(func(_, value []byte) error {
|
||||
var session ResumableSession
|
||||
if err := json.Unmarshal(value, &session); err != nil {
|
||||
return err
|
||||
}
|
||||
if !session.ExpiresAt.After(now) || session.Status != ResumableStatusUploading {
|
||||
candidates = append(candidates, session)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
})
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
for _, session := range candidates {
|
||||
if err := os.RemoveAll(s.resumableSessionDir(session.ID)); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
}
|
||||
err = s.db.Update(func(tx *bbolt.Tx) error {
|
||||
bucket := tx.Bucket(resumableUploadsBucket)
|
||||
if bucket == nil {
|
||||
return nil
|
||||
}
|
||||
for _, session := range candidates {
|
||||
if err := bucket.Delete([]byte(session.ID)); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return len(candidates), err
|
||||
}
|
||||
|
||||
func (s *UploadService) saveResumableSession(session ResumableSession) error {
|
||||
if err := s.ensureResumableBucket(); err != nil {
|
||||
return err
|
||||
}
|
||||
return s.db.Update(func(tx *bbolt.Tx) error {
|
||||
data, err := json.Marshal(session)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Bucket(resumableUploadsBucket).Put([]byte(session.ID), data)
|
||||
})
|
||||
}
|
||||
|
||||
func (s *UploadService) resumableFilesFromInput(files []ResumableFileInput, opts UploadOptions, chunkSize int64, existing map[string]bool) ([]ResumableFile, error) {
|
||||
sessionFiles := make([]ResumableFile, 0, len(files))
|
||||
for _, file := range files {
|
||||
file.Name = filepath.Base(strings.TrimSpace(file.Name))
|
||||
if file.Name == "." || file.Name == "" {
|
||||
return nil, fmt.Errorf("file name is required")
|
||||
}
|
||||
if file.Size < 0 {
|
||||
return nil, fmt.Errorf("file size is invalid")
|
||||
}
|
||||
fingerprint := strings.TrimSpace(file.Fingerprint)
|
||||
key := resumableFileKey(file.Name, file.Size, fingerprint)
|
||||
if existing != nil && existing[key] {
|
||||
continue
|
||||
}
|
||||
if !opts.SkipSizeLimit {
|
||||
if err := s.ValidateSize(file.Size); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
chunks := int((file.Size + chunkSize - 1) / chunkSize)
|
||||
if chunks == 0 {
|
||||
chunks = 1
|
||||
}
|
||||
sessionFiles = append(sessionFiles, ResumableFile{
|
||||
ID: randomID(8),
|
||||
Name: file.Name,
|
||||
Size: file.Size,
|
||||
ContentType: strings.TrimSpace(file.ContentType),
|
||||
Fingerprint: fingerprint,
|
||||
ChunkCount: chunks,
|
||||
})
|
||||
if existing != nil {
|
||||
existing[key] = true
|
||||
}
|
||||
}
|
||||
return sessionFiles, nil
|
||||
}
|
||||
|
||||
func resumableFileKey(name string, size int64, fingerprint string) string {
|
||||
return strings.TrimSpace(fingerprint) + "|" + filepath.Base(strings.TrimSpace(name)) + "|" + fmt.Sprintf("%d", size)
|
||||
}
|
||||
|
||||
func (s *UploadService) assembleResumableFiles(ctx context.Context, session ResumableSession) ([]IncomingFile, func(), error) {
|
||||
assembledDir := filepath.Join(s.resumableSessionDir(session.ID), "assembled")
|
||||
if err := os.MkdirAll(assembledDir, 0o755); err != nil {
|
||||
return nil, func() {}, err
|
||||
}
|
||||
cleanup := func() {
|
||||
_ = os.RemoveAll(assembledDir)
|
||||
}
|
||||
staged := make([]IncomingFile, 0, len(session.Files))
|
||||
for _, file := range session.Files {
|
||||
if len(file.UploadedChunks) != file.ChunkCount {
|
||||
cleanup()
|
||||
return nil, func() {}, fmt.Errorf("file %s is missing chunks", file.Name)
|
||||
}
|
||||
assembledPath := filepath.Join(assembledDir, file.ID)
|
||||
target, err := os.OpenFile(assembledPath, os.O_CREATE|os.O_TRUNC|os.O_WRONLY, 0o600)
|
||||
if err != nil {
|
||||
cleanup()
|
||||
return nil, func() {}, err
|
||||
}
|
||||
var written int64
|
||||
for i := 0; i < file.ChunkCount; i++ {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
_ = target.Close()
|
||||
cleanup()
|
||||
return nil, func() {}, ctx.Err()
|
||||
default:
|
||||
}
|
||||
chunk, err := os.Open(s.resumableChunkPath(session.ID, file.ID, i))
|
||||
if err != nil {
|
||||
_ = target.Close()
|
||||
cleanup()
|
||||
return nil, func() {}, fmt.Errorf("file %s is missing chunks", file.Name)
|
||||
}
|
||||
n, copyErr := io.Copy(target, chunk)
|
||||
closeErr := chunk.Close()
|
||||
if copyErr != nil {
|
||||
_ = target.Close()
|
||||
cleanup()
|
||||
return nil, func() {}, copyErr
|
||||
}
|
||||
if closeErr != nil {
|
||||
_ = target.Close()
|
||||
cleanup()
|
||||
return nil, func() {}, closeErr
|
||||
}
|
||||
written += n
|
||||
}
|
||||
if err := target.Close(); err != nil {
|
||||
cleanup()
|
||||
return nil, func() {}, err
|
||||
}
|
||||
if written != file.Size {
|
||||
cleanup()
|
||||
return nil, func() {}, fmt.Errorf("assembled file size mismatch")
|
||||
}
|
||||
staged = append(staged, StagedUploadFile{
|
||||
Filename: file.Name,
|
||||
FileSize: file.Size,
|
||||
MIMEType: file.ContentType,
|
||||
Path: assembledPath,
|
||||
})
|
||||
}
|
||||
return staged, cleanup, nil
|
||||
}
|
||||
|
||||
func resumableSessionWritable(session ResumableSession) error {
|
||||
if session.Status != ResumableStatusUploading {
|
||||
return fmt.Errorf("upload session is not active")
|
||||
}
|
||||
if !session.ExpiresAt.After(time.Now().UTC()) {
|
||||
return fmt.Errorf("upload session expired")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func expectedChunkSize(fileSize, chunkSize int64, index int) int64 {
|
||||
offset := int64(index) * chunkSize
|
||||
remaining := fileSize - offset
|
||||
if remaining < 0 {
|
||||
return 0
|
||||
}
|
||||
if remaining > chunkSize {
|
||||
return chunkSize
|
||||
}
|
||||
return remaining
|
||||
}
|
||||
|
||||
func addChunkIndex(chunks []int, index int) []int {
|
||||
for _, chunk := range chunks {
|
||||
if chunk == index {
|
||||
return chunks
|
||||
}
|
||||
}
|
||||
chunks = append(chunks, index)
|
||||
sort.Ints(chunks)
|
||||
return chunks
|
||||
}
|
||||
|
||||
func (s *UploadService) resumableSessionDir(sessionID string) string {
|
||||
return filepath.Join(s.dataDir, "tmp", "uploads", sessionID)
|
||||
}
|
||||
|
||||
func (s *UploadService) resumableFileDir(sessionID, fileID string) string {
|
||||
return filepath.Join(s.resumableSessionDir(sessionID), fileID)
|
||||
}
|
||||
|
||||
func (s *UploadService) resumableChunkPath(sessionID, fileID string, index int) string {
|
||||
return filepath.Join(s.resumableFileDir(sessionID, fileID), fmt.Sprintf("%06d.part", index))
|
||||
}
|
||||
@@ -42,6 +42,8 @@ type UploadOptions struct {
|
||||
ExpiresInMinutes int
|
||||
MaxDownloads int
|
||||
Password string
|
||||
PasswordSalt string
|
||||
PasswordHash string
|
||||
ObfuscateMetadata bool
|
||||
OwnerID string
|
||||
CollectionID string
|
||||
@@ -50,6 +52,56 @@ type UploadOptions struct {
|
||||
StorageBackendID string
|
||||
}
|
||||
|
||||
type IncomingFile interface {
|
||||
Name() string
|
||||
Size() int64
|
||||
ContentType() string
|
||||
Open() (io.ReadCloser, error)
|
||||
}
|
||||
|
||||
type multipartIncomingFile struct {
|
||||
header *multipart.FileHeader
|
||||
}
|
||||
|
||||
func (f multipartIncomingFile) Name() string {
|
||||
return f.header.Filename
|
||||
}
|
||||
|
||||
func (f multipartIncomingFile) Size() int64 {
|
||||
return f.header.Size
|
||||
}
|
||||
|
||||
func (f multipartIncomingFile) ContentType() string {
|
||||
return f.header.Header.Get("Content-Type")
|
||||
}
|
||||
|
||||
func (f multipartIncomingFile) Open() (io.ReadCloser, error) {
|
||||
return f.header.Open()
|
||||
}
|
||||
|
||||
type StagedUploadFile struct {
|
||||
Filename string
|
||||
FileSize int64
|
||||
MIMEType string
|
||||
Path string
|
||||
}
|
||||
|
||||
func (f StagedUploadFile) Name() string {
|
||||
return f.Filename
|
||||
}
|
||||
|
||||
func (f StagedUploadFile) Size() int64 {
|
||||
return f.FileSize
|
||||
}
|
||||
|
||||
func (f StagedUploadFile) ContentType() string {
|
||||
return f.MIMEType
|
||||
}
|
||||
|
||||
func (f StagedUploadFile) Open() (io.ReadCloser, error) {
|
||||
return os.Open(f.Path)
|
||||
}
|
||||
|
||||
type Box struct {
|
||||
ID string `json:"id"`
|
||||
OwnerID string `json:"ownerId,omitempty"`
|
||||
@@ -198,6 +250,10 @@ func (s *UploadService) ValidateSize(size int64) error {
|
||||
}
|
||||
|
||||
func (s *UploadService) CreateBox(files []*multipart.FileHeader, opts UploadOptions) (UploadResult, error) {
|
||||
return s.CreateBoxFromIncoming(multipartIncomingFiles(files), opts)
|
||||
}
|
||||
|
||||
func (s *UploadService) CreateBoxFromIncoming(files []IncomingFile, opts UploadOptions) (UploadResult, error) {
|
||||
if len(files) == 0 {
|
||||
return UploadResult{}, fmt.Errorf("no files were uploaded")
|
||||
}
|
||||
@@ -232,13 +288,16 @@ func (s *UploadService) CreateBox(files []*multipart.FileHeader, opts UploadOpti
|
||||
}
|
||||
deleteToken := randomID(32)
|
||||
box.DeleteTokenHash = deleteTokenHash(box.ID, deleteToken)
|
||||
if strings.TrimSpace(opts.Password) != "" {
|
||||
if strings.TrimSpace(opts.PasswordHash) != "" {
|
||||
box.PasswordSalt = opts.PasswordSalt
|
||||
box.PasswordHash = opts.PasswordHash
|
||||
} else if strings.TrimSpace(opts.Password) != "" {
|
||||
salt, hash := hashPassword(opts.Password)
|
||||
box.PasswordSalt = salt
|
||||
box.PasswordHash = hash
|
||||
}
|
||||
|
||||
if err := s.writeFilesToBox(&box, files, opts); err != nil {
|
||||
if err := s.writeIncomingFilesToBox(&box, files, opts); err != nil {
|
||||
return UploadResult{}, err
|
||||
}
|
||||
|
||||
@@ -261,6 +320,10 @@ func (s *UploadService) CreateBox(files []*multipart.FileHeader, opts UploadOpti
|
||||
// selection into a single box). The box keeps its original expiry, password and
|
||||
// other settings; only the new files are written.
|
||||
func (s *UploadService) AppendFiles(boxID string, files []*multipart.FileHeader, opts UploadOptions) (UploadResult, error) {
|
||||
return s.AppendIncomingFiles(boxID, multipartIncomingFiles(files), opts)
|
||||
}
|
||||
|
||||
func (s *UploadService) AppendIncomingFiles(boxID string, files []IncomingFile, opts UploadOptions) (UploadResult, error) {
|
||||
if len(files) == 0 {
|
||||
return UploadResult{}, fmt.Errorf("no files were uploaded")
|
||||
}
|
||||
@@ -268,7 +331,7 @@ func (s *UploadService) AppendFiles(boxID string, files []*multipart.FileHeader,
|
||||
if err != nil {
|
||||
return UploadResult{}, err
|
||||
}
|
||||
if err := s.writeFilesToBox(&box, files, opts); err != nil {
|
||||
if err := s.writeIncomingFilesToBox(&box, files, opts); err != nil {
|
||||
return UploadResult{}, err
|
||||
}
|
||||
if err := s.SaveBox(box); err != nil {
|
||||
@@ -289,14 +352,26 @@ func (s *UploadService) AppendFiles(boxID string, files []*multipart.FileHeader,
|
||||
// appends the file metadata to box.Files. The box's StorageBackendID determines
|
||||
// where files land, so it works for both new and existing boxes.
|
||||
func (s *UploadService) writeFilesToBox(box *Box, files []*multipart.FileHeader, opts UploadOptions) error {
|
||||
return s.writeIncomingFilesToBox(box, multipartIncomingFiles(files), opts)
|
||||
}
|
||||
|
||||
func multipartIncomingFiles(files []*multipart.FileHeader) []IncomingFile {
|
||||
incoming := make([]IncomingFile, 0, len(files))
|
||||
for _, file := range files {
|
||||
incoming = append(incoming, multipartIncomingFile{header: file})
|
||||
}
|
||||
return incoming
|
||||
}
|
||||
|
||||
func (s *UploadService) writeIncomingFilesToBox(box *Box, files []IncomingFile, opts UploadOptions) error {
|
||||
backend, err := s.storage.Backend(box.StorageBackendID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, header := range files {
|
||||
for _, incoming := range files {
|
||||
if !opts.SkipSizeLimit {
|
||||
if err := s.ValidateSize(header.Size); err != nil {
|
||||
if err := s.ValidateSize(incoming.Size()); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
@@ -306,15 +381,15 @@ func (s *UploadService) writeFilesToBox(box *Box, files []*multipart.FileHeader,
|
||||
maxSize = 0
|
||||
}
|
||||
|
||||
file, err := header.Open()
|
||||
file, err := incoming.Open()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
fileID := randomID(8)
|
||||
storedName := "@each@" + fileID + strings.ToLower(filepath.Ext(header.Filename))
|
||||
storedName := "@each@" + fileID + strings.ToLower(filepath.Ext(incoming.Name()))
|
||||
objectKey := boxObjectKey(box.ID, storedName)
|
||||
contentType := header.Header.Get("Content-Type")
|
||||
contentType := incoming.ContentType()
|
||||
if contentType == "" {
|
||||
buffer := make([]byte, 512)
|
||||
n, _ := file.Read(buffer)
|
||||
@@ -324,7 +399,7 @@ func (s *UploadService) writeFilesToBox(box *Box, files []*multipart.FileHeader,
|
||||
}
|
||||
}
|
||||
|
||||
if err := s.writeUploadedObject(context.Background(), backend, objectKey, file, header.Size, maxSize, contentType); err != nil {
|
||||
if err := s.writeUploadedObject(context.Background(), backend, objectKey, file, incoming.Size(), maxSize, contentType); err != nil {
|
||||
file.Close()
|
||||
return err
|
||||
}
|
||||
@@ -332,9 +407,9 @@ func (s *UploadService) writeFilesToBox(box *Box, files []*multipart.FileHeader,
|
||||
|
||||
box.Files = append(box.Files, File{
|
||||
ID: fileID,
|
||||
Name: filepath.Base(header.Filename),
|
||||
Name: filepath.Base(incoming.Name()),
|
||||
StoredName: storedName,
|
||||
Size: header.Size,
|
||||
Size: incoming.Size(),
|
||||
ContentType: contentType,
|
||||
PreviewKind: previewKind(contentType),
|
||||
ObjectKey: objectKey,
|
||||
@@ -931,21 +1006,17 @@ func writeUploadedFile(path string, source multipart.File, maxSize int64) error
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *UploadService) writeUploadedObject(ctx context.Context, backend StorageBackend, key string, source multipart.File, size, maxSize int64, contentType string) error {
|
||||
func (s *UploadService) writeUploadedObject(ctx context.Context, backend StorageBackend, key string, source io.Reader, size, maxSize int64, contentType string) error {
|
||||
var reader io.Reader = source
|
||||
putSize := size
|
||||
if maxSize > 0 {
|
||||
reader = io.LimitReader(source, maxSize+1)
|
||||
var buffer bytes.Buffer
|
||||
written, err := io.Copy(&buffer, reader)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if written > maxSize {
|
||||
if size > maxSize {
|
||||
return fmt.Errorf("file exceeds max upload size")
|
||||
}
|
||||
return backend.Put(ctx, key, bytes.NewReader(buffer.Bytes()), written, contentType)
|
||||
reader = io.LimitReader(source, maxSize)
|
||||
putSize = size
|
||||
}
|
||||
return backend.Put(ctx, key, reader, size, contentType)
|
||||
return backend.Put(ctx, key, reader, putSize, contentType)
|
||||
}
|
||||
|
||||
func boxObjectKey(boxID, name string) string {
|
||||
|
||||
@@ -126,6 +126,166 @@ func TestLocalStorageBackendAndLegacyFallback(t *testing.T) {
|
||||
object.Body.Close()
|
||||
}
|
||||
|
||||
func TestResumableSessionUploadOutOfOrderAndComplete(t *testing.T) {
|
||||
service := newTestUploadService(t)
|
||||
session, err := service.CreateResumableSession([]ResumableFileInput{{
|
||||
Name: "note.txt",
|
||||
Size: 11,
|
||||
ContentType: "text/plain",
|
||||
Fingerprint: "sha256:first-chunk",
|
||||
}}, UploadOptions{MaxDays: 1, Password: "secret"}, 4, time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("CreateResumableSession returned error: %v", err)
|
||||
}
|
||||
if session.Options.Password != "" || session.Options.PasswordHash == "" || session.Options.PasswordSalt == "" {
|
||||
t.Fatalf("resumable session did not hash password before storage: %+v", session.Options)
|
||||
}
|
||||
if session.Files[0].ChunkCount != 3 {
|
||||
t.Fatalf("ChunkCount = %d, want 3", session.Files[0].ChunkCount)
|
||||
}
|
||||
if session.Files[0].Fingerprint != "sha256:first-chunk" {
|
||||
t.Fatalf("Fingerprint = %q", session.Files[0].Fingerprint)
|
||||
}
|
||||
for index, body := range map[int]string{2: "rld", 0: "hell", 1: "o wo"} {
|
||||
updated, err := service.PutResumableChunk(testContext(), session.ID, session.Files[0].ID, index, strings.NewReader(body))
|
||||
if err != nil {
|
||||
t.Fatalf("PutResumableChunk(%d) returned error: %v", index, err)
|
||||
}
|
||||
if len(updated.Files[0].UploadedChunks) == 0 {
|
||||
t.Fatalf("UploadedChunks was not updated")
|
||||
}
|
||||
}
|
||||
result, completed, err := service.CompleteResumableSession(testContext(), session.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("CompleteResumableSession returned error: %v", err)
|
||||
}
|
||||
if completed.Status != ResumableStatusCompleted || completed.BoxID != result.BoxID {
|
||||
t.Fatalf("completed session = %+v, result = %+v", completed, result)
|
||||
}
|
||||
box := getTestBox(t, service, result.BoxID)
|
||||
if box.PasswordHash == "" || box.PasswordSalt == "" || box.PasswordHash != session.Options.PasswordHash {
|
||||
t.Fatalf("completed box did not preserve hashed password")
|
||||
}
|
||||
object, err := service.OpenFileObject(testContext(), box, box.Files[0])
|
||||
if err != nil {
|
||||
t.Fatalf("OpenFileObject returned error: %v", err)
|
||||
}
|
||||
data, err := io.ReadAll(object.Body)
|
||||
object.Body.Close()
|
||||
if err != nil {
|
||||
t.Fatalf("ReadAll returned error: %v", err)
|
||||
}
|
||||
if string(data) != "hello world" {
|
||||
t.Fatalf("object body = %q", string(data))
|
||||
}
|
||||
if _, err := os.Stat(service.resumableSessionDir(session.ID)); !os.IsNotExist(err) {
|
||||
t.Fatalf("resumable temp dir after complete error = %v, want os.ErrNotExist", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResumableCompleteRejectsMissingChunks(t *testing.T) {
|
||||
service := newTestUploadService(t)
|
||||
session, err := service.CreateResumableSession([]ResumableFileInput{{
|
||||
Name: "note.txt",
|
||||
Size: 8,
|
||||
ContentType: "text/plain",
|
||||
}}, UploadOptions{MaxDays: 1}, 4, time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("CreateResumableSession returned error: %v", err)
|
||||
}
|
||||
if _, err := service.PutResumableChunk(testContext(), session.ID, session.Files[0].ID, 0, strings.NewReader("hell")); err != nil {
|
||||
t.Fatalf("PutResumableChunk returned error: %v", err)
|
||||
}
|
||||
if _, _, err := service.CompleteResumableSession(testContext(), session.ID); err == nil {
|
||||
t.Fatalf("CompleteResumableSession accepted missing chunks")
|
||||
}
|
||||
}
|
||||
|
||||
func TestResumableSessionCanAddFilesBeforeComplete(t *testing.T) {
|
||||
service := newTestUploadService(t)
|
||||
session, err := service.CreateResumableSession([]ResumableFileInput{{
|
||||
Name: "one.txt",
|
||||
Size: 4,
|
||||
ContentType: "text/plain",
|
||||
Fingerprint: "one",
|
||||
}}, UploadOptions{MaxDays: 1}, 4, time.Hour)
|
||||
if err != nil {
|
||||
t.Fatalf("CreateResumableSession returned error: %v", err)
|
||||
}
|
||||
if _, err := service.PutResumableChunk(testContext(), session.ID, session.Files[0].ID, 0, strings.NewReader("one!")); err != nil {
|
||||
t.Fatalf("PutResumableChunk one returned error: %v", err)
|
||||
}
|
||||
updated, err := service.AddResumableFiles(session.ID, []ResumableFileInput{{
|
||||
Name: "two.txt",
|
||||
Size: 4,
|
||||
ContentType: "text/plain",
|
||||
Fingerprint: "two",
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatalf("AddResumableFiles returned error: %v", err)
|
||||
}
|
||||
if len(updated.Files) != 2 {
|
||||
t.Fatalf("files after add = %d, want 2", len(updated.Files))
|
||||
}
|
||||
if updated.Files[0].UploadedChunks[0] != 0 {
|
||||
t.Fatalf("existing uploaded chunk was not preserved: %+v", updated.Files[0])
|
||||
}
|
||||
if _, err := service.AddResumableFiles(session.ID, []ResumableFileInput{{
|
||||
Name: "two.txt",
|
||||
Size: 4,
|
||||
ContentType: "text/plain",
|
||||
Fingerprint: "two",
|
||||
}}); err != nil {
|
||||
t.Fatalf("duplicate AddResumableFiles returned error: %v", err)
|
||||
}
|
||||
updated, err = service.GetResumableSession(session.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("GetResumableSession returned error: %v", err)
|
||||
}
|
||||
if len(updated.Files) != 2 {
|
||||
t.Fatalf("duplicate add changed file count to %d", len(updated.Files))
|
||||
}
|
||||
if _, err := service.PutResumableChunk(testContext(), session.ID, updated.Files[1].ID, 0, strings.NewReader("two!")); err != nil {
|
||||
t.Fatalf("PutResumableChunk two returned error: %v", err)
|
||||
}
|
||||
result, _, err := service.CompleteResumableSession(testContext(), session.ID)
|
||||
if err != nil {
|
||||
t.Fatalf("CompleteResumableSession returned error: %v", err)
|
||||
}
|
||||
box := getTestBox(t, service, result.BoxID)
|
||||
if len(box.Files) != 2 {
|
||||
t.Fatalf("completed box file count = %d, want 2", len(box.Files))
|
||||
}
|
||||
}
|
||||
|
||||
func TestResumableCleanupRemovesExpiredSessionsAndChunks(t *testing.T) {
|
||||
service := newTestUploadService(t)
|
||||
session, err := service.CreateResumableSession([]ResumableFileInput{{
|
||||
Name: "note.txt",
|
||||
Size: 4,
|
||||
ContentType: "text/plain",
|
||||
}}, UploadOptions{MaxDays: 1}, 4, time.Millisecond)
|
||||
if err != nil {
|
||||
t.Fatalf("CreateResumableSession returned error: %v", err)
|
||||
}
|
||||
if _, err := service.PutResumableChunk(testContext(), session.ID, session.Files[0].ID, 0, strings.NewReader("hell")); err != nil {
|
||||
t.Fatalf("PutResumableChunk returned error: %v", err)
|
||||
}
|
||||
cleaned, err := service.CleanupExpiredResumableSessions(time.Now().UTC().Add(time.Hour))
|
||||
if err != nil {
|
||||
t.Fatalf("CleanupExpiredResumableSessions returned error: %v", err)
|
||||
}
|
||||
if cleaned != 1 {
|
||||
t.Fatalf("cleaned = %d, want 1", cleaned)
|
||||
}
|
||||
if _, err := service.GetResumableSession(session.ID); !os.IsNotExist(err) {
|
||||
t.Fatalf("GetResumableSession after cleanup error = %v, want os.ErrNotExist", err)
|
||||
}
|
||||
if _, err := os.Stat(service.resumableSessionDir(session.ID)); !os.IsNotExist(err) {
|
||||
t.Fatalf("resumable temp dir after cleanup error = %v, want os.ErrNotExist", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestContaboStorageConfigAllowsDisplayNamesWithSpaces(t *testing.T) {
|
||||
service := newTestUploadService(t)
|
||||
cfg, err := service.Storage().CreateS3Backend(StorageBackendConfig{
|
||||
|
||||
Reference in New Issue
Block a user