// 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)