Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions lib/hypervisor/cloudhypervisor/cloudhypervisor.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ func CapabilitiesForVersion(v vmm.CHVersion) hypervisor.Capabilities {
case vmm.V51_1:
caps.SupportsDiskResize = true
caps.SupportsConcurrentForkPrepare = true
caps.SupportsSnapshotBaseReuse = experimentalDiffSnapshotsEnabled()
}
return caps
}
Expand Down Expand Up @@ -163,15 +164,40 @@ func (c *CloudHypervisor) Resume(ctx context.Context) error {

// Snapshot creates a VM snapshot.
func (c *CloudHypervisor) Snapshot(ctx context.Context, destPath string) error {
diff, err := prepareDiffSnapshotDestination(destPath)
if err != nil {
return fmt.Errorf("prepare diff snapshot destination: %w", err)
}

snapshotURL := "file://" + destPath
snapshotConfig := vmm.VmSnapshotConfig{DestinationUrl: &snapshotURL}
if experimentalDiffSnapshotsEnabled() {
snapshotType := vmm.Full
if diff {
snapshotType = vmm.Diff
}
snapshotConfig.SnapshotType = &snapshotType
}
resp, err := c.client.PutVmSnapshotWithResponse(ctx, snapshotConfig)
if err != nil {
return fmt.Errorf("snapshot: %w", err)
}
if resp.StatusCode() != 204 {
return fmt.Errorf("snapshot failed with status %d", resp.StatusCode())
}
if diff {
stats, err := mergeCloudHypervisorDiff(destPath)
if err != nil {
return fmt.Errorf("merge diff snapshot: %w", err)
}
logger.FromContext(ctx).InfoContext(ctx, "merged Cloud Hypervisor diff snapshot",
"snapshot_dir", destPath,
"delta_bytes", stats.DeltaBytes,
"extents", stats.ExtentCount,
"reflinked_bytes", stats.ReflinkedBytes,
"copied_bytes", stats.CopiedBytes,
)
}
optimized, err := prepareSnapshotForKernelPaging(destPath)
if err != nil {
return fmt.Errorf("prepare kernel-paged snapshot: %w", err)
Expand Down
55 changes: 55 additions & 0 deletions lib/hypervisor/cloudhypervisor/diff_snapshot.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
package cloudhypervisor

import (
"errors"
"fmt"
"os"
"path/filepath"
"strings"
)

const experimentalDiffSnapshotsEnv = "HYPEMAN_EXPERIMENTAL_CH_DIFF_SNAPSHOTS"

const cloudHypervisorDiffMemoryFile = cloudHypervisorMemoryFile + ".diff"

type diffMergeStats struct {
DeltaBytes int64
ExtentCount int
ReflinkedBytes int64
CopiedBytes int64
}

func experimentalDiffSnapshotsEnabled() bool {
switch strings.ToLower(strings.TrimSpace(os.Getenv(experimentalDiffSnapshotsEnv))) {
case "1", "true", "yes", "on":
return true
default:
return false
}
}

// prepareDiffSnapshotDestination leaves the retained memory baseline in place
// and removes metadata that Cloud Hypervisor recreates with O_EXCL.
func prepareDiffSnapshotDestination(snapshotDir string) (bool, error) {
if !experimentalDiffSnapshotsEnabled() {
return false, nil
}
memoryPath := filepath.Join(snapshotDir, cloudHypervisorMemoryFile)
if _, err := os.Stat(memoryPath); err != nil {
if errors.Is(err, os.ErrNotExist) {
return false, nil
}
return false, fmt.Errorf("stat retained snapshot memory: %w", err)
}

for _, name := range []string{
cloudHypervisorConfigFile,
cloudHypervisorStateFile,
cloudHypervisorDiffMemoryFile,
} {
if err := os.Remove(filepath.Join(snapshotDir, name)); err != nil && !errors.Is(err, os.ErrNotExist) {
return false, fmt.Errorf("remove retained snapshot file %s: %w", name, err)
}
}
return true, nil
}
165 changes: 165 additions & 0 deletions lib/hypervisor/cloudhypervisor/diff_snapshot_linux.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,165 @@
//go:build linux

package cloudhypervisor

import (
"errors"
"fmt"
"io"
"os"
"path/filepath"

"golang.org/x/sys/unix"
)

func mergeCloudHypervisorDiff(snapshotDir string) (diffMergeStats, error) {
var stats diffMergeStats
basePath := filepath.Join(snapshotDir, cloudHypervisorMemoryFile)
diffPath := filepath.Join(snapshotDir, cloudHypervisorDiffMemoryFile)

src, err := os.Open(diffPath)
if err != nil {
return stats, fmt.Errorf("open diff snapshot memory: %w", err)
}
defer src.Close()
// Resolve delayed allocation before cloning extents. FICLONERANGE can
// otherwise clone the pre-write extent state while dirty pages still live
// only in the source file's page cache.
if err := src.Sync(); err != nil {
return stats, fmt.Errorf("sync diff snapshot memory: %w", err)
}
dst, err := os.OpenFile(basePath, os.O_RDWR, 0)
if err != nil {
return stats, fmt.Errorf("open retained snapshot memory: %w", err)
}
defer dst.Close()

srcInfo, err := src.Stat()
if err != nil {
return stats, fmt.Errorf("stat diff snapshot memory: %w", err)
}
dstInfo, err := dst.Stat()
if err != nil {
return stats, fmt.Errorf("stat retained snapshot memory: %w", err)
}
if srcInfo.Size() != dstInfo.Size() {
return stats, fmt.Errorf("diff snapshot memory size %d does not match baseline %d", srcInfo.Size(), dstInfo.Size())
}

srcFD := int(src.Fd())
dstFD := int(dst.Fd())
size := srcInfo.Size()
offset := int64(0)
cloneRanges := true
buf := make([]byte, 1<<20)

for offset < size {
dataStart, err := unix.Seek(srcFD, offset, unix.SEEK_DATA)
if err != nil {
if errors.Is(err, unix.ENXIO) {
break
}
return stats, fmt.Errorf("seek diff data at %d: %w", offset, err)
}
dataEnd, err := unix.Seek(srcFD, dataStart, unix.SEEK_HOLE)
if err != nil {
if errors.Is(err, unix.ENXIO) {
dataEnd = size
} else {
return stats, fmt.Errorf("seek diff hole at %d: %w", dataStart, err)
}
}
if dataEnd > size {
dataEnd = size
}
if dataEnd <= dataStart {
return stats, fmt.Errorf("invalid diff extent [%d,%d)", dataStart, dataEnd)
}

length := dataEnd - dataStart
stats.DeltaBytes += length
stats.ExtentCount++
if cloneRanges {
clone := unix.FileCloneRange{
Src_fd: int64(srcFD),
Src_offset: uint64(dataStart),
Src_length: uint64(length),
Dest_offset: uint64(dataStart),
}
if err := unix.IoctlFileCloneRange(dstFD, &clone); err == nil {
stats.ReflinkedBytes += length
offset = dataEnd
continue
} else if !diffRangeCloneUnsupported(err) {
return stats, fmt.Errorf("reflink diff extent [%d,%d): %w", dataStart, dataEnd, err)
}
cloneRanges = false
}

if err := copyDiffExtent(srcFD, dstFD, dataStart, length, buf); err != nil {
return stats, fmt.Errorf("copy diff extent [%d,%d): %w", dataStart, dataEnd, err)
}
stats.CopiedBytes += length
offset = dataEnd
}

if err := dst.Sync(); err != nil {
return stats, fmt.Errorf("sync merged snapshot memory: %w", err)
}
if err := os.Remove(diffPath); err != nil {
return stats, fmt.Errorf("remove merged diff snapshot: %w", err)
}
dir, err := os.Open(snapshotDir)
if err != nil {
return stats, fmt.Errorf("open snapshot directory: %w", err)
}
if err := dir.Sync(); err != nil {
_ = dir.Close()
return stats, fmt.Errorf("sync snapshot directory: %w", err)
}
if err := dir.Close(); err != nil {
return stats, fmt.Errorf("close snapshot directory: %w", err)
}
return stats, nil
}

func copyDiffExtent(srcFD, dstFD int, offset, length int64, buf []byte) error {
position := offset
remaining := length
for remaining > 0 {
chunk := int64(len(buf))
if remaining < chunk {
chunk = remaining
}
n, err := unix.Pread(srcFD, buf[:int(chunk)], position)
if err != nil {
return err
}
if n == 0 {
return io.ErrUnexpectedEOF
}
written := 0
for written < n {
wn, err := unix.Pwrite(dstFD, buf[written:n], position+int64(written))
if err != nil {
return err
}
if wn == 0 {
return io.ErrShortWrite
}
written += wn
}
position += int64(n)
remaining -= int64(n)
}
return nil
}

func diffRangeCloneUnsupported(err error) bool {
return errors.Is(err, unix.EINVAL) ||
errors.Is(err, unix.ENOTSUP) ||
errors.Is(err, unix.EOPNOTSUPP) ||
errors.Is(err, unix.EXDEV) ||
errors.Is(err, unix.ENOTTY) ||
errors.Is(err, unix.ETXTBSY)
}
52 changes: 52 additions & 0 deletions lib/hypervisor/cloudhypervisor/diff_snapshot_linux_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
//go:build linux

package cloudhypervisor

import (
"bytes"
"errors"
"os"
"path/filepath"
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"golang.org/x/sys/unix"
)

func TestMergeCloudHypervisorDiff(t *testing.T) {
dir := t.TempDir()
const size = 4 << 20
base := bytes.Repeat([]byte{0x7b}, size)
basePath := filepath.Join(dir, cloudHypervisorMemoryFile)
diffPath := filepath.Join(dir, cloudHypervisorDiffMemoryFile)
require.NoError(t, os.WriteFile(basePath, base, 0600))

diff, err := os.OpenFile(diffPath, os.O_CREATE|os.O_RDWR, 0600)
require.NoError(t, err)
require.NoError(t, diff.Truncate(size))
changed := bytes.Repeat([]byte{0x2a}, 4096)
zeroed := make([]byte, 4096)
_, err = diff.WriteAt(changed, 64*4096)
require.NoError(t, err)
_, err = diff.WriteAt(zeroed, 512*4096)
require.NoError(t, err)
require.NoError(t, diff.Sync())
require.NoError(t, diff.Close())

stats, err := mergeCloudHypervisorDiff(dir)
if errors.Is(err, unix.EINVAL) || errors.Is(err, unix.ENOTSUP) || errors.Is(err, unix.EOPNOTSUPP) {
t.Skipf("filesystem does not support sparse extent discovery: %v", err)
}
require.NoError(t, err)
assert.GreaterOrEqual(t, stats.DeltaBytes, int64(2*4096))
assert.Equal(t, stats.DeltaBytes, stats.ReflinkedBytes+stats.CopiedBytes)
assert.NoFileExists(t, diffPath)

want := append([]byte(nil), base...)
copy(want[64*4096:], changed)
copy(want[512*4096:], zeroed)
got, err := os.ReadFile(basePath)
require.NoError(t, err)
assert.Equal(t, want, got)
}
39 changes: 39 additions & 0 deletions lib/hypervisor/cloudhypervisor/diff_snapshot_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
package cloudhypervisor

import (
"os"
"path/filepath"
"testing"

"github.com/kernel/hypeman/lib/vmm"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

func TestPrepareDiffSnapshotDestination(t *testing.T) {
t.Setenv(experimentalDiffSnapshotsEnv, "true")
dir := t.TempDir()
for _, name := range []string{
cloudHypervisorMemoryFile,
cloudHypervisorConfigFile,
cloudHypervisorStateFile,
cloudHypervisorDiffMemoryFile,
} {
require.NoError(t, os.WriteFile(filepath.Join(dir, name), []byte(name), 0600))
}

diff, err := prepareDiffSnapshotDestination(dir)
require.NoError(t, err)
assert.True(t, diff)
assert.FileExists(t, filepath.Join(dir, cloudHypervisorMemoryFile))
assert.NoFileExists(t, filepath.Join(dir, cloudHypervisorConfigFile))
assert.NoFileExists(t, filepath.Join(dir, cloudHypervisorStateFile))
assert.NoFileExists(t, filepath.Join(dir, cloudHypervisorDiffMemoryFile))
}

func TestExperimentalDiffSnapshotCapability(t *testing.T) {
t.Setenv(experimentalDiffSnapshotsEnv, "true")
assert.True(t, CapabilitiesForVersion(vmm.V51_1).SupportsSnapshotBaseReuse)
t.Setenv(experimentalDiffSnapshotsEnv, "false")
assert.False(t, CapabilitiesForVersion(vmm.V51_1).SupportsSnapshotBaseReuse)
}
9 changes: 9 additions & 0 deletions lib/hypervisor/cloudhypervisor/diff_snapshot_unsupported.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
//go:build !linux

package cloudhypervisor

import "errors"

func mergeCloudHypervisorDiff(string) (diffMergeStats, error) {
return diffMergeStats{}, errors.New("cloud hypervisor diff snapshots require Linux")
}
3 changes: 3 additions & 0 deletions lib/hypervisor/cloudhypervisor/process.go
Original file line number Diff line number Diff line change
Expand Up @@ -261,6 +261,9 @@ func (s *Starter) RestoreVM(ctx context.Context, p *paths.Paths, version string,
SourceUrl: sourceURL,
Prefault: ptr(false),
}
if experimentalDiffSnapshotsEnabled() {
restoreConfig.TrackDirtyPages = ptr(true)
}
resp, err := hv.client.PutVmRestoreWithResponse(ctx, restoreConfig)
if err != nil {
return 0, nil, fmt.Errorf("restore: %w", err)
Expand Down
Loading
Loading