Files
at-container-registry/pkg/hold/pds/import.go
T
Evan Jarrett e3843db9d8 Implement did:plc support for holds with the ability to import/export CARs.
did:plc Identity Support (pkg/hold/pds/did.go, pkg/hold/config.go, pkg/hold/server.go)

  The big feature — holds can now use did:plc identities instead of only did:web. This adds:
  - LoadOrCreateDID() — resolves hold DID by priority: config DID > did.txt on disk > create new
  - CreatePLCIdentity() — builds a genesis operation, signs with rotation key, submits to PLC directory
  - EnsurePLCCurrent() — on boot, compares local signing key + URL against PLC directory and auto-updates if they've drifted (requires rotation key)
  - New config fields: did_method (web/plc), did, plc_directory_url, rotation_key_path
  - GenerateDIDDocument() now uses the stored DID instead of always deriving did:web from URL
  - NewHoldServer wired up to call LoadOrCreateDID instead of GenerateDIDFromURL

  CAR Export/Import (pkg/hold/pds/export.go, pkg/hold/pds/import.go, cmd/hold/repo.go)

  New CLI subcommands for repo backup/restore:
  - atcr-hold repo export — streams the hold's repo as a CAR file to stdout
  - atcr-hold repo import <file>... — reads CAR files, upserts all records in a single atomic commit. Uses a bulkImportRecords method that opens a delta session, checks each record for
  create vs update, commits once, and fires repo events.
  - openHoldPDS() helper to spin up a HoldPDS from config for offline CLI operations

  Admin UI Fixes (pkg/hold/admin/)

  - Logout changed from GET to POST — nav template now uses a <form method=POST> instead of an <a> link (prevents CSRF on logout)
  - Removed return_to parameter from login flow — simplified redirect logic, auth middleware now redirects to /admin/auth/login without query params

  Config/Deploy

  - config-hold.example.yaml and deploy/upcloud/configs/hold.yaml.tmpl updated with the four new did:plc config fields
  - go.mod / go.sum — added github.com/did-method-plc/go-didplc dependency
2026-02-14 15:17:53 -06:00

181 lines
4.3 KiB
Go

package pds
import (
"context"
"fmt"
"io"
"strings"
"github.com/bluesky-social/indigo/repo"
"github.com/ipfs/go-cid"
"go.opentelemetry.io/otel"
)
// rawCBOR wraps raw bytes to satisfy cbg.CBORMarshaler.
// Used to pass through record bytes from a CAR without decoding.
type rawCBOR []byte
func (r rawCBOR) MarshalCBOR(w io.Writer) error {
_, err := w.Write(r)
return err
}
// bulkRecord holds a single record to import.
type bulkRecord struct {
Collection string
Rkey string
Data rawCBOR
}
// ImportResult summarizes a CAR import operation.
type ImportResult struct {
Total int
PerCollection map[string]int
}
// ImportFromCAR reads a CAR file and imports all records into the hold's repo.
// Records are upserted (overwrite on conflict) in a single atomic commit.
// The repo is initialized if it doesn't exist yet.
func (p *HoldPDS) ImportFromCAR(ctx context.Context, r io.Reader) (*ImportResult, error) {
// Ensure repo exists
head, err := p.carstore.GetUserRepoHead(ctx, p.uid)
if err != nil || !head.Defined() {
if err := p.repomgr.InitNewActor(ctx, p.uid, "", p.did, "", "", ""); err != nil {
return nil, fmt.Errorf("failed to initialize repo: %w", err)
}
}
// Parse the CAR into an in-memory repo
sourceRepo, err := repo.ReadRepoFromCar(ctx, r)
if err != nil {
return nil, fmt.Errorf("failed to read CAR: %w", err)
}
// Collect all records
var records []bulkRecord
err = sourceRepo.ForEach(ctx, "", func(k string, v cid.Cid) error {
_, recBytes, err := sourceRepo.GetRecordBytes(ctx, k)
if err != nil {
return fmt.Errorf("failed to get record bytes for %s: %w", k, err)
}
parts := strings.SplitN(k, "/", 2)
if len(parts) != 2 {
return fmt.Errorf("unexpected record path format: %s", k)
}
records = append(records, bulkRecord{
Collection: parts[0],
Rkey: parts[1],
Data: rawCBOR(*recBytes),
})
return nil
})
if err != nil {
return nil, fmt.Errorf("failed to iterate CAR records: %w", err)
}
if len(records) == 0 {
return &ImportResult{PerCollection: map[string]int{}}, nil
}
// Bulk upsert all records in a single commit
if err := p.bulkImportRecords(ctx, records); err != nil {
return nil, fmt.Errorf("failed to import records: %w", err)
}
// Build result
result := &ImportResult{
Total: len(records),
PerCollection: make(map[string]int),
}
for _, rec := range records {
result.PerCollection[rec.Collection]++
}
return result, nil
}
// bulkImportRecords writes all records in a single delta session + commit.
// Each record is upserted: created if new, updated if exists.
func (p *HoldPDS) bulkImportRecords(ctx context.Context, records []bulkRecord) error {
ctx, span := otel.Tracer("repoman").Start(ctx, "BulkImportRecords")
defer span.End()
unlock := p.repomgr.lockUser(ctx, p.uid)
defer unlock()
rev, err := p.repomgr.cs.GetUserRepoRev(ctx, p.uid)
if err != nil {
return err
}
ds, err := p.repomgr.cs.NewDeltaSession(ctx, p.uid, &rev)
if err != nil {
return err
}
head := ds.BaseCid()
r, err := repo.OpenRepo(ctx, ds, head)
if err != nil {
return err
}
ops := make([]RepoOp, 0, len(records))
for _, rec := range records {
rpath := rec.Collection + "/" + rec.Rkey
// Check if record exists to determine create vs update
_, _, getErr := r.GetRecordBytes(ctx, rpath)
recordExists := getErr == nil
var cc cid.Cid
var evtKind EventKind
if recordExists {
cc, err = r.UpdateRecord(ctx, rpath, rec.Data)
evtKind = EvtKindUpdateRecord
} else {
cc, err = r.PutRecord(ctx, rpath, rec.Data)
evtKind = EvtKindCreateRecord
}
if err != nil {
return fmt.Errorf("failed to write %s: %w", rpath, err)
}
ops = append(ops, RepoOp{
Kind: evtKind,
Collection: rec.Collection,
Rkey: rec.Rkey,
RecCid: &cc,
})
}
nroot, nrev, err := r.Commit(ctx, p.repomgr.kmgr.SignForUser)
if err != nil {
return err
}
rslice, err := ds.CloseWithRoot(ctx, nroot, nrev)
if err != nil {
return fmt.Errorf("close with root: %w", err)
}
var oldroot *cid.Cid
if head.Defined() {
oldroot = &head
}
if p.repomgr.events != nil {
p.repomgr.events(ctx, &RepoEvent{
User: p.uid,
OldRoot: oldroot,
NewRoot: nroot,
Rev: nrev,
Since: &rev,
Ops: ops,
RepoSlice: rslice,
})
}
return nil
}