diff --git a/.gitignore b/.gitignore index b2f7876a..12a8a756 100644 --- a/.gitignore +++ b/.gitignore @@ -20,3 +20,8 @@ qemu/*.efi # c object *.o + +# session + migration dev cruft +.claude/ +scripts/migration.bak/ +scripts/migration/zosmigration diff --git a/client/node.go b/client/node.go index 281a169c..afcf471e 100644 --- a/client/node.go +++ b/client/node.go @@ -107,6 +107,56 @@ func (n *NodeClient) DeploymentDelete(ctx context.Context, contractID uint64) er return n.bus.Call(ctx, n.nodeTwin, cmd, in, nil) } +// TransferItem maps a workload (a zmount or a zmachine, by name) to a presigned +// S3 URL used to move its bytes between nodes. +type TransferItem struct { + WorkloadName gridtypes.Name `json:"workload_name"` + URL string `json:"url"` + // Size is the volume size to (re)create on the target for a rootfs volume + // download; ignored for uploads and for zmount downloads. + Size gridtypes.Unit `json:"size,omitempty"` +} + +// DeploymentTransfer asks the (old) node to upload the persistent bytes of the +// given contract's workloads (zmount disks and/or the VM rootfs writable layer) +// to the provided presigned S3 URLs. The node pauses the deployment first for a +// consistent copy. +func (n *NodeClient) DeploymentTransfer(ctx context.Context, contractID uint64, uploads []TransferItem) error { + const cmd = "zos.deployment.transfer" + in := args{ + "contract_id": contractID, + "uploads": uploads, + } + + return n.bus.Call(ctx, n.nodeTwin, cmd, in, nil) +} + +// DeploymentPrepare asks the (new) node to stage the deployment without starting +// the zmachine (provision network + zmount(s)) and to pull each listed workload's +// bytes from the provided presigned S3 URLs into the new zmount disk / rootfs +// volume. The VM is started by a follow-up call once data is in place. +func (n *NodeClient) DeploymentPrepare(ctx context.Context, dl gridtypes.Deployment, downloads []TransferItem, start bool) error { + const cmd = "zos.deployment.prepare" + in := args{ + "deployment": dl, + "downloads": downloads, + "start": start, + } + + return n.bus.Call(ctx, n.nodeTwin, cmd, in, nil) +} + +// DeploymentStart boots the zmachine(s) of a deployment previously staged on the +// (new) node via DeploymentPrepare, completing a contract move. +func (n *NodeClient) DeploymentStart(ctx context.Context, contractID uint64) error { + const cmd = "zos.deployment.start" + in := args{ + "contract_id": contractID, + } + + return n.bus.Call(ctx, n.nodeTwin, cmd, in, nil) +} + // Counters (statistics) of the node type Counters struct { // Total system capacity diff --git a/docs/deployment-move-testing.md b/docs/deployment-move-testing.md new file mode 100644 index 00000000..75090fde --- /dev/null +++ b/docs/deployment-move-testing.md @@ -0,0 +1,161 @@ +# Testing a deployment move between two nodes (zmount + rootfs) + +This is a dev runbook for the `zos.deployment.transfer` / `prepare` / `start` feature +(branch `feat/deployment-transfer`). It moves a zmachine (mycelium network + zmount) +from **node A** to **node B**, preserving zmount data *and* rootfs changes, via S3. + +## The flow at a glance + +``` +node A S3 (MinIO) node B + │ transfer (PUT) ──────────▶ zmount.raw │ + │ transfer (PUT) ──────────▶ rootfs.tar │ + │ │ prepare (GET) ◀────────┤ creates net+zmount+rootfs, pulls data + │ └───────────────────────▶ │ (VM NOT booted yet) + │ start ─────────┤ boots the VM on B + ▼ cancel old contract ▼ VM running on B +``` + +`transfer` pauses the source VM first, so the copy is consistent. + +## Prerequisites + +1. **Two zos-light nodes** A and B on the same farm/network, both reachable over RMB, + both running a node image **built from this branch** (see "Build the node image"). +2. **A twin + mnemonic** you control (the deployment owner), KYC-verified on the target + chain env — `prepare` runs the same validation as `deploy` (twin match, signature, + contract hash). Dev/QA net is fine. +3. **A MinIO (or S3) bucket** reachable from both nodes, plus the `mc` client to mint + presigned URLs. +4. **A running source deployment on A**: a mycelium network + a zmount + a zmachine that + mounts the zmount. Write a marker file into both the rootfs and the zmount so you can + prove they survived the move. + +## Build the node image + +zoslight is already wired to this branch via a local replace (added during implementation): + +``` +# in the zoslight repo +grep 'zos_base =>' go.mod +# replace github.com/threefoldtech/zos_base => /home/afouda/code/github/threefoldtech/zosbase +go build ./... # sanity: compiles against the branch +``` + +Build/flash the node image the way you normally do for zoslight (the replace makes the +image include these RMB endpoints). For a throwaway test you can also run the affected +daemons (api_gateway, provisiond, storaged) from source on a dev node. + +> Before pushing the branch, drop the local `replace` and pin a real `zos_base` +> pseudo-version instead — the path replace only works on your machine. + +## Mint presigned URLs (MinIO example) + +```bash +mc alias set m http://MINIO:9000 ACCESS SECRET +mc mb m/zos-move + +# uploads (used by node A on transfer) — PUT +ZMOUNT_PUT=$(mc share upload --expire 12h m/zos-move/zmount.raw | awk '/share:/{print $2}') +ROOTFS_PUT=$(mc share upload --expire 12h m/zos-move/rootfs.tar | awk '/share:/{print $2}') +# downloads (used by node B on prepare) — GET +ZMOUNT_GET=$(mc share download --expire 12h m/zos-move/zmount.raw | awk '/share:/{print $2}') +ROOTFS_GET=$(mc share download --expire 12h m/zos-move/rootfs.tar | awk '/share:/{print $2}') +``` + +The node streams a definite Content-Length, so plain presigned PUT/GET work (no chunked). + +## Driver program (uses the client SDK added on this branch) + +`client.NodeClient` gained `DeploymentTransfer`, `DeploymentPrepare`, `DeploymentStart`. +Workload names below (`"data"`, `"vm"`) are the `Name` fields in *your* deployment. + +```go +package main + +import ( + "context" + "time" + + "github.com/threefoldtech/zos_base/client" + "github.com/threefoldtech/zos_base/pkg/gridtypes" + "github.com/threefoldtech/zos_sdk_go/rmb-sdk-go" +) + +func main() { + cl, err := rmb.Default() // or a peer-backed rmb.Client for your env + if err != nil { panic(err) } + + const ( + twin = 1234 // your twin id + nodeATwin = 1111 // node A twin id + nodeBTwin = 2222 // node B twin id + contractA = 100 // running contract on A + ) + + nodeA := client.NewNodeClient(nodeATwin, cl) + nodeB := client.NewNodeClient(nodeBTwin, cl) + ctx := context.Background() + + // 0) grab the running deployment from A (you'll reuse it for prepare on B) + dl, err := nodeA.DeploymentGet(ctx, contractA) + if err != nil { panic(err) } + + // 1) A: pause + upload zmount + rootfs to S3 + must(nodeA.DeploymentTransfer(ctx, contractA, []client.TransferItem{ + {WorkloadName: "data", URL: ZMOUNT_PUT}, // zmount + {WorkloadName: "vm", URL: ROOTFS_PUT}, // zmachine rootfs + })) + + // 2) stage a contract on B with the SAME deployment hash (see note below), + // set dl.ContractID = , dl.Version = 0. + dl.ContractID = 200 + dl.Version = 0 + + // 3) B: prepare (provision net+zmount, skip VM) and pull data back + must(nodeB.DeploymentPrepare(ctx, dl, []client.TransferItem{ + {WorkloadName: "data", URL: ZMOUNT_GET}, + {WorkloadName: "vm", URL: ROOTFS_GET, Size: rootfsSize}, // rootfs volume size + })) + + // 4) B: boot the VM + time.Sleep(2 * time.Second) + must(nodeB.DeploymentStart(ctx, dl.ContractID)) + + // 5) A: cancel the old contract on chain to decommission the source +} + +func must(err error) { if err != nil { panic(err) } } +``` + +### The on-chain contract for B (the manual bit) + +`prepare` validates against a node contract for **B** whose `DeploymentHash` matches the +deployment. In production the grid/tfchain orchestration creates it; for a manual test, +create a node contract on B for the same hash (`dl.ChallengeHash()`) using your usual tool +(tfchain client / `tfcmd` / grid client), then use that contract id as `dl.ContractID` +in step 2. Same-contract migration (updating the existing contract's node id on chain) is +a later refinement — for now use a fresh contract on B. + +## Verify + +- `nodeB.DeploymentGet(ctx, 200)` → zmount + network `StateOk`, zmachine `StateOk` after start. +- On node B host: the zmount file exists at `/mnt//vdisks/-200-data` and the + rootfs volume at `/mnt//rootfs:-200-vm/rw`. +- Console into the VM (mycelium IP is derived from the deployment seed, so it's the same as + on A) and confirm your **marker files** in both the rootfs and the mounted zmount. + +## Cleanup + +Cancel the source contract on A (normal decommission path); the node tears down the paused +deployment. Delete the S3 objects. + +## Known limitations (this phase) + +- `transfer` pauses vCPUs but does not `fsync` the guest — tiny consistency risk; for a + clean test, quiesce the app inside the VM before transferring. +- The rootfs volume quota on B is set at `prepare` time from the `Size` you pass; make it + match the VM's rootfs size. +- Same-contract (in-place node reassignment) orchestration is not implemented — use a new + contract on B as above. +``` diff --git a/pkg/provision.go b/pkg/provision.go index 75b9702b..6adaddeb 100644 --- a/pkg/provision.go +++ b/pkg/provision.go @@ -15,6 +15,12 @@ type Provision interface { // GetWorkloadStatus: returns status, bool(true if workload exits otherwise it is false), error GetWorkloadStatus(id string) (gridtypes.ResultState, bool, error) CreateOrUpdate(twin uint32, deployment gridtypes.Deployment, update bool) error + // PrepareDeployment stages a deployment (provisions all but the zmachine) for a contract move + PrepareDeployment(twin uint32, deployment gridtypes.Deployment) error + // PauseDeployment pauses a deployment (used to freeze a source VM before transferring its data) + PauseDeployment(twin uint32, id uint64) error + // StartDeployment boots the zmachine(s) of a deployment previously staged via PrepareDeployment + StartDeployment(twin uint32, id uint64) error Get(twin uint32, contractID uint64) (gridtypes.Deployment, error) List(twin uint32) ([]gridtypes.Deployment, error) Changes(twin uint32, contractID uint64) ([]gridtypes.Workload, error) diff --git a/pkg/provision/engine.go b/pkg/provision/engine.go index ba1e1387..f2264587 100644 --- a/pkg/provision/engine.go +++ b/pkg/provision/engine.go @@ -92,6 +92,14 @@ const ( opPause // opResume resumes a deployment opResume + // opPrepare provisions all workloads of a deployment EXCEPT the zmachine + // types, leaving the VM unbooted. Used to stage a deployment on a target + // node during a contract move (network + zmount + rootfs volume are created + // so their data can be pulled in before the VM is started). + opPrepare + // opStart installs ONLY the zmachine workloads of a previously prepared + // deployment, booting the VM after its data has been pulled in. + opStart // servers default timeout defaultHttpTimeout = 10 * time.Second ) @@ -368,6 +376,33 @@ func (e *NativeEngine) Provision(ctx context.Context, deployment gridtypes.Deplo return e.queue.Enqueue(&job) } +// Prepare stages a deployment on this node: it persists the whole deployment +// (so the zmachine is stored and can be started later) and provisions every +// workload EXCEPT the zmachine types, leaving the VM unbooted. This is used +// during a contract move so the target node can create the network, zmount(s) +// and rootfs volume and pull their data before starting the VM. +func (e *NativeEngine) Prepare(ctx context.Context, deployment gridtypes.Deployment) error { + log.Info(). + Uint32("twin", deployment.TwinID). + Uint64("contract", deployment.ContractID). + Msg("scheduling deployment for prepare") + + if deployment.Version != 0 { + return errors.Wrap(ErrInvalidVersion, "expected version to be 0 on deployment creation") + } + + if err := e.storage.Create(deployment); err != nil { + return err + } + + job := engineJob{ + Target: deployment, + Op: opPrepare, + } + + return e.queue.Enqueue(&job) +} + // Pause deployment func (e *NativeEngine) Pause(ctx context.Context, twin uint32, id uint64) error { deployment, err := e.storage.Get(twin, id) @@ -408,6 +443,28 @@ func (e *NativeEngine) Resume(ctx context.Context, twin uint32, id uint64) error return e.queue.Enqueue(&job) } +// Start boots the zmachine(s) of a previously prepared deployment (see opPrepare). +// It installs only the zmachine workloads that were skipped during prepare, +// reusing the network, zmount(s) and rootfs volume already staged on this node. +func (e *NativeEngine) Start(ctx context.Context, twin uint32, id uint64) error { + deployment, err := e.storage.Get(twin, id) + if err != nil { + return err + } + + log.Info(). + Uint32("twin", deployment.TwinID). + Uint64("contract", deployment.ContractID). + Msg("schedule for start") + + job := engineJob{ + Target: deployment, + Op: opStart, + } + + return e.queue.Enqueue(&job) +} + // Deprovision workload func (e *NativeEngine) Deprovision(ctx context.Context, twin uint32, id uint64, reason string) error { deployment, err := e.storage.Get(twin, id) @@ -513,6 +570,7 @@ func (e *NativeEngine) Run(root context.Context) error { // this should ONLY be done on provosion and update operation if job.Op == opProvision || job.Op == opUpdate || + job.Op == opPrepare || job.Op == opProvisionNoValidation { // otherwise, contract validation is needed ctx, err = e.validate(ctx, &job.Target, job.Op == opProvisionNoValidation) @@ -535,6 +593,10 @@ func (e *NativeEngine) Run(root context.Context) error { fallthrough case opProvision: e.installDeployment(ctx, &job.Target) + case opPrepare: + e.prepareDeployment(ctx, &job.Target) + case opStart: + e.startDeployment(ctx, &job.Target) case opDeprovision: e.uninstallDeployment(ctx, &job.Target, job.Message) case opPause: @@ -927,6 +989,46 @@ func (e *NativeEngine) installDeployment(ctx context.Context, getter gridtypes.W } } +// prepareDeployment installs every workload of a deployment EXCEPT the zmachine +// types, leaving the VM unbooted. Same iteration/order as installDeployment. +func (e *NativeEngine) prepareDeployment(ctx context.Context, getter gridtypes.WorkloadGetter) { + for _, typ := range e.order { + if typ == zos.ZMachineType || typ == zos.ZMachineLightType { + continue + } + + workloads := getter.ByType(typ) + + if typ == zos.ZMountType || typ == zos.VolumeType { + sortMountWorkloads(workloads) + } + + for _, wl := range workloads { + if err := e.installWorkload(ctx, wl); err != nil { + log.Error().Err(err).Stringer("id", wl.ID).Msg("failed to install workload during prepare") + } + } + } +} + +// startDeployment installs ONLY the zmachine workloads of a deployment. It is the +// counterpart of prepareDeployment: prepare stages everything except the VM, then +// start boots the VM once its migrated data is in place. installWorkload reads the +// stored StateInit record and provisions the VM (the real boot path). +func (e *NativeEngine) startDeployment(ctx context.Context, getter gridtypes.WorkloadGetter) { + for _, typ := range e.order { + if typ != zos.ZMachineType && typ != zos.ZMachineLightType { + continue + } + + for _, wl := range getter.ByType(typ) { + if err := e.installWorkload(ctx, wl); err != nil { + log.Error().Err(err).Stringer("id", wl.ID).Msg("failed to install workload during start") + } + } + } +} + func (e *NativeEngine) lockDeployment(ctx context.Context, getter gridtypes.WorkloadGetter) { for i := len(e.order) - 1; i >= 0; i-- { typ := e.order[i] @@ -1066,6 +1168,56 @@ func (n *NativeEngine) CreateOrUpdate(twin uint32, deployment gridtypes.Deployme return action(ctx, deployment) } +// PrepareDeployment is the zbus/RMB-facing entry to stage a deployment on this +// node during a contract move. It runs the same validation as CreateOrUpdate +// (ownership, KYC, signatures) then provisions everything except the zmachine +// via the internal Prepare (see opPrepare). The VM is started later by a +// separate call once the migrated data has been pulled in. +func (n *NativeEngine) PrepareDeployment(twin uint32, deployment gridtypes.Deployment) error { + if err := deployment.Valid(); err != nil { + return err + } + + if deployment.TwinID != twin { + return fmt.Errorf("twin id mismatch (deployment: %d, message: %d)", deployment.TwinID, twin) + } + + check := func() error { + if ok, err := isTwinVerified(twin); err != nil { + return err + } else if !ok { + return fmt.Errorf("user with twin id %d is not verified", twin) + } + return nil + } + + if err := backoff.Retry(check, backoff.WithMaxRetries(backoff.NewExponentialBackOff(), 5)); err != nil { + return err + } + + if err := deployment.Verify(n.twins); err != nil { + return err + } + + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute) + defer cancel() + + return n.Prepare(ctx, deployment) +} + +// PauseDeployment is the zbus/RMB-facing entry to pause a deployment (freezes a +// source VM so its disks/rootfs can be transferred consistently). It wraps the +// internal Pause, which scopes to the given twin's storage. +func (n *NativeEngine) PauseDeployment(twin uint32, id uint64) error { + return n.Pause(context.Background(), twin, id) +} + +// StartDeployment is the zbus/RMB-facing entry to boot the zmachine(s) of a +// deployment previously staged via PrepareDeployment. It wraps the internal Start. +func (n *NativeEngine) StartDeployment(twin uint32, id uint64) error { + return n.Start(context.Background(), twin, id) +} + func (n *NativeEngine) Get(twin uint32, contractID uint64) (gridtypes.Deployment, error) { deployment, err := n.storage.Get(twin, contractID) if errors.Is(err, ErrDeploymentNotExists) { diff --git a/pkg/provision/interface.go b/pkg/provision/interface.go index 297e0f70..6805735c 100644 --- a/pkg/provision/interface.go +++ b/pkg/provision/interface.go @@ -18,8 +18,12 @@ type Engine interface { // means that workload has been committed to storage (accepts) // and will be processes later Provision(ctx context.Context, wl gridtypes.Deployment) error + // Prepare provisions every workload except the zmachine (leaves the VM unbooted) + Prepare(ctx context.Context, wl gridtypes.Deployment) error Deprovision(ctx context.Context, twin uint32, id uint64, reason string) error Pause(ctx context.Context, twin uint32, id uint64) error + // Start boots the zmachine(s) of a previously prepared deployment + Start(ctx context.Context, twin uint32, id uint64) error Resume(ctx context.Context, twin uint32, id uint64) error Update(ctx context.Context, update gridtypes.Deployment) error Storage() Storage diff --git a/pkg/storage.go b/pkg/storage.go index a3b4ed07..d626b905 100644 --- a/pkg/storage.go +++ b/pkg/storage.go @@ -146,6 +146,21 @@ type StorageModule interface { DiskDelete(name string) error DiskList() ([]VDisk, error) + + // Migration transfer (deployment move between nodes) + + // DiskUpload streams the raw vdisk `id` to the presigned URL via HTTP PUT + DiskUpload(id string, url string) error + + // DiskDownload downloads from the presigned URL into the existing vdisk `id` + DiskDownload(id string, url string) error + + // VolumeUpload tars the subvolume `name` and streams it to the presigned URL via HTTP PUT + VolumeUpload(name string, url string) error + + // VolumeDownload ensures a subvolume `name` of `size` exists then extracts the tar from url into it + VolumeDownload(name string, size gridtypes.Unit, url string) error + // Device management //Devices list all "allocated" devices diff --git a/pkg/storage/transfer.go b/pkg/storage/transfer.go new file mode 100644 index 00000000..c36d4f03 --- /dev/null +++ b/pkg/storage/transfer.go @@ -0,0 +1,55 @@ +package storage + +import ( + "github.com/pkg/errors" + log "github.com/rs/zerolog/log" + "github.com/threefoldtech/zos_base/pkg/gridtypes" + "github.com/threefoldtech/zos_base/pkg/storagetransfer" +) + +// DiskUpload streams the raw vdisk file identified by id to the given presigned +// URL using an HTTP PUT. Used to move a zmount (or a full-VM boot disk) to another +// node during a deployment transfer. +func (s *Module) DiskUpload(id string, url string) error { + disk, err := s.DiskLookup(id) + if err != nil { + return errors.Wrapf(err, "failed to lookup disk '%s'", id) + } + log.Info().Str("id", id).Str("path", disk.Path).Msg("uploading disk") + return storagetransfer.UploadFile(disk.Path, url) +} + +// DiskDownload downloads bytes from the given presigned URL and writes them into +// the existing vdisk file identified by id. The disk must already exist (it is +// created when the zmount workload is provisioned during prepare). +func (s *Module) DiskDownload(id string, url string) error { + disk, err := s.DiskLookup(id) + if err != nil { + return errors.Wrapf(err, "failed to lookup disk '%s'", id) + } + log.Info().Str("id", id).Str("path", disk.Path).Msg("downloading disk") + return storagetransfer.DownloadToFile(disk.Path, url) +} + +// VolumeUpload tars the subvolume named `name` and streams it to the given +// presigned URL using an HTTP PUT. Used to move a VM's writable rootfs layer +// (subvolume "rootfs:") to another node. +func (s *Module) VolumeUpload(name string, url string) error { + vol, err := s.VolumeLookup(name) + if err != nil { + return errors.Wrapf(err, "failed to lookup volume '%s'", name) + } + log.Info().Str("name", name).Str("path", vol.Path).Msg("uploading volume") + return storagetransfer.UploadDir(vol.Path, url) +} + +// VolumeDownload makes sure a subvolume `name` of the given size exists, then +// downloads a tar from the presigned URL and extracts it into the subvolume. +func (s *Module) VolumeDownload(name string, size gridtypes.Unit, url string) error { + vol, err := s.VolumeCreate(name, size) + if err != nil { + return errors.Wrapf(err, "failed to create volume '%s'", name) + } + log.Info().Str("name", name).Str("path", vol.Path).Msg("downloading volume") + return storagetransfer.DownloadDir(vol.Path, url) +} diff --git a/pkg/storage_light/transfer.go b/pkg/storage_light/transfer.go new file mode 100644 index 00000000..c36d4f03 --- /dev/null +++ b/pkg/storage_light/transfer.go @@ -0,0 +1,55 @@ +package storage + +import ( + "github.com/pkg/errors" + log "github.com/rs/zerolog/log" + "github.com/threefoldtech/zos_base/pkg/gridtypes" + "github.com/threefoldtech/zos_base/pkg/storagetransfer" +) + +// DiskUpload streams the raw vdisk file identified by id to the given presigned +// URL using an HTTP PUT. Used to move a zmount (or a full-VM boot disk) to another +// node during a deployment transfer. +func (s *Module) DiskUpload(id string, url string) error { + disk, err := s.DiskLookup(id) + if err != nil { + return errors.Wrapf(err, "failed to lookup disk '%s'", id) + } + log.Info().Str("id", id).Str("path", disk.Path).Msg("uploading disk") + return storagetransfer.UploadFile(disk.Path, url) +} + +// DiskDownload downloads bytes from the given presigned URL and writes them into +// the existing vdisk file identified by id. The disk must already exist (it is +// created when the zmount workload is provisioned during prepare). +func (s *Module) DiskDownload(id string, url string) error { + disk, err := s.DiskLookup(id) + if err != nil { + return errors.Wrapf(err, "failed to lookup disk '%s'", id) + } + log.Info().Str("id", id).Str("path", disk.Path).Msg("downloading disk") + return storagetransfer.DownloadToFile(disk.Path, url) +} + +// VolumeUpload tars the subvolume named `name` and streams it to the given +// presigned URL using an HTTP PUT. Used to move a VM's writable rootfs layer +// (subvolume "rootfs:") to another node. +func (s *Module) VolumeUpload(name string, url string) error { + vol, err := s.VolumeLookup(name) + if err != nil { + return errors.Wrapf(err, "failed to lookup volume '%s'", name) + } + log.Info().Str("name", name).Str("path", vol.Path).Msg("uploading volume") + return storagetransfer.UploadDir(vol.Path, url) +} + +// VolumeDownload makes sure a subvolume `name` of the given size exists, then +// downloads a tar from the presigned URL and extracts it into the subvolume. +func (s *Module) VolumeDownload(name string, size gridtypes.Unit, url string) error { + vol, err := s.VolumeCreate(name, size) + if err != nil { + return errors.Wrapf(err, "failed to create volume '%s'", name) + } + log.Info().Str("name", name).Str("path", vol.Path).Msg("downloading volume") + return storagetransfer.DownloadDir(vol.Path, url) +} diff --git a/pkg/storagetransfer/transfer.go b/pkg/storagetransfer/transfer.go new file mode 100644 index 00000000..c00a1c3a --- /dev/null +++ b/pkg/storagetransfer/transfer.go @@ -0,0 +1,289 @@ +// Package storagetransfer provides HTTP streaming of vdisk files and subvolumes +// to/from presigned S3 URLs. It is used to move a deployment's persistent bytes +// (zmount disks and VM rootfs layers) between nodes during a contract move. +// The helpers are shared by both the storage and storage_light modules. +package storagetransfer + +import ( + "archive/tar" + "context" + "fmt" + "io" + "net/http" + "os" + "path/filepath" + "syscall" + "time" + + "github.com/pkg/errors" + log "github.com/rs/zerolog/log" +) + +// httpTimeout bounds a single upload/download. vdisks/rootfs images can be large, +// so this is generous; the caller's presigned URL usually expires sooner. +const httpTimeout = 6 * time.Hour + +// UploadFile PUTs a seekable file so net/http sets a definite Content-Length +// (presigned S3/MinIO PUT rejects chunked transfer encoding). +func UploadFile(path, url string) error { + f, err := os.Open(path) + if err != nil { + return err + } + defer f.Close() + + st, err := f.Stat() + if err != nil { + return err + } + return httpPut(url, f, st.Size()) +} + +// DownloadToFile GETs url and writes the body into an existing file, truncating it. +func DownloadToFile(path, url string) error { + body, err := httpGet(url) + if err != nil { + return err + } + defer body.Close() + + f, err := os.OpenFile(path, os.O_WRONLY|os.O_TRUNC, 0600) + if err != nil { + return err + } + defer f.Close() + + if _, err := io.Copy(f, body); err != nil { + return err + } + return f.Sync() +} + +// UploadDir streams a tar of `dir` to url. It uses two passes so it needs no +// temp file (nodes run from a small tmpfs): pass 1 tars to a byte counter to get +// the exact Content-Length (presigned PUT rejects chunked encoding), pass 2 tars +// straight into the request body via an io.Pipe. The source is a paused VM's +// rootfs, so it is static between the two passes. +func UploadDir(dir, url string) error { + var counter countWriter + if err := writeTar(&counter, dir); err != nil { + return errors.Wrap(err, "failed to size volume tar") + } + + pr, pw := io.Pipe() + go func() { + pw.CloseWithError(writeTar(pw, dir)) + }() + defer pr.Close() + + return httpPut(url, pr, counter.n) +} + +// countWriter counts bytes written to it and discards them. +type countWriter struct{ n int64 } + +func (c *countWriter) Write(p []byte) (int, error) { + c.n += int64(len(p)) + return len(p), nil +} + +// DownloadDir GETs url and extracts the tar stream into `dir`. +func DownloadDir(dir, url string) error { + body, err := httpGet(url) + if err != nil { + return err + } + defer body.Close() + return extractTar(body, dir) +} + +func httpPut(url string, body io.Reader, size int64) error { + ctx, cancel := context.WithTimeout(context.Background(), httpTimeout) + defer cancel() + + req, err := http.NewRequestWithContext(ctx, http.MethodPut, url, body) + if err != nil { + return err + } + // explicit length => no chunked transfer-encoding, which presigned PUT rejects + req.ContentLength = size + req.Header.Set("Content-Type", "application/octet-stream") + + resp, err := http.DefaultClient.Do(req) + if err != nil { + return err + } + defer resp.Body.Close() + + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return fmt.Errorf("upload failed with status %s: %s", resp.Status, string(b)) + } + _, _ = io.Copy(io.Discard, resp.Body) + return nil +} + +// httpGet returns the response body of a GET request. The returned ReadCloser +// owns the request context and cancels it on Close. +func httpGet(url string) (io.ReadCloser, error) { + ctx, cancel := context.WithTimeout(context.Background(), httpTimeout) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + cancel() + return nil, err + } + resp, err := http.DefaultClient.Do(req) + if err != nil { + cancel() + return nil, err + } + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + b, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + resp.Body.Close() + cancel() + return nil, fmt.Errorf("download failed with status %s: %s", resp.Status, string(b)) + } + return &cancelReadCloser{ReadCloser: resp.Body, cancel: cancel}, nil +} + +type cancelReadCloser struct { + io.ReadCloser + cancel context.CancelFunc +} + +func (c *cancelReadCloser) Close() error { + err := c.ReadCloser.Close() + c.cancel() + return err +} + +// writeTar walks `dir` and writes every entry (relative to dir) to w. It preserves +// regular files, directories, symlinks and character devices. Character devices +// matter because an overlayfs upper dir represents deleted files as 0/0 whiteout +// char devices; losing them would resurrect deleted files on the target node. +func writeTar(w io.Writer, dir string) error { + tw := tar.NewWriter(w) + defer tw.Close() + + return filepath.Walk(dir, func(path string, info os.FileInfo, err error) error { + if err != nil { + return err + } + + rel, err := filepath.Rel(dir, path) + if err != nil { + return err + } + if rel == "." { + return nil + } + + var link string + if info.Mode()&os.ModeSymlink != 0 { + if link, err = os.Readlink(path); err != nil { + return err + } + } + + hdr, err := tar.FileInfoHeader(info, link) + if err != nil { + return err + } + hdr.Name = rel + + // preserve device numbers for char/block devices (e.g. overlay whiteouts) + if st, ok := info.Sys().(*syscall.Stat_t); ok && + (info.Mode()&os.ModeCharDevice != 0 || info.Mode()&os.ModeDevice != 0) { + hdr.Devmajor = int64(unixMajor(uint64(st.Rdev))) + hdr.Devminor = int64(unixMinor(uint64(st.Rdev))) + } + + if err := tw.WriteHeader(hdr); err != nil { + return err + } + + if info.Mode().IsRegular() { + f, err := os.Open(path) + if err != nil { + return err + } + defer f.Close() + if _, err := io.Copy(tw, f); err != nil { + return err + } + } + return nil + }) +} + +// extractTar extracts a tar stream into `dir`, recreating files, directories, +// symlinks and character/block devices. +func extractTar(r io.Reader, dir string) error { + tr := tar.NewReader(r) + for { + hdr, err := tr.Next() + if err == io.EOF { + return nil + } + if err != nil { + return err + } + + target := filepath.Join(dir, filepath.Clean("/"+hdr.Name)) + if err := extractEntry(tr, hdr, target); err != nil { + return errors.Wrapf(err, "failed to extract '%s'", hdr.Name) + } + } +} + +func extractEntry(tr *tar.Reader, hdr *tar.Header, target string) error { + mode := hdr.FileInfo().Mode() + switch hdr.Typeflag { + case tar.TypeDir: + return os.MkdirAll(target, mode.Perm()) + case tar.TypeReg: + if err := os.MkdirAll(filepath.Dir(target), 0755); err != nil { + return err + } + f, err := os.OpenFile(target, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, mode.Perm()) + if err != nil { + return err + } + defer f.Close() + _, err = io.Copy(f, tr) + return err + case tar.TypeSymlink: + _ = os.Remove(target) + return os.Symlink(hdr.Linkname, target) + case tar.TypeLink: + return os.Link(filepath.Join(filepath.Dir(target), hdr.Linkname), target) + case tar.TypeChar, tar.TypeBlock: + kind := uint32(syscall.S_IFCHR) + if hdr.Typeflag == tar.TypeBlock { + kind = syscall.S_IFBLK + } + dev := unixMkdev(uint32(hdr.Devmajor), uint32(hdr.Devminor)) + _ = os.Remove(target) + return syscall.Mknod(target, kind|uint32(mode.Perm()), int(dev)) + default: + // fifo/socket and anything else are not expected in a rootfs overlay; skip + log.Debug().Str("name", hdr.Name).Uint8("type", hdr.Typeflag).Msg("skipping unsupported tar entry") + return nil + } +} + +// Linux dev_t helpers (glibc encoding), avoiding a dependency on x/sys/unix. +func unixMajor(dev uint64) uint32 { + return uint32((dev>>8)&0xfff) | uint32((dev>>32)&^uint64(0xfff)) +} + +func unixMinor(dev uint64) uint32 { + return uint32(dev&0xff) | uint32((dev>>12)&^uint64(0xff)) +} + +func unixMkdev(major, minor uint32) uint64 { + return (uint64(major&0xfff) << 8) | + uint64(minor&0xff) | + (uint64(minor&^uint32(0xff)) << 12) +} diff --git a/pkg/storagetransfer/transfer_test.go b/pkg/storagetransfer/transfer_test.go new file mode 100644 index 00000000..f333c5c3 --- /dev/null +++ b/pkg/storagetransfer/transfer_test.go @@ -0,0 +1,157 @@ +package storagetransfer + +import ( + "bytes" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "sync" + "testing" +) + +// objectStore is a minimal in-memory stand-in for a presigned S3 endpoint: +// PUT stores the body under the request path, GET serves it back. +func newObjectStore(t *testing.T) *httptest.Server { + t.Helper() + var mu sync.Mutex + store := map[string][]byte{} + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + defer mu.Unlock() + switch r.Method { + case http.MethodPut: + // presigned PUT must not be chunked + if r.ContentLength < 0 { + http.Error(w, "missing content-length", http.StatusNotImplemented) + return + } + b, err := io.ReadAll(r.Body) + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + store[r.URL.Path] = b + w.WriteHeader(http.StatusOK) + case http.MethodGet: + b, ok := store[r.URL.Path] + if !ok { + http.Error(w, "not found", http.StatusNotFound) + return + } + _, _ = w.Write(b) + default: + http.Error(w, "unexpected method", http.StatusMethodNotAllowed) + } + })) + t.Cleanup(srv.Close) + return srv +} + +func TestFileRoundTrip(t *testing.T) { + srv := newObjectStore(t) + dir := t.TempDir() + + data := bytes.Repeat([]byte("zmount-bytes-0123456789"), 5000) // ~115KB + src := filepath.Join(dir, "disk.raw") + if err := os.WriteFile(src, data, 0600); err != nil { + t.Fatal(err) + } + + if err := UploadFile(src, srv.URL+"/disk"); err != nil { + t.Fatalf("upload: %v", err) + } + + // DownloadToFile writes into an existing file (as a provisioned vdisk would be) + dst := filepath.Join(dir, "dst.raw") + if err := os.WriteFile(dst, nil, 0600); err != nil { + t.Fatal(err) + } + if err := DownloadToFile(dst, srv.URL+"/disk"); err != nil { + t.Fatalf("download: %v", err) + } + + got, err := os.ReadFile(dst) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got, data) { + t.Fatalf("round-trip mismatch: got %d bytes, want %d", len(got), len(data)) + } +} + +func TestDirRoundTrip(t *testing.T) { + srv := newObjectStore(t) + dir := t.TempDir() + + // mimic a rootfs subvolume: rw/ (upperdir with user changes), wd/ (workdir) + src := filepath.Join(dir, "vol") + mustMkdir(t, filepath.Join(src, "rw", "etc")) + mustMkdir(t, filepath.Join(src, "wd")) + mustWrite(t, filepath.Join(src, "rw", "etc", "hostname"), "node-b\n") + mustWrite(t, filepath.Join(src, "rw", "data.bin"), "user-modified") + if err := os.Symlink("data.bin", filepath.Join(src, "rw", "link")); err != nil { + t.Fatal(err) + } + + if err := UploadDir(src, srv.URL+"/vol"); err != nil { + t.Fatalf("upload dir: %v", err) + } + + dst := filepath.Join(dir, "voldst") + mustMkdir(t, dst) + if err := DownloadDir(dst, srv.URL+"/vol"); err != nil { + t.Fatalf("download dir: %v", err) + } + + if got := mustRead(t, filepath.Join(dst, "rw", "etc", "hostname")); got != "node-b\n" { + t.Fatalf("hostname mismatch: %q", got) + } + if got := mustRead(t, filepath.Join(dst, "rw", "data.bin")); got != "user-modified" { + t.Fatalf("data.bin mismatch: %q", got) + } + link, err := os.Readlink(filepath.Join(dst, "rw", "link")) + if err != nil || link != "data.bin" { + t.Fatalf("symlink not restored: link=%q err=%v", link, err) + } + if fi, err := os.Stat(filepath.Join(dst, "wd")); err != nil || !fi.IsDir() { + t.Fatalf("workdir not restored: err=%v", err) + } +} + +func TestDownloadErrorStatus(t *testing.T) { + srv := newObjectStore(t) + dst := filepath.Join(t.TempDir(), "x.raw") + if err := os.WriteFile(dst, nil, 0600); err != nil { + t.Fatal(err) + } + // nothing was uploaded => GET returns 404 => must error + if err := DownloadToFile(dst, srv.URL+"/missing"); err == nil { + t.Fatal("expected error on 404 download") + } +} + +func mustMkdir(t *testing.T, path string) { + t.Helper() + if err := os.MkdirAll(path, 0755); err != nil { + t.Fatal(err) + } +} + +func mustWrite(t *testing.T, path, content string) { + t.Helper() + if err := os.WriteFile(path, []byte(content), 0644); err != nil { + t.Fatal(err) + } +} + +func mustRead(t *testing.T, path string) string { + t.Helper() + b, err := os.ReadFile(path) + if err != nil { + t.Fatal(err) + } + return string(b) +} diff --git a/pkg/stubs/provision_stub.go b/pkg/stubs/provision_stub.go index 6448169f..c9ef66f9 100644 --- a/pkg/stubs/provision_stub.go +++ b/pkg/stubs/provision_stub.go @@ -176,3 +176,48 @@ func (s *ProvisionStub) ListTwins(ctx context.Context) (ret0 []uint32, ret1 erro } return } + +func (s *ProvisionStub) PauseDeployment(ctx context.Context, arg0 uint32, arg1 uint64) (ret0 error) { + args := []interface{}{arg0, arg1} + result, err := s.client.RequestContext(ctx, s.module, s.object, "PauseDeployment", args...) + if err != nil { + panic(err) + } + result.PanicOnError() + ret0 = result.CallError() + loader := zbus.Loader{} + if err := result.Unmarshal(&loader); err != nil { + panic(err) + } + return +} + +func (s *ProvisionStub) PrepareDeployment(ctx context.Context, arg0 uint32, arg1 gridtypes.Deployment) (ret0 error) { + args := []interface{}{arg0, arg1} + result, err := s.client.RequestContext(ctx, s.module, s.object, "PrepareDeployment", args...) + if err != nil { + panic(err) + } + result.PanicOnError() + ret0 = result.CallError() + loader := zbus.Loader{} + if err := result.Unmarshal(&loader); err != nil { + panic(err) + } + return +} + +func (s *ProvisionStub) StartDeployment(ctx context.Context, arg0 uint32, arg1 uint64) (ret0 error) { + args := []interface{}{arg0, arg1} + result, err := s.client.RequestContext(ctx, s.module, s.object, "StartDeployment", args...) + if err != nil { + panic(err) + } + result.PanicOnError() + ret0 = result.CallError() + loader := zbus.Loader{} + if err := result.Unmarshal(&loader); err != nil { + panic(err) + } + return +} diff --git a/pkg/stubs/storage_stub.go b/pkg/stubs/storage_stub.go index 6e6928bb..655fd3f5 100644 --- a/pkg/stubs/storage_stub.go +++ b/pkg/stubs/storage_stub.go @@ -161,6 +161,21 @@ func (s *StorageModuleStub) DiskDelete(ctx context.Context, arg0 string) (ret0 e return } +func (s *StorageModuleStub) DiskDownload(ctx context.Context, arg0 string, arg1 string) (ret0 error) { + args := []interface{}{arg0, arg1} + result, err := s.client.RequestContext(ctx, s.module, s.object, "DiskDownload", args...) + if err != nil { + panic(err) + } + result.PanicOnError() + ret0 = result.CallError() + loader := zbus.Loader{} + if err := result.Unmarshal(&loader); err != nil { + panic(err) + } + return +} + func (s *StorageModuleStub) DiskExists(ctx context.Context, arg0 string) (ret0 bool) { args := []interface{}{arg0} result, err := s.client.RequestContext(ctx, s.module, s.object, "DiskExists", args...) @@ -243,6 +258,21 @@ func (s *StorageModuleStub) DiskResize(ctx context.Context, arg0 string, arg1 gr return } +func (s *StorageModuleStub) DiskUpload(ctx context.Context, arg0 string, arg1 string) (ret0 error) { + args := []interface{}{arg0, arg1} + result, err := s.client.RequestContext(ctx, s.module, s.object, "DiskUpload", args...) + if err != nil { + panic(err) + } + result.PanicOnError() + ret0 = result.CallError() + loader := zbus.Loader{} + if err := result.Unmarshal(&loader); err != nil { + panic(err) + } + return +} + func (s *StorageModuleStub) DiskWrite(ctx context.Context, arg0 string, arg1 string) (ret0 error) { args := []interface{}{arg0, arg1} result, err := s.client.RequestContext(ctx, s.module, s.object, "DiskWrite", args...) @@ -348,6 +378,21 @@ func (s *StorageModuleStub) VolumeDelete(ctx context.Context, arg0 string) (ret0 return } +func (s *StorageModuleStub) VolumeDownload(ctx context.Context, arg0 string, arg1 gridtypes.Unit, arg2 string) (ret0 error) { + args := []interface{}{arg0, arg1, arg2} + result, err := s.client.RequestContext(ctx, s.module, s.object, "VolumeDownload", args...) + if err != nil { + panic(err) + } + result.PanicOnError() + ret0 = result.CallError() + loader := zbus.Loader{} + if err := result.Unmarshal(&loader); err != nil { + panic(err) + } + return +} + func (s *StorageModuleStub) VolumeExists(ctx context.Context, arg0 string) (ret0 bool, ret1 error) { args := []interface{}{arg0} result, err := s.client.RequestContext(ctx, s.module, s.object, "VolumeExists", args...) @@ -413,3 +458,18 @@ func (s *StorageModuleStub) VolumeUpdate(ctx context.Context, arg0 string, arg1 } return } + +func (s *StorageModuleStub) VolumeUpload(ctx context.Context, arg0 string, arg1 string) (ret0 error) { + args := []interface{}{arg0, arg1} + result, err := s.client.RequestContext(ctx, s.module, s.object, "VolumeUpload", args...) + if err != nil { + panic(err) + } + result.PanicOnError() + ret0 = result.CallError() + loader := zbus.Loader{} + if err := result.Unmarshal(&loader); err != nil { + panic(err) + } + return +} diff --git a/pkg/zos_api/deployment.go b/pkg/zos_api/deployment.go index 293dafcf..28895026 100644 --- a/pkg/zos_api/deployment.go +++ b/pkg/zos_api/deployment.go @@ -4,11 +4,33 @@ import ( "context" "encoding/json" "fmt" + "time" + "github.com/rs/zerolog/log" "github.com/threefoldtech/zos_base/pkg/gridtypes" + "github.com/threefoldtech/zos_base/pkg/gridtypes/zos" "github.com/threefoldtech/zos_sdk_go/rmb-sdk-go/peer" ) +// rootfsVolumePrefix is how the VM primitive names a container VM's writable +// rootfs subvolume (see pkg/primitives/vm/utils.go). We reuse it here to +// transfer/restore the user's rootfs changes. +const rootfsVolumePrefix = "rootfs:" + +const ( + waitWorkloadTimeout = 3 * time.Minute + waitWorkloadInterval = 2 * time.Second +) + +// transferItem maps a workload in the source deployment to a presigned URL. +type transferItem struct { + WorkloadName gridtypes.Name `json:"workload_name"` + URL string `json:"url"` + // Size is the volume size to (re)create on the target for a rootfs volume + // download; ignored for uploads and for zmount downloads. + Size gridtypes.Unit `json:"size,omitempty"` +} + func (g *ZosAPI) deploymentDeployHandler(ctx context.Context, payload []byte) (interface{}, error) { var deployment gridtypes.Deployment if err := json.Unmarshal(payload, &deployment); err != nil { @@ -58,3 +80,215 @@ func (g *ZosAPI) deploymentChangesHandler(ctx context.Context, payload []byte) ( } return g.provisionStub.Changes(ctx, peer.GetTwinID(ctx), args.ContractID) } + +// deploymentTransferHandler moves the persistent bytes of a deployment (zmount +// disks and/or the VM rootfs writable layer) to another node. It pauses the +// source deployment for a consistent copy, then uploads each requested workload +// to a caller-provided presigned S3 URL (HTTP PUT). Used on the OLD node during +// a contract move. +func (g *ZosAPI) deploymentTransferHandler(ctx context.Context, payload []byte) (interface{}, error) { + var args struct { + ContractID uint64 `json:"contract_id"` + Uploads []transferItem `json:"uploads"` + } + if err := json.Unmarshal(payload, &args); err != nil { + return nil, err + } + + twin := peer.GetTwinID(ctx) + deployment, err := g.provisionStub.Get(ctx, twin, args.ContractID) + if err != nil { + return nil, err + } + + // resolve + validate the requested workloads up front so obvious errors are + // returned synchronously, before we detach the long-running uploads + jobs, err := resolveTransferJobs(&deployment, args.Uploads) + if err != nil { + return nil, err + } + + // freeze the source VM(s) so disks/rootfs are quiescent while copied + if err := g.provisionStub.PauseDeployment(ctx, twin, args.ContractID); err != nil { + return nil, fmt.Errorf("failed to pause deployment before transfer: %w", err) + } + + // uploads take minutes — far longer than the RMB response window — so run them + // in the background (detached from the request ctx) and report progress via + // logs. Completion is observable by the objects appearing at the target URLs. + go g.runUploads(args.ContractID, jobs) + + return map[string]interface{}{"status": "started", "uploads": len(jobs)}, nil +} + +// deploymentPrepareHandler stages a deployment on this node without starting the +// zmachine: it provisions the network + zmount(s) (and skips the VM), then pulls +// each requested workload's bytes from a caller-provided presigned S3 URL (HTTP +// GET) into the freshly created zmount disk / rootfs volume. Used on the NEW node +// during a contract move; the VM is started by a follow-up call. +func (g *ZosAPI) deploymentPrepareHandler(ctx context.Context, payload []byte) (interface{}, error) { + var args struct { + Deployment gridtypes.Deployment `json:"deployment"` + Downloads []transferItem `json:"downloads"` + // Start boots the zmachine automatically once all downloads finish. + Start bool `json:"start"` + } + if err := json.Unmarshal(payload, &args); err != nil { + return nil, err + } + + twin := peer.GetTwinID(ctx) + + // resolve + validate the requested workloads up front (synchronous errors) + jobs, err := resolveTransferJobs(&args.Deployment, args.Downloads) + if err != nil { + return nil, err + } + + // provision network + zmount(s), skip the zmachine (async) + if err := g.provisionStub.PrepareDeployment(ctx, twin, args.Deployment); err != nil { + return nil, err + } + + // downloads take minutes — run them in the background. If Start is set, the VM + // is booted automatically once "all downloads complete"; the caller can then + // wait for the zmachine workload to reach the Ok state. + go g.runDownloads(twin, args.Deployment.ContractID, jobs, args.Start) + + return map[string]interface{}{"status": "started", "downloads": len(jobs)}, nil +} + +// deploymentStartHandler boots the zmachine(s) of a deployment previously staged +// via prepare. Used on the NEW node to finish a contract move once the migrated +// data (zmount disks, rootfs) is in place. +func (g *ZosAPI) deploymentStartHandler(ctx context.Context, payload []byte) (interface{}, error) { + var args struct { + ContractID uint64 `json:"contract_id"` + } + if err := json.Unmarshal(payload, &args); err != nil { + return nil, err + } + return nil, g.provisionStub.StartDeployment(ctx, peer.GetTwinID(ctx), args.ContractID) +} + +// waitWorkloadProvisioned blocks until the workload with the given global id +// reaches an OK state, or returns an error on failure/timeout. +func (g *ZosAPI) waitWorkloadProvisioned(ctx context.Context, id string) error { + ctx, cancel := context.WithTimeout(ctx, waitWorkloadTimeout) + defer cancel() + + ticker := time.NewTicker(waitWorkloadInterval) + defer ticker.Stop() + + for { + state, exists, err := g.provisionStub.GetWorkloadStatus(ctx, id) + if err == nil && exists { + if state.IsOkay() { + return nil + } + if state == gridtypes.StateError { + return fmt.Errorf("workload %q failed to provision", id) + } + } + + select { + case <-ctx.Done(): + return fmt.Errorf("timeout waiting for workload %q to be provisioned", id) + case <-ticker.C: + } + } +} + +// transferJob is a resolved upload/download unit: a workload's global id + target URL. +type transferJob struct { + name gridtypes.Name + typ gridtypes.WorkloadType + id string + url string + size gridtypes.Unit +} + +// resolveTransferJobs maps requested items to workloads in the deployment and +// rejects anything that isn't a zmount or a zmachine. +func resolveTransferJobs(dl *gridtypes.Deployment, items []transferItem) ([]transferJob, error) { + jobs := make([]transferJob, 0, len(items)) + for _, it := range items { + wl, err := dl.Get(it.WorkloadName) + if err != nil { + return nil, err + } + switch wl.Type { + case zos.ZMountType, zos.ZMachineType, zos.ZMachineLightType: + default: + return nil, fmt.Errorf("workload %q of type %q is not transferable", it.WorkloadName, wl.Type) + } + jobs = append(jobs, transferJob{ + name: it.WorkloadName, + typ: wl.Type, + id: wl.ID.String(), + url: it.URL, + size: it.Size, + }) + } + return jobs, nil +} + +// runUploads streams each job's bytes to its URL (zmount raw disk, or zmachine +// rootfs). Runs detached from any request context; progress is logged. +func (g *ZosAPI) runUploads(contractID uint64, jobs []transferJob) { + ctx := context.Background() + for _, j := range jobs { + var err error + switch j.typ { + case zos.ZMountType: + err = g.storageStub.DiskUpload(ctx, j.id, j.url) + default: // zmachine rootfs + err = g.storageStub.VolumeUpload(ctx, rootfsVolumePrefix+j.id, j.url) + } + if err != nil { + log.Error().Err(err).Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment transfer: upload failed") + return + } + log.Info().Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment transfer: upload complete") + } + log.Info().Uint64("contract", contractID).Msg("deployment transfer: all uploads complete") +} + +// runDownloads pulls each job's bytes from its URL into the freshly staged +// zmount disk / rootfs volume. Runs detached; progress is logged. The VM must be +// started only after "all downloads complete". +func (g *ZosAPI) runDownloads(twin uint32, contractID uint64, jobs []transferJob, start bool) { + ctx := context.Background() + for _, j := range jobs { + var err error + switch j.typ { + case zos.ZMountType: + // the disk is created by the async prepare; wait for it before writing + if err = g.waitWorkloadProvisioned(ctx, j.id); err == nil { + err = g.storageStub.DiskDownload(ctx, j.id, j.url) + } + default: // zmachine rootfs; VolumeDownload creates the volume + err = g.storageStub.VolumeDownload(ctx, rootfsVolumePrefix+j.id, j.size, j.url) + } + if err != nil { + log.Error().Err(err).Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment prepare: download failed") + return + } + log.Info().Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment prepare: download complete") + } + log.Info().Uint64("contract", contractID).Msg("deployment prepare: all downloads complete") + + // boot the VM now that its data is in place (the caller waits for the zmachine + // workload to reach Ok) + if start { + if err := g.provisionStub.StartDeployment(ctx, twin, contractID); err != nil { + log.Error().Err(err).Uint64("contract", contractID).Msg("deployment prepare: auto-start failed") + } else { + log.Info().Uint64("contract", contractID).Msg("deployment prepare: auto-start scheduled") + } + } +} diff --git a/pkg/zos_api/routes.go b/pkg/zos_api/routes.go index 1460645c..88523cd0 100644 --- a/pkg/zos_api/routes.go +++ b/pkg/zos_api/routes.go @@ -51,6 +51,9 @@ func (g *ZosAPI) SetupRoutes(router *peer.Router) { deployment.WithHandler("get", g.deploymentGetHandler) deployment.WithHandler("list", g.deploymentListHandler) deployment.WithHandler("changes", g.deploymentChangesHandler) + deployment.WithHandler("transfer", g.deploymentTransferHandler) + deployment.WithHandler("prepare", g.deploymentPrepareHandler) + deployment.WithHandler("start", g.deploymentStartHandler) admin := root.SubRoute("admin") admin.Use(g.authorized) diff --git a/pkg/zos_api_light/deployment.go b/pkg/zos_api_light/deployment.go index c600779d..6c2dde52 100644 --- a/pkg/zos_api_light/deployment.go +++ b/pkg/zos_api_light/deployment.go @@ -4,11 +4,35 @@ import ( "context" "encoding/json" "fmt" + "time" + "github.com/rs/zerolog/log" gridtypes "github.com/threefoldtech/zos_base/pkg/gridtypes" + "github.com/threefoldtech/zos_base/pkg/gridtypes/zos" "github.com/threefoldtech/zos_sdk_go/rmb-sdk-go/peer" ) +// rootfsVolumePrefix is how the VM primitive names a container VM's writable +// rootfs subvolume (see pkg/primitives/vm-light/utils.go). We reuse it here to +// transfer/restore the user's rootfs changes. +const rootfsVolumePrefix = "rootfs:" + +// waitWorkloadTimeout bounds how long prepare waits for an async workload +// (zmount) to be provisioned before pulling its data. +const ( + waitWorkloadTimeout = 3 * time.Minute + waitWorkloadInterval = 2 * time.Second +) + +// transferItem maps a workload in the source deployment to a presigned URL. +type transferItem struct { + WorkloadName gridtypes.Name `json:"workload_name"` + URL string `json:"url"` + // Size is the volume size to (re)create on the target for a rootfs volume + // download; ignored for uploads and for zmount downloads. + Size gridtypes.Unit `json:"size,omitempty"` +} + func (g *ZosAPI) deploymentDeployHandler(ctx context.Context, payload []byte) (interface{}, error) { var deployment gridtypes.Deployment if err := json.Unmarshal(payload, &deployment); err != nil { @@ -58,3 +82,215 @@ func (g *ZosAPI) deploymentChangesHandler(ctx context.Context, payload []byte) ( } return g.provisionStub.Changes(ctx, peer.GetTwinID(ctx), args.ContractID) } + +// deploymentTransferHandler moves the persistent bytes of a deployment (zmount +// disks and/or the VM rootfs writable layer) to another node. It pauses the +// source deployment for a consistent copy, then uploads each requested workload +// to a caller-provided presigned S3 URL (HTTP PUT). Used on the OLD node during +// a contract move. +func (g *ZosAPI) deploymentTransferHandler(ctx context.Context, payload []byte) (interface{}, error) { + var args struct { + ContractID uint64 `json:"contract_id"` + Uploads []transferItem `json:"uploads"` + } + if err := json.Unmarshal(payload, &args); err != nil { + return nil, err + } + + twin := peer.GetTwinID(ctx) + deployment, err := g.provisionStub.Get(ctx, twin, args.ContractID) + if err != nil { + return nil, err + } + + // resolve + validate the requested workloads up front so obvious errors are + // returned synchronously, before we detach the long-running uploads + jobs, err := resolveTransferJobs(&deployment, args.Uploads) + if err != nil { + return nil, err + } + + // freeze the source VM(s) so disks/rootfs are quiescent while copied + if err := g.provisionStub.PauseDeployment(ctx, twin, args.ContractID); err != nil { + return nil, fmt.Errorf("failed to pause deployment before transfer: %w", err) + } + + // uploads take minutes — far longer than the RMB response window — so run them + // in the background (detached from the request ctx) and report progress via + // logs. Completion is observable by the objects appearing at the target URLs. + go g.runUploads(args.ContractID, jobs) + + return map[string]interface{}{"status": "started", "uploads": len(jobs)}, nil +} + +// deploymentPrepareHandler stages a deployment on this node without starting the +// zmachine: it provisions the network + zmount(s) (and skips the VM), then pulls +// each requested workload's bytes from a caller-provided presigned S3 URL (HTTP +// GET) into the freshly created zmount disk / rootfs volume. Used on the NEW node +// during a contract move; the VM is started by a follow-up call. +func (g *ZosAPI) deploymentPrepareHandler(ctx context.Context, payload []byte) (interface{}, error) { + var args struct { + Deployment gridtypes.Deployment `json:"deployment"` + Downloads []transferItem `json:"downloads"` + // Start boots the zmachine automatically once all downloads finish. + Start bool `json:"start"` + } + if err := json.Unmarshal(payload, &args); err != nil { + return nil, err + } + + twin := peer.GetTwinID(ctx) + + // resolve + validate the requested workloads up front (synchronous errors) + jobs, err := resolveTransferJobs(&args.Deployment, args.Downloads) + if err != nil { + return nil, err + } + + // provision network + zmount(s), skip the zmachine (async) + if err := g.provisionStub.PrepareDeployment(ctx, twin, args.Deployment); err != nil { + return nil, err + } + + // downloads take minutes — run them in the background. If Start is set, the VM + // is booted automatically once "all downloads complete"; the caller can then + // wait for the zmachine workload to reach the Ok state. + go g.runDownloads(twin, args.Deployment.ContractID, jobs, args.Start) + + return map[string]interface{}{"status": "started", "downloads": len(jobs)}, nil +} + +// deploymentStartHandler boots the zmachine(s) of a deployment previously staged +// via prepare. Used on the NEW node to finish a contract move once the migrated +// data (zmount disks, rootfs) is in place. +func (g *ZosAPI) deploymentStartHandler(ctx context.Context, payload []byte) (interface{}, error) { + var args struct { + ContractID uint64 `json:"contract_id"` + } + if err := json.Unmarshal(payload, &args); err != nil { + return nil, err + } + return nil, g.provisionStub.StartDeployment(ctx, peer.GetTwinID(ctx), args.ContractID) +} + +// waitWorkloadProvisioned blocks until the workload with the given global id +// reaches an OK state, or returns an error on failure/timeout. +func (g *ZosAPI) waitWorkloadProvisioned(ctx context.Context, id string) error { + ctx, cancel := context.WithTimeout(ctx, waitWorkloadTimeout) + defer cancel() + + ticker := time.NewTicker(waitWorkloadInterval) + defer ticker.Stop() + + for { + state, exists, err := g.provisionStub.GetWorkloadStatus(ctx, id) + if err == nil && exists { + if state.IsOkay() { + return nil + } + if state == gridtypes.StateError { + return fmt.Errorf("workload %q failed to provision", id) + } + } + + select { + case <-ctx.Done(): + return fmt.Errorf("timeout waiting for workload %q to be provisioned", id) + case <-ticker.C: + } + } +} + +// transferJob is a resolved upload/download unit: a workload's global id + target URL. +type transferJob struct { + name gridtypes.Name + typ gridtypes.WorkloadType + id string + url string + size gridtypes.Unit +} + +// resolveTransferJobs maps requested items to workloads in the deployment and +// rejects anything that isn't a zmount or a zmachine. +func resolveTransferJobs(dl *gridtypes.Deployment, items []transferItem) ([]transferJob, error) { + jobs := make([]transferJob, 0, len(items)) + for _, it := range items { + wl, err := dl.Get(it.WorkloadName) + if err != nil { + return nil, err + } + switch wl.Type { + case zos.ZMountType, zos.ZMachineType, zos.ZMachineLightType: + default: + return nil, fmt.Errorf("workload %q of type %q is not transferable", it.WorkloadName, wl.Type) + } + jobs = append(jobs, transferJob{ + name: it.WorkloadName, + typ: wl.Type, + id: wl.ID.String(), + url: it.URL, + size: it.Size, + }) + } + return jobs, nil +} + +// runUploads streams each job's bytes to its URL (zmount raw disk, or zmachine +// rootfs). Runs detached from any request context; progress is logged. +func (g *ZosAPI) runUploads(contractID uint64, jobs []transferJob) { + ctx := context.Background() + for _, j := range jobs { + var err error + switch j.typ { + case zos.ZMountType: + err = g.storageStub.DiskUpload(ctx, j.id, j.url) + default: // zmachine (light) rootfs + err = g.storageStub.VolumeUpload(ctx, rootfsVolumePrefix+j.id, j.url) + } + if err != nil { + log.Error().Err(err).Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment transfer: upload failed") + return + } + log.Info().Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment transfer: upload complete") + } + log.Info().Uint64("contract", contractID).Msg("deployment transfer: all uploads complete") +} + +// runDownloads pulls each job's bytes from its URL into the freshly staged +// zmount disk / rootfs volume. Runs detached; progress is logged. The VM must be +// started only after "all downloads complete". +func (g *ZosAPI) runDownloads(twin uint32, contractID uint64, jobs []transferJob, start bool) { + ctx := context.Background() + for _, j := range jobs { + var err error + switch j.typ { + case zos.ZMountType: + // the disk is created by the async prepare; wait for it before writing + if err = g.waitWorkloadProvisioned(ctx, j.id); err == nil { + err = g.storageStub.DiskDownload(ctx, j.id, j.url) + } + default: // zmachine (light) rootfs; VolumeDownload creates the volume + err = g.storageStub.VolumeDownload(ctx, rootfsVolumePrefix+j.id, j.size, j.url) + } + if err != nil { + log.Error().Err(err).Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment prepare: download failed") + return + } + log.Info().Uint64("contract", contractID).Str("workload", string(j.name)). + Msg("deployment prepare: download complete") + } + log.Info().Uint64("contract", contractID).Msg("deployment prepare: all downloads complete") + + // boot the VM now that its data is in place (the caller waits for the zmachine + // workload to reach Ok) + if start { + if err := g.provisionStub.StartDeployment(ctx, twin, contractID); err != nil { + log.Error().Err(err).Uint64("contract", contractID).Msg("deployment prepare: auto-start failed") + } else { + log.Info().Uint64("contract", contractID).Msg("deployment prepare: auto-start scheduled") + } + } +} diff --git a/pkg/zos_api_light/routes.go b/pkg/zos_api_light/routes.go index 666e8bfd..4468a412 100644 --- a/pkg/zos_api_light/routes.go +++ b/pkg/zos_api_light/routes.go @@ -42,6 +42,9 @@ func (g *ZosAPI) SetupRoutes(router *peer.Router) { deployment.WithHandler("get", g.deploymentGetHandler) deployment.WithHandler("list", g.deploymentListHandler) deployment.WithHandler("changes", g.deploymentChangesHandler) + deployment.WithHandler("transfer", g.deploymentTransferHandler) + deployment.WithHandler("prepare", g.deploymentPrepareHandler) + deployment.WithHandler("start", g.deploymentStartHandler) admin := root.SubRoute("admin") admin.Use(g.authorized) diff --git a/scripts/migration/README.md b/scripts/migration/README.md new file mode 100644 index 00000000..498a7962 --- /dev/null +++ b/scripts/migration/README.md @@ -0,0 +1,73 @@ +# zosmigration + +One command to move a zos zmachine (VM + zmount + its network) between nodes. You +give it **S3 credentials**, the **source VM contract id**, and the **source + target +node ids** — it does the rest. + +Builds against this repo via `replace github.com/threefoldtech/zos_base => ../..`. + +## What it automates + +1. Creates the S3 bucket (if missing). +2. Fetches the source VM deployment over RMB and **auto-detects** the zmount(s), the + zmachine, and the rootfs size (`ZMachine.Size`) — no workload names to type. +3. **Transfer**: presigns PUT URLs and calls `zos.deployment.transfer` on the source + node (pauses the VM, uploads zmount disk(s) + rootfs to S3), then waits for the + objects to land in the bucket. +4. **Auto-detects** the network the VM attaches to (from its zmachine interface) and + finds its contract via `DeploymentList` on the source — then fetches + deploys that + network on the target. (Pass `-src-net-contract` to override, or skip auto-detect.) +5. **Prepare**: presigns GET URLs, creates the target contract, calls + `zos.deployment.prepare` on the target (provisions network+zmount and pulls the + data) — without booting the VM. +6. Prints the follow-up `-start-contract` command. + +It re-signs deployments with the owner mnemonic (DeploymentGet re-marshals the `Env` +map, shifting the challenge hash) and copies the source contract's `deployment_data` +so the dashboard shows the workload type. + +## Usage + +```bash +go build -o zosmigration . + +# full migration: transfer -> network -> prepare +MNEMONIC="word word ..." ./zosmigration \ + -s3-endpoint https://gateway.storjshare.io \ + -s3-access \ + -s3-secret \ + -src-node-id 396 -src-contract 270333 \ + -dst-node-id 371 + +# then, once the target logs "all downloads complete", boot the VM: +MNEMONIC="word word ..." ./zosmigration -dst-node-id 371 -start-contract +``` + +- `-src-net-contract` is **optional** — the network is auto-detected from the VM + contract. Pass it only to override; if the network already exists on the target and + you want to skip re-deploying it, it will still be (idempotently) re-deployed. +- S3 creds can also come from `S3_ACCESS` / `S3_SECRET` env vars. +- Defaults: bucket `zos-migration`, region `global` (Storj), devnet substrate/relay. + Override `-substrate` / `-relay` for test/main. +- `-upload-timeout` (default 2h) bounds the wait for slow uplinks. + +## Start / waiting + +By default the tool prepares with **auto-start**: the node boots the VM once its +downloads finish, and the tool **waits until the zmachine reaches `Ok`** before it +exits — so one command runs start-to-finish and stops exactly when the VM is up on +the new node. + +- `-stage-only` prepares without auto-start/wait and prints a `-start-contract` + command to boot it later (fire-and-exit). +- `-start-contract ` boots an already-prepared deployment on `-dst-node-id` and + exits (used by `-stage-only`, or to retry a boot). + +Auto-start runs in the **target node's api-gateway**, so both nodes must run a zos +build that carries the `start` flag on `zos.deployment.prepare`. + +## Prerequisites + +- Both nodes run a zos build with the transfer/prepare/start support. +- The owner **mnemonic** (the deployment's `twin_id`); it signs the RMB calls, the + on-chain contract creation, and the deployments. diff --git a/scripts/migration/go.mod b/scripts/migration/go.mod new file mode 100644 index 00000000..19a92fd8 --- /dev/null +++ b/scripts/migration/go.mod @@ -0,0 +1,79 @@ +module zosmigration + +go 1.25.0 + +require ( + github.com/minio/minio-go/v7 v7.0.10 + github.com/threefoldtech/tfchain/clients/tfchain-client-go v0.0.0-20260302124210-526158ffbc00 + github.com/threefoldtech/zos_base v1.1.2 + github.com/threefoldtech/zos_sdk_go/rmb-sdk-go v0.18.0 +) + +require ( + filippo.io/edwards25519 v1.2.0 // indirect + github.com/ChainSafe/go-schnorrkel v1.1.0 // indirect + github.com/blang/semver v3.5.1+incompatible // indirect + github.com/cenkalti/backoff v2.2.1+incompatible // indirect + github.com/centrifuge/go-substrate-rpc-client/v4 v4.2.1 // indirect + github.com/cosmos/go-bip39 v1.0.0 // indirect + github.com/dave/jennifer v1.3.0 // indirect + github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect + github.com/deckarep/golang-set v1.8.0 // indirect + github.com/decred/base58 v1.0.6 // indirect + github.com/decred/dcrd/crypto/blake256 v1.1.0 // indirect + github.com/decred/dcrd/dcrec/secp256k1/v4 v4.3.0 // indirect + github.com/ethereum/go-ethereum v1.17.1 // indirect + github.com/go-ole/go-ole v1.3.0 // indirect + github.com/golang-jwt/jwt v3.2.2+incompatible // indirect + github.com/golang/protobuf v1.5.4 // indirect + github.com/gomodule/redigo v2.0.0+incompatible // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/gorilla/websocket v1.5.3 // indirect + github.com/gtank/merlin v0.1.1 // indirect + github.com/gtank/ristretto255 v0.2.0 // indirect + github.com/hashicorp/errwrap v1.1.0 // indirect + github.com/hashicorp/go-multierror v1.1.1 // indirect + github.com/holiman/uint256 v1.3.2 // indirect + github.com/jbenet/go-base58 v0.0.0-20150317085156-6237cf65f3a6 // indirect + github.com/json-iterator/go v1.1.10 // indirect + github.com/klauspost/cpuid v1.3.1 // indirect + github.com/klauspost/cpuid/v2 v2.0.9 // indirect + github.com/mattn/go-colorable v0.1.14 // indirect + github.com/mattn/go-isatty v0.0.20 // indirect + github.com/mimoo/StrobeGo v0.0.0-20220103164710-9a04d6ca976b // indirect + github.com/minio/md5-simd v1.1.0 // indirect + github.com/minio/sha256-simd v1.0.0 // indirect + github.com/mitchellh/go-homedir v1.1.0 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.1 // indirect + github.com/pierrec/xxHash v0.1.5 // indirect + github.com/pkg/errors v0.9.1 // indirect + github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect + github.com/rs/cors v1.11.1 // indirect + github.com/rs/xid v1.6.0 // indirect + github.com/rs/zerolog v1.34.0 // indirect + github.com/shirou/gopsutil v3.21.11+incompatible // indirect + github.com/stretchr/objx v0.5.2 // indirect + github.com/stretchr/testify v1.11.1 // indirect + github.com/threefoldtech/zbus v1.0.1 // indirect + github.com/tklauser/go-sysconf v0.3.12 // indirect + github.com/tklauser/numcpus v0.6.1 // indirect + github.com/vedhavyas/go-subkey v1.0.3 // indirect + github.com/vedhavyas/go-subkey/v2 v2.0.0 // indirect + github.com/vishvananda/netlink v1.2.1-beta.2 // indirect + github.com/vishvananda/netns v0.0.0-20210104183010-2eb08e3e575f // indirect + github.com/vmihailenco/msgpack v4.0.4+incompatible // indirect + github.com/yusufpapurcu/wmi v1.2.4 // indirect + golang.org/x/crypto v0.49.0 // indirect + golang.org/x/net v0.51.0 // indirect + golang.org/x/sys v0.42.0 // indirect + golang.org/x/text v0.35.0 // indirect + gonum.org/v1/gonum v0.16.0 // indirect + google.golang.org/appengine v1.6.7 // indirect + google.golang.org/protobuf v1.36.11 // indirect + gopkg.in/ini.v1 v1.57.0 // indirect + gopkg.in/natefinch/npipe.v2 v2.0.0-20160621034901-c1b8fa8bdcce // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect +) + +replace github.com/threefoldtech/zos_base => ../.. diff --git a/scripts/migration/go.sum b/scripts/migration/go.sum new file mode 100644 index 00000000..f00ebaa0 --- /dev/null +++ b/scripts/migration/go.sum @@ -0,0 +1,254 @@ +filippo.io/edwards25519 v1.2.0 h1:crnVqOiS4jqYleHd9vaKZ+HKtHfllngJIiOpNpoJsjo= +filippo.io/edwards25519 v1.2.0/go.mod h1:xzAOLCNug/yB62zG1bQ8uziwrIqIuxhctzJT18Q77mc= +github.com/ChainSafe/go-schnorrkel v1.1.0 h1:rZ6EU+CZFCjB4sHUE1jIu8VDoB/wRKZxoe1tkcO71Wk= +github.com/ChainSafe/go-schnorrkel v1.1.0/go.mod h1:ABkENxiP+cvjFiByMIZ9LYbRoNNLeBLiakC1XeTFxfE= +github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= +github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= +github.com/ProjectZKM/Ziren/crates/go-runtime/zkvm_runtime v0.0.0-20251001021608-1fe7b43fc4d6 h1:1zYrtlhrZ6/b6SAjLSfKzWtdgqK0U+HtH/VcBWh1BaU= +github.com/ProjectZKM/Ziren/crates/go-runtime/zkvm_runtime v0.0.0-20251001021608-1fe7b43fc4d6/go.mod h1:ioLG6R+5bUSO1oeGSDxOV3FADARuMoytZCSX6MEMQkI= +github.com/blang/semver v3.5.1+incompatible h1:cQNTCjp13qL8KC3Nbxr/y2Bqb63oX6wdnnjpJbkM4JQ= +github.com/blang/semver v3.5.1+incompatible/go.mod h1:kRBLl5iJ+tD4TcOOxsy/0fnwebNt5EWlYSAyrTnjyyk= +github.com/btcsuite/btcutil v1.0.3-0.20201208143702-a53e38424cce h1:YtWJF7RHm2pYCvA5t0RPmAaLUhREsKuKd+SLhxFbFeQ= +github.com/btcsuite/btcutil v1.0.3-0.20201208143702-a53e38424cce/go.mod h1:0DVlHczLPewLcPGEIeUEzfOJhqGPQ0mJJRDBtD307+o= +github.com/cenkalti/backoff v2.2.1+incompatible h1:tNowT99t7UNflLxfYYSlKYsBpXdEet03Pg2g16Swow4= +github.com/cenkalti/backoff v2.2.1+incompatible/go.mod h1:90ReRw6GdpyfrHakVjL/QHaoyV4aDUVVkXQJJJ3NXXM= +github.com/centrifuge/go-substrate-rpc-client/v4 v4.2.1 h1:io49TJ8IOIlzipioJc9pJlrjgdJvqktpUWYxVY5AUjE= +github.com/centrifuge/go-substrate-rpc-client/v4 v4.2.1/go.mod h1:k61SBXqYmnZO4frAJyH3iuqjolYrYsq79r8EstmklDY= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/coreos/go-systemd v0.0.0-20190321100706-95778dfbb74e/go.mod h1:F5haX7vjVVG0kc13fIWeqUViNPyEJxv/OmvnBo0Yme4= +github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc= +github.com/cosmos/go-bip39 v1.0.0 h1:pcomnQdrdH22njcAatO0yWojsUnCO3y2tNoV1cb6hHY= +github.com/cosmos/go-bip39 v1.0.0/go.mod h1:RNJv0H/pOIVgxw6KS7QeX2a0Uo0aKUlfhZ4xuwvCdJw= +github.com/dave/jennifer v1.3.0 h1:p3tl41zjjCZTNBytMwrUuiAnherNUZktlhPTKoF/sEk= +github.com/dave/jennifer v1.3.0/go.mod h1:fIb+770HOpJ2fmN9EPPKOqm1vMGhB+TwXKMZhrIygKg= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc h1:U9qPSI2PIWSS1VwoXQT9A3Wy9MM3WgvqSxFWenqJduM= +github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/deckarep/golang-set v1.8.0 h1:sk9/l/KqpunDwP7pSjUg0keiOOLEnOBHzykLrsPppp4= +github.com/deckarep/golang-set v1.8.0/go.mod h1:5nI87KwE7wgsBU1F4GKAw2Qod7p5kyS383rP6+o6qqo= +github.com/deckarep/golang-set/v2 v2.6.0 h1:XfcQbWM1LlMB8BsJ8N9vW5ehnnPVIw0je80NsVHagjM= +github.com/deckarep/golang-set/v2 v2.6.0/go.mod h1:VAky9rY/yGXJOLEDv3OMci+7wtDpOF4IN+y82NBOac4= +github.com/decred/base58 v1.0.6 h1:NXndBcO+ubGZORV3EulvqeBcMuQM7doqVGa7pBhMOs4= +github.com/decred/base58 v1.0.6/go.mod h1:KR7Oh9njDPXTagD4P67KJZwroL8jT653u8CffkYqhcQ= +github.com/decred/dcrd/crypto/blake256 v1.1.0 h1:zPMNGQCm0g4QTY27fOCorQW7EryeQ/U0x++OzVrdms8= +github.com/decred/dcrd/crypto/blake256 v1.1.0/go.mod h1:2OfgNZ5wDpcsFmHmCK5gZTPcCXqlm2ArzUIkw9czNJo= +github.com/decred/dcrd/dcrec/secp256k1/v4 v4.3.0 h1:rpfIENRNNilwHwZeG5+P150SMrnNEcHYvcCuK6dPZSg= +github.com/decred/dcrd/dcrec/secp256k1/v4 v4.3.0/go.mod h1:v57UDF4pDQJcEfFUCRop3lJL149eHGSe9Jvczhzjo/0= +github.com/ethereum/go-ethereum v1.17.1 h1:IjlQDjgxg2uL+GzPRkygGULPMLzcYWncEI7wbaizvho= +github.com/ethereum/go-ethereum v1.17.1/go.mod h1:7UWOVHL7K3b8RfVRea022btnzLCaanwHtBuH1jUCH/I= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0= +github.com/go-ole/go-ole v1.3.0 h1:Dt6ye7+vXGIKZ7Xtk4s6/xVdGDQynvom7xCFEdWr6uE= +github.com/go-ole/go-ole v1.3.0/go.mod h1:5LS6F96DhAwUc7C+1HLexzMXY1xGRSryjyPPKW6zv78= +github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= +github.com/golang-jwt/jwt v3.2.2+incompatible h1:IfV12K8xAKAnZqdXVzCZ+TOjboZ2keLg81eXfW3O+oY= +github.com/golang-jwt/jwt v3.2.2+incompatible/go.mod h1:8pz2t5EyA70fFQQSrl6XZXzqecmYZeUEB8OUGHkxJ+I= +github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc= +github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs= +github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= +github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= +github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= +github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/gomodule/redigo v1.8.9/go.mod h1:7ArFNvsTjH8GMMzB4uy1snslv2BwmginuMs06a1uzZE= +github.com/gomodule/redigo v2.0.0+incompatible h1:K/R+8tc58AaqLkqG2Ol3Qk+DR/TlNuhuh457pBFPtt0= +github.com/gomodule/redigo v2.0.0+incompatible/go.mod h1:B4C85qUVwatsJoIUNIfCRsp7qO0iAmpGFZ4EELWSbC4= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/google/gofuzz v1.2.0 h1:xRy4A+RhZaiKjJ1bPfwQ8sedCA+YS2YcCHW6ec7JMi0= +github.com/google/gofuzz v1.2.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= +github.com/google/uuid v1.1.1/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/gopherjs/gopherjs v0.0.0-20181017120253-0766667cb4d1 h1:EGx4pi6eqNxGaHF6qqu48+N2wcFQ5qg5FXgOdqsJ5d8= +github.com/gopherjs/gopherjs v0.0.0-20181017120253-0766667cb4d1/go.mod h1:wJfORRmW1u3UXTncJ5qlYoELFm8eSnnEO6hX4iZ3EWY= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +github.com/gtank/merlin v0.1.1 h1:eQ90iG7K9pOhtereWsmyRJ6RAwcP4tHTDBHXNg+u5is= +github.com/gtank/merlin v0.1.1/go.mod h1:T86dnYJhcGOh5BjZFCJWTDeTK7XW8uE+E21Cy/bIQ+s= +github.com/gtank/ristretto255 v0.2.0 h1:LeOuWr6giplWkkMizx2emfG03SRPJqKt1nfIHLVHQ/0= +github.com/gtank/ristretto255 v0.2.0/go.mod h1:OJ1ox/dWcp7sJ5grYDcZ+kkHYuj5nelW5aaL7ESVXBw= +github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4= +github.com/hashicorp/errwrap v1.1.0 h1:OxrOeh75EUXMY8TBjag2fzXGZ40LB6IKw45YeGUDY2I= +github.com/hashicorp/errwrap v1.1.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4= +github.com/hashicorp/go-multierror v1.1.1 h1:H5DkEtf6CXdFp0N0Em5UCwQpXMWke8IA0+lD48awMYo= +github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM= +github.com/holiman/uint256 v1.3.2 h1:a9EgMPSC1AAaj1SZL5zIQD3WbwTuHrMGOerLjGmM/TA= +github.com/holiman/uint256 v1.3.2/go.mod h1:EOMSn4q6Nyt9P6efbI3bueV4e1b3dGlUCXeiRV4ng7E= +github.com/jbenet/go-base58 v0.0.0-20150317085156-6237cf65f3a6 h1:4zOlv2my+vf98jT1nQt4bT/yKWUImevYPJ2H344CloE= +github.com/jbenet/go-base58 v0.0.0-20150317085156-6237cf65f3a6/go.mod h1:r/8JmuR0qjuCiEhAolkfvdZgmPiHTnJaG0UXCSeR1Zo= +github.com/json-iterator/go v1.1.10 h1:Kz6Cvnvv2wGdaG/V8yMvfkmNiXq9Ya2KUv4rouJJr68= +github.com/json-iterator/go v1.1.10/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/uOdHXbAo4= +github.com/jtolds/gls v4.20.0+incompatible h1:xdiiI2gbIgH/gLH7ADydsJ1uDOEzR8yvV7C0MuV77Wo= +github.com/jtolds/gls v4.20.0+incompatible/go.mod h1:QJZ7F/aHp+rZTRtaJ1ow/lLfFfVYBRgL+9YlvaHOwJU= +github.com/klauspost/cpuid v1.2.3/go.mod h1:Pj4uuM528wm8OyEC2QMXAi2YiTZ96dNQPGgoMS4s3ek= +github.com/klauspost/cpuid v1.3.1 h1:5JNjFYYQrZeKRJ0734q51WCEEn2huer72Dc7K+R/b6s= +github.com/klauspost/cpuid v1.3.1/go.mod h1:bYW4mA6ZgKPob1/Dlai2LviZJO7KGI3uoWLd42rAQw4= +github.com/klauspost/cpuid/v2 v2.0.4/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= +github.com/klauspost/cpuid/v2 v2.0.9 h1:lgaqFMSdTdQYdZ04uHyN2d/eKdOMyi2YLSvlQIBFYa4= +github.com/klauspost/cpuid/v2 v2.0.9/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= +github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= +github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/mattn/go-colorable v0.1.13/go.mod h1:7S9/ev0klgBDR4GtXTXX8a3vIGJpMovkB8vQcUbaXHg= +github.com/mattn/go-colorable v0.1.14 h1:9A9LHSqF/7dyVVX6g0U9cwm9pG3kP9gSzcuIPHPsaIE= +github.com/mattn/go-colorable v0.1.14/go.mod h1:6LmQG8QLFO4G5z1gPvYEzlUgJ2wF+stgPZH1UqBm1s8= +github.com/mattn/go-isatty v0.0.16/go.mod h1:kYGgaQfpe5nmfYZH+SKPsOc2e4SrIfOl2e/yFXSvRLM= +github.com/mattn/go-isatty v0.0.19/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= +github.com/mattn/go-isatty v0.0.20/go.mod h1:W+V8PltTTMOvKvAeJH7IuucS94S2C6jfK/D7dTCTo3Y= +github.com/mimoo/StrobeGo v0.0.0-20181016162300-f8f6d4d2b643/go.mod h1:43+3pMjjKimDBf5Kr4ZFNGbLql1zKkbImw+fZbw3geM= +github.com/mimoo/StrobeGo v0.0.0-20220103164710-9a04d6ca976b h1:QrHweqAtyJ9EwCaGHBu1fghwxIPiopAHV06JlXrMHjk= +github.com/mimoo/StrobeGo v0.0.0-20220103164710-9a04d6ca976b/go.mod h1:xxLb2ip6sSUts3g1irPVHyk/DGslwQsNOo9I7smJfNU= +github.com/minio/md5-simd v1.1.0 h1:QPfiOqlZH+Cj9teu0t9b1nTBfPbyTl16Of5MeuShdK4= +github.com/minio/md5-simd v1.1.0/go.mod h1:XpBqgZULrMYD3R+M28PcmP0CkI7PEMzB3U77ZrKZ0Gw= +github.com/minio/minio-go/v7 v7.0.10 h1:1oUKe4EOPUEhw2qnPQaPsJ0lmVTYLFu03SiItauXs94= +github.com/minio/minio-go/v7 v7.0.10/go.mod h1:td4gW1ldOsj1PbSNS+WYK43j+P1XVhX/8W8awaYlBFo= +github.com/minio/sha256-simd v0.1.1/go.mod h1:B5e1o+1/KgNmWrSQK08Y6Z1Vb5pwIktudl0J58iy0KM= +github.com/minio/sha256-simd v1.0.0 h1:v1ta+49hkWZyvaKwrQB8elexRqm6Y0aMLjCNsrYxo6g= +github.com/minio/sha256-simd v1.0.0/go.mod h1:OuYzVNI5vcoYIAmbIvHPl3N3jUzVedXbKy5RFepssQM= +github.com/mitchellh/go-homedir v1.1.0 h1:lukF9ziXFxDFPkA1vsr5zpc1XuPDn/wFntq5mG+4E0Y= +github.com/mitchellh/go-homedir v1.1.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0= +github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= +github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= +github.com/modern-go/reflect2 v0.0.0-20180701023420-4b7aa43c6742/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0= +github.com/modern-go/reflect2 v1.0.1 h1:9f412s+6RmYXLWZSEzVVgPGK7C2PphHj5RJrvfx9AWI= +github.com/modern-go/reflect2 v1.0.1/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0= +github.com/pierrec/xxHash v0.1.5 h1:n/jBpwTHiER4xYvK3/CdPVnLDPchj8eTJFFLUb4QHBo= +github.com/pierrec/xxHash v0.1.5/go.mod h1:w2waW5Zoa/Wc4Yqe0wgrIYAGKqRMf7czn2HNKXmuL+I= +github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= +github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= +github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U= +github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/rs/cors v1.11.1 h1:eU3gRzXLRK57F5rKMGMZURNdIG4EoAmX8k94r9wXWHA= +github.com/rs/cors v1.11.1/go.mod h1:XyqrcTp5zjWr1wsJ8PIRZssZ8b/WMcMf71DJnit4EMU= +github.com/rs/xid v1.2.1/go.mod h1:+uKXf+4Djp6Md1KODXJxgGQPKngRmWyn10oCKFzNHOQ= +github.com/rs/xid v1.6.0 h1:fV591PaemRlL6JfRxGDEPl69wICngIQ3shQtzfy2gxU= +github.com/rs/xid v1.6.0/go.mod h1:7XoLgs4eV+QndskICGsho+ADou8ySMSjJKDIan90Nz0= +github.com/rs/zerolog v1.14.3/go.mod h1:3WXPzbXEEliJ+a6UFE4vhIxV8qR1EML6ngzP9ug4eYg= +github.com/rs/zerolog v1.34.0 h1:k43nTLIwcTVQAncfCw4KZ2VY6ukYoZaBPNOE8txlOeY= +github.com/rs/zerolog v1.34.0/go.mod h1:bJsvje4Z08ROH4Nhs5iH600c3IkWhwp44iRc54W6wYQ= +github.com/shirou/gopsutil v3.21.11+incompatible h1:+1+c1VGhc88SSonWP6foOcLhvnKlUeu/erjjvaPEYiI= +github.com/shirou/gopsutil v3.21.11+incompatible/go.mod h1:5b4v6he4MtMOwMlS0TUMTu2PcXUg8+E1lC7eC3UO/RA= +github.com/smartystreets/assertions v0.0.0-20180927180507-b2de0cb4f26d h1:zE9ykElWQ6/NYmHa3jpm/yHnI4xSofP+UP6SpjHcSeM= +github.com/smartystreets/assertions v0.0.0-20180927180507-b2de0cb4f26d/go.mod h1:OnSkiWE9lh6wB0YB77sQom3nweQdgAjqCqsofrRNTgc= +github.com/smartystreets/goconvey v1.6.4 h1:fv0U8FUIMPNf1L9lnHLvLhgicrIVChEkdzIKYqbNC9s= +github.com/smartystreets/goconvey v1.6.4/go.mod h1:syvi0/a8iFYH4r/RixwvyeAJjdLS9QV7WQ/tjFTllLA= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY= +github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= +github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/threefoldtech/tfchain/clients/tfchain-client-go v0.0.0-20260302124210-526158ffbc00 h1:Q/vkWlTfZBCo4SclrbpSV35sSg/EF2MIIi0LgxPdzCI= +github.com/threefoldtech/tfchain/clients/tfchain-client-go v0.0.0-20260302124210-526158ffbc00/go.mod h1:cOL5YgHUmDG5SAXrsZxFjUECRQQuAqOoqvXhZG5sEUw= +github.com/threefoldtech/zbus v1.0.1 h1:3KaEpyOiDYAw+lrAyoQUGIvY9BcjVRXlQ1beBRqhRNk= +github.com/threefoldtech/zbus v1.0.1/go.mod h1:E/v/xEvG/l6z/Oj0aDkuSUXFm/1RVJkhKBwDTAIdsHo= +github.com/threefoldtech/zos_sdk_go/rmb-sdk-go v0.18.0 h1:CD8/yd0pAZINF93PHDODn2e+M6d5f6x6ECfpQac4kgw= +github.com/threefoldtech/zos_sdk_go/rmb-sdk-go v0.18.0/go.mod h1:OI3omGvjs88MXvOKSfYULhxweQWSz6jvQuUHDTI98gw= +github.com/tklauser/go-sysconf v0.3.12 h1:0QaGUFOdQaIVdPgfITYzaTegZvdCjmYO52cSFAEVmqU= +github.com/tklauser/go-sysconf v0.3.12/go.mod h1:Ho14jnntGE1fpdOqQEEaiKRpvIavV0hSfmBq8nJbHYI= +github.com/tklauser/numcpus v0.6.1 h1:ng9scYS7az0Bk4OZLvrNXNSAO2Pxr1XXRAPyjhIx+Fk= +github.com/tklauser/numcpus v0.6.1/go.mod h1:1XfjsgE2zo8GVw7POkMbHENHzVg3GzmoZ9fESEdAacY= +github.com/vedhavyas/go-subkey v1.0.3 h1:iKR33BB/akKmcR2PMlXPBeeODjWLM90EL98OrOGs8CA= +github.com/vedhavyas/go-subkey v1.0.3/go.mod h1:CloUaFQSSTdWnINfBRFjVMkWXZANW+nd8+TI5jYcl6Y= +github.com/vedhavyas/go-subkey/v2 v2.0.0 h1:LemDIsrVtRSOkp0FA8HxP6ynfKjeOj3BY2U9UNfeDMA= +github.com/vedhavyas/go-subkey/v2 v2.0.0/go.mod h1:95aZ+XDCWAUUynjlmi7BtPExjXgXxByE0WfBwbmIRH4= +github.com/vishvananda/netlink v1.2.1-beta.2 h1:Llsql0lnQEbHj0I1OuKyp8otXp0r3q0mPkuhwHfStVs= +github.com/vishvananda/netlink v1.2.1-beta.2/go.mod h1:twkDnbuQxJYemMlGd4JFIcuhgX83tXhKS2B/PRMpOho= +github.com/vishvananda/netns v0.0.0-20200728191858-db3c7e526aae/go.mod h1:DD4vA1DwXk04H54A1oHXtwZmA0grkVMdPxx/VGLCah0= +github.com/vishvananda/netns v0.0.0-20210104183010-2eb08e3e575f h1:p4VB7kIXpOQvVn1ZaTIVp+3vuYAXFe3OJEvjbUYJLaA= +github.com/vishvananda/netns v0.0.0-20210104183010-2eb08e3e575f/go.mod h1:DD4vA1DwXk04H54A1oHXtwZmA0grkVMdPxx/VGLCah0= +github.com/vmihailenco/msgpack v4.0.3+incompatible/go.mod h1:fy3FlTQTDXWkZ7Bh6AcGMlsjHatGryHQYUTf1ShIgkk= +github.com/vmihailenco/msgpack v4.0.4+incompatible h1:dSLoQfGFAo3F6OoNhwUmLwVgaUXK79GlxNBwueZn0xI= +github.com/vmihailenco/msgpack v4.0.4+incompatible/go.mod h1:fy3FlTQTDXWkZ7Bh6AcGMlsjHatGryHQYUTf1ShIgkk= +github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= +github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= +github.com/zenazn/goji v0.9.0/go.mod h1:7S9M489iMyHBNxwZnk9/EHS098H4/F6TATF2mIxtB1Q= +go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= +go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/otel v1.39.0 h1:8yPrr/S0ND9QEfTfdP9V+SiwT4E0G7Y5MO7p85nis48= +go.opentelemetry.io/otel v1.39.0/go.mod h1:kLlFTywNWrFyEdH0oj2xK0bFYZtHRYUdv1NklR/tgc8= +go.opentelemetry.io/otel/metric v1.39.0 h1:d1UzonvEZriVfpNKEVmHXbdf909uGTOQjA0HF0Ls5Q0= +go.opentelemetry.io/otel/metric v1.39.0/go.mod h1:jrZSWL33sD7bBxg1xjrqyDjnuzTUB0x1nBERXd7Ftcs= +go.opentelemetry.io/otel/trace v1.39.0 h1:2d2vfpEDmCJ5zVYz7ijaJdOF59xLomrvj7bjt6/qCJI= +go.opentelemetry.io/otel/trace v1.39.0/go.mod h1:88w4/PnZSazkGzz/w84VHpQafiU4EtqqlVdxWy+rNOA= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +go.uber.org/mock v0.5.2 h1:LbtPTcP8A5k9WPXj54PPPbjcI4Y6lhyOZXn+VS7wNko= +go.uber.org/mock v0.5.2/go.mod h1:wLlUxC2vVTPTaE3UD51E0BGOAElKrILxhVSDYQLld5o= +golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/crypto v0.0.0-20200709230013-948cd5f35899/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/crypto v0.0.0-20200728195943-123391ffb6de/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= +golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= +golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= +golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= +golang.org/x/net v0.0.0-20190514140710-3ec191127204/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= +golang.org/x/net v0.0.0-20190603091049-60506f45cf65/go.mod h1:HSz+uSET+XFnRR8LxR5pz3Of3rY3CfYBVs4xY44aLks= +golang.org/x/net v0.0.0-20200707034311-ab3426394381/go.mod h1:/O7V0waA8r7cgGh81Ro3o1hOxt32SMVPicZroKQ2sZA= +golang.org/x/net v0.51.0 h1:94R/GTO7mt3/4wIKpcR5gkGmRLOuE/2hNGeWq/GBIFo= +golang.org/x/net v0.51.0/go.mod h1:aamm+2QF5ogm02fjy5Bb7CQ0WMt1/WVM7FtyaTLlA9Y= +golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= +golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20200217220822-9197077df867/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20200625212154-ddb9806d33ae/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20200728102440-3e129f6d46b1/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.11.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo= +golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= +golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= +golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= +golang.org/x/text v0.35.0 h1:JOVx6vVDFokkpaq1AEptVzLTpDe9KGpj5tR4/X+ybL8= +golang.org/x/text v0.35.0/go.mod h1:khi/HExzZJ2pGnjenulevKNX1W67CUy0AsXcNubPGCA= +golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20190328211700-ab21143f2384/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs= +golang.org/x/tools v0.0.0-20190425163242-31fd60d6bfdc/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= +gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk= +gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E= +google.golang.org/appengine v1.5.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4= +google.golang.org/appengine v1.6.7 h1:FZR1q0exgwxzPzp/aF+VccGrSfxfPpkBqjIIEq3ru6c= +google.golang.org/appengine v1.6.7/go.mod h1:8WjMMxjGQR8xUklV/ARdw2HLXBOI7O7uCIDZVag1xfc= +google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= +google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= +gopkg.in/ini.v1 v1.57.0 h1:9unxIsFcTt4I55uWluz+UmL95q4kdJ0buvQ1ZIqVQww= +gopkg.in/ini.v1 v1.57.0/go.mod h1:pNLf8WUiyNEtQjuu5G5vTm06TEv9tsIgeAvK8hOrP4k= +gopkg.in/natefinch/npipe.v2 v2.0.0-20160621034901-c1b8fa8bdcce h1:+JknDZhAj8YMt7GC73Ei8pv4MzjDUNPHgQWJdtMAaDU= +gopkg.in/natefinch/npipe.v2 v2.0.0-20160621034901-c1b8fa8bdcce/go.mod h1:5AcXVHNjg+BDxry382+8OKon8SEWiKktQR07RKPsv1c= +gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= +gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/scripts/migration/main.go b/scripts/migration/main.go new file mode 100644 index 00000000..7dba6f88 --- /dev/null +++ b/scripts/migration/main.go @@ -0,0 +1,600 @@ +// Command zosmigration automates moving a zos zmachine deployment (VM + its zmount, +// plus the separate network deployment it references) from one node to another. +// +// You give it: S3 credentials, the source VM contract id, and the source + target +// node ids. It then does the whole flow: +// +// 1. creates the S3 bucket (if missing) +// 2. fetches the source VM deployment over RMB and auto-detects its zmount(s), +// the zmachine, and the rootfs size (ZMachine.Size) +// 3. TRANSFER: presigns PUT URLs and calls zos.deployment.transfer on the source +// node (pauses the VM, uploads the zmount disk(s) + the VM rootfs to S3), then +// waits for the objects to appear in the bucket +// 4. (optional) fetches + deploys the referenced network deployment on the target +// 5. PREPARE: presigns GET URLs, creates the target contract, and calls +// zos.deployment.prepare on the target node (provisions network+zmount and pulls +// the data), without booting the VM +// 6. prints the follow-up `-start-contract` command to boot the VM once downloads +// finish (downloads run in the node background with no client-visible signal) +// +// Deployments are re-signed with the owner mnemonic (DeploymentGet re-marshals the +// Env map, shifting the challenge hash), and the source contract's deployment_data +// is copied so the dashboard shows the workload type. ContractID and signatures are +// excluded from the challenge hash, so this stays consistent. +// +// Example: +// +// MNEMONIC="word word ..." go run . \ +// -s3-endpoint https://gateway.storjshare.io \ +// -s3-access -s3-secret \ +// -src-node-id 396 -src-contract 270333 -src-net-contract 270332 \ +// -dst-node-id 371 +// +// then, once "all downloads complete" on the target: +// +// MNEMONIC="..." go run . -dst-node-id 371 -start-contract +package main + +import ( + "context" + "encoding/hex" + "encoding/json" + "flag" + "fmt" + "log" + "net/url" + "os" + "strings" + "time" + + "github.com/minio/minio-go/v7" + "github.com/minio/minio-go/v7/pkg/credentials" + substrate "github.com/threefoldtech/tfchain/clients/tfchain-client-go" + "github.com/threefoldtech/zos_base/client" + "github.com/threefoldtech/zos_base/pkg/gridtypes" + "github.com/threefoldtech/zos_base/pkg/gridtypes/zos" + "github.com/threefoldtech/zos_sdk_go/rmb-sdk-go/peer" +) + +type config struct { + // s3 + s3Endpoint, s3Access, s3Secret, s3Bucket, s3Region string + s3Secure bool + presignExpiry time.Duration + // rmb / chain + mnemonic, keyType, substrateURL, relayURL string + timeout, uploadTimeout time.Duration + // migration + srcNodeID, srcNodeTwin, dstNodeID, dstNodeTwin uint32 + srcContract, srcNetContract uint64 + // start-only + startContract uint64 + // behavior + stageOnly bool +} + +func main() { + if err := run(); err != nil { + log.Fatalf("error: %v", err) + } +} + +func run() error { + cfg, err := parseFlags() + if err != nil { + return err + } + + ctx := context.Background() + mgr := substrate.NewManager(cfg.substrateURL) + rpc, err := peer.NewRpcClient(ctx, cfg.mnemonic, mgr, + peer.WithKeyType(cfg.keyType), peer.WithRelay(cfg.relayURL), peer.WithSession("zosmigration")) + if err != nil { + return fmt.Errorf("create rmb client: %w", err) + } + + // start-only mode: boot a prepared deployment + if cfg.startContract != 0 { + twin, err := resolveTwin(mgr, cfg.dstNodeTwin, cfg.dstNodeID) + if err != nil { + return err + } + node := client.NewNodeClient(twin, rpc) + if err := call(ctx, cfg.timeout, func(c context.Context) error { return node.DeploymentStart(c, cfg.startContract) }); err != nil { + return fmt.Errorf("start: %w", err) + } + log.Printf("start accepted for contract %d — the VM will boot", cfg.startContract) + return nil + } + + return migrate(ctx, cfg, mgr, rpc) +} + +func migrate(ctx context.Context, cfg config, mgr substrate.Manager, rpc *peer.RpcClient) error { + if cfg.srcContract == 0 || cfg.srcNodeID == 0 && cfg.srcNodeTwin == 0 || cfg.dstNodeID == 0 && cfg.dstNodeTwin == 0 { + return fmt.Errorf("-src-contract, a source node (-src-node-id/twin) and a target node (-dst-node-id/twin) are required") + } + + s3, err := newS3(cfg) + if err != nil { + return fmt.Errorf("s3 client: %w", err) + } + if err := ensureBucket(ctx, s3, cfg.s3Bucket, cfg.s3Region); err != nil { + return fmt.Errorf("ensure bucket: %w", err) + } + log.Printf("bucket %q ready on %s", cfg.s3Bucket, cfg.s3Endpoint) + + srcTwin, err := resolveTwin(mgr, cfg.srcNodeTwin, cfg.srcNodeID) + if err != nil { + return fmt.Errorf("resolve source: %w", err) + } + dstTwin, err := resolveTwin(mgr, cfg.dstNodeTwin, cfg.dstNodeID) + if err != nil { + return fmt.Errorf("resolve target: %w", err) + } + src := client.NewNodeClient(srcTwin, rpc) + dst := client.NewNodeClient(dstTwin, rpc) + + // fetch the VM deployment and figure out what to move + vmDl, err := getDeployment(ctx, src, cfg.srcContract, cfg.timeout) + if err != nil { + return fmt.Errorf("fetch VM deployment %d: %w", cfg.srcContract, err) + } + items, err := inspect(vmDl, cfg.srcContract) + if err != nil { + return err + } + log.Printf("VM deployment %d: %d transferable workload(s) detected", cfg.srcContract, len(items)) + for _, it := range items { + log.Printf(" - %s (%s) -> s3://%s/%s", it.name, it.typ, cfg.s3Bucket, it.key) + } + + // ---- TRANSFER (source node) ---- + uploads := make([]client.TransferItem, 0, len(items)) + for _, it := range items { + // remove any stale object so we can detect the fresh upload's completion + _ = s3.RemoveObject(ctx, cfg.s3Bucket, it.key, minio.RemoveObjectOptions{}) + put, err := s3.PresignedPutObject(ctx, cfg.s3Bucket, it.key, cfg.presignExpiry) + if err != nil { + return fmt.Errorf("presign put %s: %w", it.key, err) + } + uploads = append(uploads, client.TransferItem{WorkloadName: it.name, URL: put.String()}) + } + log.Printf("calling transfer on source node %d (this pauses the VM)...", cfg.srcNodeID) + if err := call(ctx, cfg.timeout, func(c context.Context) error { + return src.DeploymentTransfer(c, cfg.srcContract, uploads) + }); err != nil { + // the ack often does not survive a flaky link; the node keeps uploading. + log.Printf("note: transfer ack failed (%v) — continuing; will confirm via S3", err) + } + log.Printf("waiting for %d object(s) to finish uploading to S3 (timeout %s)...", len(items), cfg.uploadTimeout) + if err := waitObjects(ctx, s3, cfg.s3Bucket, items, cfg.uploadTimeout); err != nil { + return fmt.Errorf("transfer did not complete: %w", err) + } + log.Printf("all objects uploaded to S3") + + // ---- NETWORK (deploy on target) ---- + // Auto-detect the network deployment from the VM's network reference if the + // caller didn't pass -src-net-contract. + netContract := cfg.srcNetContract + if netContract == 0 { + if names := networkNames(vmDl); len(names) > 0 { + if c, n, err := findNetworkContract(ctx, src, cfg.timeout, names); err == nil { + netContract = c + log.Printf("auto-detected network %q in contract %d", n, c) + } else { + log.Printf("could not auto-detect network contract (%v); assuming it already exists on the target", err) + } + } + } + if netContract != 0 { + if err := migrateNetwork(ctx, cfg, mgr, src, dst, netContract); err != nil { + return err + } + } + + // ---- PREPARE (target node) ---- + downloads := make([]client.TransferItem, 0, len(items)) + for _, it := range items { + get, err := s3.PresignedGetObject(ctx, cfg.s3Bucket, it.key, cfg.presignExpiry, url.Values{}) + if err != nil { + return fmt.Errorf("presign get %s: %w", it.key, err) + } + downloads = append(downloads, client.TransferItem{WorkloadName: it.name, URL: get.String(), Size: it.size}) + } + + if err := prepZeroVersion(&vmDl); err != nil { + return fmt.Errorf("VM deployment: %w", err) + } + if err := reSign(vmDl.TwinID, cfg, &vmDl); err != nil { + return fmt.Errorf("re-sign VM: %w", err) + } + body := bestBody(mgr, cfg.srcContract, vmDl.Metadata) + vmContract, err := createNodeContract(mgr, cfg, hashHex(&vmDl), body) + if err != nil { + return fmt.Errorf("create VM contract: %w", err) + } + vmDl.ContractID = vmContract + autoStart := !cfg.stageOnly + log.Printf("preparing VM as contract %d on node %d (downloads=%d, auto-start=%v)", vmContract, cfg.dstNodeID, len(downloads), autoStart) + if err := call(ctx, cfg.timeout, func(c context.Context) error { + return dst.DeploymentPrepare(c, vmDl, downloads, autoStart) + }); err != nil { + log.Printf("note: prepare ack failed (%v) — the node keeps working; continuing", err) + } + + if cfg.stageOnly { + log.Printf("") + log.Printf("=== migration staged: VM contract %d on node %d ===", vmContract, cfg.dstNodeID) + log.Printf("the target is pulling the data in the background; boot it later with:") + log.Printf(" MNEMONIC=... go run . -dst-node-id %d -start-contract %d", cfg.dstNodeID, vmContract) + return nil + } + + // wait until the target finishes downloading AND boots the VM: the zmachine + // workload only reaches Ok after the auto-start that follows the downloads. + vmName := zmachineName(items) + log.Printf("waiting for the target to finish downloading and boot the VM (up to %s)...", cfg.uploadTimeout) + if err := waitWorkloadOk(ctx, dst, vmContract, vmName, cfg.uploadTimeout); err != nil { + return fmt.Errorf("VM did not come up on the target: %w", err) + } + log.Printf("") + log.Printf("=== migration complete: VM %q is up on node %d (contract %d) ===", vmName, cfg.dstNodeID, vmContract) + return nil +} + +func zmachineName(items []transferItem) gridtypes.Name { + for _, it := range items { + if it.typ == zos.ZMachineType || it.typ == zos.ZMachineLightType { + return it.name + } + } + return "" +} + +func migrateNetwork(ctx context.Context, cfg config, mgr substrate.Manager, src, dst *client.NodeClient, srcNetContract uint64) error { + netDl, err := getDeployment(ctx, src, srcNetContract, cfg.timeout) + if err != nil { + return fmt.Errorf("fetch network deployment %d: %w", srcNetContract, err) + } + netName, err := firstNetworkName(netDl) + if err != nil { + return err + } + if err := prepZeroVersion(&netDl); err != nil { + return fmt.Errorf("network deployment: %w", err) + } + if err := reSign(netDl.TwinID, cfg, &netDl); err != nil { + return fmt.Errorf("re-sign network: %w", err) + } + body := bestBody(mgr, srcNetContract, netDl.Metadata) + netContract, err := createNodeContract(mgr, cfg, hashHex(&netDl), body) + if err != nil { + return fmt.Errorf("create network contract: %w", err) + } + netDl.ContractID = netContract + log.Printf("deploying network %q as contract %d on node %d", netName, netContract, cfg.dstNodeID) + if err := call(ctx, cfg.timeout, func(c context.Context) error { return dst.DeploymentDeploy(c, netDl) }); err != nil { + return fmt.Errorf("deploy network: %w", err) + } + log.Printf("waiting for network %q to provision...", netName) + if err := waitWorkloadOk(ctx, dst, netContract, netName, 3*time.Minute); err != nil { + return fmt.Errorf("network did not provision: %w", err) + } + log.Printf("network %q is up on node %d", netName, cfg.dstNodeID) + return nil +} + +// transferItem is a resolved workload to move: its S3 object key and (for a +// zmachine) the rootfs volume size to recreate on the target. +type transferItem struct { + name gridtypes.Name + typ gridtypes.WorkloadType + key string + size gridtypes.Unit +} + +// inspect finds the zmount(s) and the zmachine in the deployment and derives the +// rootfs size from ZMachine.Size. +func inspect(dl gridtypes.Deployment, srcContract uint64) ([]transferItem, error) { + var items []transferItem + for i := range dl.Workloads { + wl := &dl.Workloads[i] + switch wl.Type { + case zos.ZMountType: + items = append(items, transferItem{ + name: wl.Name, typ: wl.Type, + key: fmt.Sprintf("mig-%d/%s.raw", srcContract, wl.Name), + }) + case zos.ZMachineType, zos.ZMachineLightType: + var zm zos.ZMachine + if err := json.Unmarshal(wl.Data, &zm); err != nil { + return nil, fmt.Errorf("decode zmachine %q: %w", wl.Name, err) + } + items = append(items, transferItem{ + name: wl.Name, typ: wl.Type, + key: fmt.Sprintf("mig-%d/%s-rootfs.tar", srcContract, wl.Name), + size: zm.Size, + }) + } + } + if len(items) == 0 { + return nil, fmt.Errorf("no zmount/zmachine workloads found in deployment %d", srcContract) + } + return items, nil +} + +// networkNames returns the znet names the VM's zmachine(s) attach to. +func networkNames(dl gridtypes.Deployment) []gridtypes.Name { + var names []gridtypes.Name + seen := map[gridtypes.Name]bool{} + for i := range dl.Workloads { + wl := &dl.Workloads[i] + if wl.Type != zos.ZMachineType && wl.Type != zos.ZMachineLightType { + continue + } + var zm zos.ZMachine + if json.Unmarshal(wl.Data, &zm) != nil { + continue + } + for _, iface := range zm.Network.Interfaces { + if iface.Network != "" && !seen[iface.Network] { + seen[iface.Network] = true + names = append(names, iface.Network) + } + } + } + return names +} + +// findNetworkContract lists the twin's deployments on the source node and returns +// the contract that holds a network workload matching one of the given names. +func findNetworkContract(ctx context.Context, node *client.NodeClient, timeout time.Duration, names []gridtypes.Name) (uint64, gridtypes.Name, error) { + c, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + deps, err := node.DeploymentList(c) + if err != nil { + return 0, "", err + } + want := map[gridtypes.Name]bool{} + for _, n := range names { + want[n] = true + } + for _, d := range deps { + for _, wl := range d.Workloads { + if (wl.Type == zos.NetworkType || wl.Type == zos.NetworkLightType) && want[wl.Name] { + return d.ContractID, wl.Name, nil + } + } + } + return 0, "", fmt.Errorf("no network deployment found for %v on the source node", names) +} + +// ---- S3 helpers ---- + +func newS3(cfg config) (*minio.Client, error) { + u, err := url.Parse(cfg.s3Endpoint) + if err != nil { + return nil, err + } + host := u.Host + if host == "" { + host = cfg.s3Endpoint + } + return minio.New(host, &minio.Options{ + Creds: credentials.NewStaticV4(cfg.s3Access, cfg.s3Secret, ""), + Secure: cfg.s3Secure, + Region: cfg.s3Region, + }) +} + +func ensureBucket(ctx context.Context, s3 *minio.Client, bucket, region string) error { + err := s3.MakeBucket(ctx, bucket, minio.MakeBucketOptions{Region: region}) + if err == nil { + return nil + } + if exists, e := s3.BucketExists(ctx, bucket); e == nil && exists { + return nil + } + return err +} + +func waitObjects(ctx context.Context, s3 *minio.Client, bucket string, items []transferItem, timeout time.Duration) error { + deadline := time.Now().Add(timeout) + done := make(map[string]bool) + for { + for _, it := range items { + if done[it.key] { + continue + } + if info, err := s3.StatObject(ctx, bucket, it.key, minio.StatObjectOptions{}); err == nil { + done[it.key] = true + log.Printf(" uploaded: %s (%d bytes)", it.key, info.Size) + } + } + if len(done) == len(items) { + return nil + } + if time.Now().After(deadline) { + return fmt.Errorf("timed out with %d/%d objects uploaded", len(done), len(items)) + } + time.Sleep(15 * time.Second) + } +} + +// ---- deployment / chain helpers ---- + +func getDeployment(ctx context.Context, node *client.NodeClient, contract uint64, timeout time.Duration) (gridtypes.Deployment, error) { + c, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + return node.DeploymentGet(c, contract) +} + +func call(ctx context.Context, timeout time.Duration, fn func(context.Context) error) error { + c, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + return fn(c) +} + +func prepZeroVersion(dl *gridtypes.Deployment) error { + if dl.Version != 0 { + return fmt.Errorf("deployment version is %d, but a fresh prepare/deploy requires version 0", dl.Version) + } + return dl.Valid() +} + +func hashHex(dl *gridtypes.Deployment) string { + h, _ := dl.ChallengeHash() + return hex.EncodeToString(h) +} + +// bestBody returns the source contract's on-chain deployment_data (so the dashboard +// shows the type), falling back to the deployment metadata. +func bestBody(mgr substrate.Manager, srcContract uint64, metadata string) string { + sub, err := mgr.Substrate() + if err != nil { + return metadata + } + defer sub.Close() + c, err := sub.GetContract(srcContract) + if err != nil { + log.Printf("warning: could not read deployment_data of contract %d (dashboard label may be blank): %v", srcContract, err) + return metadata + } + if d := c.ContractType.NodeContract.DeploymentData; d != "" { + return d + } + return metadata +} + +func createNodeContract(mgr substrate.Manager, cfg config, hashHex, body string) (uint64, error) { + id, err := identityFrom(cfg.keyType, cfg.mnemonic) + if err != nil { + return 0, err + } + sub, err := mgr.Substrate() + if err != nil { + return 0, fmt.Errorf("connect substrate: %w", err) + } + defer sub.Close() + return sub.CreateNodeContract(id, cfg.dstNodeID, body, hashHex, 0, nil) +} + +func reSign(twin uint32, cfg config, dl *gridtypes.Deployment) error { + id, err := identityFrom(cfg.keyType, cfg.mnemonic) + if err != nil { + return err + } + return dl.Sign(twin, identitySigner{id: id, keyType: cfg.keyType}) +} + +func firstNetworkName(dl gridtypes.Deployment) (gridtypes.Name, error) { + for _, wl := range dl.Workloads { + if wl.Type == zos.NetworkType || wl.Type == zos.NetworkLightType { + return wl.Name, nil + } + } + return "", fmt.Errorf("no network workload found in the network deployment") +} + +func waitWorkloadOk(ctx context.Context, node *client.NodeClient, contract uint64, name gridtypes.Name, timeout time.Duration) error { + deadline := time.Now().Add(timeout) + for { + if dl, err := getDeployment(ctx, node, contract, 20*time.Second); err == nil { + if wl, err := dl.Get(name); err == nil { + if wl.Result.State.IsOkay() { + return nil + } + if wl.Result.State == gridtypes.StateError { + return fmt.Errorf("workload %q errored: %s", name, wl.Result.Error) + } + } + } + if time.Now().After(deadline) { + return fmt.Errorf("timeout waiting for %q", name) + } + time.Sleep(2 * time.Second) + } +} + +func resolveTwin(mgr substrate.Manager, twin, nodeID uint32) (uint32, error) { + if twin != 0 { + return twin, nil + } + if nodeID == 0 { + return 0, fmt.Errorf("a node twin or node id is required") + } + sub, err := mgr.Substrate() + if err != nil { + return 0, fmt.Errorf("connect substrate: %w", err) + } + defer sub.Close() + node, err := sub.GetNode(nodeID) + if err != nil { + return 0, fmt.Errorf("resolve node %d: %w", nodeID, err) + } + return uint32(node.TwinID), nil +} + +func identityFrom(keyType, mnemonic string) (substrate.Identity, error) { + if keyType == "ed25519" { + return substrate.NewIdentityFromEd25519Phrase(mnemonic) + } + return substrate.NewIdentityFromSr25519Phrase(mnemonic) +} + +type identitySigner struct { + id substrate.Identity + keyType string +} + +func (s identitySigner) Sign(msg []byte) ([]byte, error) { return s.id.Sign(msg) } +func (s identitySigner) Type() string { return s.keyType } + +// ---- flags ---- + +func parseFlags() (config, error) { + var cfg config + var ( + srcNodeID = flag.Uint("src-node-id", 0, "source node id") + srcNodeTwin = flag.Uint("src-node-twin", 0, "source node twin (else -src-node-id)") + dstNodeID = flag.Uint("dst-node-id", 0, "target node id") + dstNodeTwin = flag.Uint("dst-node-twin", 0, "target node twin (else -dst-node-id)") + ) + flag.StringVar(&cfg.mnemonic, "mnemonic", os.Getenv("MNEMONIC"), "owner twin mnemonic (or MNEMONIC env)") + flag.StringVar(&cfg.keyType, "keytype", "sr25519", "mnemonic key type: sr25519 or ed25519") + flag.StringVar(&cfg.substrateURL, "substrate", "wss://tfchain.dev.grid.tf/ws", "substrate websocket url") + flag.StringVar(&cfg.relayURL, "relay", "wss://relay.dev.grid.tf", "rmb relay url") + flag.DurationVar(&cfg.timeout, "timeout", 60*time.Second, "per-RMB-call timeout") + flag.DurationVar(&cfg.uploadTimeout, "upload-timeout", 2*time.Hour, "how long to wait for S3 uploads to finish") + + flag.StringVar(&cfg.s3Endpoint, "s3-endpoint", "", "S3 endpoint url (e.g. https://gateway.storjshare.io)") + flag.StringVar(&cfg.s3Access, "s3-access", os.Getenv("S3_ACCESS"), "S3 access key (or S3_ACCESS env)") + flag.StringVar(&cfg.s3Secret, "s3-secret", os.Getenv("S3_SECRET"), "S3 secret key (or S3_SECRET env)") + flag.StringVar(&cfg.s3Bucket, "s3-bucket", "zos-migration", "S3 bucket to use") + flag.StringVar(&cfg.s3Region, "s3-region", "global", "S3 region for signing") + flag.BoolVar(&cfg.s3Secure, "s3-secure", true, "use https for S3") + flag.DurationVar(&cfg.presignExpiry, "presign-expiry", 24*time.Hour, "presigned URL lifetime") + + flag.Uint64Var(&cfg.srcContract, "src-contract", 0, "source VM contract id") + flag.Uint64Var(&cfg.srcNetContract, "src-net-contract", 0, "source network contract id (optional if the network already exists on target)") + flag.Uint64Var(&cfg.startContract, "start-contract", 0, "start-only: boot this prepared contract on -dst-node-id and exit") + flag.BoolVar(&cfg.stageOnly, "stage-only", false, "prepare without auto-start/wait (print the start command and exit)") + flag.Parse() + + cfg.srcNodeID, cfg.srcNodeTwin = uint32(*srcNodeID), uint32(*srcNodeTwin) + cfg.dstNodeID, cfg.dstNodeTwin = uint32(*dstNodeID), uint32(*dstNodeTwin) + + if cfg.mnemonic == "" { + return cfg, fmt.Errorf("-mnemonic (or MNEMONIC env) is required") + } + if cfg.startContract == 0 { // full migration needs S3 + if cfg.s3Endpoint == "" || cfg.s3Access == "" || cfg.s3Secret == "" { + return cfg, fmt.Errorf("-s3-endpoint, -s3-access and -s3-secret are required") + } + if !strings.Contains(cfg.s3Endpoint, "://") { + cfg.s3Endpoint = "https://" + cfg.s3Endpoint + } + } + return cfg, nil +}