mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
benchmark(atelet, ateom-microvm): log per-actor checkpoint and restor… (#1941)
## Benchamark(atelet, ateom-microvm): log per-actor checkpoint and
restore phase breakdowns
### Why
`SuspendActor` and `ResumeActor` latency is only attributable down to
the `ate.actor.{checkpoint,restore}.duration` histogram phases, and the
biggest buckets — `ateom_checkpoint`, `ateom_restore` — are opaque. When
a large memory benchmark suspend takes 6 s we want to able to tell where
the time is been spent during the suspend.
### What
One joinable, developer-facing log record per operation per layer, with
the full actor identity (allowed in logs, barred from metric labels),
the snapshot scope, and one float-seconds field per phase that ran.
**atelet**
- `Checkpoint` now writes a `Checkpoint timing breakdown` record, the
same way `Restore` has written `Restore timing breakdown` since #1364.
Phases: `sandbox_assets`, `ateom_checkpoint`, `persist`, `total`. One
slice feeds both the histogram and the record, so they cannot disagree.
- Both records carry `error.type` when the operation failed (the gRPC
code; context errors map to `DeadlineExceeded` / `Canceled`). The record
is written on the way out of a failure too, so its completed phases are
kept, and the marker lets a reader exclude a timed-out restore from a
latency distribution. This restores what the record lost when the
`ateerrors` taxonomy was deleted (#1817), using the `error.type`
convention the ateapi instruments already follow.
**ateom-microvm**
- New `phaselog.go`. `CheckpointWorkload` and `RestoreWorkload` emit
records with the same two messages, under
`ateom.actor.checkpoint.duration.<phase>` and
`ateom.actor.restore.duration.<phase>`, decomposing atelet's
`ateom_checkpoint` / `ateom_restore` buckets:
- checkpoint: `pause`, `snapshot`, `durable_dir`, `rootfs_upper`,
`teardown`, `total` — the three captures run concurrently on the paused
guest, so the paused window costs their max, not their sum.
- restore: `prep`, `bundles`, `upper_join`, `lowers`, `tap`,
`vmm_launch`, `vm_restore`, `resume`, `wakeup_probe`, `total` —
sequential; they partition the total. A Data-scope cold boot records
`total` only.
- These timings already existed as ad-hoc `slog.Duration` fields on the
`Actor checkpointed` / `Actor restore phases` lines; those lines are
kept. The record adds stable keys, identity, and seconds (the
histograms' unit).
- The phase names are deliberately private to the binary rather than
added to `internal/ateattr`, so they cannot be mistaken for
`ate.snapshot.phase` metric values. They are micro-VM specific;
ateom-gvisor is unchanged.
**docs/observability.md** is updated: the Restore record is no longer
the only per-actor latency record, and the ateom records are described.
No new instruments, no registry changes, no behavior change.
### Example
```json
{"msg":"Checkpoint timing breakdown","ate.actor.uid":"8f2a…","ate.template.name":"glutton",
"ate.snapshot.scope":"full",
"ateom.actor.checkpoint.duration.pause":0.003,
"ateom.actor.checkpoint.duration.snapshot":0.846,
"ateom.actor.checkpoint.duration.rootfs_upper":0.022,
"ateom.actor.checkpoint.duration.teardown":0.232,
"ateom.actor.checkpoint.duration.total":1.081}
```
Joined with atelet's record for the same actor, a run of the glutton
workload (1 GiB resident, microVM) attributes a 6.0 s p50 suspend as 75%
`persist`, 19% `ateom_checkpoint` (of which the CH `snapshot` is 0.85 s
and `teardown` 0.23 s), and a 5.3 s p50 resume as 79% `download`, 18%
`ateom_restore` (of which `vm_restore` is 0.73 s). The consumer that
produces those tables from pod logs is a separate
`benchmarking/analysis` PR.
### Testing
- `go test ./cmd/atelet/...` and `./cmd/ateom-microvm/...` pass; new
unit tests cover the record shape (seconds, identity keys, zero phases
absent, no duplicate keys), the scope mapping, and `error.type` for
gRPC, context and plain errors.
- `GOOS=linux go vet` clean for both binaries; boilerplate and gofmt
clean.
This commit is contained in:
+19
-6
@@ -598,12 +598,25 @@ func (s *AteomHerder) Checkpoint(ctx context.Context, req *ateletpb.CheckpointRe
|
||||
kind: checkpointSnapshotKind(req),
|
||||
scope: ateattr.SnapshotScopeValue(req.GetScope()),
|
||||
}
|
||||
attribution := resources.ActorAttribution{
|
||||
Ref: actorRef,
|
||||
UID: actorUID,
|
||||
TemplateAtespace: req.GetActorTemplateAtespace(),
|
||||
TemplateName: req.GetActorTemplateName(),
|
||||
}
|
||||
defer func() {
|
||||
s.instruments.recordCheckpoint(ctx, op,
|
||||
phase{ateattr.SnapshotPhaseSandboxAssets, dAssets},
|
||||
phase{ateattr.SnapshotPhaseAteomCheckpoint, dAteom},
|
||||
phase{ateattr.SnapshotPhasePersist, dPersist},
|
||||
phase{ateattr.SnapshotPhaseTotal, time.Since(tStart)})
|
||||
// Use the same phase values for metrics and logs so their durations stay
|
||||
// consistent. The log also includes actor identity, which is intentionally
|
||||
// excluded from metric labels because of cardinality.
|
||||
phases := []phase{
|
||||
{ateattr.SnapshotPhaseSandboxAssets, dAssets},
|
||||
{ateattr.SnapshotPhaseAteomCheckpoint, dAteom},
|
||||
{ateattr.SnapshotPhasePersist, dPersist},
|
||||
{ateattr.SnapshotPhaseTotal, time.Since(tStart)},
|
||||
}
|
||||
s.instruments.recordCheckpoint(ctx, op, phases...)
|
||||
slog.LogAttrs(ctx, slog.LevelInfo, "Checkpoint timing breakdown",
|
||||
snapshotLogAttrs(attribution, op, checkpointDurationMetric, err, phases)...)
|
||||
}()
|
||||
|
||||
// Checkpoint requests no longer carry the sandbox config; recover the
|
||||
@@ -1060,7 +1073,7 @@ func (s *AteomHerder) Restore(ctx context.Context, req *ateletpb.RestoreRequest)
|
||||
}
|
||||
s.instruments.recordRestore(ctx, op, phases...)
|
||||
slog.LogAttrs(ctx, slog.LevelInfo, "Restore timing breakdown",
|
||||
snapshotLogAttrs(attribution, op, restoreDurationMetric, phases)...)
|
||||
snapshotLogAttrs(attribution, op, restoreDurationMetric, err, phases)...)
|
||||
}()
|
||||
|
||||
// Not crashing the actor, because terminal errors here indicate problems with atelet,
|
||||
|
||||
@@ -19,6 +19,8 @@ import (
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
// snapshotLogAttrs renders what recordPhases measures as a per-actor record. The
|
||||
@@ -30,7 +32,13 @@ import (
|
||||
// the values are seconds and not the nanoseconds slog.Duration writes: that
|
||||
// instrument declares unit s. Identity stays out of snapshotOp so no edit here
|
||||
// can route it into a datapoint.
|
||||
func snapshotLogAttrs(a resources.ActorAttribution, op snapshotOp, durationKey string, phases []phase) []slog.Attr {
|
||||
//
|
||||
// err is the operation's outcome. The record is written on the way out of a
|
||||
// failed operation too, so that its completed phases are not lost; error.type
|
||||
// (the gRPC code, as on the ateapi instruments) marks it so a reader does not
|
||||
// average a timed-out download in with the successful ones. Absence means
|
||||
// success, as on the instruments.
|
||||
func snapshotLogAttrs(a resources.ActorAttribution, op snapshotOp, durationKey string, err error, phases []phase) []slog.Attr {
|
||||
attrs := ateattr.ActorLogAttrs(a)
|
||||
|
||||
// Borrow the metric's dimensions rather than rebuild them: op.attrs already
|
||||
@@ -48,6 +56,17 @@ func snapshotLogAttrs(a resources.ActorAttribution, op snapshotOp, durationKey s
|
||||
attrs = append(attrs, slog.String(string(kv.Key), kv.Value.String()))
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
// Only an error that came back over gRPC carries a status; atelet's
|
||||
// own failures are plain errors, and the ones worth telling apart are
|
||||
// the context ones (a timed-out download, a cancelled restore).
|
||||
code := status.Code(err)
|
||||
if code == codes.Unknown {
|
||||
code = status.FromContextError(err).Code()
|
||||
}
|
||||
attrs = append(attrs, slog.String(string(ateattr.ErrorTypeKey), code.String()))
|
||||
}
|
||||
|
||||
// There is no ate.snapshot.phase key: on a datapoint it names the one step
|
||||
// timed, and this record carries them all.
|
||||
for _, p := range phases {
|
||||
|
||||
@@ -18,12 +18,16 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -105,7 +109,7 @@ func TestSnapshotLogAttrsSeconds(t *testing.T) {
|
||||
scope: ateattr.SnapshotScopeFull,
|
||||
sandboxClass: "gvisor",
|
||||
}
|
||||
attrs := snapshotLogAttrs(testAttribution(), op, restoreDurationMetric,
|
||||
attrs := snapshotLogAttrs(testAttribution(), op, restoreDurationMetric, nil,
|
||||
[]phase{
|
||||
{ateattr.SnapshotPhaseDownload, 2500 * time.Millisecond},
|
||||
{ateattr.SnapshotPhaseTotal, 3 * time.Second},
|
||||
@@ -128,6 +132,7 @@ func TestSnapshotLogAttrs(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
op snapshotOp
|
||||
err error
|
||||
phases []phase
|
||||
wantStrings map[string]string
|
||||
wantNumbers map[string]float64
|
||||
@@ -163,7 +168,52 @@ func TestSnapshotLogAttrs(t *testing.T) {
|
||||
},
|
||||
// The phase key names the one step a datapoint timed; this record has
|
||||
// them all, so borrowing it here would give one key two meanings.
|
||||
wantAbsent: []string{"ate.snapshot.phase", "actor", "total", "download"},
|
||||
// Absent error.type is success, as on the instruments.
|
||||
wantAbsent: []string{"ate.snapshot.phase", "actor", "total", "download", "error.type"},
|
||||
},
|
||||
{
|
||||
name: "a failed restore keeps the phases it completed and is marked with the gRPC code",
|
||||
op: fullOp,
|
||||
err: status.Error(codes.DeadlineExceeded, "download timed out"),
|
||||
phases: []phase{
|
||||
{ateattr.SnapshotPhaseDownload, 30 * time.Second},
|
||||
{ateattr.SnapshotPhaseTotal, 30 * time.Second},
|
||||
},
|
||||
wantStrings: map[string]string{
|
||||
"error.type": "DeadlineExceeded",
|
||||
},
|
||||
wantNumbers: map[string]float64{
|
||||
"ate.actor.restore.duration.download": 30,
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "a plain error is a bounded Unknown, not its message",
|
||||
op: fullOp,
|
||||
err: errors.New("disk on fire at /var/lib/ate"),
|
||||
phases: []phase{{ateattr.SnapshotPhaseTotal, time.Second}},
|
||||
wantStrings: map[string]string{
|
||||
"error.type": "Unknown",
|
||||
},
|
||||
},
|
||||
{
|
||||
// atelet's own downloads and uploads fail with wrapped context
|
||||
// errors, not status errors; a timeout must not read as Unknown.
|
||||
name: "a wrapped context deadline reports DeadlineExceeded",
|
||||
op: fullOp,
|
||||
err: fmt.Errorf("while downloading memory-ranges: %w", context.DeadlineExceeded),
|
||||
phases: []phase{{ateattr.SnapshotPhaseTotal, 30 * time.Second}},
|
||||
wantStrings: map[string]string{
|
||||
"error.type": "DeadlineExceeded",
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "a wrapped context cancellation reports Canceled",
|
||||
op: fullOp,
|
||||
err: fmt.Errorf("while restoring: %w", context.Canceled),
|
||||
phases: []phase{{ateattr.SnapshotPhaseTotal, time.Second}},
|
||||
wantStrings: map[string]string{
|
||||
"error.type": "Canceled",
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "a phase that never ran is absent, not zero",
|
||||
@@ -209,7 +259,7 @@ func TestSnapshotLogAttrs(t *testing.T) {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
attrs := snapshotLogAttrs(testAttribution(), tt.op, restoreDurationMetric, tt.phases)
|
||||
attrs := snapshotLogAttrs(testAttribution(), tt.op, restoreDurationMetric, tt.err, tt.phases)
|
||||
|
||||
// json.Unmarshal keeps the last of a repeated key, so a duplicate is
|
||||
// invisible in the rendered record and has to be caught on the slice.
|
||||
|
||||
@@ -60,7 +60,7 @@ import (
|
||||
//
|
||||
// Allow checkpointing even if the pod is shutting down. This will allow actors
|
||||
// (or the harness) to suspend on shutdown.
|
||||
func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.CheckpointWorkloadRequest) (*ateompb.CheckpointWorkloadResponse, error) {
|
||||
func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.CheckpointWorkloadRequest) (_ *ateompb.CheckpointWorkloadResponse, err error) {
|
||||
if err := validateActorDirs(req.GetActorDirs()); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -73,7 +73,28 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
|
||||
defer cancel()
|
||||
defer s.inFlight.Add(req.GetActorUid(), rpcCheckpointWorkload, nil)()
|
||||
|
||||
// Per-phase timing, recorded on the way out so a failed checkpoint still
|
||||
// reports the phases it completed, and the failing step its elapsed time.
|
||||
// Phases left at zero never ran. The snapshot, durable_dir and rootfs_upper
|
||||
// captures run concurrently on the paused guest, so those three are
|
||||
// independent observations rather than a partition of the total.
|
||||
tStart := time.Now()
|
||||
var dPrep, dPause, dSnapshot, dDurable, dUpper, dTeardown time.Duration
|
||||
attribution := ateomstats.ActorAttributionFromRequest(req)
|
||||
scope := req.GetScope()
|
||||
defer func() {
|
||||
logSnapshotPhases(ctx, "Checkpoint timing breakdown", attribution, scope,
|
||||
checkpointDurationKey, err, []phase{
|
||||
{phasePrep, dPrep},
|
||||
{phasePause, dPause},
|
||||
{phaseSnapshot, dSnapshot},
|
||||
{phaseDurableDir, dDurable},
|
||||
{phaseRootfsUpper, dUpper},
|
||||
{phaseTeardown, dTeardown},
|
||||
{phaseTotal, time.Since(tStart)},
|
||||
})
|
||||
}()
|
||||
|
||||
if err := s.deactivateActorNetworking(ctx, attribution); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -93,7 +114,6 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
|
||||
// here as plain DATA) and lands in the default rejection.
|
||||
durable := hasDurableVolumes(req.GetSpec().GetContainers())
|
||||
csi := hasCsiVolumes(req.GetSpec().GetContainers())
|
||||
scope := req.GetScope()
|
||||
switch scope {
|
||||
case ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL:
|
||||
case ateompb.SnapshotScope_SNAPSHOT_SCOPE_DATA:
|
||||
@@ -119,10 +139,12 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
|
||||
}
|
||||
|
||||
tPause := time.Now()
|
||||
if err := client.Pause(ctx); err != nil {
|
||||
return nil, fmt.Errorf("while pausing guest: %w", err)
|
||||
dPrep = tPause.Sub(tStart)
|
||||
pauseErr := client.Pause(ctx)
|
||||
dPause = time.Since(tPause)
|
||||
if pauseErr != nil {
|
||||
return nil, fmt.Errorf("while pausing guest: %w", pauseErr)
|
||||
}
|
||||
dPause := time.Since(tPause)
|
||||
|
||||
checkpointDir := actorDirs.GetCheckpointDir()
|
||||
// Start from a clean dir so CH's snapshot files are the only contents.
|
||||
@@ -149,33 +171,32 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
|
||||
// - Rootfs upper tar (Full only): host-backed like the durable volumes —
|
||||
// the memory snapshot does not carry rootfs writes. Under Data the
|
||||
// workload cold-starts on restore, discarding rootfs state.
|
||||
var dSnapshot, dDurable, dUpper time.Duration
|
||||
g, gctx := errgroup.WithContext(ctx)
|
||||
if scope == ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL {
|
||||
g.Go(func() error {
|
||||
var err error
|
||||
dSnapshot, err = s.snapshotVMState(gctx, client, ra, actorUID, checkpointDir)
|
||||
t := time.Now()
|
||||
d, err := s.snapshotVMState(gctx, client, ra, actorUID, checkpointDir)
|
||||
if err != nil {
|
||||
d = time.Since(t)
|
||||
}
|
||||
dSnapshot = d
|
||||
return err
|
||||
})
|
||||
}
|
||||
if durable {
|
||||
g.Go(func() error {
|
||||
t := time.Now()
|
||||
if err := tarDurableVolumes(gctx, actorDirs.GetDurableDirVolumeMountsDir(), checkpointDir); err != nil {
|
||||
return err
|
||||
}
|
||||
err := tarDurableVolumes(gctx, actorDirs.GetDurableDirVolumeMountsDir(), checkpointDir)
|
||||
dDurable = time.Since(t)
|
||||
return nil
|
||||
return err
|
||||
})
|
||||
}
|
||||
if scope == ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL {
|
||||
g.Go(func() error {
|
||||
t := time.Now()
|
||||
if err := tarRootfsUpper(gctx, rootfsUpperDir(actorDirs), checkpointDir); err != nil {
|
||||
return err
|
||||
}
|
||||
err := tarRootfsUpper(gctx, rootfsUpperDir(actorDirs), checkpointDir)
|
||||
dUpper = time.Since(t)
|
||||
return nil
|
||||
return err
|
||||
})
|
||||
}
|
||||
if err := g.Wait(); err != nil {
|
||||
@@ -199,7 +220,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
|
||||
slog.String("actorUID", actorUID),
|
||||
slog.Any("err", err))
|
||||
}
|
||||
dTeardown := time.Since(tTeardown)
|
||||
dTeardown = time.Since(tTeardown)
|
||||
|
||||
s.actorLogger.EmitLifecycleLog(ctx, "Actor checkpointed", attribution)
|
||||
slog.InfoContext(ctx, "Actor checkpointed", slog.String("id", actorUID), slog.Any("snapshot_files", snapshotFiles),
|
||||
|
||||
@@ -0,0 +1,131 @@
|
||||
//go:build linux
|
||||
|
||||
// Copyright 2026 Google LLC
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"time"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/proto/ateompb"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
// The keys the per-phase durations are logged under. They are named like
|
||||
// instruments but deliberately are not ones: the phases are implementation
|
||||
// details of this binary, so they stay a developer-facing log record rather
|
||||
// than becoming metric API. They decompose the ateom_restore /
|
||||
// ateom_checkpoint phases of atelet's ate.actor.*.duration histograms, whose
|
||||
// names they extend, and are joined per actor by benchmarking tooling.
|
||||
const (
|
||||
restoreDurationKey = "ateom.actor.restore.duration"
|
||||
checkpointDurationKey = "ateom.actor.checkpoint.duration"
|
||||
)
|
||||
|
||||
// The phase names, suffixed onto the duration keys. Kept out of ateattr on
|
||||
// purpose: that package's SnapshotPhase* values are the ate.snapshot.phase
|
||||
// metric enum, and these are not values of it. "total" is shared so the two
|
||||
// layers' records agree on the denominator.
|
||||
//
|
||||
// Checkpoint: snapshot, durable_dir and rootfs_upper run concurrently on the
|
||||
// paused guest, so the paused window costs their max, not their sum; prep,
|
||||
// pause and teardown are sequential around them. Restore: every phase is
|
||||
// sequential and the phases partition the total.
|
||||
const (
|
||||
phasePause = "pause"
|
||||
phaseSnapshot = "snapshot"
|
||||
phaseDurableDir = "durable_dir"
|
||||
phaseRootfsUpper = "rootfs_upper"
|
||||
phaseTeardown = "teardown"
|
||||
|
||||
phasePrep = "prep"
|
||||
phaseBundles = "bundles"
|
||||
phaseUpperJoin = "upper_join"
|
||||
phaseLowers = "lowers"
|
||||
phaseTap = "tap"
|
||||
phaseVMMLaunch = "vmm_launch"
|
||||
phaseVMRestore = "vm_restore"
|
||||
phaseResume = "resume"
|
||||
phaseWakeupProbe = "wakeup_probe"
|
||||
|
||||
phaseTotal = ateattr.SnapshotPhaseTotal
|
||||
)
|
||||
|
||||
// phase is one timed step of a snapshot operation. A zero duration means the
|
||||
// phase never ran (a Data-scope checkpoint captures no guest) and is skipped
|
||||
// rather than logged as instant.
|
||||
type phase struct {
|
||||
name string
|
||||
d time.Duration
|
||||
}
|
||||
|
||||
// scopeLogValue maps the ateom wire enum onto the shared scope label values,
|
||||
// the same way ateattr.SnapshotScopeValue does for the atelet enum. An
|
||||
// unrecognized scope reports as unknown rather than stringified, so no wire
|
||||
// value can widen the value set readers key on.
|
||||
func scopeLogValue(scope ateompb.SnapshotScope) string {
|
||||
switch scope {
|
||||
case ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL:
|
||||
return ateattr.SnapshotScopeFull
|
||||
case ateompb.SnapshotScope_SNAPSHOT_SCOPE_DATA:
|
||||
return ateattr.SnapshotScopeData
|
||||
case ateompb.SnapshotScope_SNAPSHOT_SCOPE_DATA_ON_GOLDEN:
|
||||
return ateattr.SnapshotScopeDataOnGolden
|
||||
default:
|
||||
return ateattr.SnapshotScopeUnknown
|
||||
}
|
||||
}
|
||||
|
||||
// snapshotPhaseAttrs renders one joinable per-actor record for a checkpoint
|
||||
// or restore, mirroring atelet's snapshotLogAttrs: full actor identity, the
|
||||
// scope, and one float-seconds attr per non-zero phase under
|
||||
// durationKey.<phase>. One record carries every phase of the operation, so a
|
||||
// reader never joins two half-records that disagree about the same actor.
|
||||
//
|
||||
// err is the operation's outcome. A failed operation still records the phases
|
||||
// it completed, marked with error.type (the gRPC code, context errors as
|
||||
// DeadlineExceeded / Canceled) so a reader can leave it out of a latency
|
||||
// distribution. Absence means success.
|
||||
func snapshotPhaseAttrs(a resources.ActorAttribution, scope ateompb.SnapshotScope, durationKey string, err error, phases []phase) []slog.Attr {
|
||||
attrs := ateattr.ActorLogAttrs(a)
|
||||
attrs = append(attrs, slog.String(string(ateattr.SnapshotScopeKey), scopeLogValue(scope)))
|
||||
if err != nil {
|
||||
code := status.Code(err)
|
||||
if code == codes.Unknown {
|
||||
code = status.FromContextError(err).Code()
|
||||
}
|
||||
attrs = append(attrs, slog.String(string(ateattr.ErrorTypeKey), code.String()))
|
||||
}
|
||||
for _, p := range phases {
|
||||
if p.d == 0 {
|
||||
continue
|
||||
}
|
||||
// Seconds and not slog.Duration's nanoseconds: the values sit under
|
||||
// the atelet histograms' unit (s), so readers compare them without a
|
||||
// per-key unit table.
|
||||
attrs = append(attrs, slog.Float64(durationKey+"."+p.name, p.d.Seconds()))
|
||||
}
|
||||
return attrs
|
||||
}
|
||||
|
||||
// logSnapshotPhases emits the snapshotPhaseAttrs record under msg.
|
||||
func logSnapshotPhases(ctx context.Context, msg string, a resources.ActorAttribution, scope ateompb.SnapshotScope, durationKey string, err error, phases []phase) {
|
||||
slog.LogAttrs(ctx, slog.LevelInfo, msg, snapshotPhaseAttrs(a, scope, durationKey, err, phases)...)
|
||||
}
|
||||
@@ -0,0 +1,182 @@
|
||||
//go:build linux
|
||||
|
||||
// Copyright 2026 Google LLC
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/proto/ateompb"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
)
|
||||
|
||||
func phaseLogAttribution() resources.ActorAttribution {
|
||||
return resources.ActorAttribution{
|
||||
Ref: resources.ActorRef{Atespace: "team-a", Name: "support-agent-42"},
|
||||
UID: "uid-abc",
|
||||
TemplateAtespace: "templates",
|
||||
TemplateName: "support-agent",
|
||||
}
|
||||
}
|
||||
|
||||
// renderPhaseRecord logs attrs through a local JSON handler, so the test sees
|
||||
// the record a collector would parse rather than the slice that built it.
|
||||
func renderPhaseRecord(t *testing.T, attrs []slog.Attr) map[string]any {
|
||||
t.Helper()
|
||||
var buf bytes.Buffer
|
||||
slog.New(slog.NewJSONHandler(&buf, nil)).LogAttrs(context.Background(), slog.LevelInfo, "Checkpoint timing breakdown", attrs...)
|
||||
var rec map[string]any
|
||||
if err := json.Unmarshal(buf.Bytes(), &rec); err != nil {
|
||||
t.Fatalf("unmarshal record %q: %v", buf.String(), err)
|
||||
}
|
||||
return rec
|
||||
}
|
||||
|
||||
func TestSnapshotPhaseAttrs(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
attrs := snapshotPhaseAttrs(phaseLogAttribution(), ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL,
|
||||
checkpointDurationKey, nil, []phase{
|
||||
{phasePrep, 40 * time.Millisecond},
|
||||
{phasePause, 3 * time.Millisecond},
|
||||
{phaseSnapshot, 850 * time.Millisecond},
|
||||
// A capture that did not run stays off the record.
|
||||
{phaseDurableDir, 0},
|
||||
{phaseTeardown, 230 * time.Millisecond},
|
||||
{phaseTotal, 1250 * time.Millisecond},
|
||||
})
|
||||
|
||||
// json.Unmarshal keeps the last of a repeated key, so a duplicate is
|
||||
// invisible in the rendered record and has to be caught on the slice.
|
||||
seen := make(map[string]bool, len(attrs))
|
||||
for _, attr := range attrs {
|
||||
if seen[attr.Key] {
|
||||
t.Errorf("key %s emitted twice", attr.Key)
|
||||
}
|
||||
seen[attr.Key] = true
|
||||
}
|
||||
|
||||
rec := renderPhaseRecord(t, attrs)
|
||||
for k, want := range map[string]string{
|
||||
"ate.atespace": "team-a",
|
||||
"ate.actor.name": "support-agent-42",
|
||||
"ate.actor.uid": "uid-abc",
|
||||
"ate.template.atespace": "templates",
|
||||
"ate.template.name": "support-agent",
|
||||
"ate.snapshot.scope": ateattr.SnapshotScopeFull,
|
||||
} {
|
||||
if got, ok := rec[k]; !ok {
|
||||
t.Errorf("missing %s", k)
|
||||
} else if got != want {
|
||||
t.Errorf("%s = %v, want %q", k, got, want)
|
||||
}
|
||||
}
|
||||
// Seconds, not slog.Duration's nanoseconds: the keys extend the atelet
|
||||
// histograms' names, and those declare unit s.
|
||||
for k, want := range map[string]float64{
|
||||
"ateom.actor.checkpoint.duration.prep": 0.04,
|
||||
"ateom.actor.checkpoint.duration.pause": 0.003,
|
||||
"ateom.actor.checkpoint.duration.snapshot": 0.85,
|
||||
"ateom.actor.checkpoint.duration.teardown": 0.23,
|
||||
"ateom.actor.checkpoint.duration.total": 1.25,
|
||||
} {
|
||||
got, ok := rec[k].(float64)
|
||||
if !ok {
|
||||
t.Errorf("%s = %v (%T), want a number", k, rec[k], rec[k])
|
||||
} else if got != want {
|
||||
t.Errorf("%s = %v, want %v", k, got, want)
|
||||
}
|
||||
}
|
||||
for _, k := range []string{
|
||||
"ateom.actor.checkpoint.duration.durable_dir",
|
||||
// The phase key names the one step a datapoint timed; this record
|
||||
// carries them all.
|
||||
"ate.snapshot.phase",
|
||||
// Absent error.type is success, as on the instruments.
|
||||
"error.type",
|
||||
} {
|
||||
if v, ok := rec[k]; ok {
|
||||
t.Errorf("%s is present with %v, want absent", k, v)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestSnapshotPhaseAttrsFailure: a checkpoint that died keeps the phases it
|
||||
// completed and is marked, so a reader can exclude it from percentiles.
|
||||
func TestSnapshotPhaseAttrsFailure(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
err error
|
||||
want string
|
||||
}{
|
||||
{"a wrapped context deadline reports DeadlineExceeded",
|
||||
fmt.Errorf("while snapshotting guest: %w", context.DeadlineExceeded), "DeadlineExceeded"},
|
||||
{"a wrapped context cancellation reports Canceled",
|
||||
fmt.Errorf("while pausing guest: %w", context.Canceled), "Canceled"},
|
||||
{"a plain error is a bounded Unknown, not its message",
|
||||
fmt.Errorf("while clearing checkpoint dir: disk on fire"), "Unknown"},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
rec := renderPhaseRecord(t, snapshotPhaseAttrs(phaseLogAttribution(),
|
||||
ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL, checkpointDurationKey, tt.err,
|
||||
[]phase{{phasePrep, 40 * time.Millisecond}, {phasePause, 3 * time.Millisecond}, {phaseTotal, 30 * time.Second}}))
|
||||
if got := rec["error.type"]; got != tt.want {
|
||||
t.Errorf("error.type = %v, want %q", got, tt.want)
|
||||
}
|
||||
// The phases that ran before the failure are still on the record;
|
||||
// the ones that never started are not.
|
||||
if _, ok := rec["ateom.actor.checkpoint.duration.pause"]; !ok {
|
||||
t.Error("missing ateom.actor.checkpoint.duration.pause")
|
||||
}
|
||||
if v, ok := rec["ateom.actor.checkpoint.duration.snapshot"]; ok {
|
||||
t.Errorf("snapshot phase present with %v, want absent", v)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestScopeLogValue pins the mapping onto the shared scope values, so the
|
||||
// ateom and atelet records of one operation agree on the scope they carry.
|
||||
func TestScopeLogValue(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
scope ateompb.SnapshotScope
|
||||
want string
|
||||
}{
|
||||
{ateompb.SnapshotScope_SNAPSHOT_SCOPE_FULL, ateattr.SnapshotScopeFull},
|
||||
{ateompb.SnapshotScope_SNAPSHOT_SCOPE_DATA, ateattr.SnapshotScopeData},
|
||||
{ateompb.SnapshotScope_SNAPSHOT_SCOPE_DATA_ON_GOLDEN, ateattr.SnapshotScopeDataOnGolden},
|
||||
{ateompb.SnapshotScope_SNAPSHOT_SCOPE_UNSPECIFIED, ateattr.SnapshotScopeUnknown},
|
||||
{ateompb.SnapshotScope(99), ateattr.SnapshotScopeUnknown},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
if got := scopeLogValue(tt.scope); got != tt.want {
|
||||
t.Errorf("scopeLogValue(%v) = %q, want %q", tt.scope, got, tt.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -181,7 +181,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
|
||||
// files, and the untar above re-materialized the ACTOR's durable-dir
|
||||
// data, so resuming the golden guest picks up the actor's data through
|
||||
// the durable virtio-fs share.
|
||||
if err := s.restoreFullScope(ctx, p, restoreDir, tStart); err != nil {
|
||||
if err := s.restoreFullScope(ctx, p, scope, restoreDir, tStart); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
case ateompb.SnapshotScope_SNAPSHOT_SCOPE_DATA:
|
||||
@@ -191,8 +191,13 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
|
||||
if err := s.coldBootActorRetrying(ctx, p); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
dTotal := time.Since(tStart)
|
||||
slog.InfoContext(ctx, "Actor restored (durable-dir volumes, cold boot)",
|
||||
slog.String("id", p.actorUID), slog.Duration("total", time.Since(tStart)))
|
||||
slog.String("id", p.actorUID), slog.Duration("total", dTotal))
|
||||
// A cold boot has none of the full-scope phases, so the total is the
|
||||
// only observation on its record.
|
||||
logSnapshotPhases(ctx, "Restore timing breakdown", attribution, scope,
|
||||
restoreDurationKey, nil, []phase{{phaseTotal, dTotal}})
|
||||
default:
|
||||
return nil, status.Errorf(codes.InvalidArgument, "unsupported snapshot scope: %v", scope)
|
||||
}
|
||||
@@ -213,7 +218,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
|
||||
// and resume. Guest RAM — the actor's in-memory state and the frozen network config —
|
||||
// comes back from the memory snapshot; the durable-dir volumes were restored by the
|
||||
// caller from their tar.
|
||||
func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams, restoreDir string, tStart time.Time) (retErr error) {
|
||||
func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams, scope ateompb.SnapshotScope, restoreDir string, tStart time.Time) (retErr error) {
|
||||
actorUID := p.actorUID
|
||||
|
||||
rr := s.resolveRuntime(p.assetPaths)
|
||||
@@ -433,6 +438,8 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
|
||||
// which hid that a first (cold) restore and a later (warm) one differ by more
|
||||
// than 5x on the same actor. upper/lowers is the host reassembling the rootfs;
|
||||
// vm_restore is cloud-hypervisor reading guest RAM back.
|
||||
dWakeupProbe := time.Since(tResume)
|
||||
dTotal := time.Since(tStart)
|
||||
slog.InfoContext(ctx, "Actor restore phases", slog.String("id", actorUID),
|
||||
slog.Duration("prep", tPrep.Sub(tStart)),
|
||||
slog.Duration("bundles", tBundles.Sub(tPrep)),
|
||||
@@ -443,8 +450,24 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
|
||||
slog.Duration("vmm_launch", tLaunch.Sub(tTap)),
|
||||
slog.Duration("vm_restore", tVMRestore.Sub(tLaunch)),
|
||||
slog.Duration("resume", tResume.Sub(tVMRestore)),
|
||||
slog.Duration("wakeup_probe", time.Since(tResume)),
|
||||
slog.Duration("total", time.Since(tStart)))
|
||||
slog.Duration("wakeup_probe", dWakeupProbe),
|
||||
slog.Duration("total", dTotal))
|
||||
// The joinable per-actor record the benchmarking tooling aggregates. The
|
||||
// durable delta is not carried: tDurable is pinned to tLowers today, so it
|
||||
// would always be the zero a record skips.
|
||||
logSnapshotPhases(ctx, "Restore timing breakdown", p.actorAttribution(), scope,
|
||||
restoreDurationKey, nil, []phase{
|
||||
{phasePrep, tPrep.Sub(tStart)},
|
||||
{phaseBundles, tBundles.Sub(tPrep)},
|
||||
{phaseUpperJoin, tUpper.Sub(tBundles)},
|
||||
{phaseLowers, tLowers.Sub(tUpper)},
|
||||
{phaseTap, tTap.Sub(tDurable)},
|
||||
{phaseVMMLaunch, tLaunch.Sub(tTap)},
|
||||
{phaseVMRestore, tVMRestore.Sub(tLaunch)},
|
||||
{phaseResume, tResume.Sub(tVMRestore)},
|
||||
{phaseWakeupProbe, dWakeupProbe},
|
||||
{phaseTotal, dTotal},
|
||||
})
|
||||
|
||||
// An eager restore has read the whole snapshot into guest memory, and nothing
|
||||
// merges against it afterwards, so the staged copy is dead weight from here on —
|
||||
|
||||
@@ -123,7 +123,7 @@ An actor's **own** lines carry trace context only if the actor emits these field
|
||||
|
||||
A component's own `slog` output can also be about a specific actor. Those records take the identity keys from [`internal/ateattr`](../internal/ateattr) too, flat at the top level rather than inside a label group: a component writes no envelope, so a collector lifts the keys straight onto the log record's attributes. `ateattr.ActorLogAttrs` and `ateattr.ActorLogLabels` return the same five keys for this reason, and a test holds them together. Filtering on `ate.actor.uid` therefore finds a component record and an actor's own output alike.
|
||||
|
||||
atelet's `Restore timing breakdown` is the first of these, and the only unsampled per-actor latency record Substrate produces. It is emitted once per restore, whether the restore succeeded or failed:
|
||||
atelet's `Restore timing breakdown` and `Checkpoint timing breakdown` are the unsampled per-actor latency records Substrate produces. Each is emitted once per operation, whether it succeeded or failed; a failed one also carries `error.type` (the gRPC code, `DeadlineExceeded` and `Canceled` for context errors), so a reader can leave it out of a latency distribution:
|
||||
|
||||
```json
|
||||
{"time":"…","level":"INFO","msg":"Restore timing breakdown",
|
||||
@@ -141,6 +141,8 @@ The duration keys are the [`ate.actor.restore.duration`](#the-metric-registry) i
|
||||
|
||||
This is the record to use for a per-actor wake-up distribution. The histogram cannot answer that question at all, because actor identity is barred from metric labels; traces can, but the data plane is head-sampled at 1%.
|
||||
|
||||
ateom-microvm writes records with the same two messages, from inside the `ateom_restore` and `ateom_checkpoint` phases, under `ateom.actor.restore.duration.<phase>` and `ateom.actor.checkpoint.duration.<phase>` (for example `vm_restore`, `wakeup_probe`, `prep`, `pause`, `snapshot`, `teardown`). Those keys are not instruments: the phases are implementation details of one runtime, so they stay a developer-facing record. The same identity keys make the two layers' records joinable per actor. The checkpoint record is written on failure too, with `error.type` and the elapsed time of the step that failed; the restore record is written on success only. The restore phases are sequential and partition the total; a checkpoint's `snapshot`, `durable_dir` and `rootfs_upper` run concurrently on the paused guest, so the paused window costs their maximum, not their sum.
|
||||
|
||||
ateapi's `Actor state changed` is written once per committed actor state transition. ateapi owns the state machine, so this is where an actor's state and the time it reached it come from:
|
||||
|
||||
```json
|
||||
|
||||
Reference in New Issue
Block a user