- internal/upload owns chunked-upload domain logic (fingerprint resume, part tmp+rename, assemble, session sweep) and the shared UniquePath helper used by both single-file and chunked completion paths - handlers/uploads.go is now HTTP-only: bind params, call upload.U, map sentinel errors to the unchanged status/code/message contract - B16: session sweep moved out of the UploadInit request path onto the scanner ticker via a Sweeper hook (upload.U satisfies it) - unit tests for the package without PG/Redis; contract pinned by the existing handler integration tests (all green)
358 lines
10 KiB
Go
358 lines
10 KiB
Go
// Package upload owns the resumable chunked-upload subsystem and the shared
|
|
// unique-path placement helper used by both chunked and single-file uploads.
|
|
//
|
|
// A session lives at <BooksDir>/.uploads/<uid>/ (meta.json + parts/N). The uid
|
|
// is a fingerprint of (libID, name, size, chunkSize), so re-initialising the
|
|
// same file resumes the existing session instead of restarting it. Expired
|
|
// sessions are swept by the scanner ticker (B16), not on the request path.
|
|
//
|
|
// This package is HTTP-free: it returns sentinel errors that the handlers map
|
|
// onto status/code/message tuples. The prior client-facing contract is
|
|
// preserved exactly.
|
|
package upload
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
|
|
"booklib/internal/bookfile"
|
|
"booklib/internal/ports"
|
|
)
|
|
|
|
const (
|
|
maxChunkBytes = 32 << 20
|
|
defaultChunk = 8 << 20
|
|
uploadSessTTL = 24 * time.Hour
|
|
uploadSessionIn = ".uploads"
|
|
)
|
|
|
|
// Sentinel errors returned to handlers for status/code mapping. The four shared
|
|
// with ports (ErrTooLarge/ErrNotFound/ErrIncomplete/ErrSizeMismatch) live in the
|
|
// ports package so the interface contract and the implementation agree.
|
|
var (
|
|
ErrBadName = errors.New("bad name")
|
|
ErrBadFormat = errors.New("bad format")
|
|
ErrBadSize = errors.New("bad size")
|
|
ErrBadChunk = errors.New("bad chunk size")
|
|
ErrBadUploadID = errors.New("bad upload id")
|
|
ErrBadIndex = errors.New("bad part index")
|
|
ErrPartTooBig = errors.New("part exceeds declared size")
|
|
ErrPartSizeMismatch = errors.New("part size mismatch")
|
|
ErrCorrupt = errors.New("corrupt session")
|
|
)
|
|
|
|
// OpError wraps an internal filesystem/IO failure with the short operation label
|
|
// that the client-facing 500 message uses, preserving the prior contract strings
|
|
// ("create session", "write meta", "create part", ...).
|
|
type OpError struct {
|
|
Op string
|
|
Err error
|
|
}
|
|
|
|
func (e *OpError) Error() string { return e.Op + ": " + e.Err.Error() }
|
|
func (e *OpError) Unwrap() error { return e.Err }
|
|
|
|
func opErr(op string, err error) error { return &OpError{Op: op, Err: err} }
|
|
|
|
type uploadMeta struct {
|
|
Name string `json:"name"`
|
|
Size int64 `json:"size"`
|
|
ChunkSize int64 `json:"chunkSize"`
|
|
LibraryID int64 `json:"libraryId"`
|
|
}
|
|
|
|
// U is the upload subsystem. It is stateless beyond the filesystem session dir.
|
|
type U struct {
|
|
booksDir string
|
|
uploadMaxMB int64
|
|
}
|
|
|
|
// New builds the upload subsystem. booksDir is the storage root (sessions live
|
|
// under booksDir/.uploads); uploadMaxMB caps the total declared file size.
|
|
func New(booksDir string, uploadMaxMB int64) *U {
|
|
return &U{booksDir: booksDir, uploadMaxMB: uploadMaxMB}
|
|
}
|
|
|
|
// compile-time proof that *U satisfies the consumer-side interface.
|
|
var _ ports.UploadSessions = (*U)(nil)
|
|
|
|
// ---------- pure helpers ----------
|
|
|
|
func validUploadID(s string) bool {
|
|
if len(s) != 32 {
|
|
return false
|
|
}
|
|
for _, r := range s {
|
|
if !((r >= '0' && r <= '9') || (r >= 'a' && r <= 'f')) {
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
func uploadIDFor(libID int64, name string, size, chunk int64) string {
|
|
h := sha256.Sum256([]byte(fmt.Sprintf("%d|%s|%d|%d", libID, name, size, chunk)))
|
|
return hex.EncodeToString(h[:16])
|
|
}
|
|
|
|
func (u *U) uploadDir(uid string) string {
|
|
return filepath.Join(filepath.Clean(u.booksDir), uploadSessionIn, uid)
|
|
}
|
|
|
|
func chunkRange(m uploadMeta, i int64) (int64, int64) {
|
|
lo := i * m.ChunkSize
|
|
hi := min(lo+m.ChunkSize, m.Size)
|
|
return lo, hi
|
|
}
|
|
|
|
func numParts(m uploadMeta) int64 {
|
|
return (m.Size + m.ChunkSize - 1) / m.ChunkSize
|
|
}
|
|
|
|
// ---------- session meta ----------
|
|
|
|
// loadMeta validates uid + reads meta.json. Returns ErrBadUploadID,
|
|
// ports.ErrNotFound, or ErrCorrupt on failure.
|
|
func (u *U) loadMeta(uid string) (uploadMeta, string, error) {
|
|
if !validUploadID(uid) {
|
|
return uploadMeta{}, "", ErrBadUploadID
|
|
}
|
|
dir := u.uploadDir(uid)
|
|
b, e := os.ReadFile(filepath.Join(dir, "meta.json"))
|
|
if e != nil {
|
|
return uploadMeta{}, "", ports.ErrNotFound
|
|
}
|
|
var m uploadMeta
|
|
if json.Unmarshal(b, &m) != nil {
|
|
return uploadMeta{}, "", ErrCorrupt
|
|
}
|
|
return m, dir, nil
|
|
}
|
|
|
|
// LibraryID returns the target library id recorded in the session, so the
|
|
// handler can resolve+validate the library root before calling Complete.
|
|
func (u *U) LibraryID(ctx context.Context, uid string) (int64, error) {
|
|
m, _, e := u.loadMeta(uid)
|
|
if e != nil {
|
|
return 0, e
|
|
}
|
|
return m.LibraryID, nil
|
|
}
|
|
|
|
// ---------- public API (satisfies ports.UploadSessions) ----------
|
|
|
|
// Init validates the declared upload, then creates or resumes a session. The
|
|
// returned uid is deterministic for a given (libID, safeName, size, chunkSize),
|
|
// so a re-init of the same file resumes; a fingerprint collision with different
|
|
// content restarts the session.
|
|
func (u *U) Init(ctx context.Context, libID int64, name string, size, chunkSize int64) (string, error) {
|
|
safe := bookfile.SafeName(name)
|
|
if safe == "" {
|
|
return "", ErrBadName
|
|
}
|
|
if bookfile.FormatFromExt(safe) == "" {
|
|
return "", ErrBadFormat
|
|
}
|
|
if size <= 0 {
|
|
return "", ErrBadSize
|
|
}
|
|
if size > u.uploadMaxMB<<20 {
|
|
return "", ports.ErrTooLarge
|
|
}
|
|
if chunkSize == 0 {
|
|
chunkSize = defaultChunk
|
|
}
|
|
if chunkSize > maxChunkBytes {
|
|
return "", ErrBadChunk
|
|
}
|
|
uid := uploadIDFor(libID, safe, size, chunkSize)
|
|
dir := u.uploadDir(uid)
|
|
meta := uploadMeta{Name: safe, Size: size, ChunkSize: chunkSize, LibraryID: libID}
|
|
if b, e := os.ReadFile(filepath.Join(dir, "meta.json")); e == nil {
|
|
var old uploadMeta
|
|
if json.Unmarshal(b, &old) == nil && old == meta { // same fingerprint → resume
|
|
return uid, nil
|
|
}
|
|
os.RemoveAll(dir) // fingerprint collided but content differs → restart
|
|
}
|
|
if e := os.MkdirAll(filepath.Join(dir, "parts"), 0o755); e != nil {
|
|
return "", opErr("create session", e)
|
|
}
|
|
b, _ := json.Marshal(meta)
|
|
if e := os.WriteFile(filepath.Join(dir, "meta.json"), b, 0o644); e != nil {
|
|
return "", opErr("write meta", e)
|
|
}
|
|
return uid, nil
|
|
}
|
|
|
|
// Status returns the sorted indices of parts already received.
|
|
func (u *U) Status(ctx context.Context, uid string) ([]int64, error) {
|
|
_, dir, e := u.loadMeta(uid)
|
|
if e != nil {
|
|
return nil, e
|
|
}
|
|
recv := []int64{}
|
|
es, e := os.ReadDir(filepath.Join(dir, "parts"))
|
|
if e == nil {
|
|
for _, en := range es {
|
|
if i, e := strconv.ParseInt(en.Name(), 10, 64); e == nil {
|
|
recv = append(recv, i)
|
|
}
|
|
}
|
|
}
|
|
sort.Slice(recv, func(i, j int) bool { return recv[i] < recv[j] })
|
|
return recv, nil
|
|
}
|
|
|
|
// PutPart writes one part via tmp+rename (B4: a truncated part is never reported
|
|
// as received). body is read up to the declared part size; reading past it yields
|
|
// ErrPartTooBig, a short read yields ErrPartSizeMismatch. maxSize, when > 0, is a
|
|
// defensive ceiling on bytes read.
|
|
func (u *U) PutPart(ctx context.Context, uid string, index int64, body io.Reader, maxSize int64) error {
|
|
m, dir, e := u.loadMeta(uid)
|
|
if e != nil {
|
|
return e
|
|
}
|
|
if index < 0 || index >= numParts(m) {
|
|
return ErrBadIndex
|
|
}
|
|
lo, hi := chunkRange(m, index)
|
|
want := hi - lo
|
|
limit := want + 1
|
|
if maxSize > 0 && maxSize+1 < limit {
|
|
limit = maxSize + 1
|
|
}
|
|
p := filepath.Join(dir, "parts", strconv.FormatInt(index, 10))
|
|
tmp := p + ".tmp"
|
|
f, e := os.OpenFile(tmp, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0o644)
|
|
if e != nil {
|
|
return opErr("create part", e)
|
|
}
|
|
n, copyErr := io.Copy(f, io.LimitReader(body, limit))
|
|
f.Close()
|
|
if copyErr != nil {
|
|
os.Remove(tmp)
|
|
return ErrPartSizeMismatch
|
|
}
|
|
if n > want {
|
|
os.Remove(tmp)
|
|
return ErrPartTooBig
|
|
}
|
|
if n != want {
|
|
os.Remove(tmp)
|
|
return ErrPartSizeMismatch
|
|
}
|
|
if e := os.Rename(tmp, p); e != nil {
|
|
os.Remove(tmp)
|
|
return opErr("rename part", e)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Complete verifies every part is present and correctly sized, assembles them
|
|
// into a tmp file, then atomically renames it into root under a unique name.
|
|
// It returns the path relative to root. The session dir is removed on success.
|
|
func (u *U) Complete(ctx context.Context, uid, root string) (string, error) {
|
|
m, dir, e := u.loadMeta(uid)
|
|
if e != nil {
|
|
return "", e
|
|
}
|
|
var total int64
|
|
for i := int64(0); i < numParts(m); i++ {
|
|
lo, hi := chunkRange(m, i)
|
|
fi, e := os.Stat(filepath.Join(dir, "parts", strconv.FormatInt(i, 10)))
|
|
if e != nil || fi.Size() != hi-lo {
|
|
return "", ports.ErrIncomplete
|
|
}
|
|
total += fi.Size()
|
|
}
|
|
if total != m.Size {
|
|
return "", ports.ErrSizeMismatch
|
|
}
|
|
dst, e := u.UniquePath(root, m.Name)
|
|
if e != nil {
|
|
return "", e // os.ErrInvalid / os.ErrExist → handler maps to 403
|
|
}
|
|
tmp := filepath.Join(dir, "assembled")
|
|
out, e := os.OpenFile(tmp, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0o644)
|
|
if e != nil {
|
|
return "", opErr("create tmp", e)
|
|
}
|
|
for i := int64(0); i < numParts(m); i++ {
|
|
pf, e := os.Open(filepath.Join(dir, "parts", strconv.FormatInt(i, 10)))
|
|
if e != nil {
|
|
out.Close()
|
|
return "", opErr("open part", e)
|
|
}
|
|
_, copyErr := io.Copy(out, pf)
|
|
pf.Close()
|
|
if copyErr != nil {
|
|
out.Close()
|
|
os.Remove(tmp)
|
|
return "", opErr("assemble", copyErr)
|
|
}
|
|
}
|
|
out.Close()
|
|
if e := os.Rename(tmp, dst); e != nil { // atomic placement; scanner picks it up
|
|
os.Remove(tmp)
|
|
return "", opErr("rename", e)
|
|
}
|
|
os.RemoveAll(dir)
|
|
return strings.TrimPrefix(dst, root+string(os.PathSeparator)), nil
|
|
}
|
|
|
|
// Sweep removes session dirs untouched for longer than the TTL. Best-effort:
|
|
// a missing/unreadable base dir is not an error. Run from the scanner ticker
|
|
// (B16) instead of the request path.
|
|
func (u *U) Sweep(ctx context.Context) error {
|
|
base := filepath.Join(filepath.Clean(u.booksDir), uploadSessionIn)
|
|
es, e := os.ReadDir(base)
|
|
if e != nil {
|
|
return nil // no sessions yet
|
|
}
|
|
for _, en := range es {
|
|
if fi, e := en.Info(); e == nil && time.Since(fi.ModTime()) > uploadSessTTL {
|
|
os.RemoveAll(filepath.Join(base, en.Name()))
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// UniquePath returns a path under root for name that does not yet exist,
|
|
// appending " (n)" on collision. The cleaned name must stay inside root.
|
|
// Shared by single-file and chunked upload completion.
|
|
func (u *U) UniquePath(root, name string) (string, error) {
|
|
ext := filepath.Ext(name)
|
|
base := strings.TrimSuffix(name, ext)
|
|
for i := 0; ; i++ {
|
|
cand := base + ext
|
|
if i > 0 {
|
|
cand = base + " (" + strconv.Itoa(i) + ")" + ext
|
|
}
|
|
p := filepath.Join(root, cand)
|
|
if filepath.Clean(p) != filepath.Join(root, filepath.Clean(cand)) ||
|
|
!strings.HasPrefix(filepath.Clean(p), root+string(os.PathSeparator)) {
|
|
return "", os.ErrInvalid
|
|
}
|
|
if _, e := os.Stat(p); os.IsNotExist(e) {
|
|
return p, nil
|
|
} else if e != nil {
|
|
return "", e
|
|
}
|
|
if i > 999 {
|
|
return "", os.ErrExist
|
|
}
|
|
}
|
|
}
|