mirror of
https://tangled.org/evan.jarrett.net/at-container-registry
synced 2026-09-02 16:26:56 +00:00
Completes the swap 0033 set up. layers and manifest_references move onto manifest_key and manifests.id is gone, which removes the last node-allocated identifier in the AppView schema. Statement order in 0034 is load-bearing. With foreign keys on, DROP TABLE performs an implicit DELETE FROM, so dropping manifests while layers still holds an ON DELETE CASCADE reference deletes every layer row. Migration 0009 did exactly that; it went unnoticed because the Jetstream backfill rebuilds layers from PDS records, so the damage healed itself. PRAGMA foreign_keys is no help: it is a no-op inside a transaction and migrations run in one. So the new children are built pointing at manifests_new, the old children are dropped first, and only then is the old manifests table dropped, by which point nothing references it. Verified both behaviors before relying on them. manifest_key is declared NOT NULL as well as PRIMARY KEY, because in SQLite a PRIMARY KEY column still accepts NULL unless it is INTEGER PRIMARY KEY. That constraint immediately caught four test helpers inserting manifests without one. Five queries used MAX(id) as "the newest manifest in this repo", which I had previously reported as absent after grepping only for ORDER BY. A derived key has no ordering, so recency now comes from created_at with manifest_key as a deterministic tiebreak. This is a real behavior change, and a fix: the two disagree whenever a manifest is indexed out of order, which the backfill does routinely, and created_at is the push time these queries always wanted. Both directions are tested, including that ties resolve the same way every run. InsertManifest and BatchInsertManifests no longer read anything back. The key is derived from (did, repository, digest), so the writer knows it before the statement runs: the select-back, its per-DID IN list, and the "manifest missing id after batch insert" branch all go away, along with the UNIQUE-conflict fallback that existed only to recover a rowid. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
530 lines
16 KiB
Go
530 lines
16 KiB
Go
package jetstream
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log/slog"
|
|
"strings"
|
|
"time"
|
|
|
|
"atcr.io/pkg/appview/db"
|
|
"atcr.io/pkg/atproto"
|
|
)
|
|
|
|
// batchManifests decodes all manifest records for a repo and writes them as
|
|
// a small set of multi-row INSERTs: one per table (manifests, layers,
|
|
// manifest_references, repository_annotations). This replaces the previous
|
|
// per-record chunked-transaction loop, which exceeded Bunny Database's
|
|
// remote transaction timeout once chunks grew large.
|
|
//
|
|
// Returns the number of manifest records that were successfully decoded and
|
|
// included in the batch. Decode/validation failures are logged and skipped.
|
|
func (b *BackfillWorker) batchManifests(ctx context.Context, did string, records []atproto.Record) (int, error) {
|
|
if len(records) == 0 {
|
|
return 0, nil
|
|
}
|
|
|
|
type decoded struct {
|
|
manifestRecord atproto.ManifestRecord
|
|
manifest db.Manifest
|
|
}
|
|
|
|
decodedRecords := make([]decoded, 0, len(records))
|
|
for i := range records {
|
|
r := &records[i]
|
|
var mr atproto.ManifestRecord
|
|
if err := json.Unmarshal(r.Value, &mr); err != nil {
|
|
slog.Warn("Backfill skipping invalid manifest record", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
if mr.Digest == "" || mr.Repository == "" {
|
|
slog.Warn("Backfill skipping manifest with missing fields", "uri", r.URI)
|
|
continue
|
|
}
|
|
|
|
// Resolve holdDID the same way the single-record path does.
|
|
holdDID := mr.HoldDID
|
|
if holdDID == "" && mr.HoldEndpoint != "" {
|
|
if resolved, err := atproto.ResolveHoldDID(ctx, mr.HoldEndpoint); err == nil {
|
|
holdDID = resolved
|
|
}
|
|
}
|
|
|
|
isList := len(mr.Manifests) > 0
|
|
artifactType := "container-image"
|
|
if !isList && mr.Config != nil {
|
|
artifactType = db.GetArtifactType(mr.Config.MediaType)
|
|
}
|
|
|
|
m := db.Manifest{
|
|
DID: did,
|
|
Repository: mr.Repository,
|
|
Digest: mr.Digest,
|
|
MediaType: mr.MediaType,
|
|
SchemaVersion: mr.SchemaVersion,
|
|
HoldEndpoint: holdDID,
|
|
ArtifactType: artifactType,
|
|
CreatedAt: mr.CreatedAt,
|
|
}
|
|
if !isList && mr.Config != nil {
|
|
m.ConfigDigest = mr.Config.Digest
|
|
m.ConfigSize = mr.Config.Size
|
|
}
|
|
if mr.Subject != nil {
|
|
m.SubjectDigest = mr.Subject.Digest
|
|
}
|
|
decodedRecords = append(decodedRecords, decoded{mr, m})
|
|
}
|
|
|
|
if len(decodedRecords) == 0 {
|
|
return 0, nil
|
|
}
|
|
|
|
// Phase 1: upsert all manifests in one batch, fetch ids.
|
|
manifests := make([]db.Manifest, len(decodedRecords))
|
|
for i, d := range decodedRecords {
|
|
manifests[i] = d.manifest
|
|
}
|
|
if err := db.BatchInsertManifests(b.db, manifests); err != nil {
|
|
return 0, fmt.Errorf("batch insert manifests: %w", err)
|
|
}
|
|
|
|
// Phase 2: derive layers, references and annotations. The manifest key is
|
|
// computed locally rather than read back from the database, so there is no
|
|
// window in which a row exists without us knowing how to reference it.
|
|
var (
|
|
layerRows []db.Layer
|
|
refRows []db.ManifestReference
|
|
)
|
|
|
|
// For annotations, we keep only the newest manifest per (did, repo) with a
|
|
// non-empty annotation set. Matches reconcileAnnotations semantics at
|
|
// backfill.go:573.
|
|
type newest struct {
|
|
createdAt time.Time
|
|
annotations map[string]string
|
|
}
|
|
newestByRepo := make(map[string]newest)
|
|
|
|
for _, d := range decodedRecords {
|
|
mkey := db.ManifestKey(did, d.manifest.Repository, d.manifest.Digest)
|
|
|
|
if len(d.manifestRecord.Manifests) > 0 {
|
|
for i, ref := range d.manifestRecord.Manifests {
|
|
var pa, po, pv, pov string
|
|
if ref.Platform != nil {
|
|
pa = ref.Platform.Architecture
|
|
po = ref.Platform.OS
|
|
pv = ref.Platform.Variant
|
|
pov = ref.Platform.OSVersion
|
|
}
|
|
isAttestation := false
|
|
if refType, ok := ref.Annotations["vnd.docker.reference.type"]; ok {
|
|
isAttestation = refType == "attestation-manifest"
|
|
}
|
|
refRows = append(refRows, db.ManifestReference{
|
|
ManifestKey: mkey,
|
|
Digest: ref.Digest,
|
|
MediaType: ref.MediaType,
|
|
Size: ref.Size,
|
|
PlatformArchitecture: pa,
|
|
PlatformOS: po,
|
|
PlatformVariant: pv,
|
|
PlatformOSVersion: pov,
|
|
IsAttestation: isAttestation,
|
|
ReferenceIndex: i,
|
|
})
|
|
}
|
|
} else {
|
|
for i, layer := range d.manifestRecord.Layers {
|
|
layerRows = append(layerRows, db.Layer{
|
|
ManifestKey: mkey,
|
|
Digest: layer.Digest,
|
|
MediaType: layer.MediaType,
|
|
Size: layer.Size,
|
|
LayerIndex: i,
|
|
Annotations: layer.Annotations,
|
|
})
|
|
}
|
|
}
|
|
|
|
if hasNonEmpty(d.manifestRecord.Annotations) {
|
|
key := d.manifest.Repository
|
|
prev, ok := newestByRepo[key]
|
|
if !ok || d.manifestRecord.CreatedAt.After(prev.createdAt) {
|
|
newestByRepo[key] = newest{d.manifestRecord.CreatedAt, d.manifestRecord.Annotations}
|
|
}
|
|
}
|
|
}
|
|
|
|
if err := db.BatchInsertLayers(b.db, layerRows); err != nil {
|
|
return 0, err
|
|
}
|
|
if err := db.BatchInsertManifestReferences(b.db, refRows); err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
// Flatten annotations into AnnotationRows.
|
|
var annotationRows []db.AnnotationRow
|
|
for repo, n := range newestByRepo {
|
|
for k, v := range n.annotations {
|
|
if v == "" {
|
|
continue
|
|
}
|
|
annotationRows = append(annotationRows, db.AnnotationRow{
|
|
DID: did,
|
|
Repository: repo,
|
|
Key: k,
|
|
Value: v,
|
|
})
|
|
}
|
|
}
|
|
if err := db.BatchUpsertRepositoryAnnotations(b.db, annotationRows); err != nil {
|
|
return 0, err
|
|
}
|
|
|
|
slog.Info("Backfill batch manifests",
|
|
"did", did,
|
|
"manifests", len(manifests),
|
|
"layers", len(layerRows),
|
|
"references", len(refRows),
|
|
"annotations", len(annotationRows))
|
|
|
|
return len(decodedRecords), nil
|
|
}
|
|
|
|
func hasNonEmpty(m map[string]string) bool {
|
|
for _, v := range m {
|
|
if v != "" {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// batchTags decodes tag records and writes them in one multi-row upsert.
|
|
func (b *BackfillWorker) batchTags(did string, records []atproto.Record) (int, error) {
|
|
tags := make([]db.Tag, 0, len(records))
|
|
for i := range records {
|
|
r := &records[i]
|
|
var tr atproto.TagRecord
|
|
if err := json.Unmarshal(r.Value, &tr); err != nil {
|
|
slog.Warn("Backfill skipping invalid tag record", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
digest, err := tr.GetManifestDigest()
|
|
if err != nil {
|
|
slog.Warn("Backfill skipping tag record without digest", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
if tr.Repository == "" || tr.Tag == "" {
|
|
continue
|
|
}
|
|
tags = append(tags, db.Tag{
|
|
DID: did,
|
|
Repository: tr.Repository,
|
|
Tag: tr.Tag,
|
|
Digest: digest,
|
|
CreatedAt: tr.UpdatedAt,
|
|
})
|
|
}
|
|
if err := db.BatchUpsertTags(b.db, tags); err != nil {
|
|
return 0, err
|
|
}
|
|
slog.Info("Backfill batch tags", "did", did, "rows", len(tags))
|
|
return len(tags), nil
|
|
}
|
|
|
|
// batchStars decodes star records and writes them in one multi-row upsert.
|
|
// Ensures star subject owners exist as users first (FK requirement).
|
|
func (b *BackfillWorker) batchStars(ctx context.Context, did string, records []atproto.Record) (int, error) {
|
|
stars := make([]db.StarInput, 0, len(records))
|
|
ownerDIDs := make(map[string]struct{})
|
|
|
|
for i := range records {
|
|
r := &records[i]
|
|
var sr atproto.StarRecord
|
|
if err := json.Unmarshal(r.Value, &sr); err != nil {
|
|
slog.Warn("Backfill skipping invalid star record", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
owner, repo, err := sr.GetSubjectDIDAndRepository()
|
|
if err != nil {
|
|
slog.Warn("Backfill skipping star with bad subject", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
ownerDIDs[owner] = struct{}{}
|
|
stars = append(stars, db.StarInput{
|
|
StarrerDID: did,
|
|
OwnerDID: owner,
|
|
Repository: repo,
|
|
CreatedAt: sr.CreatedAt,
|
|
})
|
|
}
|
|
|
|
// Ensure every star subject has a users row (FK to users.did on stars).
|
|
// These calls are idempotent and cached, so repeated owners cost nothing.
|
|
for owner := range ownerDIDs {
|
|
if err := b.processor.EnsureUserExists(ctx, owner); err != nil {
|
|
slog.Warn("Backfill failed to ensure star subject user", "owner_did", owner, "error", err)
|
|
}
|
|
}
|
|
|
|
if err := db.BatchUpsertStars(b.db, stars); err != nil {
|
|
return 0, err
|
|
}
|
|
slog.Info("Backfill batch stars", "did", did, "rows", len(stars))
|
|
return len(stars), nil
|
|
}
|
|
|
|
// batchRepoPages decodes repo page records and writes them in one upsert.
|
|
func (b *BackfillWorker) batchRepoPages(did string, records []atproto.Record) (int, error) {
|
|
pages := make([]db.RepoPage, 0, len(records))
|
|
for i := range records {
|
|
r := &records[i]
|
|
var pr atproto.RepoPageRecord
|
|
if err := json.Unmarshal(r.Value, &pr); err != nil {
|
|
slog.Warn("Backfill skipping invalid repo page", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
if pr.Repository == "" {
|
|
continue
|
|
}
|
|
avatarCID := ""
|
|
if pr.Avatar != nil && pr.Avatar.Ref.Link != "" {
|
|
avatarCID = pr.Avatar.Ref.Link
|
|
}
|
|
pages = append(pages, db.RepoPage{
|
|
DID: did,
|
|
Repository: pr.Repository,
|
|
Description: pr.Description,
|
|
AvatarCID: avatarCID,
|
|
UserEdited: pr.UserEdited,
|
|
CreatedAt: pr.CreatedAt,
|
|
UpdatedAt: pr.UpdatedAt,
|
|
})
|
|
}
|
|
if err := db.BatchUpsertRepoPages(b.db, pages); err != nil {
|
|
return 0, err
|
|
}
|
|
slog.Info("Backfill batch repo pages", "did", did, "rows", len(pages))
|
|
return len(pages), nil
|
|
}
|
|
|
|
// batchDailyStats decodes daily stats records and writes them in one upsert.
|
|
// Ensures every distinct owner exists as a user first (FK requirement).
|
|
func (b *BackfillWorker) batchDailyStats(ctx context.Context, holdDID string, records []atproto.Record) (int, error) {
|
|
stats := make([]db.DailyStats, 0, len(records))
|
|
ownerDIDs := make(map[string]struct{})
|
|
|
|
for i := range records {
|
|
r := &records[i]
|
|
var dr atproto.DailyStatsRecord
|
|
if err := json.Unmarshal(r.Value, &dr); err != nil {
|
|
slog.Warn("Backfill skipping invalid daily stats", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
if dr.OwnerDID == "" || dr.Repository == "" || dr.Date == "" {
|
|
continue
|
|
}
|
|
ownerDIDs[dr.OwnerDID] = struct{}{}
|
|
stats = append(stats, db.DailyStats{
|
|
DID: dr.OwnerDID,
|
|
Repository: dr.Repository,
|
|
Date: dr.Date,
|
|
PullCount: int(dr.PullCount),
|
|
PushCount: int(dr.PushCount),
|
|
})
|
|
}
|
|
|
|
for owner := range ownerDIDs {
|
|
if err := b.processor.EnsureUserExists(ctx, owner); err != nil {
|
|
slog.Warn("Backfill failed to ensure daily stats owner user", "owner_did", owner, "error", err)
|
|
}
|
|
}
|
|
|
|
if err := db.BatchUpsertDailyStats(b.db, stats); err != nil {
|
|
return 0, err
|
|
}
|
|
slog.Info("Backfill batch daily stats", "hold_did", holdDID, "rows", len(stats))
|
|
return len(stats), nil
|
|
}
|
|
|
|
// batchStats updates the in-memory stats cache from a hold's stats records,
|
|
// then flushes the aggregated view of every touched (owner, repo) to the
|
|
// repository_stats table in a single multi-row upsert. Aggregation is across
|
|
// all holds known to the cache, preserving the single-record semantics.
|
|
func (b *BackfillWorker) batchStats(ctx context.Context, holdDID string, records []atproto.Record) (int, error) {
|
|
type key struct{ owner, repo string }
|
|
touched := make(map[key]struct{})
|
|
ownerDIDs := make(map[string]struct{})
|
|
|
|
for i := range records {
|
|
r := &records[i]
|
|
var sr atproto.StatsRecord
|
|
if err := json.Unmarshal(r.Value, &sr); err != nil {
|
|
slog.Warn("Backfill skipping invalid stats record", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
if sr.OwnerDID == "" || sr.Repository == "" {
|
|
continue
|
|
}
|
|
|
|
var lastPull, lastPush *time.Time
|
|
if sr.LastPull != "" {
|
|
if t, err := time.Parse(time.RFC3339, sr.LastPull); err == nil {
|
|
lastPull = &t
|
|
}
|
|
}
|
|
if sr.LastPush != "" {
|
|
if t, err := time.Parse(time.RFC3339, sr.LastPush); err == nil {
|
|
lastPush = &t
|
|
}
|
|
}
|
|
|
|
b.processor.statsCache.Update(holdDID, sr.OwnerDID, sr.Repository,
|
|
sr.PullCount, sr.PushCount, lastPull, lastPush)
|
|
touched[key{sr.OwnerDID, sr.Repository}] = struct{}{}
|
|
ownerDIDs[sr.OwnerDID] = struct{}{}
|
|
}
|
|
|
|
for owner := range ownerDIDs {
|
|
if err := b.processor.EnsureUserExists(ctx, owner); err != nil {
|
|
slog.Warn("Backfill failed to ensure stats owner user", "owner_did", owner, "error", err)
|
|
}
|
|
}
|
|
|
|
// Build aggregated rows from the cache.
|
|
rows := make([]db.RepositoryStats, 0, len(touched))
|
|
for k := range touched {
|
|
totalPull, totalPush, latestPull, latestPush := b.processor.statsCache.GetAggregated(k.owner, k.repo)
|
|
rows = append(rows, db.RepositoryStats{
|
|
DID: k.owner,
|
|
Repository: k.repo,
|
|
PullCount: int(totalPull),
|
|
PushCount: int(totalPush),
|
|
LastPull: latestPull,
|
|
LastPush: latestPush,
|
|
})
|
|
}
|
|
if err := db.BatchUpsertRepositoryStats(b.db, rows); err != nil {
|
|
return 0, err
|
|
}
|
|
slog.Info("Backfill batch stats", "hold_did", holdDID, "rows", len(rows))
|
|
return len(rows), nil
|
|
}
|
|
|
|
// batchCaptains decodes captain records and writes them in one upsert.
|
|
// Records whose publishing DID does not advertise an atcr_hold service in its
|
|
// DID document are skipped — only real holds may enter the discovery cache.
|
|
func (b *BackfillWorker) batchCaptains(ctx context.Context, holdDID string, records []atproto.Record) (int, error) {
|
|
captains := make([]db.HoldCaptainRecord, 0, len(records))
|
|
now := time.Now()
|
|
verified := make(map[string]bool)
|
|
for i := range records {
|
|
r := &records[i]
|
|
var cr atproto.CaptainRecord
|
|
if err := json.Unmarshal(r.Value, &cr); err != nil {
|
|
slog.Warn("Backfill skipping invalid captain record", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
if cr.Owner == "" || !strings.HasPrefix(cr.Owner, "did:") {
|
|
slog.Warn("Backfill skipping captain with invalid owner", "uri", r.URI)
|
|
continue
|
|
}
|
|
// Captain rkey is the hold DID (collections are stored on each hold's PDS,
|
|
// so record.URI already encodes the hold DID in the authority segment).
|
|
recordHoldDID := extractDIDFromURI(r.URI)
|
|
if recordHoldDID == "" {
|
|
recordHoldDID = holdDID
|
|
}
|
|
isHold, ok := verified[recordHoldDID]
|
|
if !ok {
|
|
var err error
|
|
isHold, err = atproto.HasHoldService(ctx, recordHoldDID)
|
|
if err != nil {
|
|
slog.Warn("Backfill skipping captain, hold DID unresolvable", "uri", r.URI, "error", err)
|
|
isHold = false
|
|
} else if !isHold {
|
|
slog.Info("Backfill skipping captain from non-hold DID", "hold_did", recordHoldDID)
|
|
}
|
|
verified[recordHoldDID] = isHold
|
|
}
|
|
if !isHold {
|
|
continue
|
|
}
|
|
captains = append(captains, db.HoldCaptainRecord{
|
|
HoldDID: recordHoldDID,
|
|
OwnerDID: cr.Owner,
|
|
Public: cr.Public,
|
|
AllowAllCrew: cr.AllowAllCrew,
|
|
DeployedAt: cr.DeployedAt,
|
|
Region: cr.Region,
|
|
Successor: cr.Successor,
|
|
UpdatedAt: now,
|
|
})
|
|
}
|
|
if err := db.BatchUpsertCaptainRecords(b.db, captains); err != nil {
|
|
return 0, err
|
|
}
|
|
slog.Info("Backfill batch captains", "rows", len(captains))
|
|
return len(captains), nil
|
|
}
|
|
|
|
// batchCrew decodes crew records and writes them in one upsert.
|
|
func (b *BackfillWorker) batchCrew(holdDID string, records []atproto.Record) (int, error) {
|
|
members := make([]db.CrewMember, 0, len(records))
|
|
for i := range records {
|
|
r := &records[i]
|
|
var cr atproto.CrewRecord
|
|
if err := json.Unmarshal(r.Value, &cr); err != nil {
|
|
slog.Warn("Backfill skipping invalid crew record", "uri", r.URI, "error", err)
|
|
continue
|
|
}
|
|
if cr.Member == "" || !strings.HasPrefix(cr.Member, "did:") {
|
|
slog.Warn("Backfill skipping crew with invalid member", "uri", r.URI)
|
|
continue
|
|
}
|
|
recordHoldDID := extractDIDFromURI(r.URI)
|
|
if recordHoldDID == "" {
|
|
recordHoldDID = holdDID
|
|
}
|
|
permsJSON := ""
|
|
if len(cr.Permissions) > 0 {
|
|
if b, err := json.Marshal(cr.Permissions); err == nil {
|
|
permsJSON = string(b)
|
|
}
|
|
}
|
|
rkey := extractRkeyFromURI(r.URI)
|
|
members = append(members, db.CrewMember{
|
|
HoldDID: recordHoldDID,
|
|
MemberDID: cr.Member,
|
|
Rkey: rkey,
|
|
Role: cr.Role,
|
|
Permissions: permsJSON,
|
|
Tier: cr.Tier,
|
|
AddedAt: cr.AddedAt,
|
|
})
|
|
}
|
|
if err := db.BatchUpsertCrewMembers(b.db, members); err != nil {
|
|
return 0, err
|
|
}
|
|
slog.Info("Backfill batch crew", "hold_did", holdDID, "rows", len(members))
|
|
return len(members), nil
|
|
}
|
|
|
|
// extractDIDFromURI pulls the DID authority segment out of an AT-URI.
|
|
// Format: at://did:…/collection/rkey → "did:…".
|
|
func extractDIDFromURI(uri string) string {
|
|
const prefix = "at://"
|
|
if !strings.HasPrefix(uri, prefix) {
|
|
return ""
|
|
}
|
|
rest := uri[len(prefix):]
|
|
if before, _, ok := strings.Cut(rest, "/"); ok {
|
|
return before
|
|
}
|
|
return rest
|
|
}
|