Files
at-container-registry/pkg/appview/jetstream/backfill.go
T

658 lines
22 KiB
Go

package jetstream
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"strings"
"time"
"github.com/bluesky-social/indigo/atproto/identity"
"github.com/bluesky-social/indigo/atproto/syntax"
"atcr.io/pkg/appview/db"
"atcr.io/pkg/atproto"
)
// BackfillWorker uses com.atproto.sync.listReposByCollection to backfill historical data
type BackfillWorker struct {
db *sql.DB
client *atproto.Client
directory identity.Directory
defaultHoldDID string // Default hold DID from AppView config (e.g., "did:web:hold01.atcr.io")
testMode bool // If true, suppress warnings for external holds
}
// BackfillState tracks backfill progress
type BackfillState struct {
Collection string
RepoCursor string // Cursor for listReposByCollection
CurrentDID string // Current DID being processed
RecordCursor string // Cursor for listRecords within current DID
ProcessedRepos int
ProcessedRecords int
Completed bool
}
// NewBackfillWorker creates a backfill worker using sync API
// defaultHoldDID should be in format "did:web:hold01.atcr.io"
// To find a hold's DID, visit: https://hold-url/.well-known/did.json
func NewBackfillWorker(database *sql.DB, relayEndpoint, defaultHoldDID string, testMode bool) (*BackfillWorker, error) {
// Create client for relay - used only for listReposByCollection
client := atproto.NewClient(relayEndpoint, "", "")
return &BackfillWorker{
db: database,
client: client, // This points to the relay
directory: identity.DefaultDirectory(),
defaultHoldDID: defaultHoldDID,
testMode: testMode,
}, nil
}
// Start runs the backfill for all ATCR collections
func (b *BackfillWorker) Start(ctx context.Context) error {
fmt.Println("Backfill: Starting sync-based backfill...")
// First, query and cache the default hold's captain record
if b.defaultHoldDID != "" {
fmt.Printf("Backfill: Querying default hold captain record: %s\n", b.defaultHoldDID)
if err := b.queryCaptainRecord(ctx, b.defaultHoldDID); err != nil {
fmt.Printf("WARNING: Failed to query default hold captain record: %v\n", err)
// Don't fail the whole backfill - just warn
}
}
collections := []string{
atproto.ManifestCollection, // io.atcr.manifest
atproto.TagCollection, // io.atcr.tag
atproto.StarCollection, // io.atcr.sailor.star
atproto.SailorProfileCollection, // io.atcr.sailor.profile
}
for _, collection := range collections {
fmt.Printf("Backfill: Processing collection: %s\n", collection)
if err := b.backfillCollection(ctx, collection); err != nil {
return fmt.Errorf("failed to backfill collection %s: %w", collection, err)
}
fmt.Printf("Backfill: Completed collection: %s\n", collection)
}
fmt.Println("Backfill: All collections completed!")
return nil
}
// backfillCollection backfills a single collection
func (b *BackfillWorker) backfillCollection(ctx context.Context, collection string) error {
var repoCursor string
processedRepos := 0
processedRecords := 0
// Paginate through all repos with this collection
for {
// List repos that have records in this collection
result, err := b.client.ListReposByCollection(ctx, collection, 1000, repoCursor)
if err != nil {
return fmt.Errorf("failed to list repos: %w", err)
}
fmt.Printf("Backfill: Found %d repos with %s (cursor: %s)\n", len(result.Repos), collection, repoCursor)
// Process each repo (DID)
for _, repo := range result.Repos {
recordCount, err := b.backfillRepo(ctx, repo.DID, collection)
if err != nil {
fmt.Printf("WARNING: Failed to backfill repo %s: %v\n", repo.DID, err)
continue
}
processedRepos++
processedRecords += recordCount
if processedRepos%10 == 0 {
fmt.Printf("Backfill: Progress - %d repos, %d records\n", processedRepos, processedRecords)
}
}
// Check if there are more pages
if result.Cursor == "" {
break
}
repoCursor = result.Cursor
}
fmt.Printf("Backfill: Collection %s complete - %d repos, %d records\n", collection, processedRepos, processedRecords)
return nil
}
// backfillRepo backfills all records for a single repo/DID
func (b *BackfillWorker) backfillRepo(ctx context.Context, did, collection string) (int, error) {
// Ensure user exists in database and get their PDS endpoint
if err := b.ensureUser(ctx, did); err != nil {
return 0, fmt.Errorf("failed to ensure user: %w", err)
}
// Resolve DID to get user's PDS endpoint
didParsed, err := syntax.ParseDID(did)
if err != nil {
return 0, fmt.Errorf("invalid DID %s: %w", did, err)
}
ident, err := b.directory.LookupDID(ctx, didParsed)
if err != nil {
return 0, fmt.Errorf("failed to resolve DID to PDS: %w", err)
}
pdsEndpoint := ident.PDSEndpoint()
if pdsEndpoint == "" {
return 0, fmt.Errorf("no PDS endpoint found for DID %s", did)
}
// Create a client for this user's PDS
pdsClient := atproto.NewClient(pdsEndpoint, "", "")
var recordCursor string
recordCount := 0
// Track which records exist on the PDS for reconciliation
var foundManifestDigests []string
var foundTags []struct{ Repository, Tag string }
foundStars := make(map[string]time.Time) // key: "ownerDID/repository", value: createdAt
// Paginate through all records for this repo
for {
records, cursor, err := pdsClient.ListRecordsForRepo(ctx, did, collection, 100, recordCursor)
if err != nil {
return recordCount, fmt.Errorf("failed to list records: %w", err)
}
// Process each record
for _, record := range records {
// Track what we found for deletion reconciliation
if collection == atproto.ManifestCollection {
var manifestRecord atproto.ManifestRecord
if err := json.Unmarshal(record.Value, &manifestRecord); err == nil {
foundManifestDigests = append(foundManifestDigests, manifestRecord.Digest)
}
} else if collection == atproto.TagCollection {
var tagRecord atproto.TagRecord
if err := json.Unmarshal(record.Value, &tagRecord); err == nil {
foundTags = append(foundTags, struct{ Repository, Tag string }{
Repository: tagRecord.Repository,
Tag: tagRecord.Tag,
})
}
} else if collection == atproto.StarCollection {
var starRecord atproto.StarRecord
if err := json.Unmarshal(record.Value, &starRecord); err == nil {
key := fmt.Sprintf("%s/%s", starRecord.Subject.DID, starRecord.Subject.Repository)
foundStars[key] = starRecord.CreatedAt
}
}
if err := b.processRecord(ctx, did, collection, &record); err != nil {
fmt.Printf("WARNING: Failed to process record %s: %v\n", record.URI, err)
continue
}
recordCount++
}
// Check if there are more pages
if cursor == "" {
break
}
recordCursor = cursor
}
// Reconcile deletions - remove records from DB that no longer exist on PDS
if err := b.reconcileDeletions(did, collection, foundManifestDigests, foundTags, foundStars); err != nil {
fmt.Printf("WARNING: Failed to reconcile deletions for %s: %v\n", did, err)
}
// After processing manifests, clean up orphaned tags (tags pointing to non-existent manifests)
if collection == atproto.ManifestCollection {
if err := db.CleanupOrphanedTags(b.db, did); err != nil {
fmt.Printf("WARNING: Failed to cleanup orphaned tags for %s: %v\n", did, err)
}
}
return recordCount, nil
}
// reconcileDeletions removes records from the database that no longer exist on the PDS
func (b *BackfillWorker) reconcileDeletions(did, collection string, foundManifestDigests []string, foundTags []struct{ Repository, Tag string }, foundStars map[string]time.Time) error {
switch collection {
case atproto.ManifestCollection:
// Get current manifests in DB
dbDigests, err := db.GetManifestDigestsForDID(b.db, did)
if err != nil {
return fmt.Errorf("failed to get DB manifests: %w", err)
}
// Delete manifests not found on PDS
if err := db.DeleteManifestsNotInList(b.db, did, foundManifestDigests); err != nil {
return fmt.Errorf("failed to delete orphaned manifests: %w", err)
}
// Log deletions
deleted := len(dbDigests) - len(foundManifestDigests)
if deleted > 0 {
fmt.Printf("Backfill: Deleted %d orphaned manifests for %s\n", deleted, did)
}
case atproto.TagCollection:
// Get current tags in DB
dbTags, err := db.GetTagsForDID(b.db, did)
if err != nil {
return fmt.Errorf("failed to get DB tags: %w", err)
}
// Delete tags not found on PDS
if err := db.DeleteTagsNotInList(b.db, did, foundTags); err != nil {
return fmt.Errorf("failed to delete orphaned tags: %w", err)
}
// Log deletions
deleted := len(dbTags) - len(foundTags)
if deleted > 0 {
fmt.Printf("Backfill: Deleted %d orphaned tags for %s\n", deleted, did)
}
case atproto.StarCollection:
// Reconcile stars - delete stars that no longer exist on PDS
// Star counts will be calculated on demand from the stars table
if err := db.DeleteStarsNotInList(b.db, did, foundStars); err != nil {
return fmt.Errorf("failed to delete orphaned stars: %w", err)
}
}
return nil
}
// processRecord processes a single record and stores it in the database
func (b *BackfillWorker) processRecord(ctx context.Context, did, collection string, record *atproto.Record) error {
switch collection {
case atproto.ManifestCollection:
return b.processManifestRecord(did, record)
case atproto.TagCollection:
return b.processTagRecord(did, record)
case atproto.StarCollection:
return b.processStarRecord(did, record)
case atproto.SailorProfileCollection:
return b.processSailorProfileRecord(ctx, did, record)
default:
return fmt.Errorf("unsupported collection: %s", collection)
}
}
// processManifestRecord processes a manifest record
func (b *BackfillWorker) processManifestRecord(did string, record *atproto.Record) error {
var manifestRecord atproto.ManifestRecord
if err := json.Unmarshal(record.Value, &manifestRecord); err != nil {
return fmt.Errorf("failed to unmarshal manifest: %w", err)
}
// Extract OCI annotations from manifest
var title, description, sourceURL, documentationURL, licenses, iconURL, readmeURL string
if manifestRecord.Annotations != nil {
title = manifestRecord.Annotations["org.opencontainers.image.title"]
description = manifestRecord.Annotations["org.opencontainers.image.description"]
sourceURL = manifestRecord.Annotations["org.opencontainers.image.source"]
documentationURL = manifestRecord.Annotations["org.opencontainers.image.documentation"]
licenses = manifestRecord.Annotations["org.opencontainers.image.licenses"]
iconURL = manifestRecord.Annotations["io.atcr.icon"]
readmeURL = manifestRecord.Annotations["io.atcr.readme"]
}
// Detect manifest type
isManifestList := len(manifestRecord.Manifests) > 0
// Prepare manifest for insertion
manifest := &db.Manifest{
DID: did,
Repository: manifestRecord.Repository,
Digest: manifestRecord.Digest,
MediaType: manifestRecord.MediaType,
SchemaVersion: manifestRecord.SchemaVersion,
HoldEndpoint: manifestRecord.HoldEndpoint,
CreatedAt: manifestRecord.CreatedAt,
Title: title,
Description: description,
SourceURL: sourceURL,
DocumentationURL: documentationURL,
Licenses: licenses,
IconURL: iconURL,
ReadmeURL: readmeURL,
}
// Set config fields only for image manifests (not manifest lists)
if !isManifestList && manifestRecord.Config != nil {
manifest.ConfigDigest = manifestRecord.Config.Digest
manifest.ConfigSize = manifestRecord.Config.Size
}
// Insert manifest
manifestID, err := db.InsertManifest(b.db, manifest)
if err != nil {
// Skip if already exists
if strings.Contains(err.Error(), "UNIQUE constraint failed") {
return nil
}
return fmt.Errorf("failed to insert manifest: %w", err)
}
if isManifestList {
// Insert manifest references (for manifest lists/indexes)
for i, ref := range manifestRecord.Manifests {
platformArch := ""
platformOS := ""
platformVariant := ""
platformOSVersion := ""
if ref.Platform != nil {
platformArch = ref.Platform.Architecture
platformOS = ref.Platform.OS
platformVariant = ref.Platform.Variant
platformOSVersion = ref.Platform.OSVersion
}
if err := db.InsertManifestReference(b.db, &db.ManifestReference{
ManifestID: manifestID,
Digest: ref.Digest,
MediaType: ref.MediaType,
Size: ref.Size,
PlatformArchitecture: platformArch,
PlatformOS: platformOS,
PlatformVariant: platformVariant,
PlatformOSVersion: platformOSVersion,
ReferenceIndex: i,
}); err != nil {
// Continue on error - reference might already exist
continue
}
}
} else {
// Insert layers (for image manifests)
for i, layer := range manifestRecord.Layers {
if err := db.InsertLayer(b.db, &db.Layer{
ManifestID: manifestID,
Digest: layer.Digest,
MediaType: layer.MediaType,
Size: layer.Size,
LayerIndex: i,
}); err != nil {
// Continue on error - layer might already exist
continue
}
}
}
return nil
}
// processTagRecord processes a tag record
func (b *BackfillWorker) processTagRecord(did string, record *atproto.Record) error {
var tagRecord atproto.TagRecord
if err := json.Unmarshal(record.Value, &tagRecord); err != nil {
return fmt.Errorf("failed to unmarshal tag: %w", err)
}
// Extract digest from tag record (tries manifest field first, falls back to manifestDigest)
manifestDigest, err := tagRecord.GetManifestDigest()
if err != nil {
return fmt.Errorf("failed to get manifest digest from tag record: %w", err)
}
// Insert or update tag
return db.UpsertTag(b.db, &db.Tag{
DID: did,
Repository: tagRecord.Repository,
Tag: tagRecord.Tag,
Digest: manifestDigest,
CreatedAt: tagRecord.UpdatedAt,
})
}
// processStarRecord processes a star record
func (b *BackfillWorker) processStarRecord(did string, record *atproto.Record) error {
var starRecord atproto.StarRecord
if err := json.Unmarshal(record.Value, &starRecord); err != nil {
return fmt.Errorf("failed to unmarshal star: %w", err)
}
// Upsert the star record (idempotent - won't duplicate)
// The DID here is the starrer (user who starred)
// The subject contains the owner DID and repository
// Star count will be calculated on demand from the stars table
return db.UpsertStar(b.db, did, starRecord.Subject.DID, starRecord.Subject.Repository, starRecord.CreatedAt)
}
// processSailorProfileRecord processes a sailor profile record
// Extracts defaultHold and queries the hold's captain record to cache it
func (b *BackfillWorker) processSailorProfileRecord(ctx context.Context, did string, record *atproto.Record) error {
var profileRecord atproto.SailorProfileRecord
if err := json.Unmarshal(record.Value, &profileRecord); err != nil {
return fmt.Errorf("failed to unmarshal sailor profile: %w", err)
}
// Skip if no default hold set
if profileRecord.DefaultHold == "" {
return nil
}
// Convert hold URL/DID to canonical DID
holdDID := atproto.ResolveHoldDIDFromURL(profileRecord.DefaultHold)
if holdDID == "" {
fmt.Printf("WARNING [backfill]: Invalid hold reference in profile for %s: %s\n", did, profileRecord.DefaultHold)
return nil
}
// Query and cache the captain record
if err := b.queryCaptainRecord(ctx, holdDID); err != nil {
// In test mode, only warn about default hold (local hold)
// External/production holds may not have captain records yet (dev ahead of prod)
if b.testMode && holdDID != b.defaultHoldDID {
// Suppress warning for external holds in test mode
return nil
}
fmt.Printf("WARNING [backfill]: Failed to query captain record for hold %s: %v\n", holdDID, err)
// Don't fail the whole backfill - just skip this hold
return nil
}
return nil
}
// queryCaptainRecord queries a hold's captain record and caches it in the database
func (b *BackfillWorker) queryCaptainRecord(ctx context.Context, holdDID string) error {
// Check if we already have it cached (skip if recently updated)
existing, err := db.GetCaptainRecord(b.db, holdDID)
if err == nil && existing != nil {
// If cached within last hour, skip refresh
if time.Since(existing.UpdatedAt) < 1*time.Hour {
return nil
}
}
// Resolve hold DID to URL
// For did:web, we need to fetch .well-known/did.json
holdURL, err := resolveHoldDIDToURL(ctx, holdDID)
if err != nil {
return fmt.Errorf("failed to resolve hold DID to URL: %w", err)
}
// Create client for hold's PDS
holdClient := atproto.NewClient(holdURL, holdDID, "")
// Query captain record with retries (for Docker startup timing)
var record *atproto.Record
maxRetries := 3
for attempt := 1; attempt <= maxRetries; attempt++ {
record, err = holdClient.GetRecord(ctx, "io.atcr.hold.captain", "self")
if err == nil {
break
}
// Retry on connection errors (hold service might still be starting)
if attempt < maxRetries && strings.Contains(err.Error(), "connection refused") {
fmt.Printf("Backfill: Hold not ready (attempt %d/%d), retrying in 2s...\n", attempt, maxRetries)
time.Sleep(2 * time.Second)
continue
}
return fmt.Errorf("failed to get captain record: %w", err)
}
// Parse captain record from the record's Value field
var captainRecord struct {
Owner string `json:"owner"`
Public bool `json:"public"`
AllowAllCrew bool `json:"allowAllCrew"`
DeployedAt string `json:"deployedAt"`
Region string `json:"region"`
Provider string `json:"provider"`
}
if err := json.Unmarshal(record.Value, &captainRecord); err != nil {
return fmt.Errorf("failed to parse captain record: %w", err)
}
// Cache in database
dbRecord := &db.HoldCaptainRecord{
HoldDID: holdDID,
OwnerDID: captainRecord.Owner,
Public: captainRecord.Public,
AllowAllCrew: captainRecord.AllowAllCrew,
DeployedAt: captainRecord.DeployedAt,
Region: captainRecord.Region,
Provider: captainRecord.Provider,
UpdatedAt: time.Now(),
}
if err := db.UpsertCaptainRecord(b.db, dbRecord); err != nil {
return fmt.Errorf("failed to cache captain record: %w", err)
}
fmt.Printf("Backfill: Cached captain record for hold %s (owner: %s)\n", holdDID, captainRecord.Owner)
return nil
}
// resolveHoldDIDToURL resolves a hold DID to its service endpoint URL
// Fetches the DID document and returns both the canonical DID and service endpoint
func resolveHoldDIDToURL(ctx context.Context, inputDID string) (string, error) {
// For did:web, construct the .well-known URL
if !strings.HasPrefix(inputDID, "did:web:") {
return "", fmt.Errorf("only did:web is supported, got: %s", inputDID)
}
// Extract hostname from did:web:hostname[:port]
hostname := strings.TrimPrefix(inputDID, "did:web:")
// Try HTTP first (for local Docker), then HTTPS
var serviceEndpoint string
for _, scheme := range []string{"http", "https"} {
testURL := fmt.Sprintf("%s://%s/.well-known/did.json", scheme, hostname)
// Fetch DID document (use NewClient to initialize httpClient)
client := atproto.NewClient("", "", "")
didDoc, err := client.FetchDIDDocument(ctx, testURL)
if err == nil && didDoc != nil {
// Extract service endpoint from DID document
for _, service := range didDoc.Service {
if service.Type == "AtprotoPersonalDataServer" || service.Type == "AtcrHoldService" {
serviceEndpoint = service.ServiceEndpoint
break
}
}
if serviceEndpoint != "" {
fmt.Printf("DEBUG [backfill]: Resolved %s → canonical DID: %s, endpoint: %s\n",
inputDID, didDoc.ID, serviceEndpoint)
return serviceEndpoint, nil
}
}
}
// Fallback: assume the hold service is at the root of the hostname
// Try HTTP first for local development
url := fmt.Sprintf("http://%s", hostname)
fmt.Printf("WARNING [backfill]: Failed to fetch DID document for %s, using fallback URL: %s\n", inputDID, url)
return url, nil
}
// ensureUser resolves and upserts a user by DID
func (b *BackfillWorker) ensureUser(ctx context.Context, did string) error {
// Check if user already exists
existingUser, err := db.GetUserByDID(b.db, did)
if err == nil && existingUser != nil {
// Update last seen
existingUser.LastSeen = time.Now()
return db.UpsertUser(b.db, existingUser)
}
// Resolve DID to get handle and PDS endpoint
didParsed, err := syntax.ParseDID(did)
if err != nil {
// Fallback: use DID as handle
user := &db.User{
DID: did,
Handle: did,
PDSEndpoint: "https://bsky.social",
LastSeen: time.Now(),
}
return db.UpsertUser(b.db, user)
}
ident, err := b.directory.LookupDID(ctx, didParsed)
if err != nil {
// Fallback: use DID as handle
user := &db.User{
DID: did,
Handle: did,
PDSEndpoint: "https://bsky.social",
LastSeen: time.Now(),
}
return db.UpsertUser(b.db, user)
}
resolvedDID := ident.DID.String()
handle := ident.Handle.String()
pdsEndpoint := ident.PDSEndpoint()
// If handle is invalid or PDS is missing, use defaults
if handle == "handle.invalid" || handle == "" {
handle = resolvedDID
}
if pdsEndpoint == "" {
pdsEndpoint = "https://bsky.social"
}
// Fetch user's Bluesky profile (including avatar)
// Use public Bluesky AppView API (doesn't require auth for public profiles)
avatar := ""
publicClient := atproto.NewClient("https://public.api.bsky.app", "", "")
profile, err := publicClient.GetActorProfile(ctx, resolvedDID)
if err != nil {
fmt.Printf("WARNING [backfill]: Failed to fetch profile for DID %s: %v\n", resolvedDID, err)
// Continue without avatar
} else {
avatar = profile.Avatar
}
// Upsert to database
user := &db.User{
DID: resolvedDID,
Handle: handle,
PDSEndpoint: pdsEndpoint,
Avatar: avatar,
LastSeen: time.Now(),
}
return db.UpsertUser(b.db, user)
}