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 }