mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Remove unnecessary fsync during ResumeActor (#1883)
While benchmarking pause actor I ran into multiple bottlenecks. I'm currently importing https://github.com/agent-substrate/substrate/pull/935 into my branch but https://github.com/agent-substrate/substrate/pull/1876 is another alternative to solve the first bottleneck. The 2nd bottleneck (primary target of this PR) is due to fsyncs being called in writeFileAtomic() for sandbox-assets.json and SystemInfoVolume files. Because these are new files that require ext4 filesystem metadata updates, and we have ext4 mounted in ordered mode, calling fsync on the file + directory causes all other pending writes (such as other checkpoints) to be flushed. So the file creation ends up being blocked behind other actor's checkpoint syncs. These files do not actually need to be flushed to disk. SystemInfo is recreatable. sandbox-assets.json is written during Run/Resume(), and read during Checkpoint(). But once the checkpoint is finished, the manifest is saved with the snapshot. If the node crashes between Run/Resume and Checkpoint, the actor is crashed anyway. ### Single-Node Pause/Resume Benchmark Summary **Configuration:** 17 workers, 15 actors, `glutton` (`512 MiB` RAM capacity, `32 MiB` churn per cycle), `120s` run, `1.0s` wait, `100%` pause (`0%` suspend) | Metric | Baseline (`upstream/main`) | + Commit 1 (`os.Link` local checkpoint) | + Commit 2 (remove `fsync` on ephemeral volumes) | | :--- | :---: | :---: | :---: | | **Completed Cycles** (`120s`) | `106.5` (`112` pauses / `101` resumes) | `160.0` (`165` pauses / `155` resumes) | **`212.5`** (`214` pauses / `211` resumes) | | **Cycle Throughput** | `53.3 cycles/min` (`0.89/s`) | `80.0 cycles/min` (`1.33/s`, **+50.2%**) | **`106.3 cycles/min`** (`1.77/s`, **2.00x**) | | **`ResumeActor` p50** | `6,000 ms` | `4,100 ms` | **`210 ms`** (**28.6x faster**) | | **`ResumeActor` avg** | `7,023 ms` | `4,867 ms` | **`246 ms`** (**28.5x faster**) | | **`ResumeActor` p90 / p99** | `15,000 ms` / `19,000 ms` | `9,400 ms` / `20,000 ms` | **`350 ms` / `590 ms`** | - [x] Tests pass - [x] Appropriate changes to documentation are included in the PR --------- Co-authored-by: Jefftree <jeffrey.ying86@live.com>
This commit is contained in:
+34
-18
@@ -68,6 +68,7 @@ import (
|
||||
"go.opentelemetry.io/otel/metric"
|
||||
semconv "go.opentelemetry.io/otel/semconv/v1.40.0"
|
||||
"golang.org/x/sync/errgroup"
|
||||
"golang.org/x/sys/unix"
|
||||
"google.golang.org/api/option"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
@@ -684,6 +685,12 @@ func (s *AteomHerder) Checkpoint(ctx context.Context, req *ateletpb.CheckpointRe
|
||||
// and the control plane tracks only a single local snapshot, which this
|
||||
// checkpoint either overwrites (pause) or clears (suspend).
|
||||
//
|
||||
// Do not move this above CheckpointWorkload to keep MergeDeltaIntoBase on its
|
||||
// in-place path: that leaves the whole checkpoint window with no local snapshot
|
||||
// while LocalSnapshotInfo still names the pruned one, and a crash there strands
|
||||
// the actor for good (resume never falls back to object storage, RequiredNodes
|
||||
// pins it to this node, nothing clears the field).
|
||||
//
|
||||
// Best-effort: if this fail, the actor's terminate prunes again.
|
||||
if err := pruneLocalCheckpoints(ctx, actorUID); err != nil {
|
||||
slog.WarnContext(ctx, "failed to prune superseded local checkpoints", slog.Any("actor", actorRef), slog.Any("err", err))
|
||||
@@ -1320,6 +1327,26 @@ func (s *AteomHerder) copyLocalCheckpoint(ctx context.Context, snapshotName stri
|
||||
}
|
||||
src := filepath.Join(srcDir, snapshotName, fileName)
|
||||
dst := filepath.Join(dstDir, fileName)
|
||||
// Link rather than copy. The local checkpoint lives under the same actor dir
|
||||
// as the restore staging area, so this stages the memory image in constant
|
||||
// time instead of re-writing its whole working set. Nothing rewrites the
|
||||
// shared inode: CH demand-pages from the staged image read-only,
|
||||
// rewriteSnapshotSocketPaths renames its rewritten config.json into place
|
||||
// rather than truncating, and MergeDeltaIntoBase refuses its in-place overlay
|
||||
// once the image carries a second link.
|
||||
//
|
||||
// EXDEV alone falls back to copying, so an unexpected link failure surfaces
|
||||
// instead of silently reverting to the full copy this exists to remove. It
|
||||
// also keeps sparsefile.CopyFile off a dst that is already a link to src, where its
|
||||
// O_TRUNC would empty both and report a successful copy of the old size.
|
||||
switch err := linkFile(src, dst); {
|
||||
case err == nil:
|
||||
continue
|
||||
case !errors.Is(err, unix.EXDEV):
|
||||
return fmt.Errorf("failed to link %s to %s: %w", src, dst, err)
|
||||
}
|
||||
slog.WarnContext(ctx, "local checkpoint and restore dir are on different filesystems; copying instead of linking",
|
||||
slog.String("src", src), slog.String("dst", dst))
|
||||
if _, err := sparsefile.CopyFile(src, dst); err != nil {
|
||||
return fmt.Errorf("failed to copy %s to %s: %w", src, dst, err)
|
||||
}
|
||||
@@ -1328,6 +1355,10 @@ func (s *AteomHerder) copyLocalCheckpoint(ctx context.Context, snapshotName stri
|
||||
return nil
|
||||
}
|
||||
|
||||
// linkFile is os.Link, indirected so a test can force the cross-filesystem
|
||||
// fallback in copyLocalCheckpoint without mounting a second filesystem.
|
||||
var linkFile = os.Link
|
||||
|
||||
// goldenOnlyFiles returns the golden snapshot files not shadowed by the
|
||||
// actor's own snapshot: on a DATA_ON_GOLDEN restore the actor's files (the
|
||||
// durable-dir data) win name collisions, and the golden snapshot supplies
|
||||
@@ -1800,10 +1831,8 @@ func validateUploadPausedCheckpointRequest(req *ateletpb.UploadPausedCheckpointR
|
||||
}
|
||||
|
||||
// writeFileAtomic writes data to path by writing a temp file in the same
|
||||
// directory, syncing, and renaming it over the target, then syncing the
|
||||
// parent directory so the rename is durable. The identity directory is
|
||||
// bind-mounted into actors, so the file must change atomically: a reader
|
||||
// must never observe a truncated or partially written value.
|
||||
// directory and renaming it over the target so readers never observe a
|
||||
// truncated or partially written value.
|
||||
func writeFileAtomic(path string, data []byte, perm os.FileMode) error {
|
||||
f, err := os.CreateTemp(filepath.Dir(path), "."+filepath.Base(path)+".tmp-*")
|
||||
if err != nil {
|
||||
@@ -1819,23 +1848,10 @@ func writeFileAtomic(path string, data []byte, perm os.FileMode) error {
|
||||
f.Close()
|
||||
return err
|
||||
}
|
||||
if err := f.Sync(); err != nil {
|
||||
f.Close()
|
||||
return err
|
||||
}
|
||||
if err := f.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := os.Rename(f.Name(), path); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
dir, err := os.Open(filepath.Dir(path))
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer dir.Close()
|
||||
return dir.Sync()
|
||||
return os.Rename(f.Name(), path)
|
||||
}
|
||||
|
||||
// resetActorDirs empties the actor's directories and leaves them in place for
|
||||
|
||||
@@ -29,6 +29,7 @@ import (
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"syscall"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -44,6 +45,7 @@ import (
|
||||
"github.com/google/go-cmp/cmp"
|
||||
"github.com/klauspost/compress/zstd"
|
||||
"github.com/spf13/pflag"
|
||||
"golang.org/x/sys/unix"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
@@ -79,6 +81,94 @@ func TestPortFlagDefault(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestCopyLocalCheckpointLinks covers staging a local checkpoint into the restore
|
||||
// dir: the files must land as extra links to the cached snapshot rather than
|
||||
// copies, so a resume does not rewrite the image's working set. Sharing the inode
|
||||
// is what MergeDeltaIntoBase's Nlink check keys off to refuse its in-place overlay.
|
||||
func TestCopyLocalCheckpointLinks(t *testing.T) {
|
||||
const snapshot = "snap-1"
|
||||
want := []byte("checkpoint pages")
|
||||
|
||||
newDirs := func(t *testing.T) (srcDir, dstDir string) {
|
||||
t.Helper()
|
||||
root := t.TempDir()
|
||||
srcDir, dstDir = filepath.Join(root, "local-checkpoint"), filepath.Join(root, "restore-state")
|
||||
if err := os.MkdirAll(filepath.Join(srcDir, snapshot), 0o700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.MkdirAll(dstDir, 0o700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(srcDir, snapshot, "memory-ranges"), want, 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return srcDir, dstDir
|
||||
}
|
||||
inode := func(t *testing.T, path string) uint64 {
|
||||
t.Helper()
|
||||
fi, err := os.Stat(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return fi.Sys().(*syscall.Stat_t).Ino
|
||||
}
|
||||
|
||||
t.Run("links when it can", func(t *testing.T) {
|
||||
srcDir, dstDir := newDirs(t)
|
||||
s := &AteomHerder{}
|
||||
if err := s.copyLocalCheckpoint(context.Background(), snapshot, srcDir, dstDir, []string{"memory-ranges"}); err != nil {
|
||||
t.Fatalf("copyLocalCheckpoint: %v", err)
|
||||
}
|
||||
src := filepath.Join(srcDir, snapshot, "memory-ranges")
|
||||
dst := filepath.Join(dstDir, "memory-ranges")
|
||||
if got, err := os.ReadFile(dst); err != nil || !bytes.Equal(got, want) {
|
||||
t.Fatalf("dst content = %q (err %v), want %q", got, err, want)
|
||||
}
|
||||
if inode(t, src) != inode(t, dst) {
|
||||
t.Error("staged file is a copy; expected a link to the cached snapshot")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("falls back to copying across filesystems", func(t *testing.T) {
|
||||
srcDir, dstDir := newDirs(t)
|
||||
// EXDEV stands in for the mount boundary a unit test cannot produce.
|
||||
orig := linkFile
|
||||
linkFile = func(string, string) error { return unix.EXDEV }
|
||||
t.Cleanup(func() { linkFile = orig })
|
||||
|
||||
s := &AteomHerder{}
|
||||
if err := s.copyLocalCheckpoint(context.Background(), snapshot, srcDir, dstDir, []string{"memory-ranges"}); err != nil {
|
||||
t.Fatalf("copyLocalCheckpoint: %v", err)
|
||||
}
|
||||
dst := filepath.Join(dstDir, "memory-ranges")
|
||||
if got, err := os.ReadFile(dst); err != nil || !bytes.Equal(got, want) {
|
||||
t.Fatalf("dst content = %q (err %v), want %q", got, err, want)
|
||||
}
|
||||
if inode(t, filepath.Join(srcDir, snapshot, "memory-ranges")) == inode(t, dst) {
|
||||
t.Error("expected a copy on the fallback path, got a link")
|
||||
}
|
||||
})
|
||||
|
||||
// Only EXDEV may fall back. Copying on any other link failure would silently
|
||||
// undo this optimization, and would hand sparsefile.CopyFile a dst that may already be a
|
||||
// link to src, where its O_TRUNC empties both before the copy reads a byte.
|
||||
t.Run("other link failures are fatal", func(t *testing.T) {
|
||||
srcDir, dstDir := newDirs(t)
|
||||
dst := filepath.Join(dstDir, "memory-ranges")
|
||||
if err := os.Link(filepath.Join(srcDir, snapshot, "memory-ranges"), dst); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
s := &AteomHerder{}
|
||||
// dst already exists, so os.Link fails with EEXIST.
|
||||
if err := s.copyLocalCheckpoint(context.Background(), snapshot, srcDir, dstDir, []string{"memory-ranges"}); err == nil {
|
||||
t.Fatal("copyLocalCheckpoint accepted a non-EXDEV link failure, want an error")
|
||||
}
|
||||
if got, err := os.ReadFile(dst); err != nil || !bytes.Equal(got, want) {
|
||||
t.Fatalf("dst content = %q (err %v), want the untouched %q", got, err, want)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func TestSnapshotManifestActorMetadata(t *testing.T) {
|
||||
rec := sandboxAssetsRecord{
|
||||
Atespace: "team-a",
|
||||
|
||||
@@ -255,23 +255,10 @@ func writeFileAtomicRoot(root *os.Root, relPath string, data []byte, perm os.Fil
|
||||
f.Close()
|
||||
return err
|
||||
}
|
||||
if err := f.Sync(); err != nil {
|
||||
f.Close()
|
||||
return err
|
||||
}
|
||||
if err := f.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := root.Rename(tmp, relPath); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
d, err := root.Open(dir)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer d.Close()
|
||||
return d.Sync()
|
||||
return root.Rename(tmp, relPath)
|
||||
}
|
||||
|
||||
// eventHandler enqueues the bundle names an event touches.
|
||||
|
||||
@@ -23,6 +23,7 @@ import (
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"syscall"
|
||||
|
||||
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/reaper"
|
||||
"golang.org/x/sys/unix"
|
||||
@@ -112,7 +113,8 @@ func MergeSparseOverlay(ctx context.Context, baseFile, deltaFile, outFile string
|
||||
// baseFile and deltaFile are siblings under the actor dir (restore-state/ and
|
||||
// checkpoint-state/), so the renames are same-filesystem (metadata-only). If they
|
||||
// straddle a mount boundary (EXDEV) it falls back to the copying MergeSparseOverlay
|
||||
// (baseFile is untouched until the first rename succeeds).
|
||||
// (baseFile is untouched until the first rename succeeds), as it does when baseFile
|
||||
// carries a second link and so cannot be overlaid in place.
|
||||
func MergeDeltaIntoBase(ctx context.Context, baseFile, deltaFile string) error {
|
||||
bi, err := os.Stat(baseFile)
|
||||
if err != nil {
|
||||
@@ -128,6 +130,18 @@ func MergeDeltaIntoBase(ctx context.Context, baseFile, deltaFile string) error {
|
||||
return fmt.Errorf("MergeDeltaIntoBase: size mismatch base=%d delta=%d", bi.Size(), di.Size())
|
||||
}
|
||||
|
||||
// The fast path below MUTATES baseFile's inode in place, which is only safe while
|
||||
// baseFile is the sole name for it. atelet stages restore-state by linking from
|
||||
// this actor's local pause snapshot rather than copying it, so a second link means
|
||||
// the overlay would rewrite that cached snapshot too, corrupting the actor's only
|
||||
// local restore point. Copy in that case, which is what MergeSparseOverlay does.
|
||||
// atelet holds that snapshot for the whole checkpoint, so this is the normal path
|
||||
// for a locally staged restore; the in-place path below is for one staged from
|
||||
// object storage.
|
||||
if st, ok := bi.Sys().(*syscall.Stat_t); ok && st.Nlink > 1 {
|
||||
return MergeSparseOverlay(ctx, baseFile, deltaFile, deltaFile)
|
||||
}
|
||||
|
||||
// Move baseFile (with its already-on-disk working set) next to deltaFile. If this
|
||||
// fails with EXDEV the two are on different filesystems and baseFile is still
|
||||
// intact, so fall back to the copying merge.
|
||||
|
||||
@@ -183,6 +183,54 @@ func TestMergeSparseOverlayNewOutFile(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestMergeDeltaIntoBaseHardlinkedBase asserts the fast path is refused when base
|
||||
// has a second link. The fast path overlays base's inode in place, so a shared
|
||||
// inode (atelet stages restore-state by linking from the actor's local pause
|
||||
// snapshot) would have that snapshot silently rewritten behind the merge.
|
||||
func TestMergeDeltaIntoBaseHardlinkedBase(t *testing.T) {
|
||||
const size = 8 << 20 // 8 MiB logical
|
||||
baseRegions := []region{{off: 0, data: fill(1, 4096)}}
|
||||
deltaRegions := []region{{off: 4 << 20, data: fill(42, 12345)}}
|
||||
want := make([]byte, size)
|
||||
copy(want[baseRegions[0].off:], baseRegions[0].data)
|
||||
copy(want[deltaRegions[0].off:], deltaRegions[0].data)
|
||||
// What the shared inode must still hold afterwards: base, with no delta in it.
|
||||
baseOnly := make([]byte, size)
|
||||
copy(baseOnly[baseRegions[0].off:], baseRegions[0].data)
|
||||
|
||||
dir := t.TempDir()
|
||||
base := filepath.Join(dir, "base")
|
||||
delta := filepath.Join(dir, "delta")
|
||||
cached := filepath.Join(dir, "cached-snapshot")
|
||||
writeSparse(t, base, size, baseRegions)
|
||||
writeSparse(t, delta, size, deltaRegions)
|
||||
if err := os.Link(base, cached); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if err := MergeDeltaIntoBase(context.Background(), base, delta); err != nil {
|
||||
t.Fatalf("MergeDeltaIntoBase: %v", err)
|
||||
}
|
||||
got, err := os.ReadFile(delta)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !bytes.Equal(got, want) {
|
||||
t.Fatalf("merged result != expected (len got=%d want=%d)", len(got), len(want))
|
||||
}
|
||||
shared, err := os.ReadFile(cached)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !bytes.Equal(shared, baseOnly) {
|
||||
t.Error("the hardlinked base was mutated; the cached snapshot behind it is now corrupt")
|
||||
}
|
||||
// The copying merge leaves base alone, unlike the fast path which consumes it.
|
||||
if _, err := os.Stat(base); err != nil {
|
||||
t.Errorf("base should survive a copying merge: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestMergeDeltaIntoBaseSizeMismatch verifies a base/delta size mismatch is
|
||||
// refused (misaligned overlay would corrupt the image) rather than silently
|
||||
// producing garbage.
|
||||
|
||||
@@ -501,8 +501,23 @@ func rewriteSnapshotSocketPaths(snapshotDir, id string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if err := os.WriteFile(cfgPath, out, 0o600); err != nil {
|
||||
// Temp file + rename, not os.WriteFile: atelet stages a local checkpoint by
|
||||
// hard-linking it into this dir, so config.json can share an inode with the
|
||||
// actor's cached pause snapshot. O_TRUNC would write straight through that link
|
||||
// and rewrite the snapshot, and a crash mid-write would leave the actor's only
|
||||
// local restore point holding a truncated config. Renaming replaces the name
|
||||
// here and leaves the linked original whole.
|
||||
tmp, err := os.CreateTemp(snapshotDir, ".config.json.tmp-*")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
defer os.Remove(tmp.Name()) // no-op once the rename succeeds
|
||||
if _, err := tmp.Write(out); err != nil {
|
||||
tmp.Close()
|
||||
return err
|
||||
}
|
||||
if err := tmp.Close(); err != nil {
|
||||
return err
|
||||
}
|
||||
return os.Rename(tmp.Name(), cfgPath)
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
@@ -131,4 +132,46 @@ func TestRewriteSnapshotSocketPaths(t *testing.T) {
|
||||
t.Errorf("serial file = %q, want %q", cfg.Serial.File, want)
|
||||
}
|
||||
})
|
||||
|
||||
// atelet stages a local checkpoint into the restore dir by hard-linking it, so
|
||||
// the config.json rewritten here can share an inode with the actor's cached
|
||||
// pause snapshot. Writing in place would rewrite that snapshot's config too.
|
||||
t.Run("a hardlinked config is not written through", func(t *testing.T) {
|
||||
dir := writeSnapshotConfig(t, []map[string]any{
|
||||
{"tag": kata.FsTag, "socket": "/run/vc/vm/golden/virtiofsd.sock"},
|
||||
})
|
||||
cfgPath := filepath.Join(dir, "config.json")
|
||||
before, err := os.ReadFile(cfgPath)
|
||||
if err != nil {
|
||||
t.Fatalf("reading config.json: %v", err)
|
||||
}
|
||||
cached := filepath.Join(t.TempDir(), "config.json")
|
||||
if err := os.Link(cfgPath, cached); err != nil {
|
||||
t.Fatalf("linking config.json: %v", err)
|
||||
}
|
||||
|
||||
if err := rewriteSnapshotSocketPaths(dir, id); err != nil {
|
||||
t.Fatalf("rewriteSnapshotSocketPaths: %v", err)
|
||||
}
|
||||
if got, want := readFsSockets(t, dir)[kata.FsTag], kata.VirtiofsdSocketPath(id); got != want {
|
||||
t.Errorf("%s socket = %q, want %q", kata.FsTag, got, want)
|
||||
}
|
||||
got, err := os.ReadFile(cached)
|
||||
if err != nil {
|
||||
t.Fatalf("reading the linked config.json: %v", err)
|
||||
}
|
||||
if !bytes.Equal(got, before) {
|
||||
t.Errorf("the linked config.json was rewritten: got %q, want the original %q", got, before)
|
||||
}
|
||||
// Nothing may be left behind for atelet to ship as part of the snapshot.
|
||||
entries, err := os.ReadDir(dir)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, e := range entries {
|
||||
if e.Name() != "config.json" {
|
||||
t.Errorf("stray file left in the restore dir: %q", e.Name())
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user