hold/gc: sweep orphaned scan and image config records

A user who deletes a manifest record directly on their PDS was invisible
to the hold. purgeManifest only fires on appview-driven deletes, and the
GC orphan sweep only walked io.atcr.hold.layer, so scan and image config
records survived forever and the user's storage number never dropped.

analyzeRecords now also walks io.atcr.hold.scan and
io.atcr.hold.image.config, applying the same orphan test the layer sweep
uses: a record dies only when its owning PDS was reachable and
demonstrably lacks the manifest. Unreachable PDS, unparseable URI,
unparseable timestamp, and immature records all keep.

Splitting the grace period was required to make this useful, not
cosmetic. Records only need to outlast a push (blobs and layer records
are written before the manifest reaches the user's PDS), so they age out
at 24h. Blobs stay on the 7 day window, but since records now age out
faster than the blobs they name, a record can no longer serve as its
blob's clock. Blob age comes from S3 LastModified instead, which
WalkBlobs already received from ListObjectsV2 and was discarding.

Net effect: pruning an old manifest frees the user's storage on the next
nightly run rather than never, while the bytes are still reclaimed on
the same best-effort schedule as before.

DeleteManifestAuxRecord is deliberately narrow, accepting only the scan
and image config collections, so a bug in the sweep can't reach captain,
crew, or layer records.

Not addressed: the SBOM and vuln report blobs those scan records point
at. They share the /repos/{holdDID}/blobs/ prefix with the hold's avatar
and OG images, so collecting them needs a referenced-CID set built from
the hold's own records rather than a prefix walk.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Evan Jarrett
2026-08-02 19:20:52 -05:00
co-authored by Claude Opus 4.8
parent 3298797603
commit f12d5e054c
7 changed files with 664 additions and 37 deletions
@@ -40,22 +40,24 @@
<!-- Orphaned Records table -->
{{if .Preview.OrphanedRecords}}
<div class="collapse collapse-arrow bg-base-100 shadow-sm">
<input type="checkbox" aria-label="Show orphaned layer records" />
<input type="checkbox" aria-label="Show orphaned records" />
<div class="collapse-title font-medium">
{{ icon "file-x" "size-4" }} Orphaned Layer Records ({{len .Preview.OrphanedRecords}})
{{ icon "file-x" "size-4" }} Orphaned Records ({{len .Preview.OrphanedRecords}})
</div>
<div class="collapse-content">
<div class="overflow-x-auto">
<table class="table table-sm table-fixed">
<caption class="sr-only">Orphaned layer records</caption>
<caption class="sr-only">Orphaned layer, scan, and image config records</caption>
<colgroup>
<col style="width: 15%">
<col style="width: 25%">
<col style="width: 45%">
<col style="width: 15%">
<col style="width: 18%">
<col style="width: 14%">
<col style="width: 20%">
<col style="width: 36%">
<col style="width: 12%">
</colgroup>
<thead>
<tr>
<th scope="col">Collection</th>
<th scope="col">RKey</th>
<th scope="col">Digest</th>
<th scope="col">Manifest</th>
@@ -65,6 +67,7 @@
<tbody>
{{range .Preview.OrphanedRecords}}
<tr>
<td><code class="text-xs font-mono truncate block max-w-full" title="{{.Collection}}">{{.Collection}}</code></td>
<td><code class="text-xs font-mono truncate block max-w-full" title="{{.Rkey}}">{{.Rkey}}</code></td>
<td><code class="text-xs font-mono truncate block max-w-full" title="{{.Digest}}">{{truncate .Digest 24}}</code></td>
<td><code class="text-xs font-mono truncate block max-w-full" title="{{.ManifestURI}}">{{truncate .ManifestURI 50}}</code></td>
@@ -171,7 +174,7 @@
hx-post="/admin/api/gc/delete-records"
hx-target="#gc-results"
hx-swap="innerHTML"
hx-confirm="Delete {{len .Preview.OrphanedRecords}} orphaned layer records?">
hx-confirm="Delete {{len .Preview.OrphanedRecords}} orphaned records?">
{{ icon "file-x" "size-4" }}
Delete {{len .Preview.OrphanedRecords}} Orphaned Records
</button>
+297
View File
@@ -0,0 +1,297 @@
package gc
import (
"bytes"
"testing"
"time"
"atcr.io/pkg/atproto"
)
// pastGrace is comfortably older than gcRecordGracePeriod; inGrace is newer.
var (
pastGrace = time.Now().Add(-gcRecordGracePeriod - 24*time.Hour)
inGrace = time.Now().Add(-1 * time.Minute)
)
func TestAuxRecordOrphaned(t *testing.T) {
const (
knownURI = "at://did:plc:alice/io.atcr.manifest/abc123"
deletedURI = "at://did:plc:alice/io.atcr.manifest/def456"
offlineURI = "at://did:plc:bob/io.atcr.manifest/ghi789"
)
knownManifests := map[string]*manifestInfo{
knownURI: {URI: knownURI, UserDID: "did:plc:alice"},
}
// alice's PDS answered this run; bob's did not.
fetchedUsers := map[string]bool{"did:plc:alice": true}
tests := []struct {
name string
manifestURI string
createdAt time.Time
wantOrphan bool
wantDID string
}{
{
name: "manifest deleted on reachable PDS is orphaned",
manifestURI: deletedURI,
createdAt: pastGrace,
wantOrphan: true,
wantDID: "did:plc:alice",
},
{
name: "manifest still present is kept",
manifestURI: knownURI,
createdAt: pastGrace,
wantOrphan: false,
},
{
name: "record inside grace window is kept",
manifestURI: deletedURI,
createdAt: inGrace,
wantOrphan: false,
},
{
name: "unreachable PDS is kept even though manifest is unknown",
manifestURI: offlineURI,
createdAt: pastGrace,
wantOrphan: false,
},
{
name: "zero timestamp is kept",
manifestURI: deletedURI,
createdAt: time.Time{},
wantOrphan: false,
},
{
name: "unparseable manifest URI is kept",
manifestURI: "not-an-at-uri",
createdAt: pastGrace,
wantOrphan: false,
},
{
name: "empty manifest URI is kept",
manifestURI: "",
createdAt: pastGrace,
wantOrphan: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
parts, orphan := auxRecordOrphaned(tt.manifestURI, tt.createdAt, knownManifests, fetchedUsers)
if orphan != tt.wantOrphan {
t.Errorf("auxRecordOrphaned(%q) = %v, want %v", tt.manifestURI, orphan, tt.wantOrphan)
}
if !orphan {
return
}
if parts == nil {
t.Fatal("expected parsed AT-URI parts for an orphaned record, got nil")
}
if parts.DID != tt.wantDID {
t.Errorf("parts.DID = %q, want %q", parts.DID, tt.wantDID)
}
})
}
}
// A record exactly at the grace boundary must be collected, not kept — the
// boundary test is `< gcRecordGracePeriod`, so equal-or-older is eligible.
func TestAuxRecordOrphanedGraceBoundary(t *testing.T) {
const uri = "at://did:plc:alice/io.atcr.manifest/gone"
fetchedUsers := map[string]bool{"did:plc:alice": true}
// A hair past the boundary, to avoid flaking on the clock ticking forward
// between constructing the timestamp and evaluating time.Since.
justPast := time.Now().Add(-gcRecordGracePeriod - time.Second)
if _, orphan := auxRecordOrphaned(uri, justPast, nil, fetchedUsers); !orphan {
t.Error("record just past the grace period should be orphaned")
}
justInside := time.Now().Add(-gcRecordGracePeriod + time.Minute)
if _, orphan := auxRecordOrphaned(uri, justInside, nil, fetchedUsers); orphan {
t.Error("record just inside the grace period should be kept")
}
}
// Blobs age on their own clock, independent of the records that name them.
// Their window is deliberately longer than the record window, so a blob whose
// records were already collected still survives until its own age is reached.
func TestBlobPastGrace(t *testing.T) {
tests := []struct {
name string
lastModified time.Time
want bool
}{
{
name: "blob older than the blob grace period is collectable",
lastModified: time.Now().Add(-gcBlobGracePeriod - time.Hour),
want: true,
},
{
name: "blob inside the blob grace period is kept",
lastModified: time.Now().Add(-gcBlobGracePeriod + time.Hour),
want: false,
},
{
name: "freshly uploaded blob is kept",
lastModified: time.Now(),
want: false,
},
{
name: "unknown modification time is kept",
lastModified: time.Time{},
want: false,
},
{
name: "blob past the record window but not the blob window is kept",
// The case that motivates the split: records for this blob are
// already eligible for collection, the bytes are not.
lastModified: time.Now().Add(-gcRecordGracePeriod - time.Hour),
want: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := blobPastGrace(tt.lastModified); got != tt.want {
t.Errorf("blobPastGrace(%v) = %v, want %v", tt.lastModified, got, tt.want)
}
})
}
}
// The split only holds if blobs outlive records. If these are ever reordered,
// a blob could be deleted while a live record still names it.
func TestBlobGraceOutlivesRecordGrace(t *testing.T) {
if gcBlobGracePeriod <= gcRecordGracePeriod {
t.Errorf("gcBlobGracePeriod (%v) must exceed gcRecordGracePeriod (%v): blobs must outlive the records naming them",
gcBlobGracePeriod, gcRecordGracePeriod)
}
}
func TestDecodeAuxRecordBytes(t *testing.T) {
const (
manifestURI = "at://did:plc:alice/io.atcr.manifest/abc123"
scannedAt = "2026-03-04T05:06:07Z"
createdAt = "2026-01-02T03:04:05Z"
)
scan := &atproto.ScanRecord{
Type: atproto.ScanCollection,
Manifest: manifestURI,
UserDID: "did:plc:alice",
ScannedAt: scannedAt,
}
var scanBuf bytes.Buffer
if err := scan.MarshalCBOR(&scanBuf); err != nil {
t.Fatalf("marshal scan record: %v", err)
}
cfg := &atproto.ImageConfigRecord{
Type: atproto.ImageConfigCollection,
Manifest: manifestURI,
ConfigJSON: `{"architecture":"amd64"}`,
CreatedAt: createdAt,
}
var cfgBuf bytes.Buffer
if err := cfg.MarshalCBOR(&cfgBuf); err != nil {
t.Fatalf("marshal image config record: %v", err)
}
tests := []struct {
name string
collection string
data []byte
wantURI string
wantTime string // RFC3339, or "" for the zero time
wantErr bool
}{
{
name: "scan record",
collection: atproto.ScanCollection,
data: scanBuf.Bytes(),
wantURI: manifestURI,
wantTime: scannedAt,
},
{
name: "image config record",
collection: atproto.ImageConfigCollection,
data: cfgBuf.Bytes(),
wantURI: manifestURI,
wantTime: createdAt,
},
{
name: "unsupported collection is rejected",
collection: atproto.LayerCollection,
data: scanBuf.Bytes(),
wantErr: true,
},
{
name: "malformed CBOR is an error",
collection: atproto.ScanCollection,
data: []byte("not cbor at all"),
wantErr: true,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
uri, ts, err := decodeAuxRecordBytes(tt.collection, tt.data)
if tt.wantErr {
if err == nil {
t.Fatal("expected an error, got nil")
}
return
}
if err != nil {
t.Fatalf("decodeAuxRecordBytes: %v", err)
}
if uri != tt.wantURI {
t.Errorf("manifest URI = %q, want %q", uri, tt.wantURI)
}
want, parseErr := time.Parse(time.RFC3339, tt.wantTime)
if parseErr != nil {
t.Fatalf("bad want time in test case: %v", parseErr)
}
if !ts.Equal(want) {
t.Errorf("timestamp = %v, want %v", ts, want)
}
})
}
}
// A record whose timestamp doesn't parse must decode without error but yield
// the zero time, so auxRecordOrphaned keeps it rather than collecting it.
func TestDecodeAuxRecordBytesBadTimestampIsKept(t *testing.T) {
const manifestURI = "at://did:plc:alice/io.atcr.manifest/abc123"
scan := &atproto.ScanRecord{
Type: atproto.ScanCollection,
Manifest: manifestURI,
ScannedAt: "yesterday-ish",
}
var buf bytes.Buffer
if err := scan.MarshalCBOR(&buf); err != nil {
t.Fatalf("marshal scan record: %v", err)
}
uri, ts, err := decodeAuxRecordBytes(atproto.ScanCollection, buf.Bytes())
if err != nil {
t.Fatalf("decodeAuxRecordBytes: %v", err)
}
if uri != manifestURI {
t.Errorf("manifest URI = %q, want %q", uri, manifestURI)
}
if !ts.IsZero() {
t.Errorf("timestamp = %v, want the zero time", ts)
}
fetchedUsers := map[string]bool{"did:plc:alice": true}
if _, orphan := auxRecordOrphaned(uri, ts, nil, fetchedUsers); orphan {
t.Error("record with an unparseable timestamp should be kept, not orphaned")
}
}
+22 -5
View File
@@ -1,6 +1,11 @@
// Package gc implements garbage collection for the hold service.
// It periodically cleans up orphaned blobs from S3 storage based on
// layer records in the hold's embedded PDS.
//
// It collects two things on separate clocks: records in the hold's embedded
// PDS (layer, scan, and image config) whose manifest no longer exists in the
// owning user's PDS, and blobs in S3 that no live manifest references. Records
// go first and blobs follow later, so a user who deletes a manifest record
// directly on their PDS stops being billed for it well before the bytes are
// actually reclaimed.
package gc
import "time"
@@ -10,9 +15,21 @@ const (
// gcInterval is how often GC runs (nightly)
gcInterval = 24 * time.Hour
// gcGracePeriod is how old a layer record must be before it's considered for GC.
// Records created in the last 7 days are skipped (GDPR/CCPA compliant).
gcGracePeriod = 7 * 24 * time.Hour
// gcBlobGracePeriod is how old a blob must be, by its own S3 modification
// time, before GC will delete it. Blob pruning is best-effort: reclaiming
// bytes a week late costs storage, reclaiming them early can destroy
// content a client is still pushing or a takedown may yet reverse.
gcBlobGracePeriod = 7 * 24 * time.Hour
// gcRecordGracePeriod is how old a record must be before GC treats a
// missing manifest as an intentional deletion rather than a race.
//
// Records are metadata, not content, so they don't need the blob window.
// What they do need is to outlast a push: blobs and layer records are
// written before the manifest reaches the user's PDS, so a record younger
// than this may name a manifest that simply hasn't landed yet. A day is
// far longer than any push and still collects on the next nightly run.
gcRecordGracePeriod = 24 * time.Hour
)
// Config holds GC configuration
+234 -22
View File
@@ -23,8 +23,11 @@ import (
// maxPreviewItems caps per-category detail slices to prevent memory/HTML bloat
const maxPreviewItems = 10000
// OrphanedRecordDetail holds info about a single orphaned layer record
// OrphanedRecordDetail holds info about a single orphaned record. Collection
// names which record type it is — layer records carry a digest, media type,
// and size; scan and image-config records leave those empty.
type OrphanedRecordDetail struct {
Collection string `json:"collection"`
Rkey string `json:"rkey"`
Digest string `json:"digest"`
ManifestURI string `json:"manifestUri"`
@@ -33,6 +36,12 @@ type OrphanedRecordDetail struct {
Size int64 `json:"size"`
}
// orphanRef is the minimal address needed to delete an orphaned record.
type orphanRef struct {
Collection string
Rkey string
}
// OrphanedBlobDetail holds info about a single orphaned blob in S3
type OrphanedBlobDetail struct {
Digest string `json:"digest"`
@@ -175,7 +184,7 @@ type manifestInfo struct {
// analysisResult holds intermediate data from record analysis, shared between Run and Preview
type analysisResult struct {
referenced map[string]bool
orphanedRkeys []string // rkeys for deletion in Run
orphanedRefs []orphanRef // collection+rkey for deletion in Run
orphanedDetails []OrphanedRecordDetail // details for display in Preview
missingDetails []MissingRecordDetail // details for creation in Run / display in Preview
usersChecked int64
@@ -404,13 +413,13 @@ func (gc *GarbageCollector) doRun(ctx context.Context) (*GCResult, error) {
return nil, fmt.Errorf("phase 1 (analyze records) failed: %w", err)
}
result.OrphanedRecords = int64(len(analysis.orphanedRkeys))
result.OrphanedRecords = int64(len(analysis.orphanedRefs))
result.UsersChecked = analysis.usersChecked
result.ManifestsChecked = analysis.manifestsChecked
gc.logger.Info("Phase 1 complete",
"referenced", len(analysis.referenced),
"orphanedRecords", len(analysis.orphanedRkeys),
"orphanedRecords", len(analysis.orphanedRefs),
"missingRecords", len(analysis.missingDetails))
// Reconcile: create missing layer records
@@ -421,7 +430,7 @@ func (gc *GarbageCollector) doRun(ctx context.Context) (*GCResult, error) {
// Phase 2: Delete orphaned layer records
gc.setProgress("deleting", "Deleting orphaned records...", "run")
if err := gc.deleteOrphanedRecords(ctx, analysis.orphanedRkeys, result); err != nil {
if err := gc.deleteOrphanedRecords(ctx, analysis.orphanedRefs, result); err != nil {
gc.logger.Error("Phase 2 (delete orphaned records) failed", "error", err)
// Continue to phase 3 - we can still clean up blobs
}
@@ -568,13 +577,13 @@ func (gc *GarbageCollector) doDeleteOrphanedRecords(ctx context.Context) (*GCRes
OrphanedRecords: int64(len(preview.OrphanedRecords)),
}
rkeys := make([]string, len(preview.OrphanedRecords))
refs := make([]orphanRef, len(preview.OrphanedRecords))
for i, r := range preview.OrphanedRecords {
rkeys[i] = r.Rkey
refs[i] = orphanRef{Collection: r.Collection, Rkey: r.Rkey}
}
gc.logger.Info("Deleting orphaned records", "count", len(rkeys))
if err := gc.deleteOrphanedRecords(ctx, rkeys, result); err != nil {
gc.logger.Info("Deleting orphaned records", "count", len(refs))
if err := gc.deleteOrphanedRecords(ctx, refs, result); err != nil {
return nil, fmt.Errorf("delete orphaned records: %w", err)
}
result.Duration = time.Since(start)
@@ -760,9 +769,10 @@ func (gc *GarbageCollector) analyzeRecords(ctx context.Context) (*analysisResult
pairKey := layer.Manifest + "|" + layer.Digest
coveredPairs[pairKey] = true
// Grace period: skip recent records
// Too young to judge: the manifest may still be in flight. Keep
// the record and protect its blob until the record matures.
recordTime := tidToTime(rec.Rkey)
if time.Since(recordTime) < gcGracePeriod {
if time.Since(recordTime) < gcRecordGracePeriod {
result.referenced[layer.Digest] = true
continue
}
@@ -778,9 +788,13 @@ func (gc *GarbageCollector) analyzeRecords(ctx context.Context) (*analysisResult
result.referenced[layer.Digest] = true
} else {
// User's PDS was reachable but manifest not found — orphaned
result.orphanedRkeys = append(result.orphanedRkeys, rec.Rkey)
result.orphanedRefs = append(result.orphanedRefs, orphanRef{
Collection: atproto.LayerCollection,
Rkey: rec.Rkey,
})
if len(result.orphanedDetails) < maxPreviewItems {
result.orphanedDetails = append(result.orphanedDetails, OrphanedRecordDetail{
Collection: atproto.LayerCollection,
Rkey: rec.Rkey,
Digest: layer.Digest,
ManifestURI: layer.Manifest,
@@ -809,6 +823,18 @@ func (gc *GarbageCollector) analyzeRecords(ctx context.Context) (*analysisResult
gc.logger.Info("Scanned layer records", "total", result.totalRecords, "coveredPairs", len(coveredPairs))
// Step 3b: Scan scan and image-config records. These are keyed by manifest
// digest rather than TID and hold no blob references of their own, so the
// only question is whether their manifest still exists. Without this sweep
// they survive forever when a user deletes a manifest record directly on
// their PDS — the purgeManifest XRPC only fires on appview-driven deletes.
gc.setProgress("records", "Scanning scan and image config records...", gc.operationType)
for _, collection := range []string{atproto.ScanCollection, atproto.ImageConfigCollection} {
if err := gc.scanAuxRecords(ctx, collection, knownManifests, fetchedUsers, result); err != nil {
return nil, fmt.Errorf("scan %s records: %w", collection, err)
}
}
// Step 4: Identify missing layer records (uncovered manifest+layer pairs)
for _, m := range knownManifests {
for _, layer := range m.Record.Layers {
@@ -831,6 +857,162 @@ func (gc *GarbageCollector) analyzeRecords(ctx context.Context) (*analysisResult
return result, nil
}
// scanAuxRecords walks one auxiliary manifest collection (scan or image
// config) and appends orphans to result. A record is orphaned when its
// manifest is absent from the owning user's PDS and that PDS was reachable —
// the same test the layer sweep applies, so an unreachable PDS never causes a
// deletion.
//
// Grace is taken from the record body (scannedAt / createdAt) rather than the
// rkey, because these collections use deterministic digest rkeys, not TIDs.
// An unparseable or absent timestamp is treated as in-grace, so a malformed
// record is kept rather than collected.
func (gc *GarbageCollector) scanAuxRecords(
ctx context.Context,
collection string,
knownManifests map[string]*manifestInfo,
fetchedUsers map[string]bool,
result *analysisResult,
) error {
recordsIndex := gc.pds.RecordsIndex()
if recordsIndex == nil {
return fmt.Errorf("records index not available")
}
cursor := ""
scanned, orphaned := 0, 0
for {
records, nextCursor, err := recordsIndex.ListRecords(collection, 1000, cursor, true)
if err != nil {
return fmt.Errorf("list records: %w", err)
}
for _, rec := range records {
scanned++
result.totalRecords++
manifestURI, createdAt, err := gc.decodeAuxRecord(ctx, collection, rec)
if err != nil {
gc.logger.Warn("Failed to decode record",
"collection", collection, "rkey", rec.Rkey, "error", err)
continue
}
parts, isOrphan := auxRecordOrphaned(manifestURI, createdAt, knownManifests, fetchedUsers)
if !isOrphan {
continue
}
orphaned++
result.orphanedRefs = append(result.orphanedRefs, orphanRef{
Collection: collection,
Rkey: rec.Rkey,
})
if len(result.orphanedDetails) < maxPreviewItems {
result.orphanedDetails = append(result.orphanedDetails, OrphanedRecordDetail{
Collection: collection,
Rkey: rec.Rkey,
ManifestURI: manifestURI,
UserDID: parts.DID,
})
}
gc.logger.Debug("Found orphaned record",
"collection", collection, "rkey", rec.Rkey, "manifest", manifestURI)
}
if nextCursor == "" {
break
}
cursor = nextCursor
}
gc.logger.Info("Scanned auxiliary manifest records",
"collection", collection, "scanned", scanned, "orphaned", orphaned)
return nil
}
// auxRecordOrphaned decides whether a scan or image-config record should be
// collected, given the manifest it names and its creation time. It returns the
// parsed AT-URI alongside the verdict so the caller can record the owning DID.
//
// Every uncertain case resolves to "keep": a record still inside the grace
// window, one whose timestamp couldn't be parsed (zero time), one whose
// manifest URI is unparseable, and one whose owning PDS was unreachable this
// run. Only a reachable PDS that demonstrably lacks the manifest orphans it.
func auxRecordOrphaned(
manifestURI string,
createdAt time.Time,
knownManifests map[string]*manifestInfo,
fetchedUsers map[string]bool,
) (*atURIParts, bool) {
if createdAt.IsZero() || time.Since(createdAt) < gcRecordGracePeriod {
return nil, false
}
if _, known := knownManifests[manifestURI]; known {
return nil, false
}
parts := parseATURI(manifestURI)
if parts == nil || !fetchedUsers[parts.DID] {
return nil, false
}
return parts, true
}
// decodeAuxRecord reads a scan or image-config record from the carstore and
// returns the manifest AT-URI it belongs to along with its creation time.
func (gc *GarbageCollector) decodeAuxRecord(ctx context.Context, collection string, rec pds.Record) (string, time.Time, error) {
recordPath := rec.Collection + "/" + rec.Rkey
_, recBytes, err := gc.pds.GetRecordBytes(ctx, recordPath)
if err != nil {
return "", time.Time{}, fmt.Errorf("get record bytes: %w", err)
}
return decodeAuxRecordBytes(collection, *recBytes)
}
// decodeAuxRecordBytes decodes raw CBOR for one of the auxiliary manifest
// collections. An unparseable timestamp yields the zero time rather than an
// error, which auxRecordOrphaned treats as in-grace, so a record with a
// malformed date is kept instead of collected.
func decodeAuxRecordBytes(collection string, data []byte) (string, time.Time, error) {
var manifestURI, timestamp string
switch collection {
case atproto.ScanCollection:
var scan atproto.ScanRecord
if err := scan.UnmarshalCBOR(bytes.NewReader(data)); err != nil {
return "", time.Time{}, fmt.Errorf("unmarshal CBOR: %w", err)
}
manifestURI, timestamp = scan.Manifest, scan.ScannedAt
case atproto.ImageConfigCollection:
var cfg atproto.ImageConfigRecord
if err := cfg.UnmarshalCBOR(bytes.NewReader(data)); err != nil {
return "", time.Time{}, fmt.Errorf("unmarshal CBOR: %w", err)
}
manifestURI, timestamp = cfg.Manifest, cfg.CreatedAt
default:
return "", time.Time{}, fmt.Errorf("unsupported collection %q", collection)
}
parsed, err := time.Parse(time.RFC3339, timestamp)
if err != nil {
return manifestURI, time.Time{}, nil
}
return manifestURI, parsed, nil
}
// blobPastGrace reports whether a blob is old enough to delete, based on its
// own S3 modification time rather than on any record that references it.
//
// Records now age out faster than blobs, so a blob can outlive every record
// naming it. Its age is the only thing left to judge it by. A listing that
// reports no modification time yields the zero time, which is treated as
// "unknown age" and keeps the blob.
func blobPastGrace(lastModified time.Time) bool {
if lastModified.IsZero() {
return false
}
return time.Since(lastModified) >= gcBlobGracePeriod
}
// scanOrphanedBlobDetails walks S3 and returns details of unreferenced blobs.
// Read-only — no deletions. Returns orphaned blob details and total blob count.
func (gc *GarbageCollector) scanOrphanedBlobDetails(ctx context.Context, referenced map[string]bool) ([]OrphanedBlobDetail, int, error) {
@@ -838,7 +1020,7 @@ func (gc *GarbageCollector) scanOrphanedBlobDetails(ctx context.Context, referen
totalBlobs := 0
blobsPath := "/docker/registry/v2/blobs"
err := gc.s3.WalkBlobs(ctx, blobsPath, func(key string, size int64) error {
err := gc.s3.WalkBlobs(ctx, blobsPath, func(key string, size int64, lastModified time.Time) error {
if !strings.HasSuffix(key, "/data") {
return nil
}
@@ -850,7 +1032,7 @@ func (gc *GarbageCollector) scanOrphanedBlobDetails(ctx context.Context, referen
totalBlobs++
if !referenced[digest] {
if !referenced[digest] && blobPastGrace(lastModified) {
if len(orphaned) < maxPreviewItems {
orphaned = append(orphaned, OrphanedBlobDetail{
Digest: digest,
@@ -1511,18 +1693,41 @@ func (gc *GarbageCollector) checkPredecessor(ctx context.Context, holdDID string
}
// deleteOrphanedRecords removes layer records whose manifests no longer exist
func (gc *GarbageCollector) deleteOrphanedRecords(ctx context.Context, orphanedRkeys []string, result *GCResult) error {
for _, rkey := range orphanedRkeys {
if err := gc.pds.DeleteLayerRecord(ctx, rkey); err != nil {
gc.logger.Error("Failed to delete layer record", "rkey", rkey, "error", err)
continue
func (gc *GarbageCollector) deleteOrphanedRecords(ctx context.Context, refs []orphanRef, result *GCResult) error {
for _, ref := range refs {
// An empty collection means the ref came from a preview taken before
// the sweep covered scan and image-config records; those were all
// layer records.
collection := ref.Collection
if collection == "" {
collection = atproto.LayerCollection
}
if collection == atproto.LayerCollection {
if err := gc.pds.DeleteLayerRecord(ctx, ref.Rkey); err != nil {
gc.logger.Error("Failed to delete layer record", "rkey", ref.Rkey, "error", err)
continue
}
} else {
deleted, err := gc.pds.DeleteManifestAuxRecord(ctx, collection, ref.Rkey)
if err != nil {
gc.logger.Error("Failed to delete record",
"collection", collection, "rkey", ref.Rkey, "error", err)
continue
}
if !deleted {
// Already gone (e.g. a concurrent purge) — not an error, but
// don't count it as work this run did.
continue
}
}
result.RecordsDeleted++
gc.logger.Debug("Deleted orphaned layer record", "rkey", rkey)
gc.logger.Debug("Deleted orphaned record", "collection", collection, "rkey", ref.Rkey)
}
gc.logger.Info("Phase 2 complete",
"orphaned", len(orphanedRkeys),
"orphaned", len(refs),
"deleted", result.RecordsDeleted)
return nil
@@ -1532,7 +1737,7 @@ func (gc *GarbageCollector) deleteOrphanedRecords(ctx context.Context, orphanedR
func (gc *GarbageCollector) deleteOrphanedBlobs(ctx context.Context, referenced map[string]bool, result *GCResult) error {
blobsPath := "/docker/registry/v2/blobs"
err := gc.s3.WalkBlobs(ctx, blobsPath, func(key string, size int64) error {
err := gc.s3.WalkBlobs(ctx, blobsPath, func(key string, size int64, lastModified time.Time) error {
// Only process data files
if !strings.HasSuffix(key, "/data") {
return nil
@@ -1549,6 +1754,13 @@ func (gc *GarbageCollector) deleteOrphanedBlobs(ctx context.Context, referenced
return nil
}
// Unreferenced, but young enough that we'd rather keep paying for it
// than risk deleting content still being pushed. Records for this
// blob may already be gone; blob pruning runs on its own clock.
if !blobPastGrace(lastModified) {
return nil
}
result.OrphanedBlobs++
if err := gc.s3.Delete(ctx, key); err != nil {
+17
View File
@@ -72,6 +72,23 @@ func (p *HoldPDS) PurgeManifestRecords(ctx context.Context, manifestURI string)
return res, nil
}
// DeleteManifestAuxRecord deletes a single scan or image-config record by its
// deterministic digest rkey. Restricted to those two collections so callers
// can't reach captain, crew, or layer records through it — layer records have
// their own typed deletion path (DeleteLayerRecord).
//
// Idempotent: a missing record reports false with no error. Used by the GC's
// orphaned-record sweep, which addresses these records directly rather than
// going through PurgeManifestRecords (it already knows they're orphaned).
func (p *HoldPDS) DeleteManifestAuxRecord(ctx context.Context, collection, rkey string) (bool, error) {
switch collection {
case atproto.ScanCollection, atproto.ImageConfigCollection:
default:
return false, fmt.Errorf("collection %q is not an auxiliary manifest record", collection)
}
return p.tryDeleteRecord(ctx, collection, rkey), nil
}
// PurgeUserManifests purges every manifest's records for a given DID. Used by
// user-level takedowns (URI = at://<did>) where the labeler has not enumerated
// individual manifest URIs.
+73
View File
@@ -5,6 +5,7 @@ import (
"testing"
"atcr.io/pkg/atproto"
"github.com/ipfs/go-cid"
)
func TestPurgeManifestRecordsRemovesAll(t *testing.T) {
@@ -125,6 +126,78 @@ func TestPurgeUserManifestsCollectsAcrossManifests(t *testing.T) {
}
}
func TestDeleteManifestAuxRecord(t *testing.T) {
pds := setupTestPDSWithIndex(t, "did:plc:owner")
ctx := sharedCtx
const manifestDigest = "sha256:e692418e4cbaf90ca69d05a66403747baa33ee08806650b51fab815ad7fc331f"
manifestURI := atproto.BuildManifestURI("did:plc:alice", manifestDigest)
rkey := atproto.ScanRecordKey(manifestDigest)
scanRec := &atproto.ScanRecord{
Type: atproto.ScanCollection,
Manifest: manifestURI,
ScannedAt: "2026-05-02T00:00:00Z",
}
if _, _, _, err := pds.repomgr.UpsertRecord(ctx, pds.uid, atproto.ScanCollection, rkey, scanRec); err != nil {
t.Fatalf("create scan record: %v", err)
}
// A layer record for the same manifest must survive — this entry point
// addresses only the deterministic-rkey auxiliary records.
mustCreateLayer(t, pds, manifestURI, "sha256:layer-a", 1024)
deleted, err := pds.DeleteManifestAuxRecord(ctx, atproto.ScanCollection, rkey)
if err != nil {
t.Fatalf("DeleteManifestAuxRecord: %v", err)
}
if !deleted {
t.Error("deleted = false, want true for an existing scan record")
}
if _, _, err := pds.repomgr.GetRecord(ctx, pds.uid, atproto.ScanCollection, rkey, cid.Undef); err == nil {
t.Error("scan record still present after deletion")
}
rkeys, err := pds.listLayerRkeysForManifest(ctx, manifestURI)
if err != nil {
t.Fatalf("list layers: %v", err)
}
if len(rkeys) != 1 {
t.Errorf("layer rkeys = %d, want 1 (aux delete must not touch layer records)", len(rkeys))
}
// Deleting again is not an error, it just reports nothing was removed.
deleted, err = pds.DeleteManifestAuxRecord(ctx, atproto.ScanCollection, rkey)
if err != nil {
t.Fatalf("second DeleteManifestAuxRecord: %v", err)
}
if deleted {
t.Error("deleted = true on second call, want false")
}
}
func TestDeleteManifestAuxRecordRejectsOtherCollections(t *testing.T) {
pds := setupTestPDSWithIndex(t, "did:plc:owner")
// Layer, captain, and crew records must not be reachable through this
// entry point — a GC bug shouldn't be able to delete the hold's identity.
for _, collection := range []string{
atproto.LayerCollection,
atproto.CaptainCollection,
atproto.CrewCollection,
} {
t.Run(collection, func(t *testing.T) {
deleted, err := pds.DeleteManifestAuxRecord(sharedCtx, collection, "self")
if err == nil {
t.Errorf("DeleteManifestAuxRecord(%q) returned nil error, want rejection", collection)
}
if deleted {
t.Errorf("DeleteManifestAuxRecord(%q) reported a deletion", collection)
}
})
}
}
// mustCreateLayer is a tiny helper for purge tests — keeps the table-driven
// TestPurge body focused on the assertions.
func mustCreateLayer(t *testing.T, pds *HoldPDS, manifestURI, digest string, size int64) {
+10 -2
View File
@@ -334,7 +334,11 @@ func (s *S3Service) Delete(ctx context.Context, blobPath string) error {
// WalkBlobs paginates ListObjectsV2 under prefix and calls fn for each object.
// Keys passed to fn have the PathPrefix stripped (same format as BlobPath output).
func (s *S3Service) WalkBlobs(ctx context.Context, prefix string, fn func(key string, size int64) error) error {
//
// lastModified is the object's S3 modification time, which callers use to age
// blobs independently of any record that references them. It is the zero time
// if the listing didn't report one.
func (s *S3Service) WalkBlobs(ctx context.Context, prefix string, fn func(key string, size int64, lastModified time.Time) error) error {
s3Prefix := s.s3Key(prefix)
if !strings.HasSuffix(s3Prefix, "/") {
s3Prefix += "/"
@@ -367,7 +371,11 @@ func (s *S3Service) WalkBlobs(ctx context.Context, prefix string, fn func(key st
if obj.Size != nil {
size = *obj.Size
}
if err := fn(key, size); err != nil {
var lastModified time.Time
if obj.LastModified != nil {
lastModified = *obj.LastModified
}
if err := fn(key, size, lastModified); err != nil {
return err
}
}