actorevent: define the actor usage event and the activation epoch (#1881)

Step 3 of #1748. Vocabulary only: nothing emits the event or sets the
epoch yet.

**Event.** `ate.actor.usage_sampled` in `actorevent`, sixteen required
keys via `ateattr`: identity, pool, sandbox class, source,
`ate.stats.kind` (periodic, initial, final), the four measurements,
observation time, epoch. Measurement names follow the
`ate.actor.stats.*` instruments; units are in the registry briefs. Body
is `Actor usage sampled`, the accepted stdout shape change.

**Scope.** The usage stream has its own instrumentation scope, the
lifecycle scope with `/usage` appended, so a consumer can sample or drop
it without touching lifecycle records. `NewEmitter` takes the scope; the
package-level `Log` keeps the lifecycle one, so ateapi is unchanged.

**Epoch.** `ate.actor.epoch` and `WorkloadStatsSample.epoch_unix_nano`:
the unix-nano time an activation began. CPU time is per activation for
every source; lifetime is the sum over epochs of each epoch's highest
value. Memory usage and working set are absolute; peak is as the source
reports it. An ateom that leaves the epoch at zero still sends the raw
guest counter. The measurement comments on fields 8 to 11 say all of
this.

**Registry.** Event group in `events.yaml`, seven attributes in
`metrics.yaml`, `ate.actor.epoch` in the no-actor-identity forbidden
list. The one-record-one-place rule and the collector guide now name
both scopes.

Left for step 4: the event's code anchor moves to the ateom emit site,
atelet's poller comments and the guide's stdout usage section get
rewritten, and the controller starts passing `OTEL_LOGS_EXPORTER` to
worker pods.
This commit is contained in:
Tim Bai
2026-09-25 19:54:15 +00:00
committed by GitHub
parent d3a556aae3
commit 03ba5821ec
11 changed files with 417 additions and 95 deletions
+10 -7
View File
@@ -433,8 +433,9 @@ for a worked example.
### A Note on Logs
**Substrate exports one thing over OTLP: the actor lifecycle events**, from
ateapi, through `serverboot.InitLogging`. They are off unless
**Substrate exports one thing over OTLP: the actor events**, ateapi's
lifecycle records and the ateoms' usage samples, through
`serverboot.InitLogging`. They are off unless
`OTEL_LOGS_EXPORTER=otlp` is set, which only the kind overlay does today. See
[the same records over OTLP](../../observability.md#the-same-records-over-otlp).
@@ -461,11 +462,13 @@ onto the log record's own trace fields with no transformation.
Those logs are collected by whatever agent already reads container stdout on
your nodes. The collector is not in that path.
If you do add such an agent, note that ateapi writes the lifecycle records to
both stdout and OTLP. Reading its stdout as well gives you each record twice:
either exclude the `ate-system` namespace from the agent, or drop the records
whose instrumentation scope is
`github.com/agent-substrate/substrate/internal/actorevent`.
If you do add such an agent, note that the actor events, ateapi's lifecycle
records and the ateoms' usage samples, go to both stdout and OTLP. Reading
stdout as well gives you each record twice. ateapi runs in `ate-system` and the
ateoms run in the worker pool namespaces, so a namespace exclusion does not
cover both: drop the records by `event.name` (`ate.actor.state_changed`,
`ate.actor.crashed`, `ate.actor.usage_sampled`), or exclude the
`ate-api-server` and `ateom` containers by name.
The logs pipeline also serves **your** workloads: actors or services
instrumented with the OpenTelemetry logs SDK that push OTLP log records to the
+71 -6
View File
@@ -15,14 +15,16 @@
# The log events that Substrate emits over OTLP. Weaver reads this file with
# metrics.yaml, which holds every attribute the groups below refer to.
#
# An event name promises a fixed set of attributes. The list in each group is
# the reviewable form of the Keys field on the matching internal/actorevent
# Event. TestEventsMatchTheRegistry in that package holds the two in step both
# ways, so an event has to be declared here before it can ship.
# An event name promises a set of attributes: a required one is on every
# record, and a conditionally required one is present exactly when its
# condition holds. The lists in each group are the reviewable form of the Keys
# and Conditional fields on the matching internal/actorevent Event.
# TestEventsMatchTheRegistry in that package holds the two in step both ways,
# so an event has to be declared here before it can ship.
#
# Metrics cannot answer what state an actor is in: the cardinality rule
# no-actor-identity bars actor identity from metric labels. These events carry
# the identity, thus they are the only per-actor state answer.
# no-actor-identity bars actor identity from metric labels. The lifecycle events
# carry the identity, thus they are the only per-actor state answer.
groups:
- id: event.ate.actor.state_changed
@@ -97,3 +99,66 @@ groups:
requirement_level: required
- ref: ate.actor.state
requirement_level: required
- id: event.ate.actor.usage_sampled
type: event
name: ate.actor.usage_sampled
stability: development
brief: A measurement of the CPU and the memory of one actor.
note: >
The ateom that runs the actor writes one record per actor on a timer,
plus an initial and a final record per activation. ate.stats.kind tells
them apart.
The timestamp is when the measurement was read. The measurements are
absent, not zero, while the actor is not measurable: the source is
unspecified, and the record says only that the actor is pending and,
through ate.actor.epoch, since when.
Group by ate.stats.source before you add values: the two sources do not
measure the same thing. ate.stats.cpu.time restarts at zero with each
ate.actor.epoch. ate.stats.memory.usage and ate.stats.memory.working_set
are absolute; ate.stats.memory.peak is as the source reports it.
This stream may be sampled or dropped. Select it by event name; thus a
consumer can sample or drop it without touching the lifecycle events.
annotations:
substrate:
emitted_by: [ateom-gvisor, ateom-microvm]
code_anchor: internal/actorevent/actorevent.go
cuj: How much CPU and memory does this actor use, and since when?
attributes:
- ref: ate.atespace
requirement_level: required
- ref: ate.actor.name
requirement_level: required
- ref: ate.actor.uid
requirement_level: required
- ref: ate.template.atespace
requirement_level: required
- ref: ate.template.name
requirement_level: required
- ref: ate.workerpool.namespace
requirement_level: required
- ref: ate.workerpool.name
requirement_level: required
- ref: ate.sandbox.class
requirement_level: required
- ref: ate.stats.source
requirement_level: required
- ref: ate.stats.kind
requirement_level: required
- ref: ate.actor.epoch
requirement_level: required
- ref: ate.stats.memory.usage
requirement_level:
conditionally_required: ate.stats.source is not unspecified.
- ref: ate.stats.memory.peak
requirement_level:
conditionally_required: ate.stats.source is not unspecified and the source reports a peak.
- ref: ate.stats.memory.working_set
requirement_level:
conditionally_required: ate.stats.source is not unspecified.
- ref: ate.stats.cpu.time
requirement_level:
conditionally_required: ate.stats.source is not unspecified.
+63
View File
@@ -69,6 +69,19 @@ groups:
docs/metrics/substrate.yaml. It separates two actors that held the
same name at different times.
examples: [8f2a1c4e6b0d47f1]
- id: ate.actor.epoch
stability: development
type: int
brief: >
The activation of the actor that a record belongs to, as the unix
time in nanoseconds when it began.
note: >
Logs only, by the cardinality rule no-actor-identity in
docs/metrics/substrate.yaml. A Run or a Restore starts an activation.
ate.stats.cpu.time restarts at zero with each one. The epoch is the
StartTimeUnixNano of a cumulative metric point, carried as an
attribute because a log record has no field for it.
examples: [1699999990000000000]
- id: ate.actor.state
stability: development
brief: The state the actor is in now.
@@ -534,6 +547,56 @@ groups:
stability: development
value: guest-agent
brief: The measurement comes from the guest agent. It includes only the containers of the workload.
- id: ate.stats.kind
stability: development
brief: Why the sample was taken.
note: >
Logs only. It is bounded, but it is only ever recorded beside actor
identity, which no metric may carry.
type:
members:
- id: periodic
stability: development
value: periodic
brief: The timer of the ateom took the sample.
- id: initial
stability: development
value: initial
brief: The initial sample of an activation, when the workload became measurable.
- id: final
stability: development
value: final
brief: The last sample of an activation, taken in the checkpoint before the workload is torn down.
- id: ate.stats.memory.usage
stability: development
type: int
brief: The memory in use, in bytes. Absolute, not per activation.
note: Logs only. The log field of the instrument ate.actor.stats.memory.usage.
examples: [41943040]
- id: ate.stats.memory.peak
stability: development
type: int
brief: >
The high-water mark of ate.stats.memory.usage, in bytes, as the
source reports it.
note: Logs only. Absent when the source cannot report one. The cgroup source restarts it on a restore; the guest-agent source does not.
examples: [50331648]
- id: ate.stats.memory.working_set
stability: development
type: int
brief: The memory in use less the reclaimable page cache, in bytes.
note: Logs only. The log field of the instrument ate.actor.stats.memory.working_set.
examples: [33554432]
- id: ate.stats.cpu.time
stability: development
type: double
brief: The CPU time used since the activation began, in seconds.
note: >
Logs only. Restarts at zero with each ate.actor.epoch. The lifetime
figure is the sum over the epochs of the highest value in each. The
unit is that of the instrument ate.actor.stats.cpu.time, thus one
query serves both signals.
examples: [1.5]
- id: registry.ate.deviation
type: attribute_group
+7 -5
View File
@@ -147,6 +147,7 @@ cardinality_rules:
- ate.atespace
- ate.actor.version
- ate.actor.container.name
- ate.actor.epoch
- id: bounded-or-catalog-scoped
brief: >
Each ate.* metric label has a list of permitted values, or it names an
@@ -196,15 +197,16 @@ lint_exceptions:
metric: none
rule: One record goes to one place.
brief: >
ateapi writes each actor lifecycle record twice: to stdout, and as an OTLP
event. Nothing reads both today. No collector config in this repository
declares a filelog receiver, and neither enterprise collector is a
ateapi writes each actor lifecycle record, and the ateoms each usage
sample, twice: to stdout, and as an OTLP event. Nothing reads both
today. No collector config in this repository declares a filelog
receiver, and neither enterprise collector is a
DaemonSet, thus no agent reads the stdout copy into the same pipeline.
The stdout copy is what kubectl logs shows. The OTLP copy is what a
backend queries and needs no knowledge of the stdout envelope.
An agent that reads pod stdout makes this a duplicate. Then drop one:
leave ate-system out of the globs of that agent, or drop the records with
the scope github.com/agent-substrate/substrate/internal/actorevent.
drop the records by event name, or leave the ate-api-server and ateom
containers out of the globs of that agent.
# -----------------------------------------------------------------------------
# The subsystems with no metrics. An agent that reads only the metrics gives the
+9 -6
View File
@@ -177,22 +177,25 @@ The counter carries no actor identity, so this record is the only way to attribu
#### The same records over OTLP
Both records also go out as OTLP log events, so a collector reads them without knowing substrate's stdout envelope. Set `OTEL_LOGS_EXPORTER=otlp` to turn it on; unset means `none`, which is what every environment but kind uses today. ateapi is the only emitter today. The ateoms have a LoggerProvider on the same switch and export through [the ateom relay](#the-ateom-otlp-relay), which carries logs, traces, and metrics. atecontroller does not pass `OTEL_LOGS_EXPORTER` to worker pods, so the kind ConfigMap turns on ateapi only; the ateoms stay at `none` until the controller propagates it.
These records also go out as OTLP log events, so a collector reads them without knowing substrate's stdout envelope. Set `OTEL_LOGS_EXPORTER=otlp` to turn it on; unset means `none`, which is what every environment but kind uses today. ateapi is the only emitter today. The ateoms have a LoggerProvider on the same switch and export through [the ateom relay](#the-ateom-otlp-relay), which carries logs, traces, and metrics. atecontroller does not pass `OTEL_LOGS_EXPORTER` to worker pods, so the kind ConfigMap turns on ateapi only; the ateoms stay at `none` until the controller propagates it.
Two `event.name` values, which is the OTLP LogRecord's own field rather than an attribute:
Three `event.name` values, which is the OTLP LogRecord's own field rather than an attribute:
| `event.name` | Body | Severity | Attributes |
|---|---|---|---|
| `ate.actor.state_changed` | `Actor state changed` | 9 | the five identity keys, `ate.actor.operation.name`, `ate.actor.state` |
| `ate.actor.crashed` | `Actor crashed` | 17 | the same keys |
| `ate.actor.usage_sampled` | `Actor usage sampled` | 9 | the five identity keys, `ate.workerpool.*`, `ate.sandbox.class`, `ate.stats.*`, `ate.actor.epoch` |
A crash is its own name because an event name promises a fixed set of attributes and a crash has a different severity and shape. There is no name per state: `ate.actor.state` already says which transition happened, so a consumer still selects on that one attribute and needs no map from a name to a state. Both names are in [`docs/metrics/registry/events.yaml`](metrics/registry/events.yaml), which `make verify` checks.
`ate.actor.usage_sampled` is the ateoms' record: one per actor per sampling period, plus an `initial` and a `final` per activation, told apart by `ate.stats.kind`. Its timestamp is when the measurement was read. Its measurements are named after the `ate.actor.stats.*` instruments and share their units, so `ate.stats.cpu.time` is seconds. They are absent, not zero, while the actor is not measurable, which the record says with `ate.stats.source` unspecified. `ate.stats.cpu.time` restarts at zero with each `ate.actor.epoch`, the unix-nano time the activation began, so a lifetime figure is the sum over epochs of each epoch's highest value; `ate.stats.memory.usage` and `ate.stats.memory.working_set` are absolute, and `ate.stats.memory.peak` is as the source reports it. The same measurements ride `WorkloadStatsSample` on the stats RPCs, with the epoch beside them.
The attributes are the same flat `ate.*` keys as the stdout copy, so they arrive as real log attributes with no transform in front of them. Trace context is not among them: it goes on the record's own `TraceId` and `SpanId` fields, where the stdout copy's top-level `trace_id`/`span_id` would be mapped to anyway. The instrumentation scope is `github.com/agent-substrate/substrate/internal/actorevent`, which is how you select this stream, or exclude it.
A crash is its own name because an event name promises a set of attributes and a crash has a different severity and shape. There is no name per state: `ate.actor.state` already says which transition happened, so a consumer still selects on that one attribute and needs no map from a name to a state. All three names are in [`docs/metrics/registry/events.yaml`](metrics/registry/events.yaml), which `make verify` checks.
**Never sample or filter this stream.** A consumer takes the last event for an actor's uid, so one dropped record reports a stale state with no sign that anything is missing. This is the one stream where a sampling policy is a correctness bug rather than a cost trade.
The attributes are the same flat `ate.*` keys as the stdout copy, so they arrive as real log attributes with no transform in front of them. Trace context is not among them: it goes on the record's own `TraceId` and `SpanId` fields, where the stdout copy's top-level `trace_id`/`span_id` would be mapped to anyway. The instrumentation scope is `github.com/agent-substrate/substrate/internal/actorevent` for every actor event. Select or drop a stream by `event.name`: the lifecycle events must never be sampled, the usage samples may be.
**Both copies exist on purpose.** No substrate or enterprise collector reads pod stdout today, so nothing is duplicated: the stdout copy is what `kubectl logs` shows and what keeps the component's bootstrap and crash output readable, and the OTLP copy is what a backend queries. If a `filelog` DaemonSet is ever added, drop one of the two — exclude ate-system from its include globs, or drop records whose scope is the one above.
**Never sample or filter the lifecycle stream.** A consumer takes the last event for an actor's uid, so one dropped record reports a stale state with no sign that anything is missing. This is the one stream where a sampling policy is a correctness bug rather than a cost trade.
**Both copies exist on purpose.** No substrate or enterprise collector reads pod stdout today, so nothing is duplicated: the stdout copy is what `kubectl logs` shows and what keeps the component's bootstrap and crash output readable, and the OTLP copy is what a backend queries. If a `filelog` DaemonSet is ever added, drop one of the two — drop the records by `event.name`, or exclude the `ate-api-server` and `ateom` containers by name.
### Per-Actor Usage Events
+66 -22
View File
@@ -12,10 +12,10 @@
// See the License for the specific language governing permissions and
// limitations under the License.
// Package actorevent emits the actor lifecycle events. Log writes both copies of
// a record, the stdout one and the OTLP one, from a single call, so nothing
// about a record is kept in step by hand. Ordinary component logs stay on
// stdout.
// Package actorevent emits the actor events: the lifecycle events and the usage
// samples. Log writes both copies of a record, the stdout one and the OTLP one,
// from a single call, so nothing about a record is kept in step by hand.
// Ordinary component logs stay on stdout.
//
// This is not an slog bridge. A bridge would put every component record on the
// wire, cannot set EventName, and would loop, because serverboot routes OTel SDK
@@ -29,6 +29,7 @@ package actorevent
import (
"context"
"log/slog"
"math"
"sync"
"time"
@@ -38,20 +39,25 @@ import (
"github.com/agent-substrate/substrate/internal/ateattr"
)
// ScopeName is how a consumer selects this stream.
// ScopeName is the instrumentation scope of every actor event. A consumer
// selects a stream by event name.
const ScopeName = "github.com/agent-substrate/substrate/internal/actorevent"
// Event is one name in the closed vocabulary. Name is the LogRecord's own event
// name field, not an attribute. Body and Severity live here rather than at a
// call site, so the two copies of a record cannot differ.
//
// Keys is the attribute set the name promises. An event name means a fixed
// shape, so the tests hold the two in step and a caller cannot widen the record.
// Keys is the attribute set the name promises on every record. Conditional
// holds the keys that are present exactly when the registry's condition holds,
// and absent otherwise: absence means not measured, zero means measured as
// zero. The tests hold both lists in step with the registry, and a caller cannot
// widen the record.
type Event struct {
Name string
Body string
Severity log.Severity
Keys []string
Name string
Body string
Severity log.Severity
Keys []string
Conditional []string
}
// Level is the stdout level for this event. slog has four levels to OTel's
@@ -78,8 +84,9 @@ var identityKeys = []string{
string(ateattr.TemplateNameKey),
}
// Two names, because a crash has a different shape and severity. Only two,
// because ate.actor.state already says which transition happened.
// Three names. A crash is its own because it has a different severity; there
// is no name per state, because ate.actor.state already says which transition
// happened. UsageSampled is the ateoms' measurement record.
var (
StateChanged = Event{
Name: "ate.actor.state_changed",
@@ -98,11 +105,32 @@ var (
string(ateattr.ActorOperationNameKey),
string(ateattr.ActorStateKey)),
}
UsageSampled = Event{
Name: "ate.actor.usage_sampled",
Body: "Actor usage sampled",
Severity: log.SeverityInfo,
Keys: append(append([]string{}, identityKeys...),
string(ateattr.WorkerPoolNamespaceKey),
string(ateattr.WorkerPoolNameKey),
string(ateattr.SandboxClassKey),
string(ateattr.StatsSourceKey),
string(ateattr.StatsKindKey),
string(ateattr.ActorEpochKey)),
// Absent while the actor is not measurable, and peak also when the
// source cannot report one.
Conditional: []string{
string(ateattr.StatsMemoryUsageKey),
string(ateattr.StatsMemoryPeakKey),
string(ateattr.StatsMemoryWorkingSetKey),
string(ateattr.StatsCPUTimeKey),
},
}
)
// events is the whole vocabulary, which the registry test walks. An event left
// out of it is never checked against docs/metrics/registry/events.yaml.
var events = []Event{StateChanged, Crashed}
var events = []Event{StateChanged, Crashed, UsageSampled}
// BuildRecord turns the stdout record into its OTLP form. Attributes carry
// everything machine-readable, so the body stays the display string.
@@ -132,7 +160,13 @@ func logValue(v slog.Value) log.Value {
case slog.KindInt64:
return log.Int64Value(v.Int64())
case slog.KindUint64:
return log.Int64Value(int64(v.Uint64()))
// OTel has no unsigned kind. Clamp rather than wrap, as the metric
// side does in addSat.
u := v.Uint64()
if u > math.MaxInt64 {
return log.Int64Value(math.MaxInt64)
}
return log.Int64Value(int64(u))
case slog.KindFloat64:
return log.Float64Value(v.Float64())
case slog.KindBool:
@@ -151,27 +185,32 @@ type Emitter struct {
logger log.Logger
}
// NewEmitter writes under ScopeName.
func NewEmitter(lp log.LoggerProvider) *Emitter {
return &Emitter{logger: lp.Logger(ScopeName)}
}
// Log writes both copies of ev from one call, off one time.Now(), so a consumer
// can join them on an exact timestamp. That is why the stdout record is built
// LogAt writes both copies of ev with t as their timestamp, so a consumer can
// join them on an exact time. t is when the thing happened: for a usage sample
// the read time, not the write time. That is why the stdout record is built
// here rather than through slog.LogAttrs, which would take its own reading.
//
// --log-level=warn silences the stdout copy of an info event while the OTLP copy
// still ships.
func (e *Emitter) Log(ctx context.Context, ev Event, attrs []slog.Attr) {
now := time.Now()
func (e *Emitter) LogAt(ctx context.Context, ev Event, t time.Time, attrs []slog.Attr) {
level := ev.Level()
if l := slog.Default(); l.Enabled(ctx, level) {
rec := slog.NewRecord(now, level, ev.Body, 0)
rec := slog.NewRecord(t, level, ev.Body, 0)
rec.AddAttrs(attrs...)
_ = l.Handler().Handle(ctx, rec)
}
e.emit(ctx, ev, now, attrs)
e.emit(ctx, ev, t, attrs)
}
// Log is LogAt with time.Now(), for an event that happens as it is written.
func (e *Emitter) Log(ctx context.Context, ev Event, attrs []slog.Attr) {
e.LogAt(ctx, ev, time.Now(), attrs)
}
// emit writes the OTLP copy. It is a no-op, and cheap, until InitLogging
@@ -194,3 +233,8 @@ var defaultEmitter = sync.OnceValue(func() *Emitter {
func Log(ctx context.Context, ev Event, attrs []slog.Attr) {
defaultEmitter().Log(ctx, ev, attrs)
}
// LogAt records ev at t through the process-wide provider.
func LogAt(ctx context.Context, ev Event, t time.Time, attrs []slog.Attr) {
defaultEmitter().LogAt(ctx, ev, t, attrs)
}
+93 -3
View File
@@ -18,6 +18,7 @@ import (
"context"
"log/slog"
"maps"
"math"
"slices"
"testing"
"time"
@@ -60,6 +61,33 @@ func crashedAttrs() []slog.Attr {
slog.String(string(ateattr.ActorStateKey), ateattr.ActorStateCrashed))
}
// usageSampledAttrs carries every UsageSampled key with a plausible value.
func usageSampledAttrs() []slog.Attr {
return append(ateattr.ActorLogAttrs(testAttribution()),
slog.String(string(ateattr.WorkerPoolNamespaceKey), "ate-system"),
slog.String(string(ateattr.WorkerPoolNameKey), "default"),
slog.String(string(ateattr.SandboxClassKey), "gvisor"),
slog.String(string(ateattr.StatsSourceKey), ateattr.StatsSourceCgroup),
slog.String(string(ateattr.StatsKindKey), ateattr.StatsKindPeriodic),
slog.Uint64(string(ateattr.StatsMemoryUsageKey), 40<<20),
slog.Uint64(string(ateattr.StatsMemoryPeakKey), 48<<20),
slog.Uint64(string(ateattr.StatsMemoryWorkingSetKey), 32<<20),
slog.Float64(string(ateattr.StatsCPUTimeKey), 1.5),
slog.Int64(string(ateattr.ActorEpochKey), 1_699_999_990_000_000_000))
}
// usagePendingAttrs is the record for an actor that is not measurable yet: the
// required keys only, with source unspecified.
func usagePendingAttrs() []slog.Attr {
return append(ateattr.ActorLogAttrs(testAttribution()),
slog.String(string(ateattr.WorkerPoolNamespaceKey), "ate-system"),
slog.String(string(ateattr.WorkerPoolNameKey), "default"),
slog.String(string(ateattr.SandboxClassKey), "gvisor"),
slog.String(string(ateattr.StatsSourceKey), ateattr.StatsSourceUnspecified),
slog.String(string(ateattr.StatsKindKey), ateattr.StatsKindPeriodic),
slog.Int64(string(ateattr.ActorEpochKey), 1_699_999_990_000_000_000))
}
func recordAttrs(rec log.Record) map[string]string {
got := make(map[string]string, rec.AttributesLen())
rec.WalkAttributes(func(kv log.KeyValue) bool {
@@ -82,6 +110,8 @@ func TestBuildRecord(t *testing.T) {
wantBody string
wantSev log.Severity
wantVals map[string]string
// wantAbsent are conditional keys this record must not carry.
wantAbsent []string
}{
{
name: "state changed",
@@ -118,6 +148,32 @@ func TestBuildRecord(t *testing.T) {
string(ateattr.ActorStateKey): ateattr.ActorStateCrashed,
},
},
{
name: "usage sampled",
event: UsageSampled,
attrs: usageSampledAttrs(),
wantName: "ate.actor.usage_sampled",
wantBody: "Actor usage sampled",
wantSev: log.SeverityInfo,
wantVals: map[string]string{
string(ateattr.StatsKindKey): ateattr.StatsKindPeriodic,
string(ateattr.ActorEpochKey): "1699999990000000000",
string(ateattr.StatsCPUTimeKey): "1.5",
string(ateattr.ActorUIDKey): testActorUID,
},
},
{
name: "usage sampled while pending carries no measurements",
event: UsageSampled,
attrs: usagePendingAttrs(),
wantName: "ate.actor.usage_sampled",
wantBody: "Actor usage sampled",
wantSev: log.SeverityInfo,
wantVals: map[string]string{
string(ateattr.StatsSourceKey): ateattr.StatsSourceUnspecified,
},
wantAbsent: UsageSampled.Conditional,
},
}
for _, tt := range tests {
@@ -150,18 +206,23 @@ func TestBuildRecord(t *testing.T) {
}
}
// The event name promises a fixed shape, so the emitted set and the
// declared set must match both ways.
// The event name promises a shape: every required key is present,
// and nothing outside the required and conditional sets.
for _, key := range tt.event.Keys {
if _, ok := got[key]; !ok {
t.Errorf("declared key %q is missing from the record", key)
}
}
for key := range got {
if !slices.Contains(tt.event.Keys, key) {
if !slices.Contains(tt.event.Keys, key) && !slices.Contains(tt.event.Conditional, key) {
t.Errorf("record carries %q, which %s does not declare", key, tt.event.Name)
}
}
for _, key := range tt.wantAbsent {
if _, ok := got[key]; ok {
t.Errorf("record carries %q, which must be absent here", key)
}
}
})
}
}
@@ -184,6 +245,7 @@ func TestEventLevel(t *testing.T) {
{"a sub-level keeps its range", log.SeverityInfo3, slog.LevelInfo},
{"the state_changed event", StateChanged.Severity, slog.LevelInfo},
{"the crashed event", Crashed.Severity, slog.LevelError},
{"the usage_sampled event", UsageSampled.Severity, slog.LevelInfo},
}
for _, tt := range tests {
@@ -237,6 +299,7 @@ func TestLogWritesBothCopies(t *testing.T) {
}{
{"state changed", StateChanged, stateChangedAttrs(ateattr.ActorStateRunning)},
{"crashed", Crashed, crashedAttrs()},
{"usage sampled", UsageSampled, usageSampledAttrs()},
}
for _, tt := range tests {
@@ -356,6 +419,7 @@ func TestBuildRecordKeepsValueKinds(t *testing.T) {
{"string", slog.String("k", "v"), log.StringValue("v")},
{"int", slog.Int64("k", 7), log.Int64Value(7)},
{"uint", slog.Uint64("k", 7), log.Int64Value(7)},
{"uint above int64 clamps", slog.Uint64("k", math.MaxUint64), log.Int64Value(math.MaxInt64)},
{"float", slog.Float64("k", 1.5), log.Float64Value(1.5)},
{"bool", slog.Bool("k", true), log.BoolValue(true)},
{"duration is nanoseconds, as in the stdout copy", slog.Duration("k", 1500*time.Millisecond), log.Int64Value(1_500_000_000)},
@@ -381,3 +445,29 @@ func TestBuildRecordKeepsValueKinds(t *testing.T) {
})
}
}
// TestLogAtStampsBothCopies pins that the caller's time, not the write time, is
// the timestamp of both copies: a usage sample is dated when it was read.
func TestLogAtStampsBothCopies(t *testing.T) {
stdout := &captureHandler{}
prev := slog.Default()
slog.SetDefault(slog.New(stdout))
t.Cleanup(func() { slog.SetDefault(prev) })
exp := &memExporter{}
lp := sdklog.NewLoggerProvider(sdklog.WithProcessor(sdklog.NewSimpleProcessor(exp)))
t.Cleanup(func() { _ = lp.Shutdown(context.Background()) })
at := time.Date(2026, 9, 25, 12, 0, 0, 0, time.UTC)
NewEmitter(lp).LogAt(context.Background(), UsageSampled, at, usageSampledAttrs())
if len(stdout.records) != 1 || len(exp.records) != 1 {
t.Fatalf("wrote %d stdout and %d OTLP records, want 1 and 1", len(stdout.records), len(exp.records))
}
if got := stdout.records[0].Time; !got.Equal(at) {
t.Errorf("stdout time = %v, want %v", got, at)
}
if got := exp.records[0].Timestamp(); !got.Equal(at) {
t.Errorf("OTLP timestamp = %v, want %v", got, at)
}
}
+29 -12
View File
@@ -31,8 +31,9 @@ type registry struct {
Type string `json:"type"`
Name string `json:"name"`
Attributes []struct {
Ref string `json:"ref"`
RequirementLevel string `json:"requirement_level"`
Ref string `json:"ref"`
// "required", or {conditionally_required: <condition>}.
RequirementLevel any `json:"requirement_level"`
} `json:"attributes"`
} `json:"groups"`
}
@@ -55,33 +56,49 @@ func TestEventsMatchTheRegistry(t *testing.T) {
t.Fatalf("parse %s: %v", registryPath, err)
}
declared := map[string][]string{}
type shape struct{ required, conditional []string }
declared := map[string]shape{}
for _, g := range reg.Groups {
if g.Type != "event" {
continue
}
keys := make([]string, 0, len(g.Attributes))
var sh shape
for _, a := range g.Attributes {
if a.RequirementLevel != "required" {
t.Errorf("%s declares %s as %q; an event name promises a fixed shape, so every attribute is required",
g.Name, a.Ref, a.RequirementLevel)
switch lvl := a.RequirementLevel.(type) {
case string:
if lvl != "required" {
t.Errorf("%s declares %s as %q; an event attribute is required or conditionally required with a stated condition",
g.Name, a.Ref, lvl)
}
sh.required = append(sh.required, a.Ref)
case map[string]any:
if c, ok := lvl["conditionally_required"].(string); !ok || c == "" {
t.Errorf("%s declares %s as %v; an event attribute is required or conditionally required with a stated condition",
g.Name, a.Ref, lvl)
}
sh.conditional = append(sh.conditional, a.Ref)
default:
t.Errorf("%s declares %s with an unreadable requirement level %v", g.Name, a.Ref, lvl)
}
keys = append(keys, a.Ref)
}
declared[g.Name] = keys
declared[g.Name] = sh
}
for _, ev := range events {
keys, ok := declared[ev.Name]
sh, ok := declared[ev.Name]
if !ok {
t.Errorf("%s has no group in %s; add one", ev.Name, registryPath)
continue
}
delete(declared, ev.Name)
got, want := slices.Sorted(slices.Values(keys)), slices.Sorted(slices.Values(ev.Keys))
got, want := slices.Sorted(slices.Values(sh.required)), slices.Sorted(slices.Values(ev.Keys))
if !slices.Equal(got, want) {
t.Errorf("%s: %s declares %v, the Event declares %v", ev.Name, registryPath, got, want)
t.Errorf("%s: %s requires %v, the Event declares %v", ev.Name, registryPath, got, want)
}
got, want = slices.Sorted(slices.Values(sh.conditional)), slices.Sorted(slices.Values(ev.Conditional))
if !slices.Equal(got, want) {
t.Errorf("%s: %s conditionally requires %v, the Event declares %v", ev.Name, registryPath, got, want)
}
}
+26
View File
@@ -173,6 +173,32 @@ const (
StatsSourceKey = attribute.Key("ate.stats.source")
)
// Keys of the actor usage record. Logs only: the measurements are unbounded and
// the epoch is per activation. The measurement names follow the
// ate.actor.stats.* instruments and share their units, so a query works the
// same on both signals: the memory keys are bytes, StatsCPUTimeKey is seconds
// as a double.
//
// ActorEpochKey identifies one activation of an actor, a Run or Restore, as the
// unix-nano time it began. It is the boundary of ate.stats.cpu.time, which
// restarts at zero with each activation.
const (
ActorEpochKey = attribute.Key("ate.actor.epoch")
StatsKindKey = attribute.Key("ate.stats.kind")
StatsMemoryUsageKey = attribute.Key("ate.stats.memory.usage")
StatsMemoryPeakKey = attribute.Key("ate.stats.memory.peak")
StatsMemoryWorkingSetKey = attribute.Key("ate.stats.memory.working_set")
StatsCPUTimeKey = attribute.Key("ate.stats.cpu.time")
)
// Values for StatsKindKey. An initial or final sample brackets an activation; a
// periodic one is the timer's.
const (
StatsKindPeriodic = "periodic"
StatsKindInitial = "initial"
StatsKindFinal = "final"
)
// Values for StatsSourceKey, mirroring ateompb.StatsSource. The two sources do
// not measure the same thing (the cgroup source charges the sandbox runtime's
// overhead along with the workload, the guest-agent source sees only the
+27 -19
View File
@@ -1472,31 +1472,31 @@ type WorkloadStatsSample struct {
// Measurements. All four are zero when source is STATS_SOURCE_UNSPECIFIED,
// which means "not measured" rather than "measured as zero".
//
// Two of them accumulate -- memory_peak_bytes and cpu_usage_usec -- and both
// are scoped to the current EPOCH rather than to the actor's lifetime. An
// epoch begins wherever the accounting behind the sample begins, which is not
// the same event for every source: STATS_SOURCE_CGROUP reads a sandbox cgroup
// that a restore recreates, so both restart at zero there, while
// STATS_SOURCE_GUEST_AGENT reads counters the guest kernel keeps in its own
// RAM, which a restored guest brings back with it. A caller that wants a
// lifetime figure has to accumulate one itself, and must read a decrease as a
// new epoch rather than emit a negative delta -- but not the converse. An
// epoch can also begin at a value above the last one reported, so no
// comparison of consecutive samples detects every boundary.
// memory_current_bytes and memory_working_set_bytes are absolute.
MemoryCurrentBytes uint64 `protobuf:"varint,8,opt,name=memory_current_bytes,json=memoryCurrentBytes,proto3" json:"memory_current_bytes,omitempty"`
// High-water mark of memory_current_bytes within the current epoch. Also zero
// when the runtime cannot report a peak at all: the cgroup source reads
// memory.peak, which only exists on Linux 5.19 and later.
// High-water mark of memory_current_bytes, as the source reports it: the
// cgroup source restarts it on a restore, the guest-agent source brings the
// guest's own back with it. Also zero when the runtime cannot report a peak:
// the cgroup source reads memory.peak, which only exists on Linux 5.19 and
// later.
MemoryPeakBytes uint64 `protobuf:"varint,9,opt,name=memory_peak_bytes,json=memoryPeakBytes,proto3" json:"memory_peak_bytes,omitempty"`
// memory_current_bytes less the reclaimable page cache, floored at zero. This
// is the figure to compare against a memory limit; memory_current_bytes
// drifts upward with cache the kernel would drop for free under pressure.
MemoryWorkingSetBytes uint64 `protobuf:"varint,10,opt,name=memory_working_set_bytes,json=memoryWorkingSetBytes,proto3" json:"memory_working_set_bytes,omitempty"`
// Cumulative CPU time within the current epoch.
// Cumulative CPU time since the activation began (epoch_unix_nano), for
// every source. The cgroup source restarts on its own; the guest-agent
// source is rebased on its first read after a restore, since the guest's
// counters survive in guest RAM. An ateom that leaves epoch_unix_nano at zero
// also sends the raw guest counter. A lifetime figure is the sum over epochs
// of each epoch's highest value.
CpuUsageUsec uint64 `protobuf:"varint,11,opt,name=cpu_usage_usec,json=cpuUsageUsec,proto3" json:"cpu_usage_usec,omitempty"`
ObservedAtUnixNano int64 `protobuf:"varint,12,opt,name=observed_at_unix_nano,json=observedAtUnixNano,proto3" json:"observed_at_unix_nano,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
// The activation this sample belongs to, as the unix-nano time it began. A
// Run or Restore starts one. Zero from an ateom that does not set it.
EpochUnixNano int64 `protobuf:"varint,13,opt,name=epoch_unix_nano,json=epochUnixNano,proto3" json:"epoch_unix_nano,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
func (x *WorkloadStatsSample) Reset() {
@@ -1613,6 +1613,13 @@ func (x *WorkloadStatsSample) GetObservedAtUnixNano() int64 {
return 0
}
func (x *WorkloadStatsSample) GetEpochUnixNano() int64 {
if x != nil {
return x.EpochUnixNano
}
return 0
}
type GetWorkloadStatsResponse struct {
state protoimpl.MessageState `protogen:"open.v1"`
Sample *WorkloadStatsSample `protobuf:"bytes,1,opt,name=sample,proto3" json:"sample,omitempty"`
@@ -1863,7 +1870,7 @@ const file_ateom_proto_rawDesc = "" +
"\x0f_egress_gateway\"\x19\n" +
"\x17RestoreWorkloadResponse\"6\n" +
"\x17GetWorkloadStatsRequest\x12\x1b\n" +
"\tactor_uid\x18\x01 \x01(\tR\bactorUid\"\xab\x04\n" +
"\tactor_uid\x18\x01 \x01(\tR\bactorUid\"\xd3\x04\n" +
"\x13WorkloadStatsSample\x12\x1a\n" +
"\batespace\x18\x01 \x01(\tR\batespace\x12\x1d\n" +
"\n" +
@@ -1878,7 +1885,8 @@ const file_ateom_proto_rawDesc = "" +
"\x18memory_working_set_bytes\x18\n" +
" \x01(\x04R\x15memoryWorkingSetBytes\x12$\n" +
"\x0ecpu_usage_usec\x18\v \x01(\x04R\fcpuUsageUsec\x121\n" +
"\x15observed_at_unix_nano\x18\f \x01(\x03R\x12observedAtUnixNano\"N\n" +
"\x15observed_at_unix_nano\x18\f \x01(\x03R\x12observedAtUnixNano\x12&\n" +
"\x0fepoch_unix_nano\x18\r \x01(\x03R\repochUnixNano\"N\n" +
"\x18GetWorkloadStatsResponse\x122\n" +
"\x06sample\x18\x01 \x01(\v2\x1a.ateom.WorkloadStatsSampleR\x06sample\"\x1f\n" +
"\x1dGetActiveWorkloadStatsRequest\"V\n" +
+16 -15
View File
@@ -394,30 +394,31 @@ message WorkloadStatsSample {
// Measurements. All four are zero when source is STATS_SOURCE_UNSPECIFIED,
// which means "not measured" rather than "measured as zero".
//
// Two of them accumulate -- memory_peak_bytes and cpu_usage_usec -- and both
// are scoped to the current EPOCH rather than to the actor's lifetime. An
// epoch begins wherever the accounting behind the sample begins, which is not
// the same event for every source: STATS_SOURCE_CGROUP reads a sandbox cgroup
// that a restore recreates, so both restart at zero there, while
// STATS_SOURCE_GUEST_AGENT reads counters the guest kernel keeps in its own
// RAM, which a restored guest brings back with it. A caller that wants a
// lifetime figure has to accumulate one itself, and must read a decrease as a
// new epoch rather than emit a negative delta -- but not the converse. An
// epoch can also begin at a value above the last one reported, so no
// comparison of consecutive samples detects every boundary.
// memory_current_bytes and memory_working_set_bytes are absolute.
uint64 memory_current_bytes = 8;
// High-water mark of memory_current_bytes within the current epoch. Also zero
// when the runtime cannot report a peak at all: the cgroup source reads
// memory.peak, which only exists on Linux 5.19 and later.
// High-water mark of memory_current_bytes, as the source reports it: the
// cgroup source restarts it on a restore, the guest-agent source brings the
// guest's own back with it. Also zero when the runtime cannot report a peak:
// the cgroup source reads memory.peak, which only exists on Linux 5.19 and
// later.
uint64 memory_peak_bytes = 9;
// memory_current_bytes less the reclaimable page cache, floored at zero. This
// is the figure to compare against a memory limit; memory_current_bytes
// drifts upward with cache the kernel would drop for free under pressure.
uint64 memory_working_set_bytes = 10;
// Cumulative CPU time within the current epoch.
// Cumulative CPU time since the activation began (epoch_unix_nano), for
// every source. The cgroup source restarts on its own; the guest-agent
// source is rebased on its first read after a restore, since the guest's
// counters survive in guest RAM. An ateom that leaves epoch_unix_nano at zero
// also sends the raw guest counter. A lifetime figure is the sum over epochs
// of each epoch's highest value.
uint64 cpu_usage_usec = 11;
int64 observed_at_unix_nano = 12;
// The activation this sample belongs to, as the unix-nano time it began. A
// Run or Restore starts one. Zero from an ateom that does not set it.
int64 epoch_unix_nano = 13;
}
message GetWorkloadStatsResponse {