diff --git a/.beans/ATFS-skt5--no-storage-quota-anywhere-uploads-pin-intents-and.md b/.beans/ATFS-skt5--no-storage-quota-anywhere-uploads-pin-intents-and.md index 6cb5b77..622709b 100644 --- a/.beans/ATFS-skt5--no-storage-quota-anywhere-uploads-pin-intents-and.md +++ b/.beans/ATFS-skt5--no-storage-quota-anywhere-uploads-pin-intents-and.md @@ -1,11 +1,11 @@ --- # ATFS-skt5 title: 'No storage quota anywhere: uploads, pin intents and follows can fill the volume' -status: todo +status: completed type: bug priority: normal created_at: 2026-08-12T04:12:24Z -updated_at: 2026-08-12T04:12:24Z +updated_at: 2026-08-12T07:55:25Z --- **Severity: Medium.** Nothing anywhere caps total stored bytes, blob count, @@ -65,3 +65,72 @@ is the product): Each limit should log when it bites — a silently-refusing appliance is as bad as a silently-full one. + +## Summary of Changes + +One storage ceiling, derived from the data volume rather than configured, +enforced at the three places content can arrive. + +**`internal/store` does the accounting.** `store.ErrNoRoom` is the single +error every ceiling reports. `Store.CheckRoom` refuses a write once the +volume has fallen to its reserve floor — a twentieth of the volume, never +below 512 MiB and never above 4 GiB — and `Store.Put` calls it as a backstop, +so no path that stores bytes can skip it. `Store.Reserve(n)` is the admission +gate proper: it promises n bytes for a write's duration and counts +outstanding promises against the free figure, which is exactly what a bare +free-space check can't do (N concurrent writers all see the same room and +then collectively write N times it). Free space comes from `statfs(2)` via +`golang.org/x/sys/unix` in a `linux || darwin` file with a "couldn't measure" +fallback elsewhere; a volume atfs can't measure is never one it refuses to +write on. Both build paths verified — `make check` on darwin/arm64, plus +`vet-gosd` and `CGO_ENABLED=0` linux/arm64 and linux/amd64 builds. + +**Uploads.** `uploadFile` reserves before reading a byte — the declared +Content-Length where there is one, so a nearly-full instance still takes +small uploads, and the whole blob limit for a chunked body — and answers +`InsufficientStorage` when it can't: a named lexicon error at 400, like every +other one atfs reports. + +**Pins.** `pin.Manager.Request` refuses a *new* intent once the volume is at +its floor or `maxIntents` (4096) already exist; joining or re-triggering an +existing intent is never refused. Each fetch attempt reserves its declared +size for the attempt's duration, so `fetchConcurrency` parallel fetches can't +overrun the floor between them. Running out is ambient, not definitive, so an +attempt backs off and retries rather than burning the intent. + +**Follow.** `request` now tells `ErrNoRoom` apart from the failures it +swallows: it stops the pass, logs once, and defers the rest to the next poll +rather than re-learning the same fact once per listed file. Absence counters +are cleared in their own pass over the listing first, so stopping early can't +strand a stale count that would wrongly release a file the origin still +lists. + +Every ceiling logs when it bites, both lexicons gained the +`InsufficientStorage` error, and the README has a "Storage ceiling" section. + +### Deviations from the proposed fix + +- **No separate concurrent-upload semaphore.** Reservations already bound + concurrent uploads by the resource actually at risk: each in-flight upload + holds a promise, so concurrency caps at `(free − floor) / reservation`, + adaptively. A fixed count would be either redundant (above what space + allows) or an arbitrary extra refusal (below it), and would need a second + error word for a refusal meaning the same thing. +- **A pin intent doesn't hold a reservation for its whole life**, only per + fetch attempt. Reserving a declared size at registration would let ~13 + `pinFile` calls naming 1 GiB phantom CIDs disable uploads on a 16 GiB + volume for the whole ~49h horizon — a self-inflicted DoS worse than the one + being fixed. The intent-count cap plus the floor bound the same attack + without it. +- **The floor is capped at 4 GiB** as well as floored at 512 MiB: a flat 5% + would withhold 50 GiB of a 1 TB volume for no benefit. +- **400 `InsufficientStorage`, not 503** — atfs maps named lexicon errors to + 400 (see getFile's `BlobNotFound`). + +Left undone deliberately: forgetting a failed intent once its horizon is +spent (the third point above). The intent cap bounds the same growth, and a +failed intent is still the answer to "why did my pin stop?" until its +requesters release it. + +`lexicons/` changed, so `make publish-lexicons` is owed before the next tag — +the release ritual (ATFS-bq8m) already runs it. diff --git a/README.md b/README.md index 95dcdfc..39b292b 100644 --- a/README.md +++ b/README.md @@ -278,6 +278,14 @@ It answers immediately with a state — `seeking`, `fetching`, `pinned` or `fail `POST /xrpc/dev.atfs.repo.deleteFile` (JSON body `{"cid": ""}`) releases the caller's pin on a blob; the content itself is only physically removed once every claimant has released it. The same goes for content that hasn't arrived yet: a pending `pinFile` request is a claim before the fact, so deleting it withdraws your interest, and the fetch is abandoned only when the last account waiting on it has done the same. There's no `com.atproto` alias here — atproto has no equivalent blob-deletion call to mirror, since a PDS just garbage-collects a blob once no record references it. Authorization is by pin-set membership rather than the upload allowlist, so an account can always release its own claims, even after being removed from the allowlist. +## Storage ceiling + +An instance keeps a slice of its data volume permanently free — a twentieth of it, never less than 512 MiB and never more than 4 GiB (see `internal/store`) — because a genuinely full volume takes the metadata, pin records and IPFS index down with the content, and on an SD-card appliance there's no shell to clear it from. There's nothing to configure: the figure is derived from the volume itself. + +Once what's left above that reserve is spoken for, `uploadFile` and `pinFile` answer `400 InsufficientStorage` and store nothing; retrieval, enumeration and deletion carry on untouched, and `deleteFile` is how you make room. Space is *reserved* for the duration of each upload and each pin fetch rather than merely checked, so several large transfers at once can't all be waved through on the same free space. There's a ceiling on outstanding `pinFile` requests too (thousands — far past allowlist scale), since a pin for content that never turns up costs a small file and a retry timer indefinitely. + +A follower that runs out of room stops adding for that poll and picks up where it left off an hour later. It never releases anything over it: the origin still lists what it holds, and being full is this instance's problem, not evidence the origin dropped a file. + ## Enumeration `GET /xrpc/dev.atfs.repo.listFiles?limit=&cursor=` lists this instance's committed, *directly claimed* files, a page at a time — `dev.atfs.file`-shaped entries, the same shape `pinFile` takes as input. `limit` (default 500, max 1000) bounds the page; `cursor` resumes from a previous response's `cursor` field, which is present only when the page filled (its absence marks the last page). Unauthenticated, unlike the upload/pin/delete calls — every cid it lists is already public, announced to the DHT and served at `/ipfs/`. diff --git a/go.mod b/go.mod index 8a65f71..edc695a 100644 --- a/go.mod +++ b/go.mod @@ -18,6 +18,7 @@ require ( github.com/multiformats/go-multibase v0.3.0 github.com/multiformats/go-multihash v0.2.3 golang.org/x/crypto/x509roots/fallback v0.0.0-20260811175631-f44d03d253a1 + golang.org/x/sys v0.47.0 ) require ( @@ -152,7 +153,6 @@ require ( golang.org/x/mod v0.39.0 // indirect golang.org/x/net v0.57.0 // indirect golang.org/x/sync v0.22.0 // indirect - golang.org/x/sys v0.47.0 // indirect golang.org/x/telemetry v0.0.0-20260811182544-a038080d80e5 // indirect golang.org/x/text v0.41.0 // indirect golang.org/x/time v0.15.0 // indirect diff --git a/go.sum b/go.sum index 7b3acd6..a0e8f7b 100644 --- a/go.sum +++ b/go.sum @@ -153,8 +153,6 @@ github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= -github.com/jphastings/gosd v0.5.0 h1:qQP3xtQAg7hfRnwqeeptdlxVMCYlyc9rSs4b8nPgWLI= -github.com/jphastings/gosd v0.5.0/go.mod h1:wGxGRP3kquMxr0x0ddrh5CR6KHx8aAwetki0dA2JoDU= github.com/jphastings/gosd v0.6.0 h1:/14f9x2RTtVqMldEEnKUtB3cWDmmuevfp/5jMzNo1+Y= github.com/jphastings/gosd v0.6.0/go.mod h1:wGxGRP3kquMxr0x0ddrh5CR6KHx8aAwetki0dA2JoDU= github.com/jtolds/gls v4.20.0+incompatible h1:xdiiI2gbIgH/gLH7ADydsJ1uDOEzR8yvV7C0MuV77Wo= diff --git a/internal/follow/follow.go b/internal/follow/follow.go index b7359df..8c30631 100644 --- a/internal/follow/follow.go +++ b/internal/follow/follow.go @@ -310,13 +310,25 @@ func (m *Manager) sync(ctx context.Context, t *target) { mine := m.mirrored(holder) + // Every file the origin listed is present, whatever happens to the + // additions below: clearing the absence counts in their own pass is what + // lets that loop stop early without a file it merely ran out of room for + // carrying a stale count into the release pass. + for c := range remote { + delete(t.absent, c) + } + var added int for c, ref := range remote { - delete(t.absent, c) if mine[c] { continue } - if m.request(ctx, ref, holder) { + taken, err := m.request(ctx, ref, holder) + if err != nil { + slog.Warn("follow: no room for the rest of a followed instance's catalogue, deferring it to the next poll", "uri", t.uri, "added", added, "err", err) + break + } + if taken { added++ } } @@ -378,18 +390,28 @@ func (m *Manager) mirrored(holder string) map[cid.Cid]bool { // request registers a mirror claim on one listed file, and reports whether // anything was actually taken on. A reference this instance could never pin // — over the size limit, structurally broken — is the origin's business, -// not a failure worth retrying every hour. -func (m *Manager) request(ctx context.Context, ref pin.Ref, holder string) bool { +// not a failure worth retrying every hour, and neither is a one-off +// registration failure. +// +// Running out of storage is different in kind, and the only failure returned +// rather than swallowed: it will be just as true of every remaining file in +// the listing, so reporting it per file would bury the one fact that matters +// under thousands of identical warnings, and pressing on would achieve +// nothing. The caller stops its pass — releases are unaffected, since those +// are driven by the origin's own listing and are what makes room again. +func (m *Manager) request(ctx context.Context, ref pin.Ref, holder string) (bool, error) { _, err := m.pins.Request(ctx, ref, store.Pin{DID: holder, Source: store.SourceFollow}) switch { + case errors.Is(err, store.ErrNoRoom): + return false, err case errors.Is(err, pin.ErrUnpinnable): slog.Debug("follow: skipping an unpinnable listed file", "cid", ref.CID, "err", err) - return false + return false, nil case err != nil: slog.Warn("follow: registering a mirror pin failed", "cid", ref.CID, "err", err) - return false + return false, nil } - return true + return true, nil } // release drops holder's mirror claim on one cid — both the claim-in-waiting diff --git a/internal/follow/follow_test.go b/internal/follow/follow_test.go index e71c29f..a79b78a 100644 --- a/internal/follow/follow_test.go +++ b/internal/follow/follow_test.go @@ -172,6 +172,7 @@ type fakePins struct { st *store.Store origin *origin stall bool // register intents instead of storing bytes + noRoom bool // refuse every request the way a full instance does requests []store.Pin pending map[cid.Cid][]store.Pin released []store.Pin @@ -186,6 +187,10 @@ func (f *fakePins) Request(ctx context.Context, ref pin.Ref, claim store.Pin) (p defer f.mu.Unlock() f.requests = append(f.requests, claim) + if f.noRoom { + return pin.Status{}, fmt.Errorf("%w: nothing left to promise", store.ErrNoRoom) + } + if f.stall { f.pending[ref.CID] = append(f.pending[ref.CID], claim) return pin.Status{State: pin.StateSeeking}, nil @@ -251,6 +256,18 @@ func (f *fakePins) claims() []store.Pin { return append([]store.Pin(nil), f.requests...) } +func (f *fakePins) releases() []store.Pin { + f.mu.Lock() + defer f.mu.Unlock() + return append([]store.Pin(nil), f.released...) +} + +func (f *fakePins) setNoRoom(noRoom bool) { + f.mu.Lock() + defer f.mu.Unlock() + f.noRoom = noRoom +} + type fixture struct { *Manager @@ -334,6 +351,40 @@ func TestSync_MirrorsAdditionsUnderTheServersOwnClaim(t *testing.T) { } } +// TestSync_OutOfRoomDefersTheRestOfTheCatalogue: a followed instance can +// list more than this one can hold. Running out of room stops the pass +// rather than being re-learnt once per remaining file, and — the part that +// matters — releases nothing: the origin still lists everything it holds, so +// a follower that merely has nowhere to put more must not start throwing +// away what it already mirrors. +func TestSync_OutOfRoomDefersTheRestOfTheCatalogue(t *testing.T) { + f := newFixture(t) + mirrored := f.origin.hold("mirrored while there was still room") + f.poll() + if !f.has(t, mirrored) { + t.Fatal("the first file wasn't mirrored") + } + + for i := range 20 { + f.origin.hold(fmt.Sprintf("more than there is room for (%d)", i)) + } + f.pins.setNoRoom(true) + + before := len(f.pins.claims()) + f.poll() + f.poll() + + if got := len(f.pins.claims()) - before; got != 2 { + t.Errorf("pin requests across two out-of-room polls = %d, want one per poll", got) + } + if got := f.pins.releases(); len(got) != 0 { + t.Errorf("releases = %+v, want none while the origin still lists everything", got) + } + if !f.has(t, mirrored) { + t.Error("a file the origin still lists was released while the follower was out of room") + } +} + // TestSync_ReleaseNeedsTwoConsecutiveAbsences is the hysteresis // listFiles's documented transient omissions demand: one missing poll // changes nothing, the second releases. diff --git a/internal/pin/pin.go b/internal/pin/pin.go index 8430ec0..3251781 100644 --- a/internal/pin/pin.go +++ b/internal/pin/pin.go @@ -79,6 +79,14 @@ const ( // is enough to stop a reboot with a backlog of pending pins from opening // a DHT walk and a download for every single one of them simultaneously. fetchConcurrency = 4 + // intentCap bounds how many intents can exist at once. One is a few + // hundred bytes and a timer, so the cost of any single one is trivial — + // but nothing else bounds how many a caller can create for content that + // will never resolve, and the horizon deliberately keeps a failed one + // afterwards. Thousands is far past what allowlist-scale pinning or a + // followed instance's catalogue asks for, and still small enough that + // the whole set is a couple of megabytes of an appliance's volume. + intentCap = 4096 ) // The schedule's numbers are constants — there's no config knob for any of @@ -90,6 +98,7 @@ var ( backoffCap = backoffLimit attemptLimit = attemptCap attemptTimeout = attemptAllowance + maxIntents = intentCap ) // attemptDeadline is how long one attempt gets: a fixed allowance for @@ -329,6 +338,10 @@ func (m *Manager) Close() error { // and comes back StateSeeking; a failed intent asked for again is // re-triggered from a clean slate, keeping the previous run's error in the // response so the caller learns why it had stopped. +// +// A request that would create a new intent — and only such a request — is +// refused with store.ErrNoRoom when this instance has nowhere to put what's +// being asked for; see admit. func (m *Manager) Request(ctx context.Context, ref Ref, claim store.Pin) (Status, error) { if err := m.validate(ref); err != nil { return Status{}, err @@ -361,6 +374,9 @@ func (m *Manager) Request(ctx context.Context, ref Ref, claim store.Pin) (Status it, existing := m.intents[ref.CID] if !existing { + if err := m.admit(ref); err != nil { + return Status{}, err + } it = &intent{ref: ref, createdAt: time.Now().UTC()} m.intents[ref.CID] = it } @@ -393,6 +409,28 @@ func (m *Manager) Request(ctx context.Context, ref Ref, claim store.Pin) (Status return m.status(it), nil } +// admit decides whether this manager can take on one more intent, and is +// the reason a pin can be refused for a reason that has nothing to do with +// the reference. Both ceilings answer store.ErrNoRoom, since both are the +// same statement — this instance has nowhere to put what you're asking for — +// and the endpoint turns that into one named lexicon error either way. +// +// Only a genuinely new intent is checked: joining or re-triggering an +// existing one adds a requester to a promise already made, and refusing that +// would strand a caller who can no longer even ask about a pin already +// under way. The caller must hold m.mu. +func (m *Manager) admit(ref Ref) error { + if len(m.intents) >= maxIntents { + slog.Warn("pin: refusing a new intent, the intent list is full", "cid", ref.CID, "intents", len(m.intents)) + return fmt.Errorf("%w: %d pin intents already registered", store.ErrNoRoom, len(m.intents)) + } + if err := m.store.CheckRoom(); err != nil { + slog.Warn("pin: refusing a new intent, the data volume is out of room", "cid", ref.CID, "size", ref.Size, "err", err) + return err + } + return nil +} + // pinStored runs the synchronous pin for content that's already in the // store. It reports ok=false when the attempt failed ambiently and the // caller should fall through to the intent flow; a definitive failure isn't @@ -611,6 +649,20 @@ func (m *Manager) attempt(c cid.Cid, it *intent) { return } + // The bytes this attempt is about to pull have to fit somewhere, and + // nothing on the volume reflects them until they've all landed — + // so the room is promised for the attempt's duration, and several + // concurrent fetches can't each see the same headroom and + // collectively overrun it. Refusal is ambient, not definitive: + // deleting something makes room, and waiting out the backoff is + // exactly the right response, so it settles like any failed attempt. + release, err := m.store.Reserve(ref.Size) + if err != nil { + m.settle(c, it, claims, ipfs.PinResult{}, err) + return + } + defer release() + res, err := m.fetcher.Pin(fctx, ref, m.maxSize, store.Meta{MimeType: ref.MimeType, Pins: claims}, ipfs.WithProgress(func(n int64) { it.fetched.Store(n) })) m.settle(c, it, claims, res, err) diff --git a/internal/pin/space_test.go b/internal/pin/space_test.go new file mode 100644 index 0000000..96bdfa5 --- /dev/null +++ b/internal/pin/space_test.go @@ -0,0 +1,114 @@ +package pin + +import ( + "context" + "errors" + "strings" + "testing" + "time" + + "atfs.dev/internal/store" +) + +// volume presents a fixed 20 GiB volume — whose reserve floor works out at a +// round 1 GiB — with free bytes still available, for the rest of the test. +func volume(t *testing.T, free int64) { + t.Helper() + original := store.VolumeSpace + store.VolumeSpace = func(string) (int64, int64, bool) { return free, 20 << 30, true } + t.Cleanup(func() { store.VolumeSpace = original }) +} + +// capIntents shrinks the intent ceiling so a test can reach it in two calls +// rather than four thousand — the same arrangement as schedule. +func capIntents(t *testing.T, n int) { + t.Helper() + original := maxIntents + maxIntents = n + t.Cleanup(func() { maxIntents = original }) +} + +// TestRequest_RefusesOnceTheIntentListIsFull: nothing else bounds how many +// intents a caller can create for content that will never resolve, and a +// failed one is deliberately kept — so the manager stops taking new ones, +// while still letting a second requester join, or re-ask about, one it +// already holds. +func TestRequest_RefusesOnceTheIntentListIsFull(t *testing.T) { + capIntents(t, 1) + o := newProviderOrigin(t) + st := newTestStore(t) + dir := t.TempDir() + m := newManager(t, dir, st, newOriginFetcher(st)) + + first := refFor(t, o, []byte("asked for first")) + if _, err := m.Request(context.Background(), first, accountClaim(didA)); err != nil { + t.Fatalf("first Request: %v", err) + } + + second := refFor(t, o, []byte("asked for once the list is full")) + _, err := m.Request(context.Background(), second, accountClaim(didA)) + if !errors.Is(err, store.ErrNoRoom) { + t.Fatalf("Request past the cap = %v, want store.ErrNoRoom", err) + } + if intentExists(dir, second.CID) { + t.Error("a refused request still wrote an intent file") + } + + if _, err := m.Request(context.Background(), first, accountClaim(didB)); err != nil { + t.Fatalf("joining an intent that already exists: %v", err) + } +} + +// TestRequest_RefusesWhenTheVolumeIsAtItsFloor: an intent is a promise to +// store bytes, so a volume with nothing left to promise takes on no more of +// them. +func TestRequest_RefusesWhenTheVolumeIsAtItsFloor(t *testing.T) { + o := newProviderOrigin(t) + st := newTestStore(t) + m := newManager(t, t.TempDir(), st, newOriginFetcher(st)) + volume(t, 0) + + ref := refFor(t, o, []byte("content there is nowhere to put")) + if _, err := m.Request(context.Background(), ref, accountClaim(didA)); !errors.Is(err, store.ErrNoRoom) { + t.Fatalf("Request = %v, want store.ErrNoRoom", err) + } +} + +// TestAttempt_NoRoomIsRetriedNotFailed: an attempt that can't be promised +// the room for its bytes has proved nothing about the content — deleting +// something makes room — so it backs off and tries again like any other +// ambient failure, rather than burning the intent. +func TestAttempt_NoRoomIsRetriedNotFailed(t *testing.T) { + schedule(t, 20*time.Millisecond, 20*time.Millisecond, 50) + o := newProviderOrigin(t) + st := newTestStore(t) + m := newManager(t, t.TempDir(), st, newOriginFetcher(st)) + + content := []byte("content that would fit if anything did") + o.stock(blessedCID(t, content), content) + ref := refFor(t, o, content) + + // Above the floor, so the intent itself is admitted, but not by enough + // to promise the file. + volume(t, 1<<30+1) + + if _, err := m.Request(context.Background(), ref, accountClaim(didA)); err != nil { + t.Fatalf("Request: %v", err) + } + + waitFor(t, "the attempt to be refused and rescheduled", func() bool { + s, ok := m.Status(ref.CID) + return ok && s.Attempts > 0 + }) + + s, _ := m.Status(ref.CID) + if s.State != StateSeeking { + t.Errorf("state = %q, want %q — running out of room proves nothing about the content", s.State, StateSeeking) + } + if !strings.Contains(s.LastError, "no room") { + t.Errorf("lastError = %q, want it to name the storage ceiling", s.LastError) + } + if got := o.requests(ref.CID); got != 0 { + t.Errorf("provider requests = %d, want none: no room means no fetch at all", got) + } +} diff --git a/internal/store/space.go b/internal/store/space.go new file mode 100644 index 0000000..d3255ae --- /dev/null +++ b/internal/store/space.go @@ -0,0 +1,102 @@ +package store + +import ( + "errors" + "fmt" + "sync" +) + +// ErrNoRoom refuses a write the data volume can't afford: either the volume +// has already fallen to the free space it must keep, or what's left above +// that is already promised to writes in flight (see Store.Reserve). +// +// It's the single signal every storage ceiling in atfs reports — a direct +// upload, a pin intent, a followed instance's mirror — so each caller maps +// one error to one named lexicon error rather than guessing at a 500. +var ErrNoRoom = errors.New("store: the data volume has no room to spare") + +// The volume keeps a slice of itself permanently free, so that filling it is +// never what stops an instance working: a full volume takes the metadata +// sidecars, pin intents and index artifacts down with the content, and on a +// gosd appliance there's no shell to clear it from (see CLAUDE.md). +// +// The floor is derived, never configured — bare-minimum configuration is the +// product — as a twentieth of the volume, never less than minFreeReserve on +// a small card, and never more than maxFreeReserve on a large disk, where a +// flat fraction would refuse far more capacity than an instance could ever +// need to dig itself out with. +const ( + minFreeReserve = 512 << 20 // 512 MiB + maxFreeReserve = 4 << 30 // 4 GiB + freeReserveDivisor = 20 // 5% +) + +// VolumeSpace reports the volume holding dir: the bytes an unprivileged +// writer can still use, the volume's total size, and whether this platform +// could answer at all. A false ok disables the free-space ceiling entirely — +// a platform atfs can't measure is not one it refuses to run on — leaving +// the count-based ceilings (see pin.Manager) as the only bound. +// +// It's a var so a test can present a volume that's full without needing one; +// it's exported because the packages that have to prove they refuse +// gracefully (internal/xrpc, internal/pin) aren't this one. Production never +// replaces it. +var VolumeSpace = volumeSpace + +func reserveFloor(total int64) int64 { + return min(max(total/freeReserveDivisor, minFreeReserve), maxFreeReserve) +} + +// CheckRoom refuses a write outright once the volume has fallen to its +// reserve floor. +// +// It deliberately ignores what Reserve has promised: a caller reaching here +// is either spending a reservation it already made — which would double-count +// — or making a write too small to be worth reserving, like a pin intent's +// own file. +func (s *Store) CheckRoom() error { + avail, total, ok := VolumeSpace(s.dir) + if !ok { + return nil + } + if floor := reserveFloor(total); avail < floor { + return fmt.Errorf("%w: %d bytes free, %d must stay free", ErrNoRoom, avail, floor) + } + return nil +} + +// Reserve promises n bytes of the volume to a write about to start, and +// returns the function that hands the promise back. The release is +// idempotent, so a caller can defer it and still release early. +// +// Reserving rather than merely checking the free space is what makes the +// floor hold under concurrency: several writers each checking a bare +// free-space figure all see the same room and then collectively write +// several times it. A promise is also how a caller declares a size the +// bytes haven't reached the disk under yet — an upload streaming in, a pin +// fetch that has found a source — since nothing on the volume reflects +// either until it's finished. +// +// n is a ceiling, not a measurement: a caller that ends up writing less has +// simply been pessimistic for the duration. +func (s *Store) Reserve(n int64) (release func(), err error) { + avail, total, measured := VolumeSpace(s.dir) + + s.spaceMu.Lock() + defer s.spaceMu.Unlock() + + if floor := reserveFloor(total); measured && avail-s.promised-n < floor { + return nil, fmt.Errorf("%w: %d bytes free, %d promised to writes in flight, %d wanted, %d must stay free", + ErrNoRoom, avail, s.promised, n, floor) + } + s.promised += n + + var once sync.Once + return func() { + once.Do(func() { + s.spaceMu.Lock() + defer s.spaceMu.Unlock() + s.promised -= n + }) + }, nil +} diff --git a/internal/store/space_other.go b/internal/store/space_other.go new file mode 100644 index 0000000..ee0b996 --- /dev/null +++ b/internal/store/space_other.go @@ -0,0 +1,8 @@ +//go:build !linux && !darwin + +package store + +// volumeSpace has no portable answer away from Linux and macOS, so it +// reports that it couldn't measure rather than inventing a figure — see +// VolumeSpace for what an unmeasurable volume means. +func volumeSpace(string) (avail, total int64, ok bool) { return 0, 0, false } diff --git a/internal/store/space_test.go b/internal/store/space_test.go new file mode 100644 index 0000000..a2bb4e8 --- /dev/null +++ b/internal/store/space_test.go @@ -0,0 +1,91 @@ +package store + +import ( + "errors" + "strings" + "testing" +) + +const gib = 1 << 30 + +// volume presents a fixed volume to every space check for the rest of the +// test. Its total is 20 GiB, which puts the reserve floor at a round 1 GiB +// (a twentieth), so free reads directly as "this much, of which 1 GiB must +// stay free". +func volume(t *testing.T, free int64) { + t.Helper() + original := VolumeSpace + VolumeSpace = func(string) (int64, int64, bool) { return free, 20 * gib, true } + t.Cleanup(func() { VolumeSpace = original }) +} + +// TestPut_RefusesOnceTheVolumeIsAtItsFloor covers the backstop every path +// that stores bytes goes through, whether the caller reserved first or not. +func TestPut_RefusesOnceTheVolumeIsAtItsFloor(t *testing.T) { + s := open(t) + volume(t, gib-1) + + _, err := s.Put(strings.NewReader("nowhere to put this"), pinMeta("text/plain", "did:plc:abc123")) + if !errors.Is(err, ErrNoRoom) { + t.Fatalf("Put = %v, want ErrNoRoom", err) + } + + for c, err := range s.List() { + if err != nil { + t.Fatalf("List: %v", err) + } + t.Errorf("a refused Put left %s in the store, want nothing", c) + } +} + +// TestReserve_PromisesAreCountedAgainstEachOther is what makes the floor +// hold under concurrency: two writes that would each fit on their own don't +// both get in, and the room comes back when one of them finishes. +func TestReserve_PromisesAreCountedAgainstEachOther(t *testing.T) { + s := open(t) + volume(t, gib+100) + + release, err := s.Reserve(60) + if err != nil { + t.Fatalf("first Reserve: %v", err) + } + if _, err := s.Reserve(60); !errors.Is(err, ErrNoRoom) { + t.Fatalf("second Reserve = %v, want ErrNoRoom", err) + } + + release() + if _, err := s.Reserve(60); err != nil { + t.Fatalf("Reserve once the first was released: %v", err) + } +} + +// TestReserve_AdmitsWhatCannotBeMeasured: a platform atfs can't ask about +// free space is never one it refuses to write on. +func TestReserve_AdmitsWhatCannotBeMeasured(t *testing.T) { + s := open(t) + original := VolumeSpace + VolumeSpace = func(string) (int64, int64, bool) { return 0, 0, false } + t.Cleanup(func() { VolumeSpace = original }) + + if _, err := s.Reserve(1 << 40); err != nil { + t.Fatalf("Reserve: %v", err) + } + if err := s.CheckRoom(); err != nil { + t.Fatalf("CheckRoom: %v", err) + } +} + +// TestVolumeSpace_MeasuresTheRealVolume exercises the platform's own +// statfs, whose block-size and block-count fields differ in type between +// Linux and macOS — a conversion mistake in either would show up here as a +// negative or impossible figure long before it showed up as a wrongly +// refused upload. +func TestVolumeSpace_MeasuresTheRealVolume(t *testing.T) { + avail, total, ok := VolumeSpace(t.TempDir()) + if !ok { + t.Skip("this platform can't measure free space") + } + if total <= 0 || avail < 0 || avail > total { + t.Fatalf("VolumeSpace = (avail %d, total %d), want a positive total and 0 <= avail <= total", avail, total) + } +} diff --git a/internal/store/space_unix.go b/internal/store/space_unix.go new file mode 100644 index 0000000..e58b8cd --- /dev/null +++ b/internal/store/space_unix.go @@ -0,0 +1,22 @@ +//go:build linux || darwin + +package store + +import "golang.org/x/sys/unix" + +// volumeSpace answers from statfs(2), which x/sys/unix reaches by raw +// syscall — no cgo, which the gosd arm64 cross-compile depends on (see +// CLAUDE.md). Linux and macOS are the two platforms atfs is built for: the +// appliance and container images, and the machine `make check` runs on. +// +// Bavail rather than Bfree: the blocks an unprivileged writer can actually +// use, which on ext4 excludes the root-reserved portion atfs will never be +// able to touch. +func volumeSpace(dir string) (avail, total int64, ok bool) { + var fs unix.Statfs_t + if err := unix.Statfs(dir, &fs); err != nil { + return 0, 0, false + } + blockSize := int64(fs.Bsize) + return int64(fs.Bavail) * blockSize, int64(fs.Blocks) * blockSize, true +} diff --git a/internal/store/store.go b/internal/store/store.go index 712c316..0566f69 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -269,6 +269,11 @@ type Store struct { // at appliance scale, and it's what stops a concurrent Put/Put or // Put/Unpin pair from losing a claim to a lost update. metaMu sync.Mutex + + // spaceMu guards promised: the volume bytes Reserve has handed out to + // writes that haven't finished yet (see space.go). + spaceMu sync.Mutex + promised int64 } // Open opens (creating if missing) a blob store rooted at dir. @@ -328,10 +333,19 @@ func normalizeMimeType(raw string) string { // empty-pins state that means "delete this at boot" (see ErrEmptyPins). // That's checked first, before anything touches the reader or disk, so a // rejected Put stores nothing at all. +// +// A volume already at its reserve floor fails with ErrNoRoom, before the +// temp file exists. That's a backstop, not the ceiling itself: it can only +// see space already spent, so it's Reserve — which callers take out before +// they start streaming — that actually keeps concurrent writes off the +// floor. func (s *Store) Put(r io.Reader, meta Meta) (PutResult, error) { if len(meta.Pins) == 0 { return PutResult{}, ErrEmptyPins } + if err := s.CheckRoom(); err != nil { + return PutResult{}, err + } meta.MimeType = normalizeMimeType(meta.MimeType) tmp, err := os.CreateTemp(s.tmpDir, "put-*") diff --git a/internal/xrpc/pinfile.go b/internal/xrpc/pinfile.go index 4a76eae..2e38c18 100644 --- a/internal/xrpc/pinfile.go +++ b/internal/xrpc/pinfile.go @@ -134,11 +134,14 @@ func (h *handler) pinFile(w http.ResponseWriter, r *http.Request) { status, err := h.pinner.Request(r.Context(), ref, store.Pin{DID: did.String(), Source: store.SourcePin}) if err != nil { - if errors.Is(err, pin.ErrUnpinnable) { + switch { + case errors.Is(err, pin.ErrUnpinnable): writeXRPCError(w, http.StatusBadRequest, "InvalidRequest", err.Error()) - return + case errors.Is(err, store.ErrNoRoom): + refuseNoRoom(w, pinFileLXM, did.String(), err) + default: + writeXRPCError(w, http.StatusInternalServerError, "InternalServerError", "failed to register the pin") } - writeXRPCError(w, http.StatusInternalServerError, "InternalServerError", "failed to register the pin") return } diff --git a/internal/xrpc/pinfile_test.go b/internal/xrpc/pinfile_test.go index 85f65b6..19e4d08 100644 --- a/internal/xrpc/pinfile_test.go +++ b/internal/xrpc/pinfile_test.go @@ -369,6 +369,21 @@ func TestPinFile_UnpinnableReferenceIsARequestProblem(t *testing.T) { assertXRPCError(t, rec, http.StatusBadRequest, "InvalidRequest") } +// TestPinFile_OutOfRoomIsANamedError: a pin is an upload by another name, so +// an instance with nowhere to put the content refuses it exactly as +// uploadFile does — one named lexicon error, not a 500 the caller would read +// as "try again, it might work". +func TestPinFile_OutOfRoomIsANamedError(t *testing.T) { + p := &fakePinner{err: fmt.Errorf("%w: 4096 pin intents already registered", store.ErrNoRoom)} + h, _ := newPinHandler(t, fakeVerifier{did: testUploaderDID}, p) + + blob := blessedCID(t, []byte("content there is no room for")) + rec := httptest.NewRecorder() + h.ServeHTTP(rec, pinFileHTTPRequest(t, blob, blob, 10, "")) + + assertXRPCError(t, rec, http.StatusBadRequest, "InsufficientStorage") +} + func TestPinFile_ManagerFailureIsAServerError(t *testing.T) { p := &fakePinner{err: errors.New("the intent directory is on fire")} h, _ := newPinHandler(t, fakeVerifier{did: testUploaderDID}, p) diff --git a/internal/xrpc/uploadfile.go b/internal/xrpc/uploadfile.go index ba3441b..f0df2ae 100644 --- a/internal/xrpc/uploadfile.go +++ b/internal/xrpc/uploadfile.go @@ -176,6 +176,13 @@ func (h *handler) uploadFile(lxm string) http.HandlerFunc { mimeType = "application/octet-stream" } + release, err := h.store.Reserve(reservationFor(r)) + if err != nil { + refuseNoRoom(w, lxm, did.String(), err) + return + } + defer release() + r.Body = http.MaxBytesReader(w, r.Body, maxBlobSize) result, err := h.store.Put(r.Body, store.Meta{ MimeType: mimeType, @@ -183,11 +190,14 @@ func (h *handler) uploadFile(lxm string) http.HandlerFunc { }) if err != nil { var tooLarge *http.MaxBytesError - if errors.As(err, &tooLarge) { + switch { + case errors.As(err, &tooLarge): writeXRPCError(w, http.StatusBadRequest, "BlobTooLarge", fmt.Sprintf("blob exceeds maximum size of %d bytes", maxBlobSize)) - return + case errors.Is(err, store.ErrNoRoom): + refuseNoRoom(w, lxm, did.String(), err) + default: + writeXRPCError(w, http.StatusInternalServerError, "InternalServerError", "failed to store blob") } - writeXRPCError(w, http.StatusInternalServerError, "InternalServerError", "failed to store blob") return } @@ -205,6 +215,30 @@ func (h *handler) uploadFile(lxm string) http.HandlerFunc { } } +// reservationFor is how much of the volume an upload has to be promised +// before a byte of it is read. A declared Content-Length is what the Go +// server then holds the body to, so an ordinary client reserves exactly what +// it will use and a nearly-full instance can still take small uploads; a +// chunked body declares nothing and has to be assumed the largest blob this +// instance would accept. +func reservationFor(r *http.Request) int64 { + if r.ContentLength >= 0 && r.ContentLength < maxBlobSize { + return r.ContentLength + } + return maxBlobSize +} + +// refuseNoRoom answers a storage ceiling with InsufficientStorage — a named +// lexicon error, so 400 like every other one atfs reports (see getFile's +// BlobNotFound note) rather than the 507 the HTTP status vocabulary would +// suggest. It always logs: an appliance that silently refuses is as bad as +// one that silently fills up, and this is the only place the reason for the +// refusal (which figure ran out, and by how much) survives. +func refuseNoRoom(w http.ResponseWriter, lxm, did string, err error) { + slog.Warn("xrpc: refused, the data volume is out of room", "lxm", lxm, "did", did, "err", err) + writeXRPCError(w, http.StatusBadRequest, "InsufficientStorage", "this instance has no room to store more content") +} + // resolveIPFSRoot determines the fetchable IPFS root to report for a // just-stored blob: the UnixFS root for a large blob that got indexed, or // the blob's own CID otherwise (no indexer at all, or a blob at or under diff --git a/internal/xrpc/uploadfile_test.go b/internal/xrpc/uploadfile_test.go index 902cf87..2123801 100644 --- a/internal/xrpc/uploadfile_test.go +++ b/internal/xrpc/uploadfile_test.go @@ -332,6 +332,37 @@ func TestUploadFile_OversizedBody(t *testing.T) { assertXRPCError(t, rec, http.StatusBadRequest, "BlobTooLarge") } +// fullVolume makes every storage check in the blob store report a data +// volume with nothing left above its reserve, for the rest of the test. +func fullVolume(t *testing.T) { + t.Helper() + original := store.VolumeSpace + store.VolumeSpace = func(string) (int64, int64, bool) { return 0, 20 << 30, true } + t.Cleanup(func() { store.VolumeSpace = original }) +} + +// TestUploadFile_RefusedWhenTheVolumeIsFull: the ceiling that stops an +// appliance's data volume being filled answers as a named lexicon error — +// 400 like every other one atfs reports, never a 500 — and stores nothing. +func TestUploadFile_RefusedWhenTheVolumeIsFull(t *testing.T) { + fullVolume(t) + h, s := newTestHandler(t, fakeVerifier{did: testUploaderDID}) + + req := httptest.NewRequest(http.MethodPost, "/xrpc/"+uploadFileLXM, strings.NewReader("more than this instance can hold")) + req.Header.Set("Authorization", "Bearer token") + rec := httptest.NewRecorder() + + h.ServeHTTP(rec, req) + + assertXRPCError(t, rec, http.StatusBadRequest, "InsufficientStorage") + for c, err := range s.List() { + if err != nil { + t.Fatalf("store.List: %v", err) + } + t.Errorf("a refused upload stored %s", c) + } +} + // TestUploadFile_RealVerifier wires a real auth.Verifier (backed by a fake // identity directory, not a fake Verify method) into the handler, proving // the two packages' interfaces actually line up end-to-end and not just diff --git a/lexicons/dev/atfs/repo/pinFile.json b/lexicons/dev/atfs/repo/pinFile.json index 2e28f5e..1899a77 100644 --- a/lexicons/dev/atfs/repo/pinFile.json +++ b/lexicons/dev/atfs/repo/pinFile.json @@ -63,6 +63,10 @@ { "name": "BlobTooLarge", "description": "The reference's declared size exceeds this instance's maximum blob size — checked before any fetch is attempted, the same limit dev.atfs.repo.uploadFile enforces on bytes it receives directly." + }, + { + "name": "InsufficientStorage", + "description": "This instance has no room to take on another pin: its data volume has reached the free space it keeps in reserve, or it is already holding as many pin requests at once as it will. No intent was registered, and pins already registered carry on. The same named error dev.atfs.repo.uploadFile answers with, for the same reason — a pin is an upload by another name — and releasing content with dev.atfs.repo.deleteFile is what makes room again." } ] }, diff --git a/lexicons/dev/atfs/repo/uploadFile.json b/lexicons/dev/atfs/repo/uploadFile.json index ad33502..b7a84ae 100644 --- a/lexicons/dev/atfs/repo/uploadFile.json +++ b/lexicons/dev/atfs/repo/uploadFile.json @@ -29,6 +29,10 @@ { "name": "BlobTooLarge", "description": "The request body exceeds this instance's maximum blob size." + }, + { + "name": "InsufficientStorage", + "description": "This instance has no room to store more content: its data volume has reached the free space it keeps in reserve, or what is left of that is already promised to uploads and pins in flight. Nothing was stored, and the instance keeps serving and deleting normally — releasing content with dev.atfs.repo.deleteFile is what makes room again. Retrying is reasonable, but only after something has been deleted or an upload in flight has finished." } ] }