Initial commit
This commit is contained in:
@@ -0,0 +1,406 @@
|
||||
// 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)
|
||||
Reference in New Issue
Block a user