407 lines
9.8 KiB
Go
407 lines
9.8 KiB
Go
// Package store implements the flat-file object store: one directory per
|
|
// object, holding the blob and a JSON metadata sidecar. There is no database;
|
|
// an in-memory index is rebuilt from disk at startup and kept in sync.
|
|
package store
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"hash"
|
|
"io"
|
|
"io/fs"
|
|
"os"
|
|
"path/filepath"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
const (
|
|
blobName = "blob"
|
|
partName = "blob.part"
|
|
metaName = "meta.json"
|
|
|
|
dirPerm fs.FileMode = 0o775
|
|
filePerm fs.FileMode = 0o664
|
|
|
|
// debrisMaxAge is how long an object directory with no metadata is left
|
|
// alone before being treated as the remains of a killed upload.
|
|
debrisMaxAge = 24 * time.Hour
|
|
)
|
|
|
|
var (
|
|
ErrNotFound = errors.New("object not found")
|
|
ErrExists = errors.New("name already taken")
|
|
ErrTooLarge = errors.New("upload exceeds the size limit")
|
|
)
|
|
|
|
// Store owns the data directory.
|
|
type Store struct {
|
|
dir string
|
|
objects string
|
|
|
|
// root confines every object file operation to the objects directory.
|
|
// The data directory is group-writable by design, so a symlink planted
|
|
// there must not be able to redirect a read or a write outside it.
|
|
root *os.Root
|
|
|
|
mu sync.RWMutex
|
|
index map[string]*Meta
|
|
total int64
|
|
}
|
|
|
|
// Open prepares the data directory and rebuilds the index from it.
|
|
func Open(dir string) (*Store, error) {
|
|
s := &Store{
|
|
dir: dir,
|
|
objects: filepath.Join(dir, "objects"),
|
|
index: make(map[string]*Meta),
|
|
}
|
|
if err := os.MkdirAll(s.objects, dirPerm); err != nil {
|
|
return nil, err
|
|
}
|
|
root, err := os.OpenRoot(s.objects)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
s.root = root
|
|
if err := s.load(); err != nil {
|
|
return nil, err
|
|
}
|
|
return s, nil
|
|
}
|
|
|
|
// DataDir is the directory the store was opened on.
|
|
func (s *Store) DataDir() string { return s.dir }
|
|
|
|
// Close releases the handle on the objects directory.
|
|
func (s *Store) Close() error { return s.root.Close() }
|
|
|
|
// objectDir builds an object's path for display. Actual file operations go
|
|
// through s.root instead, which cannot be walked out of.
|
|
func (s *Store) objectDir(id string) string { return filepath.Join(s.objects, id) }
|
|
|
|
// within builds a root-relative path for one of an object's files.
|
|
func within(id, name string) string { return id + "/" + name }
|
|
|
|
func (s *Store) load() error {
|
|
entries, err := os.ReadDir(s.objects)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, e := range entries {
|
|
if !e.IsDir() {
|
|
continue
|
|
}
|
|
id, err := CleanID(e.Name())
|
|
if err != nil || id != e.Name() {
|
|
// Not a name this service could have created; leave it be.
|
|
continue
|
|
}
|
|
m, err := s.readMeta(id)
|
|
if err != nil {
|
|
continue // incomplete or unreadable; the debris sweep handles it
|
|
}
|
|
// A leftover .part in a committed object is always stale at startup.
|
|
s.root.Remove(within(id, partName))
|
|
s.index[id] = m
|
|
s.total += m.Size
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Store) readMeta(id string) (*Meta, error) {
|
|
b, err := s.root.ReadFile(within(id, metaName))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
var m Meta
|
|
if err := json.Unmarshal(b, &m); err != nil {
|
|
return nil, err
|
|
}
|
|
if m.ID != id {
|
|
return nil, fmt.Errorf("store: metadata for %q claims id %q", id, m.ID)
|
|
}
|
|
return &m, nil
|
|
}
|
|
|
|
// Exists reports whether a name is currently taken, expired objects included:
|
|
// a name stays claimed until its object is actually removed.
|
|
func (s *Store) Exists(id string) bool {
|
|
_, err := s.root.Lstat(id)
|
|
return err == nil
|
|
}
|
|
|
|
// Reserve claims a name by creating its directory. os.Mkdir is atomic, so this
|
|
// is the point at which a vanity collision is detected - before any of the
|
|
// caller's body has been read.
|
|
func (s *Store) Reserve(id string) (*Upload, error) {
|
|
if err := s.root.Mkdir(id, dirPerm); err != nil {
|
|
if errors.Is(err, fs.ErrExist) {
|
|
return nil, ErrExists
|
|
}
|
|
return nil, err
|
|
}
|
|
f, err := s.root.OpenFile(within(id, partName), os.O_WRONLY|os.O_CREATE|os.O_EXCL, filePerm)
|
|
if err != nil {
|
|
s.root.RemoveAll(id)
|
|
return nil, err
|
|
}
|
|
return &Upload{s: s, id: id, f: f, h: sha256.New()}, nil
|
|
}
|
|
|
|
// ReserveRandom claims a fresh UUIDv4 name.
|
|
func (s *Store) ReserveRandom() (*Upload, error) {
|
|
for range 8 {
|
|
id, err := NewUUID()
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
u, err := s.Reserve(id)
|
|
if errors.Is(err, ErrExists) {
|
|
continue // astronomically unlikely; retry regardless
|
|
}
|
|
return u, err
|
|
}
|
|
return nil, errors.New("store: could not allocate an unused id")
|
|
}
|
|
|
|
// Upload is an in-flight object. It is an io.Writer so callers can stream a
|
|
// request body straight to disk; nothing is ever buffered in memory.
|
|
type Upload struct {
|
|
s *Store
|
|
id string
|
|
f *os.File
|
|
h hash.Hash
|
|
n int64
|
|
limit int64 // 0 means unlimited
|
|
done bool
|
|
}
|
|
|
|
func (u *Upload) ID() string { return u.id }
|
|
func (u *Upload) Size() int64 { return u.n }
|
|
|
|
// SetLimit caps the number of bytes the upload will accept. The cap is applied
|
|
// to bytes actually written, never to a declared Content-Length.
|
|
func (u *Upload) SetLimit(n int64) { u.limit = n }
|
|
|
|
func (u *Upload) Write(p []byte) (int, error) {
|
|
if u.limit > 0 && u.n+int64(len(p)) > u.limit {
|
|
return 0, ErrTooLarge
|
|
}
|
|
n, err := u.f.Write(p)
|
|
u.n += int64(n)
|
|
u.h.Write(p[:n])
|
|
return n, err
|
|
}
|
|
|
|
// Commit makes the object visible. The ordering matters: the blob is durable
|
|
// and in place before the metadata that advertises it is written, and the
|
|
// metadata is renamed into place atomically.
|
|
func (u *Upload) Commit(m *Meta) error {
|
|
if u.done {
|
|
return errors.New("store: upload already finished")
|
|
}
|
|
m.ID = u.id
|
|
m.Size = u.n
|
|
m.SHA256 = hex.EncodeToString(u.h.Sum(nil))
|
|
|
|
if err := u.f.Sync(); err != nil {
|
|
return err
|
|
}
|
|
if err := u.f.Close(); err != nil {
|
|
return err
|
|
}
|
|
if err := u.s.root.Rename(within(u.id, partName), within(u.id, blobName)); err != nil {
|
|
return err
|
|
}
|
|
if err := u.s.writeMetaAtomic(u.id, m); err != nil {
|
|
return err
|
|
}
|
|
if err := u.s.syncDir(u.id); err != nil {
|
|
return err
|
|
}
|
|
u.done = true
|
|
|
|
u.s.mu.Lock()
|
|
u.s.index[u.id] = m
|
|
u.s.total += m.Size
|
|
u.s.mu.Unlock()
|
|
return nil
|
|
}
|
|
|
|
// Abort discards an incomplete upload, releasing its name.
|
|
func (u *Upload) Abort() {
|
|
if u.done {
|
|
return
|
|
}
|
|
u.done = true
|
|
u.f.Close()
|
|
u.s.root.RemoveAll(u.id)
|
|
}
|
|
|
|
// writeMetaAtomic serialises m to a temporary file in the object's own
|
|
// directory, fsyncs it, and renames it into place. Only once this rename lands
|
|
// does the object become visible to a reader.
|
|
func (s *Store) writeMetaAtomic(id string, m *Meta) error {
|
|
b, err := json.MarshalIndent(m, "", " ")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
b = append(b, '\n')
|
|
|
|
suffix, err := NewSecret()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
tmpPath := within(id, "."+metaName+"."+suffix[:16])
|
|
|
|
tmp, err := s.root.OpenFile(tmpPath, os.O_WRONLY|os.O_CREATE|os.O_EXCL, filePerm)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer s.root.Remove(tmpPath) // no-op once the rename succeeds
|
|
|
|
if _, err := tmp.Write(b); err != nil {
|
|
tmp.Close()
|
|
return err
|
|
}
|
|
if err := tmp.Sync(); err != nil {
|
|
tmp.Close()
|
|
return err
|
|
}
|
|
if err := tmp.Close(); err != nil {
|
|
return err
|
|
}
|
|
return s.root.Rename(tmpPath, within(id, metaName))
|
|
}
|
|
|
|
// syncDir flushes a directory entry so a rename survives a power loss.
|
|
func (s *Store) syncDir(id string) error {
|
|
d, err := s.root.Open(id)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer d.Close()
|
|
return d.Sync()
|
|
}
|
|
|
|
// Get returns an object's metadata, treating an expired object as absent and
|
|
// removing it on the spot. Expiry is checked here, on every read, so a stalled
|
|
// sweeper can never serve a file past its lifetime.
|
|
func (s *Store) Get(id string, now time.Time) (*Meta, error) {
|
|
s.mu.RLock()
|
|
m, ok := s.index[id]
|
|
s.mu.RUnlock()
|
|
if !ok {
|
|
return nil, ErrNotFound
|
|
}
|
|
if m.Expired(now) {
|
|
s.Delete(id)
|
|
return nil, ErrNotFound
|
|
}
|
|
return m, nil
|
|
}
|
|
|
|
// OpenBlob returns the metadata and an open handle to the object's bytes.
|
|
func (s *Store) OpenBlob(id string, now time.Time) (*Meta, *os.File, error) {
|
|
m, err := s.Get(id, now)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
f, err := s.root.Open(within(id, blobName))
|
|
if err != nil {
|
|
// Metadata without a blob means the data directory was tampered with.
|
|
s.Delete(id)
|
|
return nil, nil, ErrNotFound
|
|
}
|
|
return m, f, nil
|
|
}
|
|
|
|
// Delete removes an object and frees its name.
|
|
func (s *Store) Delete(id string) error {
|
|
s.mu.Lock()
|
|
if m, ok := s.index[id]; ok {
|
|
s.total -= m.Size
|
|
delete(s.index, id)
|
|
}
|
|
s.mu.Unlock()
|
|
return s.root.RemoveAll(id)
|
|
}
|
|
|
|
// Total reports the number of bytes currently stored.
|
|
func (s *Store) Total() int64 {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.total
|
|
}
|
|
|
|
// Count reports the number of live objects.
|
|
func (s *Store) Count() int {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return len(s.index)
|
|
}
|
|
|
|
// Sweep removes expired objects and long-abandoned upload directories. It
|
|
// returns the number of objects removed.
|
|
func (s *Store) Sweep(now time.Time) int {
|
|
s.mu.RLock()
|
|
var expired []string
|
|
for id, m := range s.index {
|
|
if m.Expired(now) {
|
|
expired = append(expired, id)
|
|
}
|
|
}
|
|
s.mu.RUnlock()
|
|
|
|
for _, id := range expired {
|
|
s.Delete(id)
|
|
}
|
|
s.sweepDebris(now)
|
|
return len(expired)
|
|
}
|
|
|
|
// sweepDebris removes object directories that never gained metadata and are
|
|
// older than debrisMaxAge - the remains of an upload killed mid-flight.
|
|
func (s *Store) sweepDebris(now time.Time) {
|
|
entries, err := os.ReadDir(s.objects)
|
|
if err != nil {
|
|
return
|
|
}
|
|
for _, e := range entries {
|
|
if !e.IsDir() {
|
|
continue
|
|
}
|
|
s.mu.RLock()
|
|
_, live := s.index[e.Name()]
|
|
s.mu.RUnlock()
|
|
if live {
|
|
continue
|
|
}
|
|
info, err := e.Info()
|
|
if err != nil || now.Sub(info.ModTime()) < debrisMaxAge {
|
|
continue
|
|
}
|
|
if _, err := s.root.Stat(within(e.Name(), metaName)); err == nil {
|
|
continue // has metadata but is not indexed; leave it for a human
|
|
}
|
|
s.root.RemoveAll(e.Name())
|
|
}
|
|
}
|
|
|
|
// List returns every live object's metadata, for administrative use.
|
|
func (s *Store) List() []*Meta {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
out := make([]*Meta, 0, len(s.index))
|
|
for _, m := range s.index {
|
|
out = append(out, m)
|
|
}
|
|
return out
|
|
}
|
|
|
|
var _ io.Writer = (*Upload)(nil)
|