245 lines
8.9 KiB
Go
245 lines
8.9 KiB
Go
package submit
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"regexp"
|
|
"syscall"
|
|
)
|
|
|
|
// contextBlobName is the fixed object name of a submission's build context under
|
|
// its id-namespaced prefix. It is the single source of truth for both the
|
|
// derived context ref (deriveContextRef) and the on-disk write target
|
|
// (LocalContextStore), so the blob always lands exactly where Kaniko's
|
|
// --context points (build/jobspec.go).
|
|
const contextBlobName = "context.tar.gz"
|
|
|
|
// idRE re-validates a submission id at the storage boundary. The Manager only
|
|
// ever passes ids it loaded from the Store (already the crypto-hex ids newID
|
|
// mints), but LocalContextStore is a standalone component that treats the id as
|
|
// untrusted path input: a lowercase-alphanumeric-with-dashes id can contain no
|
|
// path separator and no "..", so it can never escape Base. This is the same
|
|
// defense-in-depth stance as internal/backup's zip-slip guard.
|
|
var idRE = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,127}$`)
|
|
|
|
// LocalContextStore is the filesystem-backed build-context blob store: it writes
|
|
// each submission's uploaded modpack to {Base}/{id}/context.tar.gz on a mounted
|
|
// PVC. It is one of two implemented Blobs backends — cmd/felis selects it when
|
|
// user_uploads_context is a local path and S3ContextStore when it is an s3:// base;
|
|
// a base that is neither (or an s3:// base with no credentials configured) leaves
|
|
// Manager.Blobs nil so the upload endpoint returns 503 rather than pretending to
|
|
// accept a file it cannot persist.
|
|
//
|
|
// Base MUST equal the Manager's ContextStore; cmd/felis wires both from the one
|
|
// config field (registry.user_uploads_context).
|
|
//
|
|
// The build Pod never mounts this PVC. With Manager.ContextBaseURL set (every
|
|
// installed API) the derived context ref is the internal face's
|
|
// /api/v1/internal/submissions/{id}/context route, which streams the blob out
|
|
// of this store to the build Job's `felis fetch-context` step.
|
|
type LocalContextStore struct {
|
|
// Base is the directory (uploads PVC mount) submission contexts are written
|
|
// under. Each submission gets its own {Base}/{id}/ subdirectory.
|
|
Base string
|
|
// MinFree is the share of Base's filesystem an upload must leave free; 0 uses
|
|
// DefaultUploadsMinFree.
|
|
MinFree float64
|
|
|
|
// sync flushes a file or directory to disk; nil is (*os.File).Sync. Tests
|
|
// replace it to watch or fail the flushes.
|
|
sync func(*os.File) error
|
|
}
|
|
|
|
// DefaultUploadsMinFree is the share of the uploads filesystem an upload must
|
|
// leave free. On k3s local-path the uploads PVC is a directory on the node's
|
|
// disk, beside the worlds and the database, and below about a tenth free the
|
|
// kubelet starts evicting pods (the same floor backup.MinFreeAfter keeps).
|
|
const DefaultUploadsMinFree = 0.10
|
|
|
|
// CheckRoom refuses an upload of up to need bytes that could push Base's
|
|
// filesystem below its free floor.
|
|
func (s *LocalContextStore) CheckRoom(need int64) error {
|
|
return checkRoom(s.Base, need, s.MinFree)
|
|
}
|
|
|
|
// checkRoom refuses need more bytes under dir when they could push its filesystem
|
|
// below the minFree share (DefaultUploadsMinFree when 0), wrapping ErrUploadsFull.
|
|
func checkRoom(dir string, need int64, minFree float64) error {
|
|
var st syscall.Statfs_t
|
|
if err := syscall.Statfs(dir, &st); err != nil {
|
|
return fmt.Errorf("submit: measure the uploads store %s: %w", dir, err)
|
|
}
|
|
bsize := uint64(st.Bsize) // uint32 on darwin
|
|
total, avail := uint64(st.Blocks)*bsize, uint64(st.Bavail)*bsize
|
|
if total == 0 {
|
|
return nil
|
|
}
|
|
if minFree <= 0 {
|
|
minFree = DefaultUploadsMinFree
|
|
}
|
|
floor := uint64(float64(total) * minFree)
|
|
if n := uint64(max(need, 0)); avail < n || avail-n < floor {
|
|
return fmt.Errorf("%w: %d MiB free of %d MiB, and an upload of up to %d MiB would leave less than %.0f%% free",
|
|
ErrUploadsFull, avail>>20, total>>20, n>>20, minFree*100)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// dir returns the per-submission directory, rejecting an id that could escape
|
|
// Base. Every path the store touches is rooted here.
|
|
func (s *LocalContextStore) dir(id string) (string, error) {
|
|
if !idRE.MatchString(id) {
|
|
return "", fmt.Errorf("submit: invalid submission id %q", id)
|
|
}
|
|
return filepath.Join(s.Base, id), nil
|
|
}
|
|
|
|
// Put writes r to {Base}/{id}/context.tar.gz atomically: it streams into a temp
|
|
// file in the same directory and renames it over any previous upload only on a
|
|
// fully successful copy. So a failed, truncated, or oversize upload never
|
|
// replaces a good context and never leaves a half-written blob for Kaniko to
|
|
// read; a re-upload while the submission is still pending simply supersedes the
|
|
// previous one. It returns the number of bytes stored, and once it does the blob
|
|
// survives a crash: the bytes are flushed before the rename and the directory
|
|
// entries after it.
|
|
func (s *LocalContextStore) Put(_ context.Context, id string, r io.Reader) (int64, error) {
|
|
dir, err := s.dir(id)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
if err := os.MkdirAll(dir, 0o750); err != nil {
|
|
return 0, fmt.Errorf("submit: mkdir context dir: %w", err)
|
|
}
|
|
tmp, err := os.CreateTemp(dir, contextBlobName+".*.tmp")
|
|
if err != nil {
|
|
return 0, fmt.Errorf("submit: create temp context: %w", err)
|
|
}
|
|
tmpName := tmp.Name()
|
|
n, err := io.Copy(tmp, r)
|
|
if err == nil {
|
|
// The bytes reach the disk before the rename can publish them, so a crash
|
|
// right after it finds the whole blob behind the digest the database
|
|
// already holds.
|
|
err = s.flush(tmp)
|
|
}
|
|
if err != nil {
|
|
tmp.Close()
|
|
os.Remove(tmpName)
|
|
return 0, fmt.Errorf("submit: write context blob: %w", err)
|
|
}
|
|
if err := tmp.Close(); err != nil {
|
|
os.Remove(tmpName)
|
|
return 0, fmt.Errorf("submit: close context blob: %w", err)
|
|
}
|
|
if err := os.Rename(tmpName, filepath.Join(dir, contextBlobName)); err != nil {
|
|
os.Remove(tmpName)
|
|
return 0, fmt.Errorf("submit: commit context blob: %w", err)
|
|
}
|
|
// The rename, and a first upload's new directory, are entries in their parent
|
|
// directories: flush those too, or a crash can undo them.
|
|
for _, d := range []string{dir, s.Base} {
|
|
if err := s.syncDir(d); err != nil {
|
|
return 0, fmt.Errorf("submit: sync context dir: %w", err)
|
|
}
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
func (s *LocalContextStore) flush(f *os.File) error {
|
|
if s.sync != nil {
|
|
return s.sync(f)
|
|
}
|
|
return f.Sync()
|
|
}
|
|
|
|
// syncDir flushes dir's entries to disk. A filesystem that cannot sync a
|
|
// directory has no stronger promise to give.
|
|
func (s *LocalContextStore) syncDir(dir string) error {
|
|
d, err := os.Open(dir)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer d.Close()
|
|
if err := s.flush(d); err != nil && !errors.Is(err, syscall.EINVAL) && !errors.Is(err, syscall.ENOTSUP) {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Exists reports whether a context blob has been stored for id. Approve consults
|
|
// it so a submission whose context was never uploaded is refused BEFORE the CAS,
|
|
// instead of being approved into a build Kaniko cannot pull.
|
|
func (s *LocalContextStore) Exists(_ context.Context, id string) (bool, error) {
|
|
dir, err := s.dir(id)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
switch _, err := os.Stat(filepath.Join(dir, contextBlobName)); {
|
|
case err == nil:
|
|
return true, nil
|
|
case os.IsNotExist(err):
|
|
return false, nil
|
|
default:
|
|
return false, fmt.Errorf("submit: stat context blob: %w", err)
|
|
}
|
|
}
|
|
|
|
// Open returns the stored context blob for id — the read side of the transport the
|
|
// build Pod's fetch initContainer uses. A missing blob is ErrBlobNotFound (404 on
|
|
// the route), never a bare os error, so the API keeps its status mapping.
|
|
func (s *LocalContextStore) Open(_ context.Context, id string) (io.ReadCloser, error) {
|
|
dir, err := s.dir(id)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
f, err := os.Open(filepath.Join(dir, contextBlobName))
|
|
if err != nil {
|
|
if os.IsNotExist(err) {
|
|
return nil, fmt.Errorf("%w: %v", ErrBlobNotFound, err)
|
|
}
|
|
return nil, fmt.Errorf("submit: open context blob: %w", err)
|
|
}
|
|
return f, nil
|
|
}
|
|
|
|
// Size reports the stored blob's size — the accounting read behind the per-user
|
|
// storage budget. A missing blob is (0, false, nil): absence is not an error
|
|
// here, it is simply no bytes to count (the same distinction Exists draws for
|
|
// Approve).
|
|
func (s *LocalContextStore) Size(_ context.Context, id string) (int64, bool, error) {
|
|
dir, err := s.dir(id)
|
|
if err != nil {
|
|
return 0, false, err
|
|
}
|
|
fi, err := os.Stat(filepath.Join(dir, contextBlobName))
|
|
switch {
|
|
case err == nil:
|
|
return fi.Size(), true, nil
|
|
case os.IsNotExist(err):
|
|
return 0, false, nil
|
|
default:
|
|
return 0, false, fmt.Errorf("submit: stat context blob: %w", err)
|
|
}
|
|
}
|
|
|
|
// Delete removes everything stored for id — the blob plus its id-namespaced
|
|
// directory (a stray temp file from an interrupted upload goes with it) — after
|
|
// the submission row is gone. Idempotent: nothing stored is success.
|
|
func (s *LocalContextStore) Delete(_ context.Context, id string) error {
|
|
dir, err := s.dir(id)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if err := os.RemoveAll(dir); err != nil {
|
|
return fmt.Errorf("submit: remove context blob: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Compile-time proof that the filesystem store satisfies the Blobs transport.
|
|
var _ Blobs = (*LocalContextStore)(nil)
|