feat(offsite): 异地副本加入提交上传的整合包,按内容去重加密,索引按版本保留 14 天,新增 fetch-uploads 恢复
This commit is contained in:
11 files changed
+871
-69
No files matched your search
+94
-33
@@ -25,13 +25,14 @@ import (
|
||||
|
||||
const offsiteUsage = `usage:
|
||||
felis offsite sync [-config path] [-archive-dir dir] [-db-dir dir] [-registry host:port|off]
|
||||
[-status-file path]
|
||||
[-uploads-dir dir] [-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|<bundle>
|
||||
felis offsite fetch-worlds [-config path] [-archive-dir dir]
|
||||
felis offsite fetch-images [-config path] [-registry host:port] [-at version]
|
||||
felis offsite fetch-uploads [-config path] [-uploads-dir dir] [-at version]
|
||||
felis offsite keygen
|
||||
|
||||
Every verb but keygen reads the bucket credentials and the encryption key from
|
||||
@@ -45,10 +46,10 @@ 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, 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).
|
||||
// archives, the database bundles, the registry's user images and the
|
||||
// submission uploads (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 {
|
||||
if len(args) == 0 {
|
||||
fmt.Fprint(stderr, offsiteUsage)
|
||||
@@ -71,6 +72,8 @@ func cmdOffsite(args []string, stdout, stderr io.Writer) int {
|
||||
return offsiteFetchWorlds(fs, rest, stdout, stderr)
|
||||
case "fetch-images":
|
||||
return offsiteFetchImages(fs, rest, stdout, stderr)
|
||||
case "fetch-uploads":
|
||||
return offsiteFetchUploads(fs, rest, stdout, stderr)
|
||||
case "keygen":
|
||||
k, err := offsite.NewKey()
|
||||
if err != nil {
|
||||
@@ -192,6 +195,8 @@ func offsiteSync(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int
|
||||
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)`)
|
||||
uploadsDir := fs.String("uploads-dir", "", "host directory of the submission uploads volume (default: resolved from the uploads PVC through the cluster)")
|
||||
uploadsPVC := fs.String("uploads-pvc", platform.UploadsPVCName, `the submission uploads PVC, in the control-plane namespace ("" copies no uploads)`)
|
||||
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
|
||||
@@ -208,7 +213,11 @@ 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, offsiteRegistryEndpoint(*registry, cfg.Registry), stderr)
|
||||
res, err := runOffsiteSync(cfg, env, offsiteSources{
|
||||
archiveDir: *archiveDir, backupPVC: *backupPVC, dbDir: *dbDir,
|
||||
registry: offsiteRegistryEndpoint(*registry, cfg.Registry),
|
||||
uploadsDir: *uploadsDir, uploadsPVC: *uploadsPVC,
|
||||
}, stderr)
|
||||
st.Result = res
|
||||
if err != nil {
|
||||
st.LastError = err.Error()
|
||||
@@ -218,10 +227,12 @@ 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; images copied=%d blobs=%d pruned=%d; bucket holds %d worlds (%s), %d bundles, %d images in %d repositories (%s)\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; uploads copied=%d pruned=%d; bucket holds %d worlds (%s), %d bundles, %d images in %d repositories (%s), %d uploads (%s)\n",
|
||||
res.WorldsUploaded, res.WorldsPending, len(res.WorldsMissing), res.WorldsExpired,
|
||||
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))
|
||||
res.UploadsUploaded, res.UploadObjectsPruned,
|
||||
res.RemoteWorlds, offsite.HumanBytes(res.RemoteBytes), res.RemoteDB, res.Images, res.ImageRepos, offsite.HumanBytes(res.RemoteImageBytes),
|
||||
res.Uploads, offsite.HumanBytes(res.RemoteUploadBytes))
|
||||
for _, m := range res.WorldsMissing {
|
||||
fmt.Fprintf(stderr, "felis offsite sync: recorded archive not on the volume, nothing to copy: %s\n", m)
|
||||
}
|
||||
@@ -235,7 +246,18 @@ func offsiteSync(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int
|
||||
return 0
|
||||
}
|
||||
|
||||
func runOffsiteSync(cfg *config.Config, env *offsiteEnv, archiveDir, backupPVC, dbDir, registry string, log io.Writer) (offsite.Result, error) {
|
||||
// offsiteSources is where one sync pass reads from: the world archive volume
|
||||
// (archiveDir, or the backupPVC's directory), the bundle directory, the
|
||||
// registry's loopback endpoint and the uploads volume (uploadsDir, or the
|
||||
// uploadsPVC's directory). An empty source is skipped.
|
||||
type offsiteSources struct {
|
||||
archiveDir, backupPVC string
|
||||
dbDir string
|
||||
registry string
|
||||
uploadsDir, uploadsPVC string
|
||||
}
|
||||
|
||||
func runOffsiteSync(cfg *config.Config, env *offsiteEnv, src offsiteSources, 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)
|
||||
@@ -244,13 +266,23 @@ func runOffsiteSync(cfg *config.Config, env *offsiteEnv, archiveDir, backupPVC,
|
||||
if err != nil {
|
||||
return offsite.Result{}, err
|
||||
}
|
||||
if archiveDir == "" && backupPVC != "" {
|
||||
dir, err := resolveArchiveDir(ctx, cfg.K8s.Namespace, backupPVC, false, log)
|
||||
archiveDir, uploadsDir := src.archiveDir, src.uploadsDir
|
||||
if archiveDir == "" && src.backupPVC != "" {
|
||||
dir, err := resolveVolumeDir(ctx, cfg.K8s.Namespace, src.backupPVC, archiveVolume, false, log)
|
||||
if err != nil {
|
||||
return offsite.Result{}, err
|
||||
}
|
||||
archiveDir = dir
|
||||
}
|
||||
// An s3:// uploads store is off the host already; only a local one, on
|
||||
// the uploads PVC, needs the copy.
|
||||
if uploadsDir == "" && src.uploadsPVC != "" && isLocalUploadsPath(cfg.Registry.UserUploadsContext) {
|
||||
dir, err := resolveVolumeDir(ctx, platform.DefaultControlNamespace, src.uploadsPVC, uploadsVolume, false, log)
|
||||
if err != nil {
|
||||
return offsite.Result{}, err
|
||||
}
|
||||
uploadsDir = dir
|
||||
}
|
||||
drv, err := openStore(ctx, cfg.Database.URL, false)
|
||||
if err != nil {
|
||||
return offsite.Result{}, fmt.Errorf("open database: %w", err)
|
||||
@@ -258,36 +290,45 @@ func runOffsiteSync(cfg *config.Config, env *offsiteEnv, archiveDir, backupPVC,
|
||||
defer drv.Close()
|
||||
s := &offsite.Syncer{
|
||||
Bucket: env.bucket, Catalog: offsite.PGCatalog{DB: drv.DB()}, Key: env.key,
|
||||
ArchiveDir: archiveDir, DBDir: dbDir, DBKeep: env.cfg.DBKeep, Log: log,
|
||||
ArchiveDir: archiveDir, DBDir: src.dbDir, DBKeep: env.cfg.DBKeep, UploadsDir: uploadsDir, Log: log,
|
||||
}
|
||||
if registry != "" {
|
||||
s.Images = newRegistryImages(registry)
|
||||
if src.registry != "" {
|
||||
s.Images = newRegistryImages(src.registry)
|
||||
}
|
||||
return s.Run(ctx)
|
||||
}
|
||||
|
||||
// resolveArchiveDir finds the host directory behind the world archive PVC: a
|
||||
// local-path volume is a directory on this node. A PVC still waiting for its
|
||||
// first consumer holds nothing yet: without bind that is "" (no archives),
|
||||
// with bind it is bound first, for fetch-worlds to write into.
|
||||
func resolveArchiveDir(ctx context.Context, ns, pvcName string, bind bool, log io.Writer) (string, error) {
|
||||
// volumeKind names a PVC the off-site copy reads or restores, for messages,
|
||||
// with the flag that bypasses finding it through the cluster.
|
||||
type volumeKind struct{ what, dirFlag, empty string }
|
||||
|
||||
var (
|
||||
archiveVolume = volumeKind{"archive volume", "-archive-dir", "no world has been archived"}
|
||||
uploadsVolume = volumeKind{"uploads volume", "-uploads-dir", "no modpack has been uploaded"}
|
||||
)
|
||||
|
||||
// resolveVolumeDir finds the host directory behind a PVC: a local-path volume
|
||||
// is a directory on this node. A PVC still waiting for its first consumer
|
||||
// holds nothing yet: without bind that is "" (nothing to copy), with bind it
|
||||
// is bound first, for a fetch to write into.
|
||||
func resolveVolumeDir(ctx context.Context, ns, pvcName string, kind volumeKind, bind bool, log io.Writer) (string, error) {
|
||||
if ns == "" {
|
||||
ns = platform.DefaultMinecraftNamespace
|
||||
}
|
||||
cl, err := buildSystemServerClient()
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("reach the cluster to find the archive volume (or pass -archive-dir): %w", err)
|
||||
return "", fmt.Errorf("reach the cluster to find the %s (or pass %s): %w", kind.what, kind.dirFlag, err)
|
||||
}
|
||||
var pvc corev1.PersistentVolumeClaim
|
||||
if err := cl.Get(ctx, types.NamespacedName{Namespace: ns, Name: pvcName}, &pvc); err != nil {
|
||||
return "", fmt.Errorf("archive volume %s/%s: %w", ns, pvcName, err)
|
||||
return "", fmt.Errorf("%s %s/%s: %w", kind.what, ns, pvcName, err)
|
||||
}
|
||||
if pvc.Spec.VolumeName == "" {
|
||||
if !bind {
|
||||
fmt.Fprintf(log, "felis offsite: archive volume %s/%s is not bound yet; no world has been archived\n", ns, pvcName)
|
||||
fmt.Fprintf(log, "felis offsite: %s %s/%s is not bound yet; %s\n", kind.what, ns, pvcName, kind.empty)
|
||||
return "", nil
|
||||
}
|
||||
if err := bindVolume(ctx, cl, ns, pvcName, log); err != nil {
|
||||
if err := bindVolume(ctx, cl, ns, pvcName, kind, log); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := cl.Get(ctx, types.NamespacedName{Namespace: ns, Name: pvcName}, &pvc); err != nil {
|
||||
@@ -296,7 +337,7 @@ func resolveArchiveDir(ctx context.Context, ns, pvcName string, bind bool, log i
|
||||
}
|
||||
var pv corev1.PersistentVolume
|
||||
if err := cl.Get(ctx, types.NamespacedName{Name: pvc.Spec.VolumeName}, &pv); err != nil {
|
||||
return "", fmt.Errorf("archive volume %s: %w", pvc.Spec.VolumeName, err)
|
||||
return "", fmt.Errorf("%s %s: %w", kind.what, pvc.Spec.VolumeName, err)
|
||||
}
|
||||
var dir string
|
||||
switch {
|
||||
@@ -305,10 +346,10 @@ func resolveArchiveDir(ctx context.Context, ns, pvcName string, bind bool, log i
|
||||
case pv.Spec.HostPath != nil:
|
||||
dir = pv.Spec.HostPath.Path
|
||||
default:
|
||||
return "", fmt.Errorf("archive volume %s is not a directory on a node (local or hostPath); pass -archive-dir with where it is mounted on this host", pv.Name)
|
||||
return "", fmt.Errorf("%s %s is not a directory on a node (local or hostPath); pass %s with where it is mounted on this host", kind.what, pv.Name, kind.dirFlag)
|
||||
}
|
||||
if fi, err := os.Stat(dir); err != nil || !fi.IsDir() {
|
||||
return "", fmt.Errorf("archive volume %s is %s on its node, which is not a directory here; run this on the node that holds it, or pass -archive-dir", pv.Name, dir)
|
||||
return "", fmt.Errorf("%s %s is %s on its node, which is not a directory here; run this on the node that holds it, or pass %s", kind.what, pv.Name, dir, kind.dirFlag)
|
||||
}
|
||||
return dir, nil
|
||||
}
|
||||
@@ -316,10 +357,10 @@ func resolveArchiveDir(ctx context.Context, ns, pvcName string, bind bool, log i
|
||||
// bindVolume runs a pod that mounts the PVC and exits, which is what makes a
|
||||
// WaitForFirstConsumer volume (k3s local-path) get provisioned. The pod uses
|
||||
// the control plane's own image, which every install already has.
|
||||
func bindVolume(ctx context.Context, cl client.Client, ns, pvcName string, log io.Writer) error {
|
||||
func bindVolume(ctx context.Context, cl client.Client, ns, pvcName string, kind volumeKind, log io.Writer) error {
|
||||
var api appsv1.Deployment
|
||||
if err := cl.Get(ctx, types.NamespacedName{Namespace: platform.DefaultControlNamespace, Name: "felis-api"}, &api); err != nil {
|
||||
return fmt.Errorf("find the felis image to bind the archive volume with: %w", err)
|
||||
return fmt.Errorf("find the felis image to bind the %s with: %w", kind.what, err)
|
||||
}
|
||||
if len(api.Spec.Template.Spec.Containers) == 0 {
|
||||
return errors.New("felis-api has no container to take the image from")
|
||||
@@ -327,9 +368,9 @@ func bindVolume(ctx context.Context, cl client.Client, ns, pvcName string, log i
|
||||
image := api.Spec.Template.Spec.Containers[0].Image
|
||||
pod := platform.VolumeBinderPod(ns, pvcName, image)
|
||||
if err := cl.Create(ctx, pod); err != nil {
|
||||
return fmt.Errorf("start a pod to bind the archive volume: %w", err)
|
||||
return fmt.Errorf("start a pod to bind the %s: %w", kind.what, err)
|
||||
}
|
||||
fmt.Fprintf(log, "felis offsite: binding the archive volume %s/%s (pod %s)\n", ns, pvcName, pod.Name)
|
||||
fmt.Fprintf(log, "felis offsite: binding the %s %s/%s (pod %s)\n", kind.what, ns, pvcName, pod.Name)
|
||||
defer func() {
|
||||
_ = cl.Delete(context.Background(), pod, client.PropagationPolicy(metav1.DeletePropagationBackground))
|
||||
}()
|
||||
@@ -345,7 +386,7 @@ func bindVolume(ctx context.Context, cl client.Client, ns, pvcName string, log i
|
||||
case <-time.After(2 * time.Second):
|
||||
}
|
||||
}
|
||||
return fmt.Errorf("the archive volume %s/%s did not bind within 3 minutes; see kubectl -n %s describe pod %s", ns, pvcName, ns, pod.Name)
|
||||
return fmt.Errorf("the %s %s/%s did not bind within 3 minutes; see kubectl -n %s describe pod %s", kind.what, ns, pvcName, ns, pod.Name)
|
||||
}
|
||||
|
||||
func offsiteStatus(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) int {
|
||||
@@ -360,7 +401,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, database bundles and user images exist on this machine only.")
|
||||
fmt.Fprintln(stdout, "off-site copy: not configured. World archives, database bundles, user images and uploaded modpacks exist on this machine only.")
|
||||
fmt.Fprintln(stdout, "See docs/troubleshooting.md §16, \"Keep a copy somewhere else\".")
|
||||
return 1
|
||||
}
|
||||
@@ -399,6 +440,12 @@ func offsiteStatus(fs *flag.FlagSet, args []string, stdout, stderr io.Writer) in
|
||||
} else {
|
||||
fmt.Fprintln(stdout, "images: not copied (no in-cluster registry, or no sync has reached it yet)")
|
||||
}
|
||||
if r.UploadIndex != "" {
|
||||
fmt.Fprintf(stdout, "uploads: %d submission contexts (%s), uploads index %s\n",
|
||||
r.Uploads, offsite.HumanBytes(r.RemoteUploadBytes), r.UploadIndex)
|
||||
} else {
|
||||
fmt.Fprintln(stdout, "uploads: not copied (an s3:// uploads store, or no sync has reached the volume 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)
|
||||
@@ -470,6 +517,20 @@ func printOffsiteList(env *offsiteEnv, stdout, stderr io.Writer) int {
|
||||
}
|
||||
fmt.Fprintf(stdout, " %s %d images in %d repositories\n", versions[i], x.Images(), len(x.Repositories))
|
||||
}
|
||||
uploads, err := offsite.UploadIndexes(ctx, env.bucket)
|
||||
if err != nil {
|
||||
fmt.Fprintf(stderr, "felis offsite list: %v\n", err)
|
||||
return 1
|
||||
}
|
||||
fmt.Fprintf(stdout, "uploads index versions (%d, newest first; restore one with fetch-uploads -at):\n", len(uploads))
|
||||
for i := len(uploads) - 1; i >= 0; i-- {
|
||||
x, err := offsite.LoadUploadIndex(ctx, env.bucket, env.key, uploads[i])
|
||||
if err != nil {
|
||||
fmt.Fprintf(stdout, " %s unreadable: %v\n", uploads[i], err)
|
||||
continue
|
||||
}
|
||||
fmt.Fprintf(stdout, " %s %d submission contexts (%s)\n", uploads[i], len(x.Contexts), offsite.HumanBytes(x.Bytes()))
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
@@ -566,7 +627,7 @@ func offsiteFetchWorlds(fs *flag.FlagSet, args []string, stdout, stderr io.Write
|
||||
defer cancel()
|
||||
dir := *archiveDir
|
||||
if dir == "" {
|
||||
if dir, err = resolveArchiveDir(ctx, cfg.K8s.Namespace, *backupPVC, true, stderr); err != nil {
|
||||
if dir, err = resolveVolumeDir(ctx, cfg.K8s.Namespace, *backupPVC, archiveVolume, true, stderr); err != nil {
|
||||
fmt.Fprintf(stderr, "felis offsite fetch-worlds: %v\n", err)
|
||||
return 1
|
||||
}
|
||||
|
||||
@@ -0,0 +1,59 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"flag"
|
||||
"fmt"
|
||||
"io"
|
||||
"time"
|
||||
|
||||
"felis.lolicon.best/internal/offsite"
|
||||
"felis.lolicon.best/internal/platform"
|
||||
)
|
||||
|
||||
func offsiteFetchUploads(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")
|
||||
uploadsDir := fs.String("uploads-dir", "", "host directory of the submission uploads volume (default: resolved from the uploads PVC, binding it if needed)")
|
||||
uploadsPVC := fs.String("uploads-pvc", platform.UploadsPVCName, "the submission uploads PVC, in the control-plane namespace")
|
||||
at := fs.String("at", "", "uploads 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-uploads: %v\n", err)
|
||||
return 1
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 6*time.Hour)
|
||||
defer cancel()
|
||||
stamp, idx, err := offsite.ChooseUploadIndex(ctx, env.bucket, env.key, *at)
|
||||
if err != nil {
|
||||
fmt.Fprintf(stderr, "felis offsite fetch-uploads: %v\n", err)
|
||||
return 1
|
||||
}
|
||||
dir := *uploadsDir
|
||||
if dir == "" {
|
||||
if !isLocalUploadsPath(cfg.Registry.UserUploadsContext) {
|
||||
fmt.Fprintf(stderr, "felis offsite fetch-uploads: [registry] user_uploads_context %q is not the uploads volume; pass -uploads-dir to restore into a directory anyway\n", cfg.Registry.UserUploadsContext)
|
||||
return 2
|
||||
}
|
||||
if dir, err = resolveVolumeDir(ctx, platform.DefaultControlNamespace, *uploadsPVC, uploadsVolume, true, stderr); err != nil {
|
||||
fmt.Fprintf(stderr, "felis offsite fetch-uploads: %v\n", err)
|
||||
return 1
|
||||
}
|
||||
}
|
||||
fmt.Fprintf(stdout, "felis offsite fetch-uploads: restoring uploads index %s (%d submission contexts, %s) into %s\n",
|
||||
stamp, len(idx.Contexts), offsite.HumanBytes(idx.Bytes()), dir)
|
||||
res, err := offsite.FetchUploads(ctx, env.bucket, env.key, idx, dir, platform.ControlPlaneUID, platform.ControlPlaneUID, stderr)
|
||||
fmt.Fprintf(stdout, "felis offsite fetch-uploads: %d written (%s), %d already in place, %d failed\n",
|
||||
res.Written, offsite.HumanBytes(res.Bytes), res.Present, len(res.Failures))
|
||||
for _, f := range res.Failures {
|
||||
fmt.Fprintf(stderr, "felis offsite fetch-uploads: %s\n", f)
|
||||
}
|
||||
if err != nil {
|
||||
fmt.Fprintf(stderr, "felis offsite fetch-uploads: %v (a second run writes only what is still missing)\n", err)
|
||||
return 1
|
||||
}
|
||||
return 0
|
||||
}
|
||||
+1
-1
@@ -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, database bundles and user images to an off-site bucket, and fetch them back (sync|status|list|fetch-db|fetch-worlds|fetch-images|keygen)
|
||||
offsite Copy world archives, database bundles, user images and uploads to an off-site bucket, and fetch them back (sync|status|list|fetch-db|fetch-worlds|fetch-images|fetch-uploads|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)
|
||||
|
||||
+1
-1
@@ -3166,7 +3166,7 @@ install_offsite_timer() {
|
||||
fi
|
||||
cat > "$OFFSITE_SERVICE" <<EOF
|
||||
[Unit]
|
||||
Description=Felis off-site copy (world archives, database bundles and user registry images, encrypted, to the [offsite] bucket)
|
||||
Description=Felis off-site copy (world archives, database bundles, user registry images and submission uploads, encrypted, to the [offsite] bucket)
|
||||
After=network-online.target k3s.service postgresql.service felis-db-backup.service
|
||||
Wants=network-online.target
|
||||
|
||||
|
||||
+32
-11
@@ -700,7 +700,8 @@ control namespace (or `--registry-namespace`):
|
||||
`sudo felis offsite fetch-images` pushes the user images back at their old
|
||||
digests, so servers pinned to them pull again; it pushes only what the
|
||||
registry lacks, so a second run after an interruption is cheap. Without an
|
||||
off-site copy, user images come back from their approved submissions:
|
||||
off-site copy of the images, user images come back from their approved
|
||||
submissions (whose uploads the off-site copy also carries):
|
||||
the uploaded context of an approved submission stays on the uploads PVC
|
||||
(`GET /api/v1/submissions/{id}/context`, its `context_ref` and `image_ref`
|
||||
are in `GET /api/v1/submissions`), so an admin can build it again through
|
||||
@@ -1722,7 +1723,17 @@ host yourself, plus the off-site encryption key if the copy is in the bucket.
|
||||
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:
|
||||
5. Put the submission uploads back:
|
||||
|
||||
```
|
||||
sudo felis offsite fetch-uploads
|
||||
```
|
||||
|
||||
It writes every upload of the newest upload list (`-at <stamp>` for an
|
||||
older one) into the `felis-uploads` volume, owned by the control plane's
|
||||
uid, checking each against its sha256, and leaves one already in place
|
||||
alone.
|
||||
6. Restore the database and bring the servers back:
|
||||
|
||||
```
|
||||
kubectl -n felis scale deployment felis-api felis-operator --replicas=0
|
||||
@@ -1732,7 +1743,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 -
|
||||
```
|
||||
|
||||
6. Bring the world archives back into the archive volume:
|
||||
7. Bring the world archives back into the archive volume:
|
||||
|
||||
```
|
||||
sudo felis offsite fetch-worlds
|
||||
@@ -1749,9 +1760,9 @@ host yourself, plus the off-site encryption key if the copy is in the bucket.
|
||||
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, 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:
|
||||
images in the platform registry or the modpacks users uploaded for review. The
|
||||
installer's off-site copy sends all four to an S3-compatible bucket (AWS S3,
|
||||
Cloudflare R2, Backblaze B2, MinIO, ...), encrypted on this host:
|
||||
|
||||
```
|
||||
FELIS_OFFSITE_ENDPOINT=https://<account>.r2.cloudflarestorage.com \
|
||||
@@ -1792,9 +1803,18 @@ What runs:
|
||||
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]
|
||||
- It also copies the submission uploads (`sub-*/context.tar.gz` on the
|
||||
`felis-uploads` volume, §8), the source an admin rebuilds an approved image
|
||||
from. Identical uploads are stored once; the upload list is versioned and
|
||||
kept for 14 days like the image list (`fetch-uploads -at`). A context is read
|
||||
again only when its size or modification time changed. It runs when
|
||||
`[registry] user_uploads_context` is a local path (the installer's default);
|
||||
`-uploads-dir` names the directory by hand, `-uploads-pvc ""` skips it.
|
||||
[VM-TESTED: 17 uploads in 10 objects, 200 MiB, restored byte-identical]
|
||||
- Objects are `worlds/<archive>.fenc`, `db/<bundle>.fenc`,
|
||||
`registry/blobs/<sha256>.fenc`, `registry/manifests/<sha256>.fenc` and
|
||||
`registry/index/<stamp>.json.fenc`: AES-256-GCM in 64 KiB segments, so
|
||||
`registry/blobs/<sha256>.fenc`, `registry/manifests/<sha256>.fenc`,
|
||||
`registry/index/<stamp>.json.fenc`, `uploads/blobs/<sha256>.fenc` and
|
||||
`uploads/index/<stamp>.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).
|
||||
@@ -1805,7 +1825,7 @@ Checking it:
|
||||
|
||||
```
|
||||
sudo felis offsite status # last run, errors, what the bucket holds, what waits
|
||||
sudo felis offsite list # the bundles and image lists in the bucket, newest first
|
||||
sudo felis offsite list # the bundles, image lists and upload lists in the bucket, newest first
|
||||
sudo journalctl -u felis-offsite -n 50 --no-pager
|
||||
sudo systemctl start felis-offsite.service # run one now
|
||||
```
|
||||
@@ -1823,8 +1843,9 @@ 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/`, and the
|
||||
registry's images under the `registry` volume's (`*_felis_registry`).
|
||||
`felis-backups` volume's directory in `/var/lib/rancher/k3s/storage/`, the
|
||||
registry's images under the `registry` volume's (`*_felis_registry`) and the
|
||||
uploads under the `felis-uploads` volume's (`*_felis_felis-uploads`).
|
||||
|
||||
### `FELIS_PRE_MIGRATE_BACKUP=0`
|
||||
|
||||
|
||||
+26
-17
@@ -455,19 +455,7 @@ func sameImages(a, b *ImageIndex) bool {
|
||||
// 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)
|
||||
}
|
||||
}
|
||||
keep, drop := retire(versions, s.now().Add(-ImageHistory))
|
||||
live := map[string]bool{}
|
||||
mark := func(x *ImageIndex) {
|
||||
for d, m := range x.Manifests {
|
||||
@@ -557,7 +545,7 @@ func (d *digestReader) Read(p []byte) (int, error) {
|
||||
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)
|
||||
return n, fmt.Errorf("blob %s is longer than the %d bytes recorded for it", d.digest, d.size)
|
||||
}
|
||||
if err == io.EOF {
|
||||
if d.n != d.size {
|
||||
@@ -570,11 +558,32 @@ func (d *digestReader) Read(p []byte) (int, error) {
|
||||
return n, err
|
||||
}
|
||||
|
||||
// imageIndexStamps lists the index versions among keys, oldest first.
|
||||
func imageIndexStamps(keys map[string]int64) []string {
|
||||
// retire splits the index versions before the newest (oldest first) into the
|
||||
// ones still kept and the ones replaced before cutoff. The newest is in
|
||||
// neither: it is the state the pass just recorded.
|
||||
func retire(versions []string, cutoff time.Time) (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)
|
||||
}
|
||||
}
|
||||
return keep, drop
|
||||
}
|
||||
|
||||
// imageIndexStamps lists the registry index versions among keys, oldest first.
|
||||
func imageIndexStamps(keys map[string]int64) []string { return indexStamps(keys, imageIndexDir) }
|
||||
|
||||
// indexStamps lists the index versions under dir among keys, oldest first.
|
||||
func indexStamps(keys map[string]int64, dir string) []string {
|
||||
var out []string
|
||||
for key := range keys {
|
||||
stamp, ok := strings.CutPrefix(key, imageIndexDir)
|
||||
stamp, ok := strings.CutPrefix(key, dir)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -1,7 +1,8 @@
|
||||
// Package offsite keeps a second copy of what a lost node would take with it:
|
||||
// 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
|
||||
// bundles (internal/dbbackup), the user images in the platform registry
|
||||
// (images.go) and the submission uploads (uploads.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.
|
||||
@@ -94,6 +95,9 @@ type Syncer struct {
|
||||
// Images is the platform registry whose user images are copied (images.go);
|
||||
// nil copies none.
|
||||
Images ImageSource
|
||||
// UploadsDir is the host directory of the uploads volume, whose submission
|
||||
// contexts are copied (uploads.go); empty copies none.
|
||||
UploadsDir string
|
||||
Now func() time.Time
|
||||
Log io.Writer
|
||||
}
|
||||
@@ -127,6 +131,13 @@ type Result struct {
|
||||
// ImagesIncomplete are manifests the registry lists without holding all
|
||||
// of them, so there was nothing whole to copy.
|
||||
ImagesIncomplete []string `json:"images_incomplete,omitempty"`
|
||||
// UploadIndex is the newest uploads index version in the bucket, which
|
||||
// names Uploads submission contexts.
|
||||
UploadIndex string `json:"upload_index,omitempty"`
|
||||
Uploads int `json:"uploads"`
|
||||
UploadsUploaded int `json:"uploads_uploaded"`
|
||||
UploadObjectsPruned int `json:"upload_objects_pruned"`
|
||||
RemoteUploadBytes int64 `json:"remote_upload_bytes"`
|
||||
Errors []string `json:"errors,omitempty"`
|
||||
}
|
||||
|
||||
@@ -143,8 +154,8 @@ func (s *Syncer) logf(format string, args ...any) {
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
// Run does one pass: world archives, database bundles, registry images,
|
||||
// submission uploads, 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
|
||||
@@ -161,6 +172,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.syncUploads(ctx, &res, fail)
|
||||
s.expireWorlds(ctx, remoteWorlds, &res, fail)
|
||||
|
||||
for _, size := range remoteWorlds {
|
||||
|
||||
@@ -0,0 +1,433 @@
|
||||
package offsite
|
||||
|
||||
// Submission build contexts: the modpacks users uploaded, kept on the uploads
|
||||
// volume as <id>/context.tar.gz (internal/submit.LocalContextStore). A reviewer
|
||||
// downloads one to inspect it, and building the submission again needs it; the
|
||||
// image already built from it is in the registry copy (images.go).
|
||||
//
|
||||
// The copy is content-addressed like the images: uploads/blobs/<sha256>.fenc
|
||||
// holds each context once, and uploads/index/<stamp>.json.fenc versions which
|
||||
// submission holds which. A context deleted here (a withdrawn submission, a
|
||||
// rejected one reaped) leaves the newest version but stays in the versions
|
||||
// before it for UploadHistory. A volume that comes back empty, as it does on a
|
||||
// rebuilt host before its restore, therefore takes nothing out of the bucket
|
||||
// for two weeks.
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/fs"
|
||||
"maps"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"slices"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
uploadsDir = "uploads/"
|
||||
uploadBlobsDir = uploadsDir + "blobs/"
|
||||
uploadIndexDir = uploadsDir + "index/"
|
||||
maxUploadIndexBytes = 16 << 20
|
||||
|
||||
// UploadContextFile is the one file of a submission's directory on the
|
||||
// uploads volume (internal/submit's contextBlobName).
|
||||
UploadContextFile = "context.tar.gz"
|
||||
)
|
||||
|
||||
// UploadHistory is how long an uploads index version is kept after a newer one
|
||||
// replaced it, with the contexts only it names.
|
||||
const UploadHistory = ImageHistory
|
||||
|
||||
// uploadIDRE is internal/submit's idRE: a submission id can hold no path
|
||||
// separator and no "..".
|
||||
var uploadIDRE = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,127}$`)
|
||||
|
||||
// UploadIndex is one version of the uploads copy: which submission holds which
|
||||
// context.
|
||||
type UploadIndex struct {
|
||||
Created time.Time `json:"created"`
|
||||
Contexts map[string]UploadContext `json:"contexts"`
|
||||
}
|
||||
|
||||
// UploadContext is one submission's context as the sync last saw it. Size and
|
||||
// ModTime let the next pass skip hashing a file that has not changed.
|
||||
type UploadContext struct {
|
||||
Digest string `json:"digest"`
|
||||
Size int64 `json:"size"`
|
||||
ModTime time.Time `json:"mtime"`
|
||||
}
|
||||
|
||||
// Bytes is the total size of the contexts x names.
|
||||
func (x *UploadIndex) Bytes() int64 {
|
||||
var n int64
|
||||
for _, c := range x.Contexts {
|
||||
n += c.Size
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
func newUploadIndex(at time.Time) *UploadIndex {
|
||||
return &UploadIndex{Created: at, Contexts: map[string]UploadContext{}}
|
||||
}
|
||||
|
||||
func uploadBlobKey(digest string) string {
|
||||
return uploadBlobsDir + strings.TrimPrefix(digest, "sha256:") + objExt
|
||||
}
|
||||
|
||||
func uploadIndexKey(stamp string) string { return uploadIndexDir + stamp + imageIndexExt }
|
||||
|
||||
func (s *Syncer) syncUploads(ctx context.Context, res *Result, fail func(string, ...any)) {
|
||||
if s.UploadsDir == "" {
|
||||
return
|
||||
}
|
||||
remote, err := s.listSizes(ctx, uploadsDir)
|
||||
if err != nil {
|
||||
fail("list %s in the bucket: %v", uploadsDir, err)
|
||||
return
|
||||
}
|
||||
versions := indexStamps(remote, uploadIndexDir)
|
||||
prev := newUploadIndex(time.Time{})
|
||||
if n := len(versions); n > 0 {
|
||||
if prev, err = LoadUploadIndex(ctx, s.Bucket, s.Key, versions[n-1]); err != nil {
|
||||
fail("read the uploads index %s from the bucket: %v", versions[n-1], err)
|
||||
return
|
||||
}
|
||||
}
|
||||
entries, err := os.ReadDir(s.UploadsDir)
|
||||
if err != nil {
|
||||
fail("read the uploads volume %s: %v", s.UploadsDir, err)
|
||||
return
|
||||
}
|
||||
next := newUploadIndex(s.now().UTC())
|
||||
// failed is set by anything that leaves next short of the volume, which
|
||||
// rules out pruning in this pass.
|
||||
failed := false
|
||||
for _, e := range entries {
|
||||
id := e.Name()
|
||||
if !e.IsDir() || !uploadIDRE.MatchString(id) {
|
||||
continue
|
||||
}
|
||||
old, had := prev.Contexts[id]
|
||||
if ctx.Err() != nil {
|
||||
fail("stopped before the upload of %s: %v", id, ctx.Err())
|
||||
failed = true
|
||||
if had {
|
||||
next.Contexts[id] = old
|
||||
}
|
||||
continue
|
||||
}
|
||||
c, err := s.copyUpload(ctx, id, old, had, remote, res)
|
||||
switch {
|
||||
case errors.Is(err, fs.ErrNotExist):
|
||||
// A directory without its context: an upload still being written,
|
||||
// or one reaped between the listing and here.
|
||||
case err != nil:
|
||||
fail("copy the upload of %s: %v", id, err)
|
||||
failed = true
|
||||
if had {
|
||||
next.Contexts[id] = old
|
||||
}
|
||||
default:
|
||||
next.Contexts[id] = c
|
||||
}
|
||||
}
|
||||
|
||||
stamp := ""
|
||||
if len(versions) > 0 {
|
||||
stamp = versions[len(versions)-1]
|
||||
}
|
||||
if len(versions) == 0 || !maps.EqualFunc(prev.Contexts, next.Contexts, sameContext) {
|
||||
stamp = next.Created.Format(imageStampLayout)
|
||||
raw, err := json.Marshal(next)
|
||||
if err == nil {
|
||||
err = s.putBytes(ctx, uploadIndexKey(stamp), raw)
|
||||
}
|
||||
if err != nil {
|
||||
fail("write the uploads index: %v", err)
|
||||
return
|
||||
}
|
||||
if len(versions) == 0 || versions[len(versions)-1] != stamp {
|
||||
versions = append(versions, stamp)
|
||||
}
|
||||
s.logf("recorded uploads index %s: %d contexts (%s)", stamp, len(next.Contexts), HumanBytes(next.Bytes()))
|
||||
}
|
||||
res.UploadIndex = stamp
|
||||
res.Uploads = len(next.Contexts)
|
||||
if !failed {
|
||||
s.pruneUploads(ctx, versions, next, remote, res, fail)
|
||||
}
|
||||
for key, size := range remote {
|
||||
if strings.HasPrefix(key, uploadBlobsDir) {
|
||||
res.RemoteUploadBytes += size
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// copyUpload makes sure the bucket holds the context of submission id and
|
||||
// returns what the index records for it. A file whose size and modification
|
||||
// time match the previous index is not read again.
|
||||
func (s *Syncer) copyUpload(ctx context.Context, id string, old UploadContext, had bool, remote map[string]int64, res *Result) (UploadContext, error) {
|
||||
p := filepath.Join(s.UploadsDir, id, UploadContextFile)
|
||||
fi, err := os.Lstat(p)
|
||||
if err != nil {
|
||||
return UploadContext{}, err
|
||||
}
|
||||
if !fi.Mode().IsRegular() {
|
||||
return UploadContext{}, fmt.Errorf("%s is not a regular file", p)
|
||||
}
|
||||
c := UploadContext{Size: fi.Size(), ModTime: fi.ModTime().UTC()}
|
||||
if had && old.Size == c.Size && old.ModTime.Equal(c.ModTime) && remote[uploadBlobKey(old.Digest)] == SealedSize(old.Size) {
|
||||
return old, nil
|
||||
}
|
||||
if c.Digest, err = hashFile(p); err != nil {
|
||||
return UploadContext{}, err
|
||||
}
|
||||
key := uploadBlobKey(c.Digest)
|
||||
if remote[key] == SealedSize(c.Size) {
|
||||
return c, nil
|
||||
}
|
||||
f, err := os.Open(p)
|
||||
if err != nil {
|
||||
return UploadContext{}, err
|
||||
}
|
||||
defer f.Close()
|
||||
// Checked again on the way up: a context replaced between the hash and
|
||||
// the upload fails instead of being stored under the old digest.
|
||||
if err := s.putStream(ctx, key, newDigestReader(f, c.Digest, c.Size), c.Size); err != nil {
|
||||
return UploadContext{}, err
|
||||
}
|
||||
remote[key] = SealedSize(c.Size)
|
||||
res.UploadsUploaded++
|
||||
res.BytesUploaded += c.Size
|
||||
s.logf("copied the upload of %s (%s)", id, HumanBytes(c.Size))
|
||||
return c, nil
|
||||
}
|
||||
|
||||
func sameContext(a, b UploadContext) bool {
|
||||
return a.Digest == b.Digest && a.Size == b.Size && a.ModTime.Equal(b.ModTime)
|
||||
}
|
||||
|
||||
func hashFile(p string) (string, error) {
|
||||
f, err := os.Open(p)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
defer f.Close()
|
||||
h := sha256.New()
|
||||
if _, err := io.Copy(h, f); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return "sha256:" + hex.EncodeToString(h.Sum(nil)), nil
|
||||
}
|
||||
|
||||
// pruneUploads drops the index versions replaced more than UploadHistory ago,
|
||||
// then every context no remaining version names.
|
||||
func (s *Syncer) pruneUploads(ctx context.Context, versions []string, newest *UploadIndex, remote map[string]int64, res *Result, fail func(string, ...any)) {
|
||||
keep, drop := retire(versions, s.now().Add(-UploadHistory))
|
||||
live := map[string]bool{}
|
||||
mark := func(x *UploadIndex) {
|
||||
for _, c := range x.Contexts {
|
||||
live[uploadBlobKey(c.Digest)] = true
|
||||
}
|
||||
}
|
||||
mark(newest)
|
||||
for _, v := range keep {
|
||||
x, err := LoadUploadIndex(ctx, s.Bucket, s.Key, v)
|
||||
if err != nil {
|
||||
fail("read the uploads index %s from the bucket: %v; nothing pruned", v, err)
|
||||
return
|
||||
}
|
||||
mark(x)
|
||||
}
|
||||
for _, v := range drop {
|
||||
if err := s.Bucket.Remove(ctx, uploadIndexKey(v)); err != nil {
|
||||
fail("prune uploads index %s: %v", v, err)
|
||||
return
|
||||
}
|
||||
delete(remote, uploadIndexKey(v))
|
||||
res.UploadObjectsPruned++
|
||||
}
|
||||
for _, key := range slices.Sorted(maps.Keys(remote)) {
|
||||
if !strings.HasPrefix(key, uploadBlobsDir) || live[key] {
|
||||
continue
|
||||
}
|
||||
if err := s.Bucket.Remove(ctx, key); err != nil {
|
||||
fail("prune %s: %v", key, err)
|
||||
continue
|
||||
}
|
||||
delete(remote, key)
|
||||
res.UploadObjectsPruned++
|
||||
}
|
||||
}
|
||||
|
||||
// UploadIndexes lists the uploads index versions in the bucket, oldest first.
|
||||
func UploadIndexes(ctx context.Context, b Bucket) ([]string, error) {
|
||||
objs, err := b.List(ctx, uploadIndexDir)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
keys := make(map[string]int64, len(objs))
|
||||
for _, o := range objs {
|
||||
keys[o.Key] = o.Size
|
||||
}
|
||||
return indexStamps(keys, uploadIndexDir), nil
|
||||
}
|
||||
|
||||
// LoadUploadIndex reads one uploads index version.
|
||||
func LoadUploadIndex(ctx context.Context, b Bucket, key []byte, stamp string) (*UploadIndex, error) {
|
||||
raw, err := getSealed(ctx, b, key, uploadIndexKey(stamp), maxUploadIndexBytes)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
x := newUploadIndex(time.Time{})
|
||||
if err := json.Unmarshal(raw, x); err != nil {
|
||||
return nil, fmt.Errorf("uploads index %s: %w", stamp, err)
|
||||
}
|
||||
if x.Contexts == nil {
|
||||
x.Contexts = map[string]UploadContext{}
|
||||
}
|
||||
return x, nil
|
||||
}
|
||||
|
||||
// ChooseUploadIndex picks the version a restore uses: at when given, otherwise
|
||||
// the newest. An empty newest version while an older one names contexts is
|
||||
// what a rebuilt host's first sync records, so it is refused with the versions
|
||||
// worth choosing instead.
|
||||
func ChooseUploadIndex(ctx context.Context, b Bucket, key []byte, at string) (string, *UploadIndex, error) {
|
||||
versions, err := UploadIndexes(ctx, b)
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
if len(versions) == 0 {
|
||||
return "", nil, errors.New("the bucket holds no uploads index: no sync has copied the uploads volume yet")
|
||||
}
|
||||
if at != "" {
|
||||
if !slices.Contains(versions, at) {
|
||||
return "", nil, fmt.Errorf("the bucket holds no uploads index %s; it holds %s", at, strings.Join(versions, ", "))
|
||||
}
|
||||
x, err := LoadUploadIndex(ctx, b, key, at)
|
||||
return at, x, err
|
||||
}
|
||||
newest := versions[len(versions)-1]
|
||||
x, err := LoadUploadIndex(ctx, b, key, newest)
|
||||
if err != nil || len(x.Contexts) > 0 {
|
||||
return newest, x, err
|
||||
}
|
||||
var older []string
|
||||
for i := len(versions) - 2; i >= 0; i-- {
|
||||
o, err := LoadUploadIndex(ctx, b, key, versions[i])
|
||||
if err != nil {
|
||||
return "", nil, err
|
||||
}
|
||||
if len(o.Contexts) > 0 {
|
||||
older = append(older, fmt.Sprintf("%s (%d contexts)", versions[i], len(o.Contexts)))
|
||||
}
|
||||
}
|
||||
if len(older) == 0 {
|
||||
return newest, x, nil
|
||||
}
|
||||
return "", nil, fmt.Errorf("the newest uploads index %s lists no contexts, which is what a rebuilt host records before its restore; pick the state to restore with -at: %s",
|
||||
newest, strings.Join(older, ", "))
|
||||
}
|
||||
|
||||
// FetchUploadsResult counts what FetchUploads wrote.
|
||||
type FetchUploadsResult struct {
|
||||
Written int
|
||||
Present int
|
||||
Bytes int64
|
||||
Failures []string
|
||||
}
|
||||
|
||||
// FetchUploads writes every context x names into dir as <id>/context.tar.gz,
|
||||
// skipping one already there with the same content. Each is written to a
|
||||
// temporary file beside its place, checked against its digest, and renamed
|
||||
// into place, so a failure leaves no partial context. uid and gid own what it
|
||||
// creates (-1 leaves the caller's).
|
||||
func FetchUploads(ctx context.Context, b Bucket, key []byte, x *UploadIndex, dir string, uid, gid int, log io.Writer) (FetchUploadsResult, error) {
|
||||
var res FetchUploadsResult
|
||||
for _, id := range slices.Sorted(maps.Keys(x.Contexts)) {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return res, err
|
||||
}
|
||||
c := x.Contexts[id]
|
||||
if !uploadIDRE.MatchString(id) {
|
||||
res.Failures = append(res.Failures, fmt.Sprintf("%s: not a submission id", id))
|
||||
continue
|
||||
}
|
||||
wrote, err := fetchUpload(ctx, b, key, id, c, dir, uid, gid)
|
||||
switch {
|
||||
case err != nil:
|
||||
res.Failures = append(res.Failures, fmt.Sprintf("%s: %v", id, err))
|
||||
case wrote:
|
||||
res.Written++
|
||||
res.Bytes += c.Size
|
||||
if log != nil {
|
||||
fmt.Fprintf(log, "felis offsite: restored the upload of %s (%s)\n", id, HumanBytes(c.Size))
|
||||
}
|
||||
default:
|
||||
res.Present++
|
||||
}
|
||||
}
|
||||
if len(res.Failures) > 0 {
|
||||
return res, fmt.Errorf("%d contexts failed; first: %s", len(res.Failures), res.Failures[0])
|
||||
}
|
||||
return res, nil
|
||||
}
|
||||
|
||||
func fetchUpload(ctx context.Context, b Bucket, key []byte, id string, c UploadContext, dir string, uid, gid int) (bool, error) {
|
||||
sub := filepath.Join(dir, id)
|
||||
final := filepath.Join(sub, UploadContextFile)
|
||||
if fi, err := os.Lstat(final); err == nil && fi.Mode().IsRegular() && fi.Size() == c.Size {
|
||||
if d, err := hashFile(final); err == nil && d == c.Digest {
|
||||
return false, nil
|
||||
}
|
||||
}
|
||||
if err := os.MkdirAll(sub, 0o770); err != nil {
|
||||
return false, err
|
||||
}
|
||||
if err := chown(sub, uid, gid); err != nil {
|
||||
return false, err
|
||||
}
|
||||
rc, err := openSealed(ctx, b, key, uploadBlobKey(c.Digest))
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
defer rc.Close()
|
||||
tmp, err := os.CreateTemp(sub, UploadContextFile+".*.restore")
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
defer os.Remove(tmp.Name())
|
||||
_, err = io.Copy(tmp, newDigestReader(rc, c.Digest, c.Size))
|
||||
if err == nil {
|
||||
err = tmp.Sync()
|
||||
}
|
||||
if cerr := tmp.Close(); err == nil {
|
||||
err = cerr
|
||||
}
|
||||
if err == nil {
|
||||
err = os.Chmod(tmp.Name(), 0o660)
|
||||
}
|
||||
if err == nil {
|
||||
err = chown(tmp.Name(), uid, gid)
|
||||
}
|
||||
if err == nil {
|
||||
err = os.Rename(tmp.Name(), final)
|
||||
}
|
||||
return err == nil, err
|
||||
}
|
||||
|
||||
func chown(p string, uid, gid int) error {
|
||||
if uid < 0 && gid < 0 {
|
||||
return nil
|
||||
}
|
||||
return os.Lchown(p, uid, gid)
|
||||
}
|
||||
@@ -0,0 +1,198 @@
|
||||
package offsite
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func newUploadsSyncer(t *testing.T) (*Syncer, *memBucket, *time.Time, string) {
|
||||
t.Helper()
|
||||
dir := t.TempDir()
|
||||
b := newMemBucket()
|
||||
clock := now
|
||||
return &Syncer{
|
||||
Bucket: b, Catalog: &fakeCatalog{}, Key: testKey(t), UploadsDir: dir,
|
||||
Now: func() time.Time { return clock },
|
||||
}, b, &clock, dir
|
||||
}
|
||||
|
||||
func putContext(t *testing.T, dir, id string, data []byte, mtime time.Time) {
|
||||
t.Helper()
|
||||
sub := filepath.Join(dir, id)
|
||||
if err := os.MkdirAll(sub, 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
p := filepath.Join(sub, UploadContextFile)
|
||||
if err := os.WriteFile(p, data, 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.Chtimes(p, mtime, mtime); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func runSync(t *testing.T, s *Syncer) Result {
|
||||
t.Helper()
|
||||
res, err := s.Run(context.Background())
|
||||
if err != nil {
|
||||
t.Fatalf("Run: %v (%v)", err, res.Errors)
|
||||
}
|
||||
return res
|
||||
}
|
||||
|
||||
// TestSyncUploadsCopiesContexts: every submission's context is copied, two
|
||||
// identical ones once; a directory without its context and a name that is not
|
||||
// a submission id are left out. A second run with nothing new sends nothing.
|
||||
func TestSyncUploadsCopiesContexts(t *testing.T) {
|
||||
s, b, _, dir := newUploadsSyncer(t)
|
||||
pack := bytes.Repeat([]byte("modpack"), 5000)
|
||||
putContext(t, dir, "sub-aaa", pack, now.Add(-time.Hour))
|
||||
putContext(t, dir, "sub-bbb", pack, now.Add(-time.Hour))
|
||||
putContext(t, dir, "sub-ccc", []byte("another pack"), now.Add(-time.Hour))
|
||||
if err := os.MkdirAll(filepath.Join(dir, "sub-writing"), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.MkdirAll(filepath.Join(dir, "lost+found"), 0o755); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
res := runSync(t, s)
|
||||
if res.Uploads != 3 || res.UploadsUploaded != 2 {
|
||||
t.Fatalf("uploads=%d uploaded=%d, want 3 contexts in 2 objects", res.Uploads, res.UploadsUploaded)
|
||||
}
|
||||
if got := keysUnder(b, uploadBlobsDir); len(got) != 2 {
|
||||
t.Fatalf("blobs %v, want 2", got)
|
||||
}
|
||||
if got := keysUnder(b, uploadIndexDir); len(got) != 1 {
|
||||
t.Fatalf("index versions %v, want 1", got)
|
||||
}
|
||||
x, err := LoadUploadIndex(context.Background(), b, s.Key, res.UploadIndex)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if len(x.Contexts) != 3 || x.Contexts["sub-aaa"].Digest != x.Contexts["sub-bbb"].Digest {
|
||||
t.Fatalf("index %+v", x.Contexts)
|
||||
}
|
||||
|
||||
puts := b.puts
|
||||
res = runSync(t, s)
|
||||
if b.puts != puts || res.UploadsUploaded != 0 {
|
||||
t.Fatalf("second run put %d objects, uploaded %d; want nothing", b.puts-puts, res.UploadsUploaded)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSyncUploadsEmptyVolumeKeepsCopies: a volume that comes back empty (a
|
||||
// rebuilt host before its restore) records an empty version, keeps every
|
||||
// context for UploadHistory, and a restore refuses the empty version until
|
||||
// told which one to use. Past the history both go.
|
||||
func TestSyncUploadsEmptyVolumeKeepsCopies(t *testing.T) {
|
||||
s, b, clock, dir := newUploadsSyncer(t)
|
||||
putContext(t, dir, "sub-aaa", []byte("pack a"), now.Add(-time.Hour))
|
||||
putContext(t, dir, "sub-bbb", []byte("pack b"), now.Add(-time.Hour))
|
||||
first := runSync(t, s).UploadIndex
|
||||
|
||||
for _, id := range []string{"sub-aaa", "sub-bbb"} {
|
||||
if err := os.RemoveAll(filepath.Join(dir, id)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
*clock = clock.Add(time.Hour)
|
||||
res := runSync(t, s)
|
||||
if res.Uploads != 0 || res.UploadObjectsPruned != 0 {
|
||||
t.Fatalf("empty volume: uploads=%d pruned=%d", res.Uploads, res.UploadObjectsPruned)
|
||||
}
|
||||
if got := keysUnder(b, uploadBlobsDir); len(got) != 2 {
|
||||
t.Fatalf("blobs %v, want both kept", got)
|
||||
}
|
||||
ctx := context.Background()
|
||||
if _, _, err := ChooseUploadIndex(ctx, b, s.Key, ""); err == nil || !strings.Contains(err.Error(), first) {
|
||||
t.Fatalf("choosing the empty newest version: %v, want a refusal naming %s", err, first)
|
||||
}
|
||||
stamp, x, err := ChooseUploadIndex(ctx, b, s.Key, first)
|
||||
if err != nil || stamp != first || len(x.Contexts) != 2 {
|
||||
t.Fatalf("choose %s: %s %v %v", first, stamp, x, err)
|
||||
}
|
||||
|
||||
*clock = clock.Add(UploadHistory + time.Hour)
|
||||
res = runSync(t, s)
|
||||
if got := keysUnder(b, uploadBlobsDir); len(got) != 0 {
|
||||
t.Fatalf("blobs %v after the history, want none", got)
|
||||
}
|
||||
if got := keysUnder(b, uploadIndexDir); len(got) != 1 {
|
||||
t.Fatalf("index versions %v after the history, want the newest only", got)
|
||||
}
|
||||
if res.UploadObjectsPruned != 3 {
|
||||
t.Fatalf("pruned %d, want the old version and its 2 contexts", res.UploadObjectsPruned)
|
||||
}
|
||||
}
|
||||
|
||||
// TestSyncUploadsReplacedContext: a context rewritten in place is copied again;
|
||||
// the old bytes stay while the version that names them is kept.
|
||||
func TestSyncUploadsReplacedContext(t *testing.T) {
|
||||
s, b, clock, dir := newUploadsSyncer(t)
|
||||
putContext(t, dir, "sub-aaa", []byte("version one"), now.Add(-time.Hour))
|
||||
runSync(t, s)
|
||||
|
||||
putContext(t, dir, "sub-aaa", []byte("version two"), now)
|
||||
*clock = clock.Add(time.Hour)
|
||||
res := runSync(t, s)
|
||||
if res.UploadsUploaded != 1 {
|
||||
t.Fatalf("uploaded %d, want the new version", res.UploadsUploaded)
|
||||
}
|
||||
if got := keysUnder(b, uploadBlobsDir); len(got) != 2 {
|
||||
t.Fatalf("blobs %v, want old and new", got)
|
||||
}
|
||||
|
||||
*clock = clock.Add(UploadHistory + time.Hour)
|
||||
runSync(t, s)
|
||||
if got := keysUnder(b, uploadBlobsDir); len(got) != 1 {
|
||||
t.Fatalf("blobs %v after the history, want the new one only", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFetchUploadsRestores: a restore writes every context back in place,
|
||||
// leaves one that is already right alone, and rewrites one whose bytes differ.
|
||||
func TestFetchUploadsRestores(t *testing.T) {
|
||||
s, b, _, dir := newUploadsSyncer(t)
|
||||
a, c := bytes.Repeat([]byte("a"), 70000), []byte("pack c")
|
||||
putContext(t, dir, "sub-aaa", a, now.Add(-time.Hour))
|
||||
putContext(t, dir, "sub-ccc", c, now.Add(-time.Hour))
|
||||
stamp := runSync(t, s).UploadIndex
|
||||
|
||||
ctx := context.Background()
|
||||
x, err := LoadUploadIndex(ctx, b, s.Key, stamp)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
target := t.TempDir()
|
||||
res, err := FetchUploads(ctx, b, s.Key, x, target, -1, -1, nil)
|
||||
if err != nil || res.Written != 2 || res.Bytes != int64(len(a)+len(c)) {
|
||||
t.Fatalf("fetch: %+v %v", res, err)
|
||||
}
|
||||
for id, want := range map[string][]byte{"sub-aaa": a, "sub-ccc": c} {
|
||||
got, err := os.ReadFile(filepath.Join(target, id, UploadContextFile))
|
||||
if err != nil || !bytes.Equal(got, want) {
|
||||
t.Fatalf("%s restored as %d bytes (%v)", id, len(got), err)
|
||||
}
|
||||
}
|
||||
|
||||
if err := os.WriteFile(filepath.Join(target, "sub-ccc", UploadContextFile), []byte("pack X"), 0o644); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
res, err = FetchUploads(ctx, b, s.Key, x, target, -1, -1, nil)
|
||||
if err != nil || res.Written != 1 || res.Present != 1 {
|
||||
t.Fatalf("second fetch: %+v %v, want the damaged one rewritten", res, err)
|
||||
}
|
||||
if got, _ := os.ReadFile(filepath.Join(target, "sub-ccc", UploadContextFile)); !bytes.Equal(got, c) {
|
||||
t.Fatalf("sub-ccc is %q after the second fetch", got)
|
||||
}
|
||||
leftovers, _ := filepath.Glob(filepath.Join(target, "*", "*.restore"))
|
||||
if len(leftovers) != 0 {
|
||||
t.Fatalf("temporary files left behind: %v", leftovers)
|
||||
}
|
||||
}
|
||||
@@ -189,6 +189,15 @@ const (
|
||||
nonRootUID int64 = 1000
|
||||
)
|
||||
|
||||
// UploadsPVCName is felis-api's submission uploads PVC in the control-plane
|
||||
// namespace; the host-side off-site copy reads and restores it.
|
||||
const UploadsPVCName = uploadsPVCName
|
||||
|
||||
// ControlPlaneUID is the uid and gid the control-plane pods run as, which own
|
||||
// what they write to their volumes; a host-side restore into one of them
|
||||
// writes as it.
|
||||
const ControlPlaneUID = int(nonRootUID)
|
||||
|
||||
// APIInternalServiceName is the ClusterIP Service that fronts the felis-api
|
||||
// internal face (8081). It is SEPARATE from the external NodePort Service (SAAPI)
|
||||
// on purpose — see apiInternalService. The login pod resolves it by cross-namespace
|
||||
|
||||
@@ -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, database bundles and user images exist on this machine only",
|
||||
Summary: "异地备份从未成功同步过,世界归档、数据库备份、用户镜像与上传的整合包只在本机",
|
||||
SummaryEN: "the off-site copy has never completed; world archives, database bundles, user images and uploaded modpacks exist on this machine only",
|
||||
Hint: hint,
|
||||
}
|
||||
if st != nil && !st.LastSuccess.IsZero() {
|
||||
|
||||
Reference in new issue
Block a user