Skip to content
Merged
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
40 changes: 40 additions & 0 deletions packages/orchestrator/pkg/sandbox/sandbox.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,17 @@ var (
meter = otel.Meter("github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox")
envdInitCalls = utils.Must(telemetry.GetCounter(meter, telemetry.EnvdInitCalls))
waitForEnvdDurationHistogram = utils.Must(telemetry.GetHistogram(meter, telemetry.WaitForEnvdDurationHistogramName))

uffdStartupPagesHistogram = utils.Must(telemetry.GetHistogram(meter, telemetry.UffdStartupPagesHistogramName))
uffdStartupSourcePagesHistogram = utils.Must(telemetry.GetHistogram(meter, telemetry.UffdStartupSourcePagesHistogramName))
uffdStartupBytesHistogram = utils.Must(telemetry.GetHistogram(meter, telemetry.UffdStartupBytesHistogramName))
)

// Sandbox start types recorded on the orchestrator.sandbox.uffd.startup.*
// metrics via the start_type attribute.
const (
StartTypeCreate = "create" // cold boot (template build)
StartTypeResume = "resume" // resume from a snapshot (the common runtime path)
)

var SandboxHttpTransport = otelhttp.NewTransport(
Expand Down Expand Up @@ -256,6 +267,14 @@ type Sandbox struct {
exit *utils.ErrorOnce

stop utils.Lazy[error]

// startupStatsOnce guards the orchestrator.sandbox.uffd.startup.* recording
// so it fires only on the first WaitForEnvd — the actual sandbox start.
// ServeStats() is lifetime-cumulative on the UFFD handler, so a later
// WaitForEnvd on the same handler (e.g. the envd-binary swap + restart in a
// template build) would otherwise emit a sample inflated with post-startup
// faults rather than that init's working set.
startupStatsOnce sync.Once
}

func (s *Sandbox) LoggerMetadata() sbxlogger.SandboxMetadata {
Expand Down Expand Up @@ -911,6 +930,7 @@ func (f *Factory) ResumeSandbox(

err = sbx.WaitForEnvd(
ctx,
StartTypeResume,
f.config.EnvdTimeout,
)
if err != nil {
Expand Down Expand Up @@ -1457,6 +1477,7 @@ func (s *Sandbox) WaitForExit(ctx context.Context) error {

func (s *Sandbox) WaitForEnvd(
ctx context.Context,
startType string,
timeout time.Duration,
) (e error) {
start := time.Now()
Expand All @@ -1471,6 +1492,25 @@ func (s *Sandbox) WaitForEnvd(
attribute.Bool("success", e == nil),
))

// Record the demand-fault working set the guest needed to reach this
// point. Only on the first WaitForEnvd: it is the actual start, and
// ServeStats() is cumulative since resume, so at this instant it equals
// the startup counts. A later WaitForEnvd on the same handler (e.g. the
// envd-binary swap + restart during a template build) would otherwise
// re-report a cumulative total polluted with intervening faults.
// Recorded for both outcomes (success label) so slow/failed starts can
// be correlated with page volume.
s.startupStatsOnce.Do(func() {
stats := s.memory.ServeStats()
startupAttrs := metric.WithAttributes(
attribute.String("start_type", startType),
attribute.Bool("success", e == nil),
)
uffdStartupPagesHistogram.Record(ctx, stats.Pages, startupAttrs)
uffdStartupSourcePagesHistogram.Record(ctx, stats.SourcePages, startupAttrs)
uffdStartupBytesHistogram.Record(ctx, stats.Bytes, startupAttrs)
})

if e != nil {
return
}
Expand Down
5 changes: 5 additions & 0 deletions packages/orchestrator/pkg/sandbox/uffd/memory_backend.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (

"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block"
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/fc"
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/uffd/userfaultfd"
"github.com/e2b-dev/infra/packages/shared/pkg/storage/header"
"github.com/e2b-dev/infra/packages/shared/pkg/utils"
)
Expand All @@ -22,4 +23,8 @@ type MemoryBackend interface {
Ready() chan struct{}
Exit() *utils.ErrorOnce
Memfd(ctx context.Context) *block.Memfd
// ServeStats returns a cumulative snapshot of demand faults served so far.
// Sampled at the envd-init boundary it yields the pages/bytes a guest
// needed to start.
ServeStats() userfaultfd.ServeSnapshot
}
7 changes: 7 additions & 0 deletions packages/orchestrator/pkg/sandbox/uffd/noop.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import (

"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/block"
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/fc"
"github.com/e2b-dev/infra/packages/orchestrator/pkg/sandbox/uffd/userfaultfd"
"github.com/e2b-dev/infra/packages/shared/pkg/storage/header"
"github.com/e2b-dev/infra/packages/shared/pkg/utils"
)
Expand Down Expand Up @@ -87,3 +88,9 @@ func (m *NoopMemory) Exit() *utils.ErrorOnce {
func (m *NoopMemory) Memfd(context.Context) *block.Memfd {
return nil
}

// ServeStats returns a zero snapshot: NoopMemory has no UFFD serve loop, so no
// pages are demand-faulted through it.
func (m *NoopMemory) ServeStats() userfaultfd.ServeSnapshot {
return userfaultfd.ServeSnapshot{}
}
12 changes: 12 additions & 0 deletions packages/orchestrator/pkg/sandbox/uffd/uffd.go
Original file line number Diff line number Diff line change
Expand Up @@ -281,3 +281,15 @@ func (u *Uffd) PrefetchData(ctx context.Context) (block.PrefetchData, error) {
func (u *Uffd) Memfd(_ context.Context) *block.Memfd {
return u.memfd.Swap(nil)
}

// ServeStats returns a cumulative snapshot of demand faults served so far, or a
// zero snapshot if the handler has not been created yet (FC has not connected).
// It never blocks.
func (u *Uffd) ServeStats() userfaultfd.ServeSnapshot {
handler, err := u.handler.Result()
if err != nil {
return userfaultfd.ServeSnapshot{}
}

return handler.ServeStats()
}
45 changes: 45 additions & 0 deletions packages/orchestrator/pkg/sandbox/uffd/userfaultfd/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,51 @@ var prefaultTimer = utils.Must(telemetry.NewTimerFactory(
"UFFD prefault attempts",
))

// ServeSnapshot is a cumulative count of demand faults a handler has served,
// read at a point in time via Userfaultfd.ServeStats. Prefaults bypass the
// serve loop and are not counted. Sampling it at the envd-init boundary yields
// the pages/bytes a guest needed to start.
type ServeSnapshot struct {
// Pages is the number of demand faults resolved (installed or already
// present) — the count of distinct pages the guest needed.
Pages int64
// SourcePages is the subset of Pages installed from the source
// (page_class=new); the rest were zero-filled or already resident (a
// prefetch hit or a concurrent install).
SourcePages int64
// Bytes is the number of bytes installed into the guest by this handler
// (page_class new + zero). On a fixed page size this is Pages-minus-present
// times the page size, but it is tracked directly so it stays correct
// across mixed page sizes.
Bytes int64
}

// ServeStats returns a cumulative snapshot of the demand faults served so far.
func (u *Userfaultfd) ServeStats() ServeSnapshot {
return ServeSnapshot{
Pages: u.servedPages.Load(),
SourcePages: u.servedSourcePages.Load(),
Bytes: u.servedBytes.Load(),
}
}

// recordServeStats folds one finished serve attempt into the cumulative
// counters read by ServeStats. Only resolved faults (installed or already
// present) count as a needed page; deferred faults are re-served and counted
// when they finally resolve, so they are not double counted here.
func (u *Userfaultfd) recordServeStats(pclass pageClass, result faultResult, servedBytes int64) {
switch result {
case faultResultInstalled, faultResultPresent:
u.servedPages.Add(1)
if servedBytes > 0 {
u.servedBytes.Add(servedBytes)
if pclass == pageClassNew {
u.servedSourcePages.Add(1)
}
}
}
}

// prefaultAttrs holds a precomputed metric.MeasurementOption per faultResult.
// Prefaults have no page_class: they only ever target not-present pages (the
// Dirty/Zero pre-check records result="skipped" instead).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,36 @@ func TestPrefaultMetric(t *testing.T) {
assert.Equal(t, int64(0), bytesSum["deferred"])
}

// TestServeStats folds the same mix of finished serve attempts that
// TestServeMetric records and asserts the cumulative ServeStats() snapshot —
// the point-in-time count the startup metric samples at the envd-init boundary.
// Only resolved faults (installed or already-present) count as a needed page;
// deferred and errored attempts do not, so a deferred fault is not
// double-counted when it is later re-served and resolves.
func TestServeStats(t *testing.T) {
t.Parallel()

pageSize := int64(header.PageSize)

var u Userfaultfd
fold := func(class pageClass, result faultResult, bytes int64, n int) {
for range n {
u.recordServeStats(class, result, bytes)
}
}
fold(pageClassNew, faultResultInstalled, pageSize, 3) // +3 pages, +3 source, +3 pages of bytes
fold(pageClassZero, faultResultInstalled, pageSize, 1) // +1 page, +1 page of bytes (not source)
fold(pageClassResident, faultResultPresent, 0, 2) // +2 pages, no bytes
fold(pageClassNew, faultResultPresent, 0, 1) // +1 page (lost race), no bytes/source
fold(pageClassNew, faultResultDeferred, 0, 1) // not counted (re-served later)
fold(pageClassNew, faultResultError, 0, 1) // not counted (never resolved)

stats := u.ServeStats()
assert.Equal(t, int64(7), stats.Pages, "resolved demand faults (installed+present)")
assert.Equal(t, int64(3), stats.SourcePages, "subset installed from the source")
assert.Equal(t, 4*pageSize, stats.Bytes, "bytes installed (new+zero), present/deferred/error contribute none")
}

func key(pageClass, result string) string { return pageClass + "/" + result }

// attrKey returns "page_class/result" for serve datapoints and just "result"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,15 @@ type Userfaultfd struct {
// that has already been closed (and potentially recycled by the OS).
closed bool

// Cumulative demand-fault serve counters, read via ServeStats(). They
// mirror the orchestrator.sandbox.uffd.serve metric but as a per-handler
// snapshot, so a caller can sample "how many pages did this guest need so
// far" at a point in time (e.g. the moment envd init returns). Prefaults
// bypass the serve loop and are not counted here. See recordServeStats.
servedPages atomic.Int64 // faults resolved (installed or already-present)
servedSourcePages atomic.Int64 // subset installed from the source (page_class=new)
servedBytes atomic.Int64 // bytes installed into the guest (new + zero)

logger logger.Logger
}

Expand Down Expand Up @@ -417,7 +426,10 @@ func (u *Userfaultfd) Serve(
pclass := pageClassUnknown
result := faultResultInstalled
var servedBytes int64
defer func() { sw.RecordRaw(ctx, servedBytes, serveAttrs[pclass][result]) }()
defer func() {
sw.RecordRaw(ctx, servedBytes, serveAttrs[pclass][result])
u.recordServeStats(pclass, result, servedBytes)
}()

if h := u.testFaultHook.Load(); h != nil {
(*h)(addr, faultPhaseBeforeRLock)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,7 @@ func (cs *CreateSandbox) Sandbox(

err = sbx.WaitForEnvd(
ctx,
sandbox.StartTypeCreate,
waitEnvdTimeout,
)
if err != nil {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -235,6 +235,7 @@ func (lb *LayerExecutor) updateEnvdInSandbox(
// Step 4: Wait for envd to initialize
err = sbx.WaitForEnvd(
ctx,
sandbox.StartTypeCreate,
waitEnvdTimeout,
)
if err != nil {
Expand Down
15 changes: 15 additions & 0 deletions packages/shared/pkg/telemetry/meters.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,14 @@ const (
OrchestratorSandboxCreateDurationName HistogramType = "orchestrator.sandbox.create.duration"
WaitForEnvdDurationHistogramName HistogramType = "orchestrator.sandbox.envd.init.duration"

// Sandbox startup working-set histograms: demand-fault pages/bytes a guest
// needed to reach a successful envd init, recorded once per start. Sampled
// per start (not per fault), so histogram_quantile yields per-sandbox
// percentiles.
UffdStartupPagesHistogramName HistogramType = "orchestrator.sandbox.uffd.startup.pages"
UffdStartupSourcePagesHistogramName HistogramType = "orchestrator.sandbox.uffd.startup.source_pages"
UffdStartupBytesHistogramName HistogramType = "orchestrator.sandbox.uffd.startup.bytes"

// TCP Firewall histograms
TCPFirewallConnectionDurationHistogramName HistogramType = "orchestrator.tcpfirewall.connection.duration"
TCPFirewallConnectionsPerSandboxHistogramName HistogramType = "orchestrator.tcpfirewall.connections.per_sandbox"
Expand Down Expand Up @@ -382,6 +390,10 @@ var histogramDesc = map[HistogramType]string{
OrchestratorSandboxCreateDurationName: "Time taken to create a sandbox",
WaitForEnvdDurationHistogramName: "Time taken for Envd to initialize successfully",

UffdStartupPagesHistogramName: "Demand-fault pages a guest needed to reach a successful envd init, per start",
UffdStartupSourcePagesHistogramName: "Subset of startup demand-fault pages pulled from the source (e.g. GCS), per start",
UffdStartupBytesHistogramName: "Bytes faulted into a guest to reach a successful envd init, per start",

TCPFirewallConnectionDurationHistogramName: "Duration of TCP firewall proxied connections",
TCPFirewallConnectionsPerSandboxHistogramName: "Number of active TCP firewall connections per sandbox",

Expand Down Expand Up @@ -422,6 +434,9 @@ var histogramUnits = map[HistogramType]string{
BuildRootfsSizeHistogramName: "{By}",
OrchestratorSandboxCreateDurationName: "ms",
WaitForEnvdDurationHistogramName: "ms",
UffdStartupPagesHistogramName: "{page}",
UffdStartupSourcePagesHistogramName: "{page}",
UffdStartupBytesHistogramName: "{By}",
TCPFirewallConnectionDurationHistogramName: "ms",
TCPFirewallConnectionsPerSandboxHistogramName: "{connection}",

Expand Down
Loading