From e3ac9cd545ba134c2840b00f4455e76e17fe313a Mon Sep 17 00:00:00 2001 From: Lemon-miaow Date: Fri, 25 Sep 2026 04:33:23 +0800 Subject: [PATCH] =?UTF-8?q?feat(offsite):=20=E5=BC=82=E5=9C=B0=E5=89=AF?= =?UTF-8?q?=E6=9C=AC=E5=8A=A0=E5=85=A5=20registry=20=E7=94=A8=E6=88=B7?= =?UTF-8?q?=E9=95=9C=E5=83=8F=EF=BC=8C=E6=8C=89=E6=91=98=E8=A6=81=E5=8A=A0?= =?UTF-8?q?=E5=AF=86=E5=8E=BB=E9=87=8D=EF=BC=8C=E7=B4=A2=E5=BC=95=E6=8C=89?= =?UTF-8?q?=E7=89=88=E6=9C=AC=E4=BF=9D=E7=95=99=2014=20=E5=A4=A9=EF=BC=8C?= =?UTF-8?q?=E6=96=B0=E5=A2=9E=20fetch-images=20=E5=9B=9E=E6=8E=A8=E6=81=A2?= =?UTF-8?q?=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- cmd/felis/mirrortools.go | 14 +- cmd/felis/offsite.go | 50 +- cmd/felis/offsite_images.go | 167 +++++++ cmd/felis/offsite_test.go | 26 + cmd/felis/run.go | 2 +- deploy/bootstrap.sh | 2 +- docs/troubleshooting.md | 57 ++- internal/imagepush/mirror.go | 11 + internal/imagepush/push.go | 20 + internal/offsite/images.go | 815 ++++++++++++++++++++++++++++++++ internal/offsite/images_test.go | 454 ++++++++++++++++++ internal/offsite/sync.go | 55 ++- internal/watchdog/probes.go | 4 +- 13 files changed, 1626 insertions(+), 51 deletions(-) create mode 100644 cmd/felis/offsite_images.go create mode 100644 internal/offsite/images.go create mode 100644 internal/offsite/images_test.go diff --git a/cmd/felis/mirrortools.go b/cmd/felis/mirrortools.go index fb01046..571356b 100644 --- a/cmd/felis/mirrortools.go +++ b/cmd/felis/mirrortools.go @@ -14,7 +14,6 @@ import ( "felis.lolicon.best/internal/build" "felis.lolicon.best/internal/imagepush" - "felis.lolicon.best/internal/registrygate" ) // defaultBuildToolsStatus is where mirror-build-tools records its last run; the @@ -49,16 +48,9 @@ func cmdMirrorBuildTools(args []string, stdout, stderr io.Writer) int { fmt.Fprintf(stderr, "felis mirror-build-tools: %v\n", err) return 2 } - if err := loadEnvFile(*secrets); err != nil { - fmt.Fprintf(stderr, "felis mirror-build-tools: read %s: %v\n", *secrets, err) - return 1 - } - user, pass := os.Getenv("FELIS_REGISTRY_USERNAME"), os.Getenv("FELIS_REGISTRY_PASSWORD") - if pass == "" { - user, pass = registrygate.PrincipalPlatform, os.Getenv("REGISTRY_PLATFORM_TOKEN") - } - if pass == "" { - fmt.Fprintln(stderr, "felis mirror-build-tools: no registry credential: set FELIS_REGISTRY_PASSWORD or run as root on the node (REGISTRY_PLATFORM_TOKEN in /etc/felis/secrets.env)") + user, pass, err := registryWriteCredential(*secrets) + if err != nil { + fmt.Fprintf(stderr, "felis mirror-build-tools: %v\n", err) return 2 } diff --git a/cmd/felis/offsite.go b/cmd/felis/offsite.go index b852c1a..4702f26 100644 --- a/cmd/felis/offsite.go +++ b/cmd/felis/offsite.go @@ -24,12 +24,14 @@ import ( ) const offsiteUsage = `usage: - felis offsite sync [-config path] [-archive-dir dir] [-db-dir dir] [-status-file path] + felis offsite sync [-config path] [-archive-dir dir] [-db-dir dir] [-registry host:port|off] + [-status-file path] felis offsite status [-config path] [-status-file path] felis offsite list [-config path] felis offsite fetch-db [-config path | -endpoint url -bucket name [-region r] [-prefix p]] [-dir dir] latest| felis offsite fetch-worlds [-config path] [-archive-dir dir] + felis offsite fetch-images [-config path] [-registry host:port] [-at version] felis offsite keygen Every verb but keygen reads the bucket credentials and the encryption key from @@ -43,7 +45,8 @@ FELIS_OFFSITE_SECRET_KEY, FELIS_OFFSITE_KEY), taking any that are unset from const defaultOffsiteEnvFile = "/etc/felis/offsite.env" // cmdOffsite implements `felis offsite`: the off-site copy of the world -// archives and the database bundles (internal/offsite). felis-offsite.timer +// archives, the database bundles and the registry's user images +// (internal/offsite). felis-offsite.timer // runs `sync` hourly on the host; the fetch verbs are the way back after the // node is lost (docs/troubleshooting.md §16). func cmdOffsite(args []string, stdout, stderr io.Writer) int { @@ -66,6 +69,8 @@ func cmdOffsite(args []string, stdout, stderr io.Writer) int { return offsiteFetchDB(fs, rest, stdout, stderr) case "fetch-worlds": return offsiteFetchWorlds(fs, rest, stdout, stderr) + case "fetch-images": + return offsiteFetchImages(fs, rest, stdout, stderr) case "keygen": k, err := offsite.NewKey() if err != nil { @@ -186,6 +191,7 @@ func offsiteSync(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int archiveDir := fs.String("archive-dir", "", "host directory of the world archive volume (default: resolved from the backup PVC through the cluster)") backupPVC := fs.String("backup-pvc", "felis-backups", `the world archive PVC, in the [k8s] namespace ("" when backups are off)`) dbDir := fs.String("db-dir", dbbackup.DefaultDir, `database bundle directory ("" copies no bundles)`) + registry := fs.String("registry", "", `host[:port] of the registry whose user images are copied (default: the in-cluster registry's loopback hostPort; "off" copies none)`) statusFile := fs.String("status-file", offsite.DefaultStatusFile, "where the result of this run is recorded for the watchdog and `status`") if err := fs.Parse(args); err != nil { return 2 @@ -202,7 +208,7 @@ func offsiteSync(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int if prev, _ := offsite.ReadStatus(*statusFile); prev != nil { st.LastSuccess = prev.LastSuccess } - res, err := runOffsiteSync(cfg, env, *archiveDir, *backupPVC, *dbDir, stderr) + res, err := runOffsiteSync(cfg, env, *archiveDir, *backupPVC, *dbDir, offsiteRegistryEndpoint(*registry, cfg.Registry), stderr) st.Result = res if err != nil { st.LastError = err.Error() @@ -212,12 +218,16 @@ func offsiteSync(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int if werr := offsite.WriteStatus(*statusFile, st); werr != nil { fmt.Fprintf(stderr, "felis offsite sync: record status: %v\n", werr) } - fmt.Fprintf(stdout, "felis offsite sync: worlds copied=%d pending=%d missing=%d expired=%d; bundles copied=%d pruned=%d; bucket holds %d worlds (%s) and %d bundles\n", + fmt.Fprintf(stdout, "felis offsite sync: worlds copied=%d pending=%d missing=%d expired=%d; bundles copied=%d pruned=%d; images copied=%d blobs=%d pruned=%d; bucket holds %d worlds (%s), %d bundles, %d images in %d repositories (%s)\n", res.WorldsUploaded, res.WorldsPending, len(res.WorldsMissing), res.WorldsExpired, - res.DBUploaded, res.DBPruned, res.RemoteWorlds, offsite.HumanBytes(res.RemoteBytes), res.RemoteDB) + res.DBUploaded, res.DBPruned, res.ImagesUploaded, res.ImageBlobsUploaded, res.ImageObjectsPruned, + res.RemoteWorlds, offsite.HumanBytes(res.RemoteBytes), res.RemoteDB, res.Images, res.ImageRepos, offsite.HumanBytes(res.RemoteImageBytes)) for _, m := range res.WorldsMissing { fmt.Fprintf(stderr, "felis offsite sync: recorded archive not on the volume, nothing to copy: %s\n", m) } + for _, m := range res.ImagesIncomplete { + fmt.Fprintf(stderr, "felis offsite sync: registry image not whole: %s\n", m) + } if err != nil { fmt.Fprintf(stderr, "felis offsite sync: %v\n", err) return 1 @@ -225,7 +235,7 @@ func offsiteSync(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int return 0 } -func runOffsiteSync(cfg *config.Config, env *offsiteEnv, archiveDir, backupPVC, dbDir string, log io.Writer) (offsite.Result, error) { +func runOffsiteSync(cfg *config.Config, env *offsiteEnv, archiveDir, backupPVC, dbDir, registry string, log io.Writer) (offsite.Result, error) { ctx, cancel := context.WithTimeout(context.Background(), 50*time.Minute) defer cancel() checkCtx, checkCancel := context.WithTimeout(ctx, 30*time.Second) @@ -250,6 +260,9 @@ func runOffsiteSync(cfg *config.Config, env *offsiteEnv, archiveDir, backupPVC, Bucket: env.bucket, Catalog: offsite.PGCatalog{DB: drv.DB()}, Key: env.key, ArchiveDir: archiveDir, DBDir: dbDir, DBKeep: env.cfg.DBKeep, Log: log, } + if registry != "" { + s.Images = newRegistryImages(registry) + } return s.Run(ctx) } @@ -347,7 +360,7 @@ func offsiteStatus(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) in return 1 } if !cfg.Offsite.Enabled() { - fmt.Fprintln(stdout, "off-site copy: not configured. World archives and database bundles exist on this machine only.") + fmt.Fprintln(stdout, "off-site copy: not configured. World archives, database bundles and user images exist on this machine only.") fmt.Fprintln(stdout, "See docs/troubleshooting.md §16, \"Keep a copy somewhere else\".") return 1 } @@ -380,10 +393,19 @@ func offsiteStatus(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) in r := st.Result fmt.Fprintf(stdout, "bucket holds: %d world archives (%s), %d database bundles, newest %s\n", r.RemoteWorlds, offsite.HumanBytes(r.RemoteBytes), r.RemoteDB, orNone(r.NewestDB)) + if r.ImageIndex != "" { + fmt.Fprintf(stdout, "images: %d in %d repositories (%s), registry index %s\n", + r.Images, r.ImageRepos, offsite.HumanBytes(r.RemoteImageBytes), r.ImageIndex) + } else { + fmt.Fprintln(stdout, "images: not copied (no in-cluster registry, or no sync has reached it yet)") + } fmt.Fprintf(stdout, "waiting: %d world archives not yet copied\n", r.WorldsPending) for _, m := range r.WorldsMissing { fmt.Fprintf(stdout, "missing: %s is recorded but not on the volume\n", m) } + for _, m := range r.ImagesIncomplete { + fmt.Fprintf(stdout, "not whole: %s\n", m) + } if st.LastSuccess.IsZero() || now.Sub(st.LastSuccess) > offsite.StaleAfter { fmt.Fprintf(stdout, "\nThe last successful sync is older than %s: journalctl -u felis-offsite -n 50\n", dbbackup.Age(offsite.StaleAfter)) return 1 @@ -434,6 +456,20 @@ func printOffsiteList(env *offsiteEnv, stdout, stderr io.Writer) int { total += w.Size } fmt.Fprintf(stdout, "world archives: %d (%s)\n", len(worlds), offsite.HumanBytes(total)) + versions, err := offsite.ImageIndexes(ctx, env.bucket) + if err != nil { + fmt.Fprintf(stderr, "felis offsite list: %v\n", err) + return 1 + } + fmt.Fprintf(stdout, "registry index versions (%d, newest first; restore one with fetch-images -at):\n", len(versions)) + for i := len(versions) - 1; i >= 0; i-- { + x, err := offsite.LoadImageIndex(ctx, env.bucket, env.key, versions[i]) + if err != nil { + fmt.Fprintf(stdout, " %s unreadable: %v\n", versions[i], err) + continue + } + fmt.Fprintf(stdout, " %s %d images in %d repositories\n", versions[i], x.Images(), len(x.Repositories)) + } return 0 } diff --git a/cmd/felis/offsite_images.go b/cmd/felis/offsite_images.go new file mode 100644 index 0000000..20cce74 --- /dev/null +++ b/cmd/felis/offsite_images.go @@ -0,0 +1,167 @@ +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "io" + "net/http" + "os" + "strings" + "time" + + "felis.lolicon.best/internal/config" + "felis.lolicon.best/internal/imagepush" + "felis.lolicon.best/internal/offsite" + "felis.lolicon.best/internal/registrygate" + "felis.lolicon.best/internal/registryprune" +) + +// registryImages is the platform registry as the off-site copy sees it: read +// anonymously through the gate (the catalog, its manifest index, manifests and +// blobs) and, for a restore, written as the platform principal. Both go through +// the node's loopback hostPort, the way the installer pushes. +type registryImages struct { + host string + index *registryprune.Client + source *imagepush.Source + pusher *imagepush.Pusher +} + +func newRegistryImages(endpoint string) *registryImages { + return ®istryImages{ + host: endpoint, + index: ®istryprune.Client{Endpoint: "http://" + endpoint}, + source: &imagepush.Source{Scheme: "http"}, + } +} + +func (r *registryImages) Repositories(ctx context.Context) ([]string, error) { + return r.index.Repositories(ctx) +} + +func (r *registryImages) Revisions(ctx context.Context, repo string) ([]string, map[string]string, error) { + idx, err := r.index.Index(ctx, repo) + if err != nil { + return nil, nil, err + } + digests := make([]string, 0, len(idx.Revisions)) + for _, rev := range idx.Revisions { + digests = append(digests, rev.Digest) + } + return digests, idx.Tags, nil +} + +func (r *registryImages) Manifest(ctx context.Context, repo, digest string) ([]byte, string, error) { + body, mt, err := r.source.Manifest(ctx, r.host, repo, digest) + return body, mt, registryGone(err) +} + +func (r *registryImages) Blob(ctx context.Context, repo, digest string) (io.ReadCloser, error) { + rc, err := r.source.Blob(ctx, r.host, repo, digest) + return rc, registryGone(err) +} + +func (r *registryImages) PutBlob(ctx context.Context, repo, digest string, size int64, open func() (io.ReadCloser, error)) error { + return r.pusher.UploadBlob(ctx, r.host, repo, digest, size, open) +} + +func (r *registryImages) PutManifest(ctx context.Context, repo, reference, mediaType string, body []byte) error { + _, err := r.pusher.PutManifest(ctx, r.host, repo, reference, mediaType, body) + return err +} + +// registryGone marks a 404 as a manifest or blob the registry no longer holds. +func registryGone(err error) error { + var se *imagepush.StatusError + if errors.As(err, &se) && se.Code == http.StatusNotFound { + return fmt.Errorf("%w: %v", offsite.ErrImageGone, err) + } + return err +} + +// offsiteRegistryEndpoint is where the host reaches the registry whose images +// the off-site copy covers: flag when given ("off" for none), otherwise the +// loopback hostPort of the in-cluster registry [registry] url names. A +// registry outside the cluster is not this install's to copy. +func offsiteRegistryEndpoint(flag string, reg config.RegistryConfig) string { + switch flag { + case "off": + return "" + case "": + default: + return flag + } + host, _, _ := strings.Cut(reg.URL, "/") + name, _, _ := strings.Cut(host, ":") + if strings.HasSuffix(name, ".svc") || strings.HasSuffix(name, ".svc.cluster.local") { + return loopbackEndpoint(host) + } + return "" +} + +// registryWriteCredential is the credential a host-side push uses: +// FELIS_REGISTRY_USERNAME/PASSWORD when set, otherwise the platform principal +// with REGISTRY_PLATFORM_TOKEN from the installer's secrets file. +func registryWriteCredential(secrets string) (string, string, error) { + if err := loadEnvFile(secrets); err != nil { + return "", "", fmt.Errorf("read %s: %w", secrets, err) + } + user, pass := os.Getenv("FELIS_REGISTRY_USERNAME"), os.Getenv("FELIS_REGISTRY_PASSWORD") + if pass == "" { + user, pass = registrygate.PrincipalPlatform, os.Getenv("REGISTRY_PLATFORM_TOKEN") + } + if pass == "" { + return "", "", errors.New("no registry credential: set FELIS_REGISTRY_PASSWORD or run as root on the node (REGISTRY_PLATFORM_TOKEN in /etc/felis/secrets.env)") + } + return user, pass, nil +} + +func offsiteFetchImages(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int { + cfgPath := fs.String("config", "/etc/felis/felis.toml", "path to felis.toml (the host copy)") + envFile := fs.String("env-file", defaultOffsiteEnvFile, "file with the [offsite] secrets, for variables not already set") + registry := fs.String("registry", "", "host[:port] of the registry to push into (default: the loopback hostPort of the in-cluster registry)") + secrets := fs.String("secrets-env", "/etc/felis/secrets.env", "installer secrets file holding REGISTRY_PLATFORM_TOKEN, read when FELIS_REGISTRY_PASSWORD is unset") + at := fs.String("at", "", "registry index version to restore (default: the newest; `felis offsite list` shows them)") + if err := fs.Parse(args); err != nil { + return 2 + } + cfg, env, err := loadOffsite(*cfgPath, *envFile) + if err != nil { + fmt.Fprintf(stderr, "felis offsite fetch-images: %v\n", err) + return 1 + } + endpoint := offsiteRegistryEndpoint(*registry, cfg.Registry) + if endpoint == "" { + fmt.Fprintf(stderr, "felis offsite fetch-images: [registry] url %q is not the in-cluster registry; pass -registry host:port\n", cfg.Registry.URL) + return 2 + } + user, pass, err := registryWriteCredential(*secrets) + if err != nil { + fmt.Fprintf(stderr, "felis offsite fetch-images: %v\n", err) + return 2 + } + ctx, cancel := context.WithTimeout(context.Background(), 6*time.Hour) + defer cancel() + stamp, idx, err := offsite.ChooseImageIndex(ctx, env.bucket, env.key, *at) + if err != nil { + fmt.Fprintf(stderr, "felis offsite fetch-images: %v\n", err) + return 1 + } + fmt.Fprintf(stdout, "felis offsite fetch-images: restoring registry index %s (%d repositories, %d images) into %s\n", + stamp, len(idx.Repositories), idx.Images(), endpoint) + target := newRegistryImages(endpoint) + target.pusher = &imagepush.Pusher{Scheme: "http", Username: user, Password: pass} + res, err := offsite.FetchImages(ctx, env.bucket, env.key, idx, target, stderr) + fmt.Fprintf(stdout, "felis offsite fetch-images: %d of %d repositories restored, %d images, %d tags, %d blobs pushed (%s)\n", + res.Repositories, len(idx.Repositories), res.Manifests, res.Tags, res.BlobsPushed, offsite.HumanBytes(res.BytesPushed)) + for _, f := range res.Failures { + fmt.Fprintf(stderr, "felis offsite fetch-images: %s\n", f) + } + if err != nil { + fmt.Fprintf(stderr, "felis offsite fetch-images: %v (a second run pushes only what is still missing)\n", err) + return 1 + } + return 0 +} diff --git a/cmd/felis/offsite_test.go b/cmd/felis/offsite_test.go index ace1c13..4c95d0d 100644 --- a/cmd/felis/offsite_test.go +++ b/cmd/felis/offsite_test.go @@ -2,12 +2,14 @@ package main import ( "bytes" + "errors" "os" "path/filepath" "strings" "testing" "felis.lolicon.best/internal/config" + "felis.lolicon.best/internal/imagepush" "felis.lolicon.best/internal/offsite" ) @@ -94,3 +96,27 @@ func TestOffsiteFetchDBRejectsOddNames(t *testing.T) { t.Fatalf("exit %d: %s", code, errb.String()) } } + +func TestOffsiteRegistryEndpoint(t *testing.T) { + for _, tc := range []struct{ flag, url, want string }{ + {"", "registry.felis.svc:5000", "127.0.0.1:5000"}, + {"", "registry.felis.svc.cluster.local:5001", "127.0.0.1:5001"}, + {"", "ghcr.io/acme", ""}, + {"", "", ""}, + {"off", "registry.felis.svc:5000", ""}, + {"10.0.0.5:5000", "ghcr.io/acme", "10.0.0.5:5000"}, + } { + if got := offsiteRegistryEndpoint(tc.flag, config.RegistryConfig{URL: tc.url}); got != tc.want { + t.Errorf("offsiteRegistryEndpoint(%q, %q) = %q, want %q", tc.flag, tc.url, got, tc.want) + } + } +} + +func TestRegistryGoneMarksNotFound(t *testing.T) { + if err := registryGone(&imagepush.StatusError{Op: "get blob", Code: 404}); !errors.Is(err, offsite.ErrImageGone) { + t.Fatalf("404 = %v, want ErrImageGone", err) + } + if err := registryGone(&imagepush.StatusError{Op: "get blob", Code: 503}); errors.Is(err, offsite.ErrImageGone) { + t.Fatalf("503 = %v, want it kept an ordinary failure", err) + } +} diff --git a/cmd/felis/run.go b/cmd/felis/run.go index 5c4c22e..f5173bc 100644 --- a/cmd/felis/run.go +++ b/cmd/felis/run.go @@ -13,7 +13,7 @@ Usage: Commands: migrate up Apply embedded database migrations under an advisory lock (snapshots the database first) db Back up, verify, list and restore the control-plane database (backup|restore|verify|list|check) - offsite Copy world archives and database bundles to an off-site bucket, and fetch them back (sync|status|list|fetch-db|fetch-worlds|keygen) + offsite Copy world archives, database bundles and user images to an off-site bucket, and fetch them back (sync|status|list|fetch-db|fetch-worlds|fetch-images|keygen) operator Run the MinecraftServer controller-manager api Run the felis-api HTTP server nano Run the Felis-nano hasJoined multiplexer (multi-Yggdrasil, no control plane) diff --git a/deploy/bootstrap.sh b/deploy/bootstrap.sh index ff4a30d..993b19c 100644 --- a/deploy/bootstrap.sh +++ b/deploy/bootstrap.sh @@ -3166,7 +3166,7 @@ install_offsite_timer() { fi cat > "$OFFSITE_SERVICE" <` for an + older one; `felis offsite list` shows them) through the registry's loopback + port as the platform principal, verifying every manifest and layer against + its digest, and pushes only what the registry lacks. The pruner counts a + restored image as freshly pushed and keeps it for 24 hours; finish the next + step within that window so the restored servers and whitelist entries keep + naming it. +5. Restore the database and bring the servers back: ``` kubectl -n felis scale deployment felis-api felis-operator --replicas=0 @@ -1715,7 +1732,7 @@ host yourself, plus the off-site encryption key if the copy is in the bucket. tar -xOf felis-db-....tar k8s/minecraftservers.json | kubectl apply -f - ``` -5. Bring the world archives back into the archive volume: +6. Bring the world archives back into the archive volume: ``` sudo felis offsite fetch-worlds @@ -1725,16 +1742,16 @@ host yourself, plus the off-site encryption key if the copy is in the bucket. and the volume lacks, provisioning the `felis-backups` volume first if nothing has used it yet (a short-lived `felis-bind-felis-backups-*` pod). It lists any it could not find in the bucket. Restore a world from its archive - as usual (§10, §13). Custom images built on the old host are rebuilt from - their submissions (§8), or re-pushed. + as usual (§10, §13). ### Keep a copy somewhere else A bundle on the same disk as the database protects against mistakes and bad upgrades, and a world archive on the same disk as the worlds protects against -a deleted server. Neither survives losing the disk. The installer's off-site -copy sends both to an S3-compatible bucket (AWS S3, Cloudflare R2, Backblaze -B2, MinIO, ...), encrypted on this host: +a deleted server. Neither survives losing the disk, and neither do the user +images in the platform registry. The installer's off-site copy sends all three +to an S3-compatible bucket (AWS S3, Cloudflare R2, Backblaze B2, MinIO, ...), +encrypted on this host: ``` FELIS_OFFSITE_ENDPOINT=https://.r2.cloudflarestorage.com \ @@ -1763,9 +1780,22 @@ What runs: its retention (`expires_at`) has passed. An object already in the bucket at the right size is recorded without being sent again, so a run cut short resumes. [GO-TESTED: `internal/offsite`] -- Objects are `worlds/.fenc` and `db/.fenc`: AES-256-GCM in - 64 KiB segments, so truncation, reordering and a wrong key are all refused - on the way back. +- The same run copies the user images in the platform registry: every + repository outside `felis/` and `mirror/` (the installer pushes those again), + each manifest the registry's index lists and every layer it names, read + through the loopback hostPort. A layer shared by many images is stored once. + When the set changed, a new version of the image list is written; versions + replaced more than 14 days ago are dropped together with the layers only + they named, so an image deleted by mistake stays restorable for two weeks + (`fetch-images -at`). A manifest the registry lost mid-run keeps its earlier + copy and is listed under `not whole:` in `status`. The copy covers the + in-cluster registry that `[registry] url` names; `-registry host:port` points + it elsewhere, `-registry off` skips images. [VM-TESTED: 16 images, 638 MiB, + restored into an empty registry at the same digests] +- Objects are `worlds/.fenc`, `db/.fenc`, + `registry/blobs/.fenc`, `registry/manifests/.fenc` and + `registry/index/.json.fenc`: AES-256-GCM in 64 KiB segments, so + truncation, reordering and a wrong key are all refused on the way back. - The reaper deletes an idle world only after its archive is in the bucket (§10). - The watchdog mails the owners when no sync has completed for 12 hours @@ -1775,7 +1805,7 @@ Checking it: ``` sudo felis offsite status # last run, errors, what the bucket holds, what waits -sudo felis offsite list # the bundles in the bucket, newest first +sudo felis offsite list # the bundles and image lists in the bucket, newest first sudo journalctl -u felis-offsite -n 50 --no-pager sudo systemctl start felis-offsite.service # run one now ``` @@ -1793,7 +1823,8 @@ Without a bucket, copy the backup directory off the host on a schedule of your own (`rsync -a root@felis-host:/var/lib/felis/db-backups/ /backups/felis-db/`, with the `.sha256` sidecars; `sha256sum -c` on the far side proves the copy). That covers the database only; the world archives are under the -`felis-backups` volume's directory in `/var/lib/rancher/k3s/storage/`. +`felis-backups` volume's directory in `/var/lib/rancher/k3s/storage/`, and the +registry's images under the `registry` volume's (`*_felis_registry`). ### `FELIS_PRE_MIGRATE_BACKUP=0` diff --git a/internal/imagepush/mirror.go b/internal/imagepush/mirror.go index 253a1d1..edf0b61 100644 --- a/internal/imagepush/mirror.go +++ b/internal/imagepush/mirror.go @@ -452,3 +452,14 @@ func WriteMirrorStatus(path string, st MirrorStatus) error { } return os.Rename(tmp, path) } + +// Manifest fetches one manifest of host/repo by tag or digest and returns its +// bytes and media type; a digest reference is checked against the bytes. +func (s *Source) Manifest(ctx context.Context, host, repo, reference string) ([]byte, string, error) { + return s.manifest(ctx, SourceRef{Host: host, Repo: repo}, reference) +} + +// Blob opens one blob of host/repo. The caller checks the bytes against digest. +func (s *Source) Blob(ctx context.Context, host, repo, digest string) (io.ReadCloser, error) { + return s.blob(ctx, SourceRef{Host: host, Repo: repo}, digest) +} diff --git a/internal/imagepush/push.go b/internal/imagepush/push.go index f5e8916..1caa61a 100644 --- a/internal/imagepush/push.go +++ b/internal/imagepush/push.go @@ -493,3 +493,23 @@ func (w *prefixWriter) Write(p []byte) (int, error) { } return len(p), nil } + +// UploadBlob uploads the blob open returns into host/repo unless the repository +// already holds digest, retrying like Push. open is called once per attempt. +func (p *Pusher) UploadBlob(ctx context.Context, host, repo, digest string, size int64, open func() (io.ReadCloser, error)) error { + r := Ref{Host: host, Repo: repo} + return p.retry(ctx, func() error { return p.uploadBlob(ctx, r, digest, size, open) }) +} + +// PutManifest stores body in host/repo under reference, a tag or the digest body +// hashes to, retrying like Push, and returns that digest. +func (p *Pusher) PutManifest(ctx context.Context, host, repo, reference, mediaType string, body []byte) (string, error) { + r := Ref{Host: host, Repo: repo, Tag: reference} + var digest string + err := p.retry(ctx, func() error { + d, err := p.putManifest(ctx, r, mediaType, body) + digest = d + return err + }) + return digest, err +} diff --git a/internal/offsite/images.go b/internal/offsite/images.go new file mode 100644 index 0000000..2f24cd9 --- /dev/null +++ b/internal/offsite/images.go @@ -0,0 +1,815 @@ +package offsite + +// The off-site copy of the platform registry's user images. The installer +// pushes everything under felis/ and mirror/ again on any host, so those are +// left out; every other repository holds builds that exist nowhere else: a +// lost registry volume would otherwise leave each server pinned to one of them +// in ImagePullBackOff until someone rebuilds it. +// +// The bucket holds each blob and manifest once, by digest, encrypted like +// everything else, plus an index that says which repository holds which +// manifests under which tags. A new index version is written whenever that +// changes, and every version stays until ImageHistory after a newer one +// replaced it, together with every object it names. So a registry that comes +// back empty, a host being rebuilt whose timer runs before anyone restores +// anything, only adds a new version: the one before it is still there to +// restore from (FetchImages, `felis offsite fetch-images -at`). + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "hash" + "io" + "maps" + "regexp" + "slices" + "sort" + "strings" + "time" + + "felis.lolicon.best/internal/registrygate" +) + +const ( + registryDir = "registry/" + imageBlobsDir = registryDir + "blobs/" + imageManifestsDir = registryDir + "manifests/" + imageIndexDir = registryDir + "index/" + imageIndexExt = ".json" + objExt + imageStampLayout = "20060102T150405Z" + + maxImageIndexBytes = 64 << 20 + maxImageManifestBytes = 4 << 20 +) + +// ImageHistory is how long a registry state stays restorable after a newer one +// replaced it. +const ImageHistory = 14 * 24 * time.Hour + +// ErrImageGone is a manifest or blob the registry lists but no longer holds. +var ErrImageGone = errors.New("offsite: not in the registry") + +// ImageSource is the registry the sync reads, anonymously. +type ImageSource interface { + Repositories(ctx context.Context) ([]string, error) + // Revisions lists every manifest repo holds, tagged or not, and its tags. + Revisions(ctx context.Context, repo string) (digests []string, tags map[string]string, err error) + // Manifest returns the manifest stored under digest and its media type. + // A manifest the registry no longer holds is ErrImageGone. + Manifest(ctx context.Context, repo, digest string) ([]byte, string, error) + // Blob opens a blob; one the registry no longer holds is ErrImageGone. + Blob(ctx context.Context, repo, digest string) (io.ReadCloser, error) +} + +// ImageTarget is the registry FetchImages pushes into. +type ImageTarget interface { + // PutBlob uploads the blob open returns unless repo already holds it. + PutBlob(ctx context.Context, repo, digest string, size int64, open func() (io.ReadCloser, error)) error + // PutManifest stores body under reference, a tag or its digest. + PutManifest(ctx context.Context, repo, reference, mediaType string, body []byte) error +} + +// ImageIndex is one state of the registry's user images. +type ImageIndex struct { + Created time.Time `json:"created"` + Repositories map[string]ImageRepo `json:"repositories"` + Manifests map[string]ImageManifest `json:"manifests"` +} + +// ImageRepo is one repository: its manifests, and the tag → digest map. +type ImageRepo struct { + Manifests []string `json:"manifests"` + Tags map[string]string `json:"tags,omitempty"` +} + +// ImageManifest describes a stored manifest: the blobs an image manifest names +// (config and layers), or the manifests an index names. +type ImageManifest struct { + MediaType string `json:"media_type"` + Size int64 `json:"size"` + Blobs []ImageBlob `json:"blobs,omitempty"` + Children []string `json:"children,omitempty"` +} + +// ImageBlob is one blob a manifest names. +type ImageBlob struct { + Digest string `json:"digest"` + Size int64 `json:"size"` +} + +// Images counts the manifests the repositories hold, a manifest in two +// repositories counting twice, like two images. +func (x *ImageIndex) Images() int { + n := 0 + for _, r := range x.Repositories { + n += len(r.Manifests) + } + return n +} + +func newImageIndex(at time.Time) *ImageIndex { + return &ImageIndex{Created: at, Repositories: map[string]ImageRepo{}, Manifests: map[string]ImageManifest{}} +} + +var imageDigestRE = regexp.MustCompile(`^sha256:[0-9a-f]{64}$`) + +func imageBlobKey(digest string) string { + return imageBlobsDir + strings.TrimPrefix(digest, "sha256:") + objExt +} + +func imageManifestKey(digest string) string { + return imageManifestsDir + strings.TrimPrefix(digest, "sha256:") + objExt +} + +func imageIndexKey(stamp string) string { return imageIndexDir + stamp + imageIndexExt } + +func reservedRepo(repo string) bool { + root, _, _ := strings.Cut(repo, "/") + return slices.Contains(registrygate.ReservedRepoRoots, root) +} + +// Manifest media types the copy understands. +const ( + mediaOCIIndex = "application/vnd.oci.image.index.v1+json" + mediaDockerList = "application/vnd.docker.distribution.manifest.list.v2+json" + mediaOCIManifest = "application/vnd.oci.image.manifest.v1+json" + mediaDockerManifest2 = "application/vnd.docker.distribution.manifest.v2+json" +) + +// describeManifest reads the blobs or child manifests out of a manifest. +func describeManifest(body []byte, mediaType string) (ImageManifest, error) { + type desc struct { + Digest string `json:"digest"` + Size int64 `json:"size"` + } + var m struct { + MediaType string `json:"mediaType"` + Config *desc `json:"config"` + Layers []desc `json:"layers"` + Manifests []desc `json:"manifests"` + } + if err := json.Unmarshal(body, &m); err != nil { + return ImageManifest{}, fmt.Errorf("manifest: %w", err) + } + if mediaType == "" { + mediaType = m.MediaType + } + im := ImageManifest{MediaType: mediaType, Size: int64(len(body))} + check := func(d desc) error { + if !imageDigestRE.MatchString(d.Digest) || d.Size < 0 { + return fmt.Errorf("manifest names %q (size %d), not a sha256 digest", d.Digest, d.Size) + } + return nil + } + switch mediaType { + case mediaOCIIndex, mediaDockerList: + for _, d := range m.Manifests { + if err := check(d); err != nil { + return ImageManifest{}, err + } + if !slices.Contains(im.Children, d.Digest) { + im.Children = append(im.Children, d.Digest) + } + } + case mediaOCIManifest, mediaDockerManifest2: + if m.Config == nil { + return ImageManifest{}, errors.New("image manifest without a config") + } + for _, d := range append([]desc{*m.Config}, m.Layers...) { + if err := check(d); err != nil { + return ImageManifest{}, err + } + if !slices.ContainsFunc(im.Blobs, func(b ImageBlob) bool { return b.Digest == d.Digest }) { + im.Blobs = append(im.Blobs, ImageBlob{Digest: d.Digest, Size: d.Size}) + } + } + default: + return ImageManifest{}, fmt.Errorf("unsupported manifest type %q", mediaType) + } + return im, nil +} + +// imageCopy is one sync pass over the registry. +type imageCopy struct { + s *Syncer + remote map[string]int64 + prev *ImageIndex + next *ImageIndex + res *Result + fail func(string, ...any) + // failed is set by anything that leaves the new index short of the + // registry, which rules out pruning in this pass. + failed bool +} + +// errSkipped marks a manifest left out without failing the pass. +var errSkipped = errors.New("skipped") + +func (s *Syncer) syncImages(ctx context.Context, res *Result, fail func(string, ...any)) { + if s.Images == nil { + return + } + remote, err := s.listSizes(ctx, registryDir) + if err != nil { + fail("list %s in the bucket: %v", registryDir, err) + return + } + versions := imageIndexStamps(remote) + prev := newImageIndex(time.Time{}) + if n := len(versions); n > 0 { + if prev, err = LoadImageIndex(ctx, s.Bucket, s.Key, versions[n-1]); err != nil { + fail("read the registry index %s from the bucket: %v", versions[n-1], err) + return + } + } + c := &imageCopy{s: s, remote: remote, prev: prev, next: newImageIndex(s.now().UTC()), res: res, fail: fail} + repos, err := s.Images.Repositories(ctx) + if err != nil { + fail("list the registry's repositories: %v", err) + return + } + sort.Strings(repos) + listed := map[string]bool{} + for _, repo := range repos { + if reservedRepo(repo) { + continue + } + listed[repo] = true + if ctx.Err() != nil { + c.fail("stopped before %s: %v", repo, ctx.Err()) + c.failed = true + c.carry(repo) + continue + } + c.repo(ctx, repo) + } + + stamp := "" + if len(versions) > 0 { + stamp = versions[len(versions)-1] + } + if len(versions) == 0 || !sameImages(prev, c.next) { + stamp = c.next.Created.Format(imageStampLayout) + if err := s.putImageIndex(ctx, stamp, c.next); err != nil { + fail("write the registry index: %v", err) + return + } + if len(versions) == 0 || versions[len(versions)-1] != stamp { + versions = append(versions, stamp) + } + s.logf("recorded registry index %s: %d repositories, %d images", stamp, len(c.next.Repositories), c.next.Images()) + } + res.ImageIndex = stamp + res.ImageRepos = len(c.next.Repositories) + res.Images = c.next.Images() + if !c.failed { + s.pruneImages(ctx, versions, c.next, remote, res, fail) + } + for key, size := range remote { + if strings.HasPrefix(key, imageBlobsDir) || strings.HasPrefix(key, imageManifestsDir) { + res.RemoteImageBytes += size + } + } +} + +// repo copies one repository into c.next. A repository that cannot be read in +// full keeps the entry the previous index had for it. +func (c *imageCopy) repo(ctx context.Context, repo string) { + digests, tags, err := c.s.Images.Revisions(ctx, repo) + if err != nil { + c.fail("list the images of %s: %v", repo, err) + c.failed = true + c.carry(repo) + return + } + entry := ImageRepo{Tags: map[string]string{}} + short := false + for _, d := range digests { + switch err := c.manifest(ctx, repo, d); { + case errors.Is(err, errSkipped): + case err != nil: + c.fail("copy %s@%s: %v", repo, d, err) + short = true + default: + entry.Manifests = append(entry.Manifests, d) + } + } + if short { + c.failed = true + if _, ok := c.prev.Repositories[repo]; ok { + c.carry(repo) + return + } + } + sort.Strings(entry.Manifests) + for tag, d := range tags { + if slices.Contains(entry.Manifests, d) { + entry.Tags[tag] = d + } + } + if len(entry.Manifests) > 0 { + c.next.Repositories[repo] = entry + } +} + +// carry keeps the previous index's entry for repo, and what it names. +func (c *imageCopy) carry(repo string) { + entry, ok := c.prev.Repositories[repo] + if !ok { + return + } + c.next.Repositories[repo] = entry + var add func(d string) + add = func(d string) { + m, ok := c.prev.Manifests[d] + if !ok { + return + } + c.next.Manifests[d] = m + for _, child := range m.Children { + add(child) + } + } + for _, d := range entry.Manifests { + add(d) + } +} + +// manifest makes sure digest, everything it names, and then the manifest +// itself are in the bucket, and records it in c.next. A manifest or blob the +// registry lost is errSkipped, unless an earlier pass already copied the +// manifest whole: that copy is then the only complete one and stays. +func (c *imageCopy) manifest(ctx context.Context, repo, digest string) error { + if _, ok := c.next.Manifests[digest]; ok { + return nil + } + if !imageDigestRE.MatchString(digest) { + return fmt.Errorf("%q is not a sha256 digest", digest) + } + key := imageManifestKey(digest) + im, stored := c.prev.Manifests[digest] + stored = stored && c.remote[key] == SealedSize(im.Size) + var body []byte + if !stored { + b, mt, err := c.s.Images.Manifest(ctx, repo, digest) + if err != nil { + return c.gone(repo, digest, err) + } + if im, err = describeManifest(b, mt); err != nil { + return err + } + body = b + } + for _, child := range im.Children { + if err := c.manifest(ctx, repo, child); err != nil { + if errors.Is(err, errSkipped) { + return c.gone(repo, digest, fmt.Errorf("%w: child manifest %s", ErrImageGone, child)) + } + return fmt.Errorf("child manifest %s: %w", child, err) + } + } + for _, b := range im.Blobs { + bkey := imageBlobKey(b.Digest) + if c.remote[bkey] == SealedSize(b.Size) { + continue + } + if err := c.s.putImageBlob(ctx, repo, b); err != nil { + return c.gone(repo, digest, fmt.Errorf("blob %s: %w", b.Digest, err)) + } + c.remote[bkey] = SealedSize(b.Size) + c.res.ImageBlobsUploaded++ + c.res.BytesUploaded += b.Size + } + if body != nil { + if err := c.s.putBytes(ctx, key, body); err != nil { + return err + } + c.remote[key] = SealedSize(int64(len(body))) + c.res.ImagesUploaded++ + c.s.logf("copied image %s@%s", repo, digest) + } + c.next.Manifests[digest] = im + return nil +} + +// gone turns err into errSkipped when it says the registry lost the manifest +// or one of its blobs, keeping the copy an earlier pass made if there is one. +func (c *imageCopy) gone(repo, digest string, err error) error { + if !errors.Is(err, ErrImageGone) { + return err + } + if _, ok := c.prev.Manifests[digest]; ok && c.prevComplete(digest) { + c.res.ImagesIncomplete = append(c.res.ImagesIncomplete, fmt.Sprintf("%s@%s: %v; kept the off-site copy", repo, digest, err)) + c.addPrev(digest) + return nil + } + c.res.ImagesIncomplete = append(c.res.ImagesIncomplete, fmt.Sprintf("%s@%s: %v; not copied", repo, digest, err)) + return errSkipped +} + +// prevComplete reports whether every object digest needs is in the bucket. +func (c *imageCopy) prevComplete(digest string) bool { + m, ok := c.prev.Manifests[digest] + if !ok || c.remote[imageManifestKey(digest)] != SealedSize(m.Size) { + return false + } + for _, b := range m.Blobs { + if c.remote[imageBlobKey(b.Digest)] != SealedSize(b.Size) { + return false + } + } + for _, child := range m.Children { + if !c.prevComplete(child) { + return false + } + } + return true +} + +func (c *imageCopy) addPrev(digest string) { + m := c.prev.Manifests[digest] + c.next.Manifests[digest] = m + for _, child := range m.Children { + c.addPrev(child) + } +} + +// sameImages reports whether two indexes record the same state. +func sameImages(a, b *ImageIndex) bool { + ja, err1 := json.Marshal(struct { + R map[string]ImageRepo + M map[string]ImageManifest + }{a.Repositories, a.Manifests}) + jb, err2 := json.Marshal(struct { + R map[string]ImageRepo + M map[string]ImageManifest + }{b.Repositories, b.Manifests}) + return err1 == nil && err2 == nil && bytes.Equal(ja, jb) +} + +// pruneImages drops the index versions replaced more than ImageHistory ago, +// then every blob and manifest no remaining version names. +func (s *Syncer) pruneImages(ctx context.Context, versions []string, newest *ImageIndex, remote map[string]int64, res *Result, fail func(string, ...any)) { + cutoff := s.now().Add(-ImageHistory) + var keep, drop []string + for i, v := range versions { + if i == len(versions)-1 { + break + } + replaced, err := time.Parse(imageStampLayout, versions[i+1]) + if err == nil && replaced.Before(cutoff) { + drop = append(drop, v) + } else { + keep = append(keep, v) + } + } + live := map[string]bool{} + mark := func(x *ImageIndex) { + for d, m := range x.Manifests { + live[imageManifestKey(d)] = true + for _, b := range m.Blobs { + live[imageBlobKey(b.Digest)] = true + } + } + } + mark(newest) + for _, v := range keep { + x, err := LoadImageIndex(ctx, s.Bucket, s.Key, v) + if err != nil { + fail("read the registry index %s from the bucket: %v; nothing pruned", v, err) + return + } + mark(x) + } + for _, v := range drop { + if err := s.Bucket.Remove(ctx, imageIndexKey(v)); err != nil { + fail("prune registry index %s: %v", v, err) + return + } + delete(remote, imageIndexKey(v)) + res.ImageObjectsPruned++ + } + for _, key := range slices.Sorted(maps.Keys(remote)) { + if !strings.HasPrefix(key, imageBlobsDir) && !strings.HasPrefix(key, imageManifestsDir) { + continue + } + if live[key] { + continue + } + if err := s.Bucket.Remove(ctx, key); err != nil { + fail("prune %s: %v", key, err) + continue + } + delete(remote, key) + res.ImageObjectsPruned++ + } +} + +// putImageBlob streams one blob from the registry into the bucket, checking it +// against its digest and size on the way: a blob that does not match fails +// the upload instead of storing bytes a restore would push as that digest. +func (s *Syncer) putImageBlob(ctx context.Context, repo string, b ImageBlob) error { + rc, err := s.Images.Blob(ctx, repo, b.Digest) + if err != nil { + return err + } + defer rc.Close() + return s.putStream(ctx, imageBlobKey(b.Digest), newDigestReader(rc, b.Digest, b.Size), b.Size) +} + +func (s *Syncer) putBytes(ctx context.Context, key string, plain []byte) error { + var sealed bytes.Buffer + if err := Encrypt(&sealed, bytes.NewReader(plain), s.Key); err != nil { + return err + } + return s.Bucket.Put(ctx, key, &sealed, int64(sealed.Len())) +} + +func (s *Syncer) putImageIndex(ctx context.Context, stamp string, x *ImageIndex) error { + raw, err := json.Marshal(x) + if err != nil { + return err + } + return s.putBytes(ctx, imageIndexKey(stamp), raw) +} + +// digestReader passes a blob through and fails at its end when the bytes do not +// hash to digest or are not size long. +type digestReader struct { + r io.Reader + h hash.Hash + n int64 + size int64 + digest string +} + +func newDigestReader(r io.Reader, digest string, size int64) *digestReader { + return &digestReader{r: r, h: sha256.New(), size: size, digest: digest} +} + +func (d *digestReader) Read(p []byte) (int, error) { + n, err := d.r.Read(p) + d.h.Write(p[:n]) + d.n += int64(n) + if d.n > d.size { + return n, fmt.Errorf("blob %s is longer than the %d bytes its manifest records", d.digest, d.size) + } + if err == io.EOF { + if d.n != d.size { + return n, fmt.Errorf("blob %s ended after %d of %d bytes", d.digest, d.n, d.size) + } + if got := "sha256:" + hex.EncodeToString(d.h.Sum(nil)); got != d.digest { + return n, fmt.Errorf("blob %s hashes to %s", d.digest, got) + } + } + return n, err +} + +// imageIndexStamps lists the index versions among keys, oldest first. +func imageIndexStamps(keys map[string]int64) []string { + var out []string + for key := range keys { + stamp, ok := strings.CutPrefix(key, imageIndexDir) + if !ok { + continue + } + stamp, ok = strings.CutSuffix(stamp, imageIndexExt) + if _, err := time.Parse(imageStampLayout, stamp); ok && err == nil { + out = append(out, stamp) + } + } + sort.Strings(out) + return out +} + +// ImageIndexes lists the registry index versions in the bucket, oldest first. +func ImageIndexes(ctx context.Context, b Bucket) ([]string, error) { + objs, err := b.List(ctx, imageIndexDir) + if err != nil { + return nil, err + } + keys := make(map[string]int64, len(objs)) + for _, o := range objs { + keys[o.Key] = o.Size + } + return imageIndexStamps(keys), nil +} + +// LoadImageIndex reads one index version. +func LoadImageIndex(ctx context.Context, b Bucket, key []byte, stamp string) (*ImageIndex, error) { + raw, err := getSealed(ctx, b, key, imageIndexKey(stamp), maxImageIndexBytes) + if err != nil { + return nil, err + } + x := newImageIndex(time.Time{}) + if err := json.Unmarshal(raw, x); err != nil { + return nil, fmt.Errorf("registry index %s: %w", stamp, err) + } + if x.Repositories == nil { + x.Repositories = map[string]ImageRepo{} + } + if x.Manifests == nil { + x.Manifests = map[string]ImageManifest{} + } + return x, nil +} + +// ChooseImageIndex picks the version a restore uses: at when given, otherwise +// the newest. A newest version with no images while an older one has some is +// what a rebuilt host's first sync records before the restore, so it is +// refused with the versions worth choosing instead. +func ChooseImageIndex(ctx context.Context, b Bucket, key []byte, at string) (string, *ImageIndex, error) { + versions, err := ImageIndexes(ctx, b) + if err != nil { + return "", nil, err + } + if len(versions) == 0 { + return "", nil, errors.New("the bucket holds no registry index: no sync has copied images yet") + } + if at != "" { + if !slices.Contains(versions, at) { + return "", nil, fmt.Errorf("the bucket holds no registry index %s; it holds %s", at, strings.Join(versions, ", ")) + } + x, err := LoadImageIndex(ctx, b, key, at) + return at, x, err + } + newest := versions[len(versions)-1] + x, err := LoadImageIndex(ctx, b, key, newest) + if err != nil || len(x.Repositories) > 0 { + return newest, x, err + } + var older []string + for i := len(versions) - 2; i >= 0; i-- { + o, err := LoadImageIndex(ctx, b, key, versions[i]) + if err != nil { + return "", nil, err + } + if len(o.Repositories) > 0 { + older = append(older, fmt.Sprintf("%s (%d repositories, %d images)", versions[i], len(o.Repositories), o.Images())) + } + } + if len(older) == 0 { + return newest, x, nil + } + return "", nil, fmt.Errorf("the newest registry index %s lists no images, which is what a rebuilt host records before its restore; pick the state to restore with -at: %s", + newest, strings.Join(older, ", ")) +} + +// getSealed downloads and decrypts a small object. +func getSealed(ctx context.Context, b Bucket, key []byte, objKey string, limit int64) ([]byte, error) { + rc, err := b.Get(ctx, objKey) + if err != nil { + return nil, err + } + defer rc.Close() + var out bytes.Buffer + if err := Decrypt(&limitWriter{w: &out, left: limit}, rc, key); err != nil { + return nil, fmt.Errorf("%s: %w", objKey, err) + } + return out.Bytes(), nil +} + +type limitWriter struct { + w io.Writer + left int64 +} + +func (l *limitWriter) Write(p []byte) (int, error) { + if int64(len(p)) > l.left { + return 0, errors.New("object is larger than expected") + } + l.left -= int64(len(p)) + return l.w.Write(p) +} + +// FetchImagesResult is what FetchImages did. +type FetchImagesResult struct { + Repositories int + Manifests int + Tags int + // BlobsPushed and BytesPushed count the blobs the registry lacked. + BlobsPushed int + BytesPushed int64 + Failures []string +} + +// FetchImages pushes every image x records back into t, each repository's +// blobs first, then its manifests by digest (an index after the manifests it +// names), then its tags. What t already holds is skipped, so a second run +// finishes what the first could not. +func FetchImages(ctx context.Context, b Bucket, key []byte, x *ImageIndex, t ImageTarget, log io.Writer) (FetchImagesResult, error) { + var res FetchImagesResult + bodies := map[string][]byte{} + body := func(d string) ([]byte, error) { + if raw, ok := bodies[d]; ok { + return raw, nil + } + raw, err := getSealed(ctx, b, key, imageManifestKey(d), maxImageManifestBytes) + if err != nil { + return nil, err + } + sum := sha256.Sum256(raw) + if got := "sha256:" + hex.EncodeToString(sum[:]); got != d { + return nil, fmt.Errorf("the stored manifest %s hashes to %s", d, got) + } + bodies[d] = raw + return raw, nil + } + pushedBlobs := map[string]bool{} + for _, repo := range slices.Sorted(maps.Keys(x.Repositories)) { + entry := x.Repositories[repo] + done := map[string]bool{} + var push func(d string) error + push = func(d string) error { + if done[d] { + return nil + } + m, ok := x.Manifests[d] + if !ok { + return fmt.Errorf("the index names %s but does not describe it", d) + } + for _, child := range m.Children { + if err := push(child); err != nil { + return fmt.Errorf("child %s: %w", child, err) + } + } + for _, bl := range m.Blobs { + open := func() (io.ReadCloser, error) { + if !pushedBlobs[repo+"@"+bl.Digest] { + pushedBlobs[repo+"@"+bl.Digest] = true + res.BlobsPushed++ + res.BytesPushed += bl.Size + } + return openSealed(ctx, b, key, imageBlobKey(bl.Digest)) + } + if err := t.PutBlob(ctx, repo, bl.Digest, bl.Size, open); err != nil { + return fmt.Errorf("blob %s: %w", bl.Digest, err) + } + } + raw, err := body(d) + if err != nil { + return err + } + if err := t.PutManifest(ctx, repo, d, m.MediaType, raw); err != nil { + return err + } + done[d] = true + return nil + } + ok := true + for _, d := range entry.Manifests { + if err := ctx.Err(); err != nil { + return res, err + } + if err := push(d); err != nil { + res.Failures = append(res.Failures, fmt.Sprintf("%s@%s: %v", repo, d, err)) + ok = false + continue + } + res.Manifests++ + } + for _, tag := range slices.Sorted(maps.Keys(entry.Tags)) { + d := entry.Tags[tag] + if !done[d] { + continue + } + if err := t.PutManifest(ctx, repo, tag, x.Manifests[d].MediaType, bodies[d]); err != nil { + res.Failures = append(res.Failures, fmt.Sprintf("%s:%s: %v", repo, tag, err)) + ok = false + continue + } + res.Tags++ + } + if ok { + res.Repositories++ + if log != nil { + fmt.Fprintf(log, "felis offsite: restored %s (%d images, %d tags)\n", repo, len(entry.Manifests), len(entry.Tags)) + } + } + } + if len(res.Failures) > 0 { + return res, fmt.Errorf("%d images or tags failed; first: %s", len(res.Failures), res.Failures[0]) + } + return res, nil +} + +// openSealed streams the decrypted content of objKey. Decryption authenticates +// every segment; the registry checks the whole against its digest. +func openSealed(ctx context.Context, b Bucket, key []byte, objKey string) (io.ReadCloser, error) { + rc, err := b.Get(ctx, objKey) + if err != nil { + return nil, err + } + pr, pw := io.Pipe() + go func() { + err := Decrypt(pw, rc, key) + rc.Close() + pw.CloseWithError(err) + }() + return pr, nil +} diff --git a/internal/offsite/images_test.go b/internal/offsite/images_test.go new file mode 100644 index 0000000..87de372 --- /dev/null +++ b/internal/offsite/images_test.go @@ -0,0 +1,454 @@ +package offsite + +import ( + "bytes" + "context" + "crypto/rand" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "slices" + "strings" + "testing" + "time" +) + +// ---- fake registry --------------------------------------------------------- + +type fakeRepo struct { + digests []string + tags map[string]string +} + +type fakeManifest struct { + body []byte + mediaType string +} + +// fakeRegistry is the registry as the sync reads it and as a restore writes +// it: one blob and manifest store shared by every repository, as +// distribution keeps it, and per-repository links. +type fakeRegistry struct { + repos map[string]*fakeRepo + blobs map[string][]byte + manifests map[string]fakeManifest + // linked records which blobs each repository holds, for a restore. + linked map[string]map[string]bool + revErr map[string]error + gone map[string]bool + corrupt map[string]bool + blobGets int + // log is what a restore did, in order. + log []string +} + +func newFakeRegistry() *fakeRegistry { + return &fakeRegistry{ + repos: map[string]*fakeRepo{}, blobs: map[string][]byte{}, manifests: map[string]fakeManifest{}, + linked: map[string]map[string]bool{}, revErr: map[string]error{}, gone: map[string]bool{}, corrupt: map[string]bool{}, + } +} + +func digestOf(b []byte) string { + sum := sha256.Sum256(b) + return "sha256:" + hex.EncodeToString(sum[:]) +} + +func (r *fakeRegistry) blob(data []byte) map[string]any { + d := digestOf(data) + r.blobs[d] = data + return map[string]any{"mediaType": "application/octet-stream", "digest": d, "size": len(data)} +} + +func (r *fakeRegistry) link(repo, digest string, tags ...string) { + rp := r.repos[repo] + if rp == nil { + rp = &fakeRepo{tags: map[string]string{}} + r.repos[repo] = rp + } + if !slices.Contains(rp.digests, digest) { + rp.digests = append(rp.digests, digest) + } + for _, t := range tags { + rp.tags[t] = digest + } +} + +func (r *fakeRegistry) putManifest(mediaType string, m map[string]any) string { + m["schemaVersion"] = 2 + m["mediaType"] = mediaType + body, _ := json.Marshal(m) + d := digestOf(body) + r.manifests[d] = fakeManifest{body: body, mediaType: mediaType} + return d +} + +// image stores an image manifest of the given layers in repo and returns its digest. +func (r *fakeRegistry) image(repo, tag string, layers ...[]byte) string { + var ls []any + for _, l := range layers { + ls = append(ls, r.blob(l)) + } + cfg := r.blob([]byte(fmt.Sprintf(`{"architecture":"arm64","repo":%q,"tag":%q}`, repo, tag))) + d := r.putManifest(mediaDockerManifest2, map[string]any{"config": cfg, "layers": ls}) + r.link(repo, d, tag) + return d +} + +// index stores an image index over children in repo. +func (r *fakeRegistry) index(repo, tag string, children ...string) string { + var ms []any + for _, c := range children { + ms = append(ms, map[string]any{"mediaType": mediaOCIManifest, "digest": c, "size": len(r.manifests[c].body)}) + } + d := r.putManifest(mediaOCIIndex, map[string]any{"manifests": ms}) + r.link(repo, d, tag) + return d +} + +func (r *fakeRegistry) Repositories(context.Context) ([]string, error) { + var out []string + for name := range r.repos { + out = append(out, name) + } + return out, nil +} + +func (r *fakeRegistry) Revisions(_ context.Context, repo string) ([]string, map[string]string, error) { + if err := r.revErr[repo]; err != nil { + return nil, nil, err + } + rp := r.repos[repo] + return slices.Clone(rp.digests), rp.tags, nil +} + +func (r *fakeRegistry) Manifest(_ context.Context, repo, digest string) ([]byte, string, error) { + m, ok := r.manifests[digest] + if !ok || r.gone[digest] { + return nil, "", fmt.Errorf("%w: manifest %s", ErrImageGone, digest) + } + return m.body, m.mediaType, nil +} + +func (r *fakeRegistry) Blob(_ context.Context, repo, digest string) (io.ReadCloser, error) { + r.blobGets++ + b, ok := r.blobs[digest] + if !ok || r.gone[digest] { + return nil, fmt.Errorf("%w: blob %s", ErrImageGone, digest) + } + if r.corrupt[digest] { + b = append(bytes.Clone(b), 'x') + } + return io.NopCloser(bytes.NewReader(b)), nil +} + +func (r *fakeRegistry) PutBlob(_ context.Context, repo, digest string, size int64, open func() (io.ReadCloser, error)) error { + if r.linked[repo][digest] { + return nil + } + rc, err := open() + if err != nil { + return err + } + defer rc.Close() + data, err := io.ReadAll(rc) + if err != nil { + return err + } + if int64(len(data)) != size || digestOf(data) != digest { + return fmt.Errorf("blob upload does not match %s", digest) + } + r.blobs[digest] = data + if r.linked[repo] == nil { + r.linked[repo] = map[string]bool{} + } + r.linked[repo][digest] = true + r.log = append(r.log, "blob "+repo+" "+digest) + return nil +} + +func (r *fakeRegistry) PutManifest(_ context.Context, repo, reference, mediaType string, body []byte) error { + d := digestOf(body) + if strings.HasPrefix(reference, "sha256:") && reference != d { + return fmt.Errorf("manifest %s hashes to %s", reference, d) + } + // Like distribution: everything a manifest names must be in the repository. + im, err := describeManifest(body, mediaType) + if err != nil { + return err + } + for _, b := range im.Blobs { + if !r.linked[repo][b.Digest] { + return fmt.Errorf("blob unknown: %s", b.Digest) + } + } + for _, c := range im.Children { + if r.repos[repo] == nil || !slices.Contains(r.repos[repo].digests, c) { + return fmt.Errorf("manifest unknown: %s", c) + } + } + r.manifests[d] = fakeManifest{body: body, mediaType: mediaType} + if strings.HasPrefix(reference, "sha256:") { + r.link(repo, d) + } else { + r.link(repo, d, reference) + } + r.log = append(r.log, "manifest "+repo+" "+reference) + return nil +} + +func layer(n int) []byte { + b := make([]byte, n) + rand.Read(b) + return b +} + +func newImageSyncer(t *testing.T, reg *fakeRegistry) (*Syncer, *memBucket, *time.Time) { + t.Helper() + b := newMemBucket() + clock := now + return &Syncer{ + Bucket: b, Catalog: &fakeCatalog{}, Key: testKey(t), Images: reg, + Now: func() time.Time { return clock }, + }, b, &clock +} + +func keysUnder(b *memBucket, prefix string) []string { + objs, _ := b.List(context.Background(), prefix) + var out []string + for _, o := range objs { + out = append(out, o.Key) + } + return out +} + +// ---- tests ----------------------------------------------------------------- + +// TestSyncImagesCopiesUserImages: user repositories are copied blob by blob, +// each blob once however many images share it, and the platform's reserved +// repositories are left alone. A second run with nothing new sends nothing. +func TestSyncImagesCopiesUserImages(t *testing.T) { + reg := newFakeRegistry() + shared := layer(3000) + a := reg.image("user-uploads/sub-a", "latest", shared, layer(70000)) + child := reg.image("e2e/multi", "arm64", shared) + idx := reg.index("e2e/multi", "latest", child) + reg.image("felis/felis", "v1", layer(5000)) + reg.image("mirror/trivy-db", "2", layer(5000)) + + s, b, _ := newImageSyncer(t, reg) + res, err := s.Run(context.Background()) + if err != nil { + t.Fatalf("Run: %v", err) + } + if res.ImageRepos != 2 || res.Images != 3 || res.ImagesUploaded != 3 { + t.Fatalf("result = %+v, want 2 repositories, 3 images, 3 uploaded", res) + } + // shared layer once, one unique layer, two configs + if got := keysUnder(b, imageBlobsDir); len(got) != 4 || res.ImageBlobsUploaded != 4 { + t.Fatalf("blobs in bucket = %v (uploaded %d), want 4", got, res.ImageBlobsUploaded) + } + for _, d := range []string{a, child, idx} { + if _, ok := b.objs[imageManifestKey(d)]; !ok { + t.Errorf("manifest %s not in the bucket", d) + } + } + x, err := LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex) + if err != nil { + t.Fatal(err) + } + if _, ok := x.Repositories["felis/felis"]; ok { + t.Fatal("reserved repository felis/felis was copied") + } + if got := x.Repositories["e2e/multi"]; got.Tags["latest"] != idx || got.Tags["arm64"] != child { + t.Fatalf("e2e/multi tags = %v", got.Tags) + } + + puts, gets := b.puts, reg.blobGets + res, err = s.Run(context.Background()) + if err != nil { + t.Fatalf("second Run: %v", err) + } + if b.puts != puts || reg.blobGets != gets { + t.Fatalf("second run sent %d objects and read %d blobs, want none", b.puts-puts, reg.blobGets-gets) + } + if len(keysUnder(b, imageIndexDir)) != 1 { + t.Fatalf("an unchanged registry wrote another index: %v", keysUnder(b, imageIndexDir)) + } +} + +// TestSyncImagesRejectsCorruptBlob: bytes that do not hash to the digest never +// become that digest's object, the image stays out of the index, and the run +// fails so the next one tries again. +func TestSyncImagesRejectsCorruptBlob(t *testing.T) { + reg := newFakeRegistry() + bad := layer(100) + reg.image("user/a", "v1", bad) + reg.corrupt[digestOf(bad)] = true + + s, b, _ := newImageSyncer(t, reg) + res, err := s.Run(context.Background()) + if err == nil || !strings.Contains(err.Error(), "longer than") { + t.Fatalf("Run = %v, want the corrupt blob to fail it", err) + } + if _, ok := b.objs[imageBlobKey(digestOf(bad))]; ok { + t.Fatal("the corrupt blob was stored") + } + if len(keysUnder(b, imageManifestsDir)) != 0 || res.Images != 0 { + t.Fatalf("manifest stored without its blob: %v, images=%d", keysUnder(b, imageManifestsDir), res.Images) + } +} + +// TestSyncImagesSkipsWhatTheRegistryLost: a manifest whose blob the registry no +// longer has is reported, not copied, and does not fail the run; one copied +// earlier keeps its off-site copy. +func TestSyncImagesSkipsWhatTheRegistryLost(t *testing.T) { + reg := newFakeRegistry() + keptLayer := layer(100) + kept := reg.image("user/a", "v1", keptLayer) + s, b, _ := newImageSyncer(t, reg) + if _, err := s.Run(context.Background()); err != nil { + t.Fatal(err) + } + + lostLayer := layer(100) + lost := reg.image("user/a", "v2", lostLayer) + reg.gone[digestOf(lostLayer)] = true + res, err := s.Run(context.Background()) + if err != nil { + t.Fatalf("Run = %v, want a lost blob to be reported only", err) + } + if len(res.ImagesIncomplete) != 1 || !strings.Contains(res.ImagesIncomplete[0], lost) { + t.Fatalf("incomplete = %v, want %s", res.ImagesIncomplete, lost) + } + x, _ := LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex) + if got := x.Repositories["user/a"]; !slices.Equal(got.Manifests, []string{kept}) || got.Tags["v2"] != "" { + t.Fatalf("user/a = %+v, want only the complete image", got) + } +} + +// TestSyncImagesKeepsReplacedStates: a registry that comes back empty (a host +// being rebuilt) writes a new index version but prunes nothing; the objects of +// the old state go only ImageHistory after it was replaced. +func TestSyncImagesKeepsReplacedStates(t *testing.T) { + reg := newFakeRegistry() + reg.image("user/a", "v1", layer(100)) + s, b, clock := newImageSyncer(t, reg) + if _, err := s.Run(context.Background()); err != nil { + t.Fatal(err) + } + objects := len(keysUnder(b, imageBlobsDir)) + len(keysUnder(b, imageManifestsDir)) + + reg.repos = map[string]*fakeRepo{} + *clock = clock.Add(time.Hour) + res, err := s.Run(context.Background()) + if err != nil { + t.Fatal(err) + } + if res.ImageRepos != 0 || len(keysUnder(b, imageIndexDir)) != 2 || res.ImageObjectsPruned != 0 { + t.Fatalf("after emptying: repos=%d versions=%v pruned=%d", res.ImageRepos, keysUnder(b, imageIndexDir), res.ImageObjectsPruned) + } + if got := len(keysUnder(b, imageBlobsDir)) + len(keysUnder(b, imageManifestsDir)); got != objects { + t.Fatalf("objects = %d, want the old state's %d kept", got, objects) + } + // A restore now refuses the empty newest state and names the old one. + if _, _, err := ChooseImageIndex(context.Background(), b, s.Key, ""); err == nil || !strings.Contains(err.Error(), "-at") { + t.Fatalf("ChooseImageIndex = %v, want it to point at -at", err) + } + + *clock = clock.Add(ImageHistory - 2*time.Hour) + if res, _ := s.Run(context.Background()); res.ImageObjectsPruned != 0 { + t.Fatalf("pruned %d objects inside the history window", res.ImageObjectsPruned) + } + *clock = clock.Add(3 * time.Hour) + res, err = s.Run(context.Background()) + if err != nil { + t.Fatal(err) + } + if res.ImageObjectsPruned != objects+1 || len(keysUnder(b, imageBlobsDir)) != 0 || len(keysUnder(b, imageIndexDir)) != 1 { + t.Fatalf("pruned %d, left blobs %v and versions %v", res.ImageObjectsPruned, keysUnder(b, imageBlobsDir), keysUnder(b, imageIndexDir)) + } +} + +// TestSyncImagesUnreadableRepoKeepsItsEntry: a repository the registry fails to +// list keeps its last recorded state, and the failed run prunes nothing. +func TestSyncImagesUnreadableRepoKeepsItsEntry(t *testing.T) { + reg := newFakeRegistry() + a := reg.image("user/a", "v1", layer(100)) + s, b, clock := newImageSyncer(t, reg) + if _, err := s.Run(context.Background()); err != nil { + t.Fatal(err) + } + reg.revErr["user/a"] = errors.New("gate restarting") + reg.image("user/b", "v1", layer(100)) + *clock = clock.Add(ImageHistory * 2) + res, err := s.Run(context.Background()) + if err == nil { + t.Fatal("Run succeeded with an unreadable repository") + } + x, _ := LoadImageIndex(context.Background(), b, s.Key, res.ImageIndex) + if got := x.Repositories["user/a"]; !slices.Equal(got.Manifests, []string{a}) { + t.Fatalf("user/a = %+v, want its last recorded state", got) + } + if _, ok := x.Repositories["user/b"]; !ok || res.ImageObjectsPruned != 0 { + t.Fatalf("user/b missing or pruned %d", res.ImageObjectsPruned) + } +} + +// TestFetchImagesRestoresRegistry: a fresh registry gets back every user image +// with the same digests and tags, blobs before the manifests that name them and +// an index after its children; a second run pushes nothing. +func TestFetchImagesRestoresRegistry(t *testing.T) { + reg := newFakeRegistry() + shared := layer(3000) + a := reg.image("user/a", "v1", shared, layer(9000)) + child := reg.image("user/multi", "arm64", shared) + idx := reg.index("user/multi", "latest", child) + s, b, _ := newImageSyncer(t, reg) + if _, err := s.Run(context.Background()); err != nil { + t.Fatal(err) + } + + stamp, x, err := ChooseImageIndex(context.Background(), b, s.Key, "") + if err != nil || stamp == "" { + t.Fatalf("ChooseImageIndex = %q, %v", stamp, err) + } + fresh := newFakeRegistry() + res, err := FetchImages(context.Background(), b, s.Key, x, fresh, nil) + if err != nil { + t.Fatalf("FetchImages: %v", err) + } + if res.Repositories != 2 || res.Manifests != 3 || res.Tags != 3 || res.BlobsPushed != 5 { + t.Fatalf("result = %+v", res) + } + for repo, want := range map[string]map[string]string{ + "user/a": {"v1": a}, + "user/multi": {"arm64": child, "latest": idx}, + } { + for tag, d := range want { + if got := fresh.repos[repo].tags[tag]; got != d { + t.Errorf("%s:%s = %s, want %s", repo, tag, got, d) + } + } + } + if !bytes.Equal(fresh.manifests[idx].body, reg.manifests[idx].body) { + t.Fatal("restored index differs from the original") + } + pos := func(entry string) int { return slices.Index(fresh.log, entry) } + if pos("manifest user/multi "+child) > pos("manifest user/multi "+idx) { + t.Fatalf("index pushed before its child: %v", fresh.log) + } + + n := len(fresh.log) + if _, err := FetchImages(context.Background(), b, s.Key, x, fresh, nil); err != nil { + t.Fatal(err) + } + for _, e := range fresh.log[n:] { + if strings.HasPrefix(e, "blob ") { + t.Fatalf("second restore uploaded a blob again: %s", e) + } + } +} diff --git a/internal/offsite/sync.go b/internal/offsite/sync.go index a30dd70..855d21a 100644 --- a/internal/offsite/sync.go +++ b/internal/offsite/sync.go @@ -1,8 +1,10 @@ // Package offsite keeps a second copy of what a lost node would take with it: -// every world archive (world_backups) and the newest control-plane database -// bundles (internal/dbbackup), encrypted, in an S3-compatible bucket off the -// machine. `felis offsite sync` runs it from felis-offsite.timer on the host, -// which is where both the archive volume and the bundle directory live. +// every world archive (world_backups), the newest control-plane database +// bundles (internal/dbbackup) and the user images in the platform registry +// (images.go), encrypted, in an S3-compatible bucket off the machine. `felis +// offsite sync` runs it from felis-offsite.timer on the host, which is where +// the archive volume and the bundle directory live and where the registry +// answers on its loopback hostPort. // // The database records the copy: world_backups.offsite_at is set once an // archive's object is in the bucket, and with [offsite] configured the reaper @@ -89,6 +91,9 @@ type Syncer struct { // bucket keeps. DBDir string DBKeep int + // Images is the platform registry whose user images are copied (images.go); + // nil copies none. + Images ImageSource Now func() time.Time Log io.Writer } @@ -110,7 +115,19 @@ type Result struct { DBPruned int `json:"db_pruned"` RemoteDB int `json:"remote_db"` NewestDB string `json:"newest_db,omitempty"` - Errors []string `json:"errors,omitempty"` + // ImageIndex is the newest registry index version in the bucket, which + // lists ImageRepos repositories holding Images images. + ImageIndex string `json:"image_index,omitempty"` + ImageRepos int `json:"image_repos"` + Images int `json:"images"` + ImagesUploaded int `json:"images_uploaded"` + ImageBlobsUploaded int `json:"image_blobs_uploaded"` + ImageObjectsPruned int `json:"image_objects_pruned"` + RemoteImageBytes int64 `json:"remote_image_bytes"` + // ImagesIncomplete are manifests the registry lists without holding all + // of them, so there was nothing whole to copy. + ImagesIncomplete []string `json:"images_incomplete,omitempty"` + Errors []string `json:"errors,omitempty"` } func (s *Syncer) now() time.Time { @@ -126,9 +143,9 @@ func (s *Syncer) logf(format string, args ...any) { } } -// Run does one pass: world archives, then database bundles, then expiry. A -// failure on one item is recorded and the pass carries on; the returned error -// is non-nil when anything failed. +// Run does one pass: world archives, database bundles, registry images, then +// expiry. A failure on one item is recorded and the pass carries on; the +// returned error is non-nil when anything failed. func (s *Syncer) Run(ctx context.Context) (Result, error) { var res Result fail := func(format string, args ...any) { @@ -143,6 +160,7 @@ func (s *Syncer) Run(ctx context.Context) (Result, error) { } s.syncWorlds(ctx, remoteWorlds, &res, fail) s.syncDB(ctx, &res, fail) + s.syncImages(ctx, &res, fail) s.expireWorlds(ctx, remoteWorlds, &res, fail) for _, size := range remoteWorlds { @@ -327,23 +345,28 @@ func (s *Syncer) expireWorlds(ctx context.Context, remote map[string]int64, res } } -// putFile encrypts the file at p into key. The sealed size is known in -// advance, so the upload streams: nothing larger than one part is buffered. +// putFile encrypts the file at p into key. func (s *Syncer) putFile(ctx context.Context, key, p string, size int64) error { f, err := os.Open(p) if err != nil { return err } defer f.Close() + // A file that changed size under us would not match the declared length; + // LimitReader keeps the stream to the size we announced and the bucket's + // length check catches a short one. + return s.putStream(ctx, key, io.LimitReader(f, size), size) +} + +// putStream encrypts size bytes of src into key. The sealed size is known in +// advance, so the upload streams: nothing larger than one part is buffered. +// An error from src fails the upload. +func (s *Syncer) putStream(ctx context.Context, key string, src io.Reader, size int64) error { pr, pw := io.Pipe() go func() { - // A file that changed size under us would not match the declared - // length; LimitReader keeps the stream to the size we announced and - // the length check below catches a short one. - err := Encrypt(pw, io.LimitReader(f, size), s.Key) - pw.CloseWithError(err) + pw.CloseWithError(Encrypt(pw, src, s.Key)) }() - err = s.Bucket.Put(ctx, key, pr, SealedSize(size)) + err := s.Bucket.Put(ctx, key, pr, SealedSize(size)) pr.CloseWithError(errors.New("upload finished")) return err } diff --git a/internal/watchdog/probes.go b/internal/watchdog/probes.go index 53f9d78..5a95d9b 100644 --- a/internal/watchdog/probes.go +++ b/internal/watchdog/probes.go @@ -361,8 +361,8 @@ func OffsiteFinding(statusFile string, now time.Time) *Finding { } f := &Finding{ Key: "offsite", Severity: Warning, For: backupFor, - Summary: "异地备份从未成功同步过,世界归档与数据库备份只在本机", - SummaryEN: "the off-site copy has never completed; world archives and database bundles exist on this machine only", + Summary: "异地备份从未成功同步过,世界归档、数据库备份与用户镜像只在本机", + SummaryEN: "the off-site copy has never completed; world archives, database bundles and user images exist on this machine only", Hint: hint, } if st != nil && !st.LastSuccess.IsZero() {