Files
Felis/internal/fileedit/stage.go
T

298 lines
9.6 KiB
Go

package fileedit
import (
"crypto/rand"
"crypto/sha256"
"crypto/subtle"
"encoding/hex"
"errors"
"fmt"
"hash"
"io"
"os"
"sync"
"syscall"
"time"
)
// Stage holds uploads between the request that brought them and the Job that
// lands them. The bytes cannot ride the Job spec the way an edit does (etcd caps
// an object near 1.5 MiB, and execve one environment string at 128 KiB), and
// felis-api cannot mount the world volume, so felis-api keeps the upload on its
// own disk and serves it once, on its internal face, to the Job it created for it
// (cmd/felis files, fetchUpload).
//
// Each staged upload is opened by an unguessable id in the URL plus a token the
// Job carries in its environment. Only a digest of the token is kept, compared in
// constant time, and the first successful Open spends it: the Job never retries,
// so a second Open could only be someone else. The release func Put returns
// deletes the file once the Job has answered, whatever it answered. A file too
// big for one request arrives in parts instead (session.go) and is fetched the
// same way.
//
// Nothing here outlives the process: the index is in memory, so Sweep empties
// Dir at startup of whatever a previous process left behind.
type Stage struct {
// Dir holds the staged files. cmd/felis puts it on the uploads volume, which
// has a real capacity; the pod's /tmp is the node's own disk.
Dir string
// MinFree is the share of Dir's filesystem an upload must leave free
// (DefaultStageMinFree when zero), so a burst of uploads cannot fill the disk
// the submission store shares.
MinFree float64
// Now is the clock sessions are aged by (time.Now when nil).
Now func() time.Time
mu sync.Mutex
items map[string]*stagedFile
sessions map[string]*session
reserved int64
}
// DefaultStageMinFree matches the submission store's own floor
// (submit.DefaultUploadsMinFree): the two share the uploads volume.
const DefaultStageMinFree = 0.10
type stagedFile struct {
path string
tokenHash [sha256.Size]byte
size int64
used bool
}
// Staged is one upload on the stage: where the Job fetches it and what it must
// be.
type Staged struct {
ID string
Token string
SHA256 string
Size int64
}
var (
// ErrStageFull is an upload that would leave less than MinFree of the staging
// filesystem free, or that ran it out of space outright.
ErrStageFull = errors.New("fileedit: no room to stage the upload")
// ErrShortUpload is a body that ended, or broke, before its declared length.
ErrShortUpload = errors.New("fileedit: the upload ended before its declared length")
// ErrNotStaged is an Open with an unknown id, a wrong token, or a spent one.
// They are one error on purpose: the internal face answers all three the same.
ErrNotStaged = errors.New("fileedit: no such staged upload")
// ErrNoDigest is bytes sent without the SHA-256 the client computed over
// them, so what arrived cannot be told apart from what was sent.
ErrNoDigest = errors.New("fileedit: the upload carries no SHA-256 digest")
// ErrDigestMismatch is bytes that do not hash to the digest they were sent
// with: they were changed on the way.
ErrDigestMismatch = errors.New("fileedit: the bytes that arrived do not match the digest they were sent with")
)
// checkDigest compares the digest of what arrived with the one it was sent with.
func checkDigest(got, want []byte) error {
if subtle.ConstantTimeCompare(got, want) != 1 {
return fmt.Errorf("%w: they hash to %x, sent as %x", ErrDigestMismatch, got, want)
}
return nil
}
// statfs reports a filesystem's available and total bytes. A var so a test can
// stage against a disk of a chosen size.
var statfs = func(dir string) (avail, total uint64, err error) {
var st syscall.Statfs_t
if err := syscall.Statfs(dir, &st); err != nil {
return 0, 0, err
}
bsize := uint64(st.Bsize) // uint32 on darwin
return uint64(st.Bavail) * bsize, uint64(st.Blocks) * bsize, nil
}
// Sweep deletes everything under Dir. Call it once, before the first Put.
func (s *Stage) Sweep() error {
if err := os.RemoveAll(s.Dir); err != nil {
return fmt.Errorf("fileedit: clear the upload stage: %w", err)
}
return nil
}
// Put stages exactly size bytes from body and returns the handle a Job fetches
// it by, plus the func that deletes it. body must end right after size bytes (an
// HTTP body with that Content-Length does): Put reads to its end, which is also
// what tells the server the body is done.
//
// want is the SHA-256 the client computed over the bytes it sent (the request's
// Content-Digest). Bytes that hash to anything else were changed on the way and
// are refused with ErrDigestMismatch; nothing is staged without one.
func (s *Stage) Put(body io.Reader, size int64, want []byte) (Staged, func(), error) {
if size < 0 {
return Staged{}, nil, fmt.Errorf("fileedit: an upload of %d bytes", size)
}
if len(want) != sha256.Size {
return Staged{}, nil, ErrNoDigest
}
if err := os.MkdirAll(s.Dir, 0o700); err != nil {
return Staged{}, nil, fmt.Errorf("fileedit: create the upload stage: %w", err)
}
if err := s.reserve(size); err != nil {
return Staged{}, nil, err
}
defer s.unreserve(size)
f, err := os.CreateTemp(s.Dir, "upload-*")
if err != nil {
return Staged{}, nil, fmt.Errorf("fileedit: stage the upload: %w", err)
}
h := sha256.New()
src := &bodyReader{r: body}
// One byte past size, so the read that finds the end happens here.
n, copyErr := io.Copy(io.MultiWriter(f, h), io.LimitReader(src, size+1))
closeErr := f.Close()
err = stageFailure(src.err, copyErr, closeErr, n, size)
if err == nil {
err = checkDigest(h.Sum(nil), want)
}
if err != nil {
os.Remove(f.Name())
return Staged{}, nil, err
}
st, tokenHash, err := newHandle(h, size)
if err != nil {
os.Remove(f.Name())
return Staged{}, nil, err
}
s.mu.Lock()
if s.items == nil {
s.items = map[string]*stagedFile{}
}
s.items[st.ID] = &stagedFile{path: f.Name(), tokenHash: tokenHash, size: size}
s.mu.Unlock()
release := func() {
s.mu.Lock()
delete(s.items, st.ID)
s.mu.Unlock()
os.Remove(f.Name())
}
return st, release, nil
}
// stageFailure decides what a finished copy means. The body's own error, or a
// body shorter than promised, is the caller's; a body longer than promised
// cannot come through net/http, which stops at Content-Length, but a direct
// caller could send one and it is refused all the same.
func stageFailure(readErr, copyErr, closeErr error, n, size int64) error {
switch {
case readErr != nil:
return fmt.Errorf("%w: %v", ErrShortUpload, readErr)
case copyErr == nil && n < size:
return fmt.Errorf("%w: got %d of %d bytes", ErrShortUpload, n, size)
case copyErr == nil && n > size:
return fmt.Errorf("fileedit: the upload is longer than its declared %d bytes", size)
}
err := copyErr
if err == nil {
err = closeErr
}
if err == nil {
return nil
}
if errors.Is(err, syscall.ENOSPC) || errors.Is(err, syscall.EDQUOT) {
return fmt.Errorf("%w: %v", ErrStageFull, err)
}
return fmt.Errorf("fileedit: stage the upload: %w", err)
}
// newHandle mints the id and token for a staged upload whose bytes h hashed.
func newHandle(h hash.Hash, size int64) (Staged, [sha256.Size]byte, error) {
id, err := randomHex(16)
if err != nil {
return Staged{}, [sha256.Size]byte{}, fmt.Errorf("fileedit: generate an upload id: %w", err)
}
token, err := randomHex(32)
if err != nil {
return Staged{}, [sha256.Size]byte{}, fmt.Errorf("fileedit: generate an upload token: %w", err)
}
st := Staged{ID: id, Token: token, SHA256: hex.EncodeToString(h.Sum(nil)), Size: size}
return st, sha256.Sum256([]byte(st.Token)), nil
}
// randomHex is n random bytes in hex.
func randomHex(n int) (string, error) {
b := make([]byte, n)
if _, err := rand.Read(b); err != nil {
return "", err
}
return hex.EncodeToString(b), nil
}
// reserve admits an upload of size bytes if the disk keeps MinFree free after it
// and after every upload still being written. Those have not reached the disk
// yet, so statfs alone would let two of them through on room for one.
func (s *Stage) reserve(size int64) error {
avail, total, err := statfs(s.Dir)
if err != nil {
return fmt.Errorf("fileedit: measure the upload stage: %w", err)
}
minFree := s.MinFree
if minFree <= 0 {
minFree = DefaultStageMinFree
}
floor := uint64(float64(total) * minFree)
s.mu.Lock()
defer s.mu.Unlock()
need := uint64(s.reserved) + uint64(size)
if avail < need || avail-need < floor {
return fmt.Errorf("%w: %d MiB free of %d MiB, and staging %d MiB would leave less than %.0f%% free",
ErrStageFull, avail>>20, total>>20, need>>20, minFree*100)
}
s.reserved += size
return nil
}
func (s *Stage) unreserve(size int64) {
s.mu.Lock()
s.reserved -= size
s.mu.Unlock()
}
// Open spends a staged upload's token, or a sealed session's (Seal), and returns
// its file and size. Any mismatch is ErrNotStaged.
func (s *Stage) Open(id, token string) (*os.File, int64, error) {
sum := sha256.Sum256([]byte(token))
s.mu.Lock()
name, size, found, err := s.openSession(id, sum)
if !found {
it, ok := s.items[id]
if !ok || it.used || subtle.ConstantTimeCompare(sum[:], it.tokenHash[:]) != 1 {
err = ErrNotStaged
} else {
it.used = true
name, size = it.path, it.size
}
}
s.mu.Unlock()
if err != nil {
return nil, 0, err
}
f, err := os.Open(name)
if err != nil {
return nil, 0, fmt.Errorf("fileedit: open the staged upload: %w", err)
}
return f, size, nil
}
// bodyReader remembers the body's own read error, so Put can tell a client that
// went away from the disk filling up.
type bodyReader struct {
r io.Reader
err error
}
func (b *bodyReader) Read(p []byte) (int, error) {
n, err := b.r.Read(p)
if err != nil && err != io.EOF {
b.err = err
}
return n, err
}