diff --git a/.beans/ATFS-o8ov--guard-the-empty-pin-set-sweep-invariant.md b/.beans/ATFS-o8ov--guard-the-empty-pin-set-sweep-invariant.md index eb1bf2f..40d8c3c 100644 --- a/.beans/ATFS-o8ov--guard-the-empty-pin-set-sweep-invariant.md +++ b/.beans/ATFS-o8ov--guard-the-empty-pin-set-sweep-invariant.md @@ -1,7 +1,7 @@ --- # ATFS-o8ov title: Guard the empty-pin-set sweep invariant -status: todo +status: completed type: bug priority: high created_at: 2026-08-07T14:02:40Z @@ -9,3 +9,52 @@ updated_at: 2026-08-07T14:02:40Z --- SweepDeletions treats an empty pin list as "a delete crashed between Unpin and removal" and reaps the blob at boot — but Store.Put and Node.Pin accept caller-supplied Meta with zero pins, so a future pinFile stage-two caller passing empty pins would create content the next boot silently deletes. Enforce the invariant at the write boundary (reject empty-pins Meta on Put/Pin) or teach the sweep to distinguish in-flight intents. Must land before pinFile stage two wires a route to Node.Pin. + +## Summary of Changes + +pinFile stage two had already landed (internal/pin's Manager + the +pinFile XRPC handler both drive `ipfs.Node.Pin`) by the time this was +picked up, so the route the bean warned about was live — but every actual +caller (internal/pin's `Request`/`attempt`, `xrpc/uploadfile.go`) already +builds its `store.Meta` from a real requester DID, so nothing currently +in the tree could trigger the bug. Enforced the invariant at both write +boundaries anyway, defensively, per the bean's instruction: + +- `store.Store.Put` now rejects any `Meta` with an empty `Pins` list with + a new sentinel, `ErrEmptyPins`, checked before anything touches the + reader or disk — a rejected `Put` stores nothing (no temp file, no + content file, no sidecar). This applies uniformly, including to a + re-Put of content that's already stored (the dedup/merge path), not + just a brand-new blob. +- `ipfs.Node.Pin` rejects an empty `meta.Pins` up front too, alongside its + existing `ref.validate` check, as a definitive (non-retryable) failure. + This isn't redundant with `Store.Put`'s guard: `Node.Pin`'s + already-stored branch calls `store.Store.Pin` (the merge-only method), + which quietly no-ops on an empty incoming pin list rather than + erroring — so without this check, pinning already-stored content with + no claim would have reported a false success while silently recording + nothing. +- `store.Store.Pin` itself is left as a no-op on empty input (documented + why): it only ever merges into an *existing* record, so an empty + incoming list can't manufacture a new zero-pins record — the one + production caller (`Node.Pin`) now guarantees non-empty input anyway. +- `Store.Unpin`'s empty-list persistence and `SweepDeletions`' reaping are + untouched, as instructed — the invariant they encode is exactly what + the write-boundary guards now protect. + +Behavioral tests added: `TestPut_RejectsEmptyPins` (rejected, stores +nothing), `TestPut_RejectsEmptyPinsEvenForAlreadyStoredContent` (rejected +on the dedup path too, existing claims untouched), +`TestPin_RejectsEmptyPinsAndNeverDials` (definitive, no network touched), +and `TestPin_RejectsEmptyPinsForAlreadyStoredContent` (definitive, no +false success on already-stored content). + +Fixing this uncovered that a fair number of existing tests across +internal/store and internal/ipfs built fixtures via `Put(..., Meta{})` — +an empty pin list used purely as "I don't care about claims for this +test". All were updated to pass a real claim (mostly via the packages' +existing `pinMeta` test helpers), which incidentally makes those fixtures +model the real invariant more faithfully: the only way to reach a +persisted empty-pins state is still Put-then-Unpin, never a bare Put. +`make build`, `go vet ./...`, `go test ./...`, and `gofmt -l .` are all +clean. diff --git a/internal/ipfs/blockstore_test.go b/internal/ipfs/blockstore_test.go index e04147c..31a870f 100644 --- a/internal/ipfs/blockstore_test.go +++ b/internal/ipfs/blockstore_test.go @@ -28,7 +28,7 @@ func newTestBlockstore(t *testing.T) (*Blockstore, *store.Store, *Indexer) { func TestBlockstore_ServesSmallRawBlock(t *testing.T) { bs, s, _ := newTestBlockstore(t) content := []byte("a small blob, well under the raw block limit") - res, err := s.Put(bytes.NewReader(content), store.Meta{}) + res, err := s.Put(bytes.NewReader(content), pinMeta(testPinDID)) if err != nil { t.Fatalf("store.Put: %v", err) } diff --git a/internal/ipfs/fetch.go b/internal/ipfs/fetch.go index 9fe8d0e..1bdbea0 100644 --- a/internal/ipfs/fetch.go +++ b/internal/ipfs/fetch.go @@ -179,10 +179,21 @@ type PinResult struct { // Failures are classified for a retrying caller — see ErrDefinitive — and // a failure before the bytes are stored leaves nothing behind: no temp // file, no partial blob, no index entry. +// +// meta.Pins must carry at least one claim, checked up front like ref's own +// fields: a pin call recording no claim isn't a smaller version of a real +// pin, it's a caller bug, and letting it through would be worse than +// refusing it. The not-yet-stored path would eventually hit store.Put's own +// ErrEmptyPins, but the already-stored path merges through store.Pin, which +// quietly no-ops on an empty claim list (see its doc) — without this check +// that would report a false success, having recorded nothing. func (n *Node) Pin(ctx context.Context, ref FileRef, maxSize int64, meta store.Meta, opts ...PinOption) (PinResult, error) { if err := ref.validate(maxSize); err != nil { return PinResult{}, err } + if len(meta.Pins) == 0 { + return PinResult{}, definitivef("pin requires at least one claim (empty Meta.Pins)") + } var o PinOptions for _, opt := range opts { opt(&o) diff --git a/internal/ipfs/fetch_test.go b/internal/ipfs/fetch_test.go index e43cb61..2c0c0e8 100644 --- a/internal/ipfs/fetch_test.go +++ b/internal/ipfs/fetch_test.go @@ -8,6 +8,7 @@ import ( "net/http/httptest" "os" "path/filepath" + "reflect" "sync/atomic" "testing" "time" @@ -301,6 +302,67 @@ func TestPin_ReferenceOverMaxSizeIsDefinitiveAndNeverDials(t *testing.T) { assertNoStoredBlobs(t, dir) } +// TestPin_RejectsEmptyPinsAndNeverDials checks that Pin refuses a Meta with +// no claims before touching the network at all — the same definitive, +// never-retry treatment as a malformed reference (see TestPin_ReferenceOverMaxSizeIsDefinitiveAndNeverDials). +func TestPin_RejectsEmptyPinsAndNeverDials(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + content := deterministicBytes(testSmallSize) + st, dir := newTestStore(t) + n := newTestNode(t, ctx, st) + + ref := FileRef{CID: blessedCID(t, content), IPFSRoot: blessedCID(t, content), Size: int64(len(content)), Providers: []string{refusingOrigin(t)}} + _, err := n.Pin(ctx, ref, testMaxSize, store.Meta{}) + if err == nil { + t.Fatal("Pin succeeded, want a failure for an empty Meta.Pins") + } + if !IsDefinitive(err) { + t.Errorf("Pin error %v is ambient, want definitive", err) + } + assertNoStoredBlobs(t, dir) +} + +// TestPin_RejectsEmptyPinsForAlreadyStoredContent checks that Pin refuses +// an empty Meta.Pins even for content it already serves — without this, the +// call would report success (Existed=true) while silently recording no +// claim at all, since store.Store.Pin no-ops on an empty claim list. +func TestPin_RejectsEmptyPinsForAlreadyStoredContent(t *testing.T) { + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + content := deterministicBytes(testBlobSize) + st, _ := newTestStore(t) + n := newTestNode(t, ctx, st) + + res, err := st.Put(bytes.NewReader(content), store.Meta{MimeType: "application/octet-stream", Pins: []store.Pin{{DID: "did:plc:uploader", Source: "upload"}}}) + if err != nil { + t.Fatalf("store.Put: %v", err) + } + if err := n.index.EnsureIndexed(ctx, res.CID); err != nil { + t.Fatalf("EnsureIndexed: %v", err) + } + root, _, err := n.index.RootFor(res.CID) + if err != nil { + t.Fatalf("RootFor: %v", err) + } + + ref := FileRef{CID: res.CID, IPFSRoot: root, Size: res.Size, Providers: []string{refusingOrigin(t)}} + if _, err := n.Pin(ctx, ref, testMaxSize, store.Meta{}); !IsDefinitive(err) { + t.Errorf("Pin(empty pins) error = %v, want a definitive failure", err) + } + + stat, err := st.Stat(res.CID) + if err != nil { + t.Fatalf("Stat: %v", err) + } + want := []store.Pin{{DID: "did:plc:uploader", Source: "upload"}} + if !reflect.DeepEqual(stat.Pins, want) { + t.Errorf("Pins after rejected Pin = %+v, want %+v (untouched)", stat.Pins, want) + } +} + func TestPin_OverlongProviderStreamIsCutOffEarly(t *testing.T) { ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() diff --git a/internal/ipfs/unixfs_test.go b/internal/ipfs/unixfs_test.go index 803b65a..4287e1f 100644 --- a/internal/ipfs/unixfs_test.go +++ b/internal/ipfs/unixfs_test.go @@ -258,7 +258,7 @@ func TestEnsureIndexed_SmallBlobsGetNoManifest(t *testing.T) { func TestWalkAndIndex_SelfHeals(t *testing.T) { content := deterministicBytes(testBlobSize) s, _ := newTestStore(t) - res, err := s.Put(bytes.NewReader(content), store.Meta{}) + res, err := s.Put(bytes.NewReader(content), pinMeta(testPinDID)) if err != nil { t.Fatalf("store.Put: %v", err) } @@ -560,7 +560,7 @@ func TestUnindex_SharedChunkSurvivesWhenOneBlobDeleted(t *testing.T) { func TestDeleteBlob_SmallNeverIndexedBlobIsNoOp(t *testing.T) { s, _ := newTestStore(t) ix := newTestIndexer(t, s) - res, err := s.Put(bytes.NewReader([]byte("small, never chunked")), store.Meta{}) + res, err := s.Put(bytes.NewReader([]byte("small, never chunked")), pinMeta(testPinDID)) if err != nil { t.Fatalf("store.Put: %v", err) } diff --git a/internal/store/store.go b/internal/store/store.go index abf4d81..efa0034 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -39,6 +39,16 @@ var ErrNotFound = errors.New("store: blob not found") // store keys. var ErrInvalidCID = errors.New("store: not a blessed CID (cidv1, raw, sha2-256)") +// ErrEmptyPins is returned by Put when meta carries no claims. A pin list +// only ever becomes empty later, when Unpin releases a last remaining claim +// (see Meta) — that's the one legitimate route to the empty, GC-eligible +// state SweepDeletions reaps at boot (see ipfs.Indexer.SweepDeletions). A +// Put that itself supplies zero claims would manufacture that same +// on-disk state for content nobody has claimed yet, which the next boot's +// sweep can't tell apart from a crashed delete — so it's refused outright, +// before anything is written. +var ErrEmptyPins = errors.New("store: put requires at least one pin") + // ErrNotPinned is returned by Unpin when a blob is stored but the given DID // holds no claim against it — distinct from ErrNotFound, so a caller can // tell "nothing to unpin" apart from "there was never anything here at @@ -186,7 +196,17 @@ func Open(dir string) (*Store, error) { // merged into the existing claim list: a DID already holding a claim is // left alone (idempotent re-upload; its original Source wins), and any DID // not yet present is appended. Either way PutResult.Existed is true. +// +// meta.Pins must carry at least one claim — a Put is always some caller +// claiming content, and nothing else is what turns a stored blob into the +// 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. func (s *Store) Put(r io.Reader, meta Meta) (PutResult, error) { + if len(meta.Pins) == 0 { + return PutResult{}, ErrEmptyPins + } + tmp, err := os.CreateTemp(s.tmpDir, "put-*") if err != nil { return PutResult{}, fmt.Errorf("store: put: %w", err) @@ -284,7 +304,16 @@ func mergePins(current, incoming []Pin) (merged []Pin, changed bool) { // to a temp file before it can compute the CID that tells it the content // was already there — a full-size write to flash to append one DID. // -// It fails with ErrNotFound if c isn't stored. +// It fails with ErrNotFound if c isn't stored. Unlike Put, an empty pins +// isn't rejected outright — mergePins reports no change and Pin quietly +// no-ops — because merging nothing in is never what manufactures the +// empty-pins state Put guards against (see ErrEmptyPins): a blob already +// holding claims is untouched, and one already at zero stays exactly the +// GC-eligible state it already was. Callers meaning to record a claim must +// still pass one; ipfs.Node.Pin is the one production caller, and it +// rejects an empty Meta.Pins itself before reaching here, for exactly the +// same reason Put does — so this no-op path only ever fires when there's +// truly nothing new to add. func (s *Store) Pin(c cid.Cid, pins []Pin) error { s.metaMu.Lock() defer s.metaMu.Unlock() diff --git a/internal/store/store_test.go b/internal/store/store_test.go index 1fb5e61..d2683d7 100644 --- a/internal/store/store_test.go +++ b/internal/store/store_test.go @@ -70,6 +70,64 @@ func TestPutGet_RoundTrip(t *testing.T) { } } +// TestPut_RejectsEmptyPins checks that Put refuses a Meta carrying no +// claims — see ErrEmptyPins — and stores nothing: no content file, no +// sidecar, no leftover temp file. +func TestPut_RejectsEmptyPins(t *testing.T) { + s := open(t) + + _, err := s.Put(strings.NewReader("nobody claims this"), Meta{MimeType: "text/plain"}) + if !errors.Is(err, ErrEmptyPins) { + t.Fatalf("Put(empty pins) error = %v, want ErrEmptyPins", err) + } + + entries, err := os.ReadDir(dirOf(s)) + if err != nil { + t.Fatalf("ReadDir: %v", err) + } + for _, e := range entries { + if e.Name() == ".tmp" { + continue + } + t.Errorf("unexpected entry after rejected Put: %s", e.Name()) + } + tmpEntries, err := os.ReadDir(filepath.Join(dirOf(s), ".tmp")) + if err != nil { + t.Fatalf("ReadDir .tmp: %v", err) + } + if len(tmpEntries) != 0 { + t.Errorf(".tmp has %d leftover entries, want 0", len(tmpEntries)) + } +} + +// TestPut_RejectsEmptyPinsEvenForAlreadyStoredContent checks that the +// rejection also covers a re-Put of content that's already in the store — +// not just the fresh-write path — since a dedup Put is still a claim on +// nothing if Pins is empty. +func TestPut_RejectsEmptyPinsEvenForAlreadyStoredContent(t *testing.T) { + s := open(t) + content := []byte("already claimed by someone") + + first, err := s.Put(bytes.NewReader(content), pinMeta("text/plain", "did:plc:original")) + if err != nil { + t.Fatalf("first Put: %v", err) + } + + if _, err := s.Put(bytes.NewReader(content), Meta{MimeType: "text/plain"}); !errors.Is(err, ErrEmptyPins) { + t.Fatalf("re-Put(empty pins) error = %v, want ErrEmptyPins", err) + } + + blob, err := s.Get(first.CID) + if err != nil { + t.Fatalf("Get: %v", err) + } + defer blob.Close() + want := []Pin{{DID: "did:plc:original", Source: "upload"}} + if !reflect.DeepEqual(blob.Meta.Pins, want) { + t.Errorf("Pins after rejected re-Put = %+v, want %+v (untouched)", blob.Meta.Pins, want) + } +} + // TestPut_BlessedCID checks the CID for known bytes against an // independently computed blessed CID: CIDv1 + raw codec (0x55) + sha2-256 // multihash (0x12, 0x20) + base32 multibase, built here by hand from @@ -78,7 +136,7 @@ func TestPut_BlessedCID(t *testing.T) { const want = "bafkreibrl5n5w5wqpdcdxcwaazheualemevr7ttxzbutiw74stdvrfhn2m" s := open(t) - res, err := s.Put(strings.NewReader("Hello, world!"), Meta{}) + res, err := s.Put(strings.NewReader("Hello, world!"), pinMeta("", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -178,7 +236,10 @@ func TestPut_ConcurrentDistinctClaimantsAllRecorded(t *testing.T) { } wg.Wait() - res, err := s.Put(bytes.NewReader(content), Meta{}) + // A re-Put by a DID that already holds a claim, purely to fetch the + // CID back for the read below — a no-op merge, so the pin count stays + // n. + res, err := s.Put(bytes.NewReader(content), pinMeta("text/plain", "did:plc:racer00")) if err != nil { t.Fatalf("final Put: %v", err) } @@ -230,7 +291,7 @@ func TestStat_MatchesGet(t *testing.T) { func TestStat_NotFound(t *testing.T) { s := open(t) - res, err := s.Put(strings.NewReader("never actually kept"), Meta{}) + res, err := s.Put(strings.NewReader("never actually kept"), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -251,7 +312,7 @@ func TestGet_NotFound(t *testing.T) { // A well-formed, never-stored CID: put content, then remove its files, // so Get sees exactly "absent", not "malformed". - res, err := s.Put(strings.NewReader("never actually kept"), Meta{}) + res, err := s.Put(strings.NewReader("never actually kept"), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -289,7 +350,7 @@ func TestGet_Has_InvalidCID(t *testing.T) { // read as a single-entry Pins list. func TestGet_DecodesLegacySidecar(t *testing.T) { s := open(t) - res, err := s.Put(strings.NewReader("legacy blob"), Meta{}) + res, err := s.Put(strings.NewReader("legacy blob"), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -313,7 +374,7 @@ func TestGet_DecodesLegacySidecar(t *testing.T) { func TestPut_UpgradesLegacySidecarOnNextWrite(t *testing.T) { s := open(t) content := []byte("legacy blob") - res, err := s.Put(bytes.NewReader(content), Meta{}) + res, err := s.Put(bytes.NewReader(content), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -420,7 +481,7 @@ func TestPin_IsIdempotentForADIDThatAlreadyHoldsAClaim(t *testing.T) { func TestPin_UnknownCID(t *testing.T) { s := open(t) - res, err := s.Put(strings.NewReader("stored, then deleted"), Meta{}) + res, err := s.Put(strings.NewReader("stored, then deleted"), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -435,7 +496,7 @@ func TestPin_UnknownCID(t *testing.T) { func TestUnpin_UnknownCID(t *testing.T) { s := open(t) - res, err := s.Put(strings.NewReader("never actually kept"), Meta{}) + res, err := s.Put(strings.NewReader("never actually kept"), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -533,7 +594,7 @@ func TestList(t *testing.T) { s := open(t) want := map[cid.Cid]bool{} for _, content := range []string{"first blob", "second blob", "third blob"} { - res, err := s.Put(strings.NewReader(content), Meta{}) + res, err := s.Put(strings.NewReader(content), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) } @@ -576,7 +637,7 @@ func (r *erroringReader) Read(p []byte) (int, error) { func TestPut_AbortedUploadLeavesNoBlob(t *testing.T) { s := open(t) - _, err := s.Put(&erroringReader{data: []byte("partial"), fail: errors.New("connection reset")}, Meta{}) + _, err := s.Put(&erroringReader{data: []byte("partial"), fail: errors.New("connection reset")}, pinMeta("text/plain", "did:plc:whoever")) if err == nil { t.Fatal("Put() = nil error, want error for aborted upload") } @@ -627,7 +688,7 @@ func TestDelete_RemovesBlobAndSidecar(t *testing.T) { func TestDelete_UnknownCID(t *testing.T) { s := open(t) - res, err := s.Put(strings.NewReader("never actually kept"), Meta{}) + res, err := s.Put(strings.NewReader("never actually kept"), pinMeta("text/plain", "did:plc:whoever")) if err != nil { t.Fatalf("Put: %v", err) }