(chore): Write both copies of an actor lifecycle event from one call (#1771)

An actor lifecycle record went to stdout and to OTLP in two separate
calls. The attributes were shared, but the severity was not, so the slog
level sat at the call site and the OTel severity sat on the `Event`.
Nothing made a caller write both copies either, so a new record could
reach stdout only and no test would notice.

Since this PR `actorevent.Log` now writes both copies. The level comes
from `Event.Severity`, so it is stored once. The OTLP only `Emit` and
the exported body constants are gone, so the dual write is the only way
out of the package.

Also normalizes `ate.actor.operation.name` on the crash event, which the
state change event already did.

Fixes #1744 

Testing

- Unit tests for the level mapping, the dual write, and the registry
check.
- End to end on a fresh kind cluster. Both event names arrive with the
right severity (9 and 17), the right attributes, and trace context on
the record fields. The stdout copies match record for record.

- [x] Tests pass
- [ ] Appropriate changes to documentation are included in the PR
This commit is contained in:
Krisztian F
2026-09-21 20:55:52 +00:00
committed by GitHub
parent aa17b3c297
commit f75e626485
9 changed files with 347 additions and 64 deletions
+2 -3
View File
@@ -123,11 +123,10 @@ func crashActor(ctx context.Context, st crashActorStore, actorRef resources.Acto
// last state an actor reached has to see this record to reach "crashed" at all.
func logActorCrashed(ctx context.Context, actor *ateapipb.Actor, opName, reason string) {
attrs := ateattr.ActorLogAttrs(resources.ActorAttributionFromActor(actor))
attrs = append(attrs, slog.String(string(ateattr.ActorOperationNameKey), opName))
attrs = append(attrs, slog.String(string(ateattr.ActorOperationNameKey), ateattr.NormalizeOperationName(opName)))
attrs = append(attrs, slog.String(string(ateattr.ActorStateKey), ateattr.ActorStateCrashed))
attrs = append(attrs, ateattr.FailureLogAttrs(reason)...)
slog.LogAttrs(ctx, slog.LevelError, actorevent.CrashedBody, attrs...)
actorevent.Emit(ctx, actorevent.Crashed, attrs)
actorevent.Log(ctx, actorevent.Crashed, attrs)
}
// crashActorStore encapsulates the subset of store operations needed to crash
+55 -13
View File
@@ -23,6 +23,7 @@ import (
"strings"
"sync"
"testing"
"time"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store"
"github.com/agent-substrate/substrate/cmd/ateapi/internal/store/storetest"
@@ -584,27 +585,36 @@ func TestCrashActorReleaseFailureLeavesWorkerReclaimable(t *testing.T) {
// crashRecords captures the "Actor crashed" records a crash emits, so a test can
// assert the identity that ate.actor.crashes is barred from carrying. crashActor
// logs through the slog default, so this swaps it and the caller cannot be parallel.
func crashRecords(t *testing.T) *[]map[string]string {
func crashRecords(t *testing.T) *[]stdoutRecord {
t.Helper()
return logRecords(t, "Actor crashed")
return logRecords(t, actorevent.Crashed.Body)
}
// logRecords captures the attributes of every record with the given message.
func logRecords(t *testing.T, msg string) *[]map[string]string {
// stdoutRecord is one captured record from the stdout copy. It keeps the level
// and the time, not just the attributes, so a test can hold those against the
// OTLP copy. The message is whatever logRecords filtered on.
type stdoutRecord struct {
level slog.Level
time time.Time
attrs map[string]string
}
// logRecords captures every record with the given message.
func logRecords(t *testing.T, msg string) *[]stdoutRecord {
t.Helper()
var records []map[string]string
var records []stdoutRecord
prev := slog.Default()
slog.SetDefault(slog.New(slogHandlerFunc(func(r slog.Record) {
if r.Message != msg {
return
}
fields := map[string]string{}
rec := stdoutRecord{level: r.Level, time: r.Time, attrs: map[string]string{}}
r.Attrs(func(a slog.Attr) bool {
fields[a.Key] = a.Value.String()
rec.attrs[a.Key] = a.Value.String()
return true
})
records = append(records, fields)
records = append(records, rec)
})))
t.Cleanup(func() { slog.SetDefault(prev) })
return &records
@@ -612,8 +622,11 @@ func logRecords(t *testing.T, msg string) *[]map[string]string {
// otlpEvent is one captured record from the OTLP copy of a log record.
type otlpEvent struct {
name string
attrs map[string]string
name string
body string
severity otellog.Severity
timestamp time.Time
attrs map[string]string
}
var (
@@ -628,7 +641,13 @@ func (otlpSinkExporter) Export(_ context.Context, records []sdklog.Record) error
otlpSinkMu.Lock()
defer otlpSinkMu.Unlock()
for _, r := range records {
e := otlpEvent{name: r.EventName(), attrs: map[string]string{}}
e := otlpEvent{
name: r.EventName(),
body: r.Body().String(),
severity: r.Severity(),
timestamp: r.Timestamp(),
attrs: map[string]string{},
}
r.WalkAttributes(func(kv otellog.KeyValue) bool {
e.attrs[kv.Key] = kv.Value.String()
return true
@@ -678,6 +697,27 @@ func (f slogHandlerFunc) Handle(_ context.Context, r slog.Record) error {
func (f slogHandlerFunc) WithAttrs([]slog.Attr) slog.Handler { return f }
func (f slogHandlerFunc) WithGroup(string) slog.Handler { return f }
// assertCopiesAgree checks the fields that no longer live at a call site. Both
// copies take their severity and body from ev, so neither can hold its own. A
// stdout body that drifted fails earlier, when logRecords matches nothing.
func assertCopiesAgree(t *testing.T, stdout stdoutRecord, otlp otlpEvent, ev actorevent.Event) {
t.Helper()
if stdout.level != ev.Level() {
t.Errorf("stdout level = %v, want %v", stdout.level, ev.Level())
}
if otlp.severity != ev.Severity {
t.Errorf("OTLP severity = %v, want %v", otlp.severity, ev.Severity)
}
if otlp.body != ev.Body {
t.Errorf("OTLP body = %q, want %q", otlp.body, ev.Body)
}
// One time.Now() serves both, so a consumer can join them on it.
if !stdout.time.Equal(otlp.timestamp) {
t.Errorf("timestamps differ: stdout %v, OTLP %v", stdout.time, otlp.timestamp)
}
}
// The crash record is the only signal carrying actor identity, so it must fire
// exactly when the counter does. A crash counted but not logged is unattributable;
// one logged but not counted double-counts on a retry.
@@ -707,7 +747,7 @@ func TestCrashActor_RecordAndCounterAgree(t *testing.T) {
t.Fatalf("got %d crash records, want 1", len(*records))
}
got := (*records)[0]
got := (*records)[0].attrs
stored, err := st.GetActor(ctx, actorRef)
if err != nil {
t.Fatalf("GetActor: %v", err)
@@ -730,7 +770,8 @@ func TestCrashActor_RecordAndCounterAgree(t *testing.T) {
t.Error("crash record carries no ate.actor.uid; it cannot survive a name reuse")
}
// The OTLP copy is the same record under an event name.
// The OTLP copy is the same record under an event name. One call writes both,
// so anything either copy holds alone is a bug in actorevent.Log.
gotEvents := events()
if len(gotEvents) != 1 {
t.Fatalf("got %d crash events, want 1: %v", len(gotEvents), gotEvents)
@@ -741,6 +782,7 @@ func TestCrashActor_RecordAndCounterAgree(t *testing.T) {
if !maps.Equal(gotEvents[0].attrs, got) {
t.Errorf("crash event attributes = %v, want the stdout record's %v", gotEvents[0].attrs, got)
}
assertCopiesAgree(t, (*records)[0], gotEvents[0], actorevent.Crashed)
// Re-crashing an already-crashed actor must move neither signal.
if err := crashActor(ctx, st, actorRef, ateattr.OperationResume, ateattr.ReasonWorkerPodGone); err != nil {
+1 -2
View File
@@ -96,8 +96,7 @@ func logActorState(ctx context.Context, actor *ateapipb.Actor, opName, state str
attrs = append(attrs,
slog.String(string(ateattr.ActorOperationNameKey), ateattr.NormalizeOperationName(opName)),
slog.String(string(ateattr.ActorStateKey), state))
slog.LogAttrs(ctx, slog.LevelInfo, actorevent.StateChangedBody, attrs...)
actorevent.Emit(ctx, actorevent.StateChanged, attrs)
actorevent.Log(ctx, actorevent.StateChanged, attrs)
}
// ActorWorkflow handles the workflows for actor's resume / suspend operations.
@@ -92,10 +92,10 @@ func TestEnsurePausedFinalized_WorkerGone(t *testing.T) {
if len(*records) != 1 {
t.Fatalf("got %d crash records, want 1", len(*records))
}
if got := (*records)[0][string(ateattr.ActorUIDKey)]; got == "" {
if got := (*records)[0].attrs[string(ateattr.ActorUIDKey)]; got == "" {
t.Error("crash record carries no ate.actor.uid")
}
if got := (*records)[0][string(ateattr.FailureDomainKey)]; got != ateattr.FailureDomainInfrastructure {
if got := (*records)[0].attrs[string(ateattr.FailureDomainKey)]; got != ateattr.FailureDomainInfrastructure {
t.Errorf("ate.failure.domain = %q, want %q", got, ateattr.FailureDomainInfrastructure)
}
}
@@ -86,7 +86,7 @@ func TestActorStateChangeRecords(t *testing.T) {
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
ctx := context.Background()
records := logRecords(t, actorevent.StateChangedBody)
records := logRecords(t, actorevent.StateChanged.Body)
events := otlpEvents(t)
persistence := newTestPersistence(t)
@@ -118,7 +118,7 @@ func TestActorStateChangeRecords(t *testing.T) {
if len(*records) != 1 {
t.Fatalf("got %d state records, want 1: %v", len(*records), *records)
}
got := (*records)[0]
got := (*records)[0].attrs
want := map[string]string{
string(ateattr.AtespaceKey): actorRef.Atespace,
string(ateattr.ActorNameKey): actorRef.Name,
@@ -137,7 +137,8 @@ func TestActorStateChangeRecords(t *testing.T) {
t.Errorf("got %d attributes, want %d: %v", len(got), len(want), got)
}
// The OTLP copy is the same record under an event name.
// The OTLP copy is the same record under an event name. One call writes
// both, so anything either copy holds alone is a bug in actorevent.Log.
gotEvents := events()
if len(gotEvents) != 1 {
t.Fatalf("got %d state events, want 1: %v", len(gotEvents), gotEvents)
@@ -148,6 +149,7 @@ func TestActorStateChangeRecords(t *testing.T) {
if !maps.Equal(gotEvents[0].attrs, got) {
t.Errorf("state event attributes = %v, want the stdout record's %v", gotEvents[0].attrs, got)
}
assertCopiesAgree(t, (*records)[0], gotEvents[0], actorevent.StateChanged)
stored, err := persistence.GetActor(ctx, actorRef)
if err != nil {
@@ -165,7 +167,7 @@ func TestActorStateChangeRecords(t *testing.T) {
// retention.
func TestActorCreatedRecord(t *testing.T) {
ctx := context.Background()
records := logRecords(t, "Actor state changed")
records := logRecords(t, actorevent.StateChanged.Body)
persistence := newTestPersistence(t)
storetest.MustCreateAtespace(t, ctx, persistence, "ns")
@@ -184,7 +186,7 @@ func TestActorCreatedRecord(t *testing.T) {
if len(*records) != 1 {
t.Fatalf("got %d state records, want 1: %v", len(*records), *records)
}
got := (*records)[0]
got := (*records)[0].attrs
if got[string(ateattr.ActorOperationNameKey)] != ateattr.OperationCreate {
t.Errorf("operation = %q, want %q", got[string(ateattr.ActorOperationNameKey)], ateattr.OperationCreate)
}
@@ -200,7 +202,7 @@ func TestActorCreatedRecord(t *testing.T) {
// finished delete from one that is stuck.
func TestActorDeletedRecord(t *testing.T) {
ctx := context.Background()
records := logRecords(t, "Actor state changed")
records := logRecords(t, actorevent.StateChanged.Body)
persistence := newTestPersistence(t)
storetest.MustCreateAtespace(t, ctx, persistence, "ns")
@@ -226,7 +228,7 @@ func TestActorDeletedRecord(t *testing.T) {
if len(*records) != 1 {
t.Fatalf("got %d state records, want 1: %v", len(*records), *records)
}
got := (*records)[0]
got := (*records)[0].attrs
if got[string(ateattr.ActorStateKey)] != ateattr.ActorStateDeleted {
t.Errorf("state = %q, want %q", got[string(ateattr.ActorStateKey)], ateattr.ActorStateDeleted)
}
@@ -246,7 +248,7 @@ func TestActorDeletedRecord(t *testing.T) {
// in the stream the store never held.
func TestActorStateChangeRecordSkippedOnConflict(t *testing.T) {
ctx := context.Background()
records := logRecords(t, "Actor state changed")
records := logRecords(t, actorevent.StateChanged.Body)
persistence := newTestPersistence(t)
storetest.MustCreateAtespace(t, ctx, persistence, "ns")
+2 -1
View File
@@ -17,7 +17,8 @@
#
# 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, and a test in that package holds the two in step.
# 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
+56 -22
View File
@@ -12,16 +12,17 @@
// See the License for the specific language governing permissions and
// limitations under the License.
// Package actorevent emits the actor lifecycle events over OTLP. Events go over
// OTLP; ordinary component logs stay on stdout. The caller passes the same
// []slog.Attr to both copies, so the two cannot drift.
// 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.
//
// 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
// errors through slog.
//
// serverboot.InitLogging pairs this with a batching processor. These records sit
// on the actor resume path, so exporting inside Emit would put a blocking gRPC
// on the actor resume path, so exporting inside Log would put a blocking gRPC
// call there and make a slow collector look like control-plane latency.
package actorevent
@@ -40,14 +41,9 @@ import (
// ScopeName is how a consumer selects this stream.
const ScopeName = "github.com/agent-substrate/substrate/internal/actorevent"
// Bodies match the stdout record's message, so both read alike.
const (
StateChangedBody = "Actor state changed"
CrashedBody = "Actor crashed"
)
// Event is one name in the closed vocabulary. Name is the LogRecord's own event
// name field, not an attribute.
// 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.
@@ -58,6 +54,21 @@ type Event struct {
Keys []string
}
// Level is the stdout level for this event. slog has four levels to OTel's
// twenty-four, so a sub-level collapses onto the range it sits in.
func (ev Event) Level() slog.Level {
switch {
case ev.Severity >= log.SeverityError:
return slog.LevelError
case ev.Severity >= log.SeverityWarn:
return slog.LevelWarn
case ev.Severity >= log.SeverityInfo:
return slog.LevelInfo
default:
return slog.LevelDebug
}
}
// identityKeys is what ateattr.ActorLogAttrs writes, in its order.
var identityKeys = []string{
string(ateattr.AtespaceKey),
@@ -72,7 +83,7 @@ var identityKeys = []string{
var (
StateChanged = Event{
Name: "ate.actor.state_changed",
Body: StateChangedBody,
Body: "Actor state changed",
Severity: log.SeverityInfo,
Keys: append(append([]string{}, identityKeys...),
string(ateattr.ActorOperationNameKey),
@@ -81,7 +92,7 @@ var (
Crashed = Event{
Name: "ate.actor.crashed",
Body: CrashedBody,
Body: "Actor crashed",
Severity: log.SeverityError,
Keys: append(append([]string{}, identityKeys...),
string(ateattr.ActorOperationNameKey),
@@ -91,6 +102,10 @@ var (
}
)
// 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}
// BuildRecord turns the stdout record into its OTLP form. Attributes carry
// everything machine-readable, so the body stays the display string.
//
@@ -132,8 +147,8 @@ func logValue(v slog.Value) log.Value {
}
}
// Emitter writes events through one log.Logger. Tests construct one directly, so
// they need no global provider and can run in parallel.
// Emitter writes the OTLP copy through one log.Logger. Tests construct one
// directly, so emit needs no global provider and can run in parallel.
type Emitter struct {
logger log.Logger
}
@@ -142,14 +157,33 @@ func NewEmitter(lp log.LoggerProvider) *Emitter {
return &Emitter{logger: lp.Logger(ScopeName)}
}
// Emit records ev. attrs is the same slice the caller gave its stdout record.
// It is a no-op, and cheap, until InitLogging installs a provider.
func (e *Emitter) Emit(ctx context.Context, ev Event, attrs []slog.Attr) {
// 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
// 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()
level := ev.Level()
if l := slog.Default(); l.Enabled(ctx, level) {
rec := slog.NewRecord(now, level, ev.Body, 0)
rec.AddAttrs(attrs...)
_ = l.Handler().Handle(ctx, rec)
}
e.emit(ctx, ev, now, attrs)
}
// emit writes the OTLP copy. It is a no-op, and cheap, until InitLogging
// installs a provider.
func (e *Emitter) emit(ctx context.Context, ev Event, t time.Time, attrs []slog.Attr) {
params := log.EnabledParameters{Severity: ev.Severity, EventName: ev.Name}
if !e.logger.Enabled(ctx, params) {
return
}
e.logger.Emit(ctx, BuildRecord(ev, time.Now(), attrs))
e.logger.Emit(ctx, BuildRecord(ev, t, attrs))
}
// The global provider delegates, so a Logger taken before InitLogging still
@@ -158,7 +192,7 @@ var defaultEmitter = sync.OnceValue(func() *Emitter {
return NewEmitter(global.GetLoggerProvider())
})
// Emit records ev through the process-wide provider.
func Emit(ctx context.Context, ev Event, attrs []slog.Attr) {
defaultEmitter().Emit(ctx, ev, attrs)
// Log records ev through the process-wide provider.
func Log(ctx context.Context, ev Event, attrs []slog.Attr) {
defaultEmitter().Log(ctx, ev, attrs)
}
+128 -13
View File
@@ -12,11 +12,12 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package actorevent_test
package actorevent
import (
"context"
"log/slog"
"maps"
"slices"
"testing"
"time"
@@ -25,7 +26,6 @@ import (
sdklog "go.opentelemetry.io/otel/sdk/log"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
"github.com/agent-substrate/substrate/internal/actorevent"
"github.com/agent-substrate/substrate/internal/ateattr"
"github.com/agent-substrate/substrate/internal/resources"
)
@@ -78,7 +78,7 @@ func TestBuildRecord(t *testing.T) {
tests := []struct {
name string
event actorevent.Event
event Event
attrs []slog.Attr
wantName string
wantBody string
@@ -87,10 +87,10 @@ func TestBuildRecord(t *testing.T) {
}{
{
name: "state changed",
event: actorevent.StateChanged,
event: StateChanged,
attrs: stateChangedAttrs(ateattr.ActorStateRunning),
wantName: "ate.actor.state_changed",
wantBody: actorevent.StateChangedBody,
wantBody: "Actor state changed",
wantSev: log.SeverityInfo,
wantVals: map[string]string{
string(ateattr.ActorStateKey): ateattr.ActorStateRunning,
@@ -100,10 +100,10 @@ func TestBuildRecord(t *testing.T) {
},
{
name: "deleted is a state, not a name of its own",
event: actorevent.StateChanged,
event: StateChanged,
attrs: stateChangedAttrs(ateattr.ActorStateDeleted),
wantName: "ate.actor.state_changed",
wantBody: actorevent.StateChangedBody,
wantBody: "Actor state changed",
wantSev: log.SeverityInfo,
wantVals: map[string]string{
string(ateattr.ActorStateKey): ateattr.ActorStateDeleted,
@@ -111,10 +111,10 @@ func TestBuildRecord(t *testing.T) {
},
{
name: "crashed",
event: actorevent.Crashed,
event: Crashed,
attrs: crashedAttrs(),
wantName: "ate.actor.crashed",
wantBody: actorevent.CrashedBody,
wantBody: "Actor crashed",
wantSev: log.SeverityError,
wantVals: map[string]string{
string(ateattr.ActorStateKey): ateattr.ActorStateCrashed,
@@ -128,7 +128,7 @@ func TestBuildRecord(t *testing.T) {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
rec := actorevent.BuildRecord(tt.event, now, tt.attrs)
rec := BuildRecord(tt.event, now, tt.attrs)
if got := rec.EventName(); got != tt.wantName {
t.Errorf("EventName() = %q, want %q", got, tt.wantName)
@@ -170,6 +170,37 @@ func TestBuildRecord(t *testing.T) {
}
}
func TestEventLevel(t *testing.T) {
t.Parallel()
tests := []struct {
name string
sev log.Severity
want slog.Level
}{
{"unset falls to the quietest level", log.SeverityUndefined, slog.LevelDebug},
{"trace", log.SeverityTrace, slog.LevelDebug},
{"debug", log.SeverityDebug, slog.LevelDebug},
{"info", log.SeverityInfo, slog.LevelInfo},
{"warn", log.SeverityWarn, slog.LevelWarn},
{"error", log.SeverityError, slog.LevelError},
{"fatal is as high as slog goes", log.SeverityFatal, slog.LevelError},
{"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},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
if got := (Event{Severity: tt.sev}).Level(); got != tt.want {
t.Errorf("Level() for severity %v = %v, want %v", tt.sev, got, tt.want)
}
})
}
}
// memExporter collects records in memory. Small enough to keep here rather than
// vendoring the SDK's test package.
type memExporter struct {
@@ -183,6 +214,90 @@ func (e *memExporter) Export(_ context.Context, records []sdklog.Record) error {
func (e *memExporter) Shutdown(context.Context) error { return nil }
func (e *memExporter) ForceFlush(context.Context) error { return nil }
// captureHandler keeps the stdout copy so a test can hold it beside the OTLP one.
type captureHandler struct {
records []slog.Record
}
func (h *captureHandler) Enabled(context.Context, slog.Level) bool { return true }
func (h *captureHandler) Handle(_ context.Context, r slog.Record) error {
h.records = append(h.records, r.Clone())
return nil
}
func (h *captureHandler) WithAttrs([]slog.Attr) slog.Handler { return h }
func (h *captureHandler) WithGroup(string) slog.Handler { return h }
// TestLogWritesBothCopies is the invariant this package exists for: one call,
// two copies, and no field either copy can hold alone.
//
// It swaps the slog default, so it cannot be parallel. Go finishes every
// non-parallel test before it resumes the parallel ones, so it does not race
// the rest of this file.
func TestLogWritesBothCopies(t *testing.T) {
tests := []struct {
name string
event Event
attrs []slog.Attr
}{
{"state changed", StateChanged, stateChangedAttrs(ateattr.ActorStateRunning)},
{"crashed", Crashed, crashedAttrs()},
}
for _, tt := range tests {
t.Run(tt.name, func(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()) })
NewEmitter(lp).Log(context.Background(), tt.event, tt.attrs)
if len(stdout.records) != 1 {
t.Fatalf("wrote %d stdout records, want 1", len(stdout.records))
}
if len(exp.records) != 1 {
t.Fatalf("exported %d OTLP records, want 1", len(exp.records))
}
stdoutRec, otlpRec := stdout.records[0], exp.records[0]
if stdoutRec.Message != tt.event.Body {
t.Errorf("stdout message = %q, want %q", stdoutRec.Message, tt.event.Body)
}
if got := otlpRec.Body().String(); got != tt.event.Body {
t.Errorf("OTLP body = %q, want %q", got, tt.event.Body)
}
if stdoutRec.Level != tt.event.Level() {
t.Errorf("stdout level = %v, want %v", stdoutRec.Level, tt.event.Level())
}
if got := otlpRec.Severity(); got != tt.event.Severity {
t.Errorf("OTLP severity = %v, want %v", got, tt.event.Severity)
}
// One time.Now() serves both, so a consumer can join them on it.
if !stdoutRec.Time.Equal(otlpRec.Timestamp()) {
t.Errorf("timestamps differ: stdout %v, OTLP %v", stdoutRec.Time, otlpRec.Timestamp())
}
stdoutAttrs := map[string]string{}
stdoutRec.Attrs(func(a slog.Attr) bool {
stdoutAttrs[a.Key] = a.Value.String()
return true
})
otlpAttrs := map[string]string{}
otlpRec.WalkAttributes(func(kv log.KeyValue) bool {
otlpAttrs[kv.Key] = kv.Value.String()
return true
})
if !maps.Equal(stdoutAttrs, otlpAttrs) {
t.Errorf("attributes differ: stdout %v, OTLP %v", stdoutAttrs, otlpAttrs)
}
})
}
}
func TestEmitCarriesTraceContext(t *testing.T) {
t.Parallel()
@@ -196,7 +311,7 @@ func TestEmitCarriesTraceContext(t *testing.T) {
ctx, span := tp.Tracer("test").Start(context.Background(), "test")
defer span.End()
actorevent.NewEmitter(lp).Emit(ctx, actorevent.StateChanged, stateChangedAttrs(ateattr.ActorStateRunning))
NewEmitter(lp).emit(ctx, StateChanged, time.Now(), stateChangedAttrs(ateattr.ActorStateRunning))
if len(exp.records) != 1 {
t.Fatalf("exported %d records, want 1", len(exp.records))
@@ -231,7 +346,7 @@ func TestEmitIsANoOpWithoutAProvider(t *testing.T) {
// The package default resolves the global provider, which no test installs.
// This asserts it does not panic rather than that it drops the record.
actorevent.Emit(context.Background(), actorevent.StateChanged, stateChangedAttrs(ateattr.ActorStateRunning))
defaultEmitter().emit(context.Background(), StateChanged, time.Now(), stateChangedAttrs(ateattr.ActorStateRunning))
}
func TestBuildRecordKeepsValueKinds(t *testing.T) {
@@ -255,7 +370,7 @@ func TestBuildRecordKeepsValueKinds(t *testing.T) {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
rec := actorevent.BuildRecord(actorevent.StateChanged, time.Now(), []slog.Attr{tt.attr})
rec := BuildRecord(StateChanged, time.Now(), []slog.Attr{tt.attr})
var got log.Value
rec.WalkAttributes(func(kv log.KeyValue) bool {
got = kv.Value
+91
View File
@@ -0,0 +1,91 @@
// 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 actorevent
import (
"os"
"slices"
"testing"
"sigs.k8s.io/yaml"
)
const registryPath = "../../docs/metrics/registry/events.yaml"
// registry is the part of a Weaver event group this test reads. Weaver has no
// field for a body or a severity, so those two are pinned in Go alone.
type registry struct {
Groups []struct {
Type string `json:"type"`
Name string `json:"name"`
Attributes []struct {
Ref string `json:"ref"`
RequirementLevel string `json:"requirement_level"`
} `json:"attributes"`
} `json:"groups"`
}
// TestEventsMatchTheRegistry holds the vocabulary and events.yaml in step both
// ways, so a new event has to be declared before it can ship.
//
// It covers every event group in the file. A component outside this package that
// starts emitting events needs a check of its own, and this one has to learn to
// skip what it does not own.
func TestEventsMatchTheRegistry(t *testing.T) {
t.Parallel()
raw, err := os.ReadFile(registryPath)
if err != nil {
t.Fatalf("read %s: %v", registryPath, err)
}
var reg registry
if err := yaml.Unmarshal(raw, &reg); err != nil {
t.Fatalf("parse %s: %v", registryPath, err)
}
declared := map[string][]string{}
for _, g := range reg.Groups {
if g.Type != "event" {
continue
}
keys := make([]string, 0, len(g.Attributes))
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)
}
keys = append(keys, a.Ref)
}
declared[g.Name] = keys
}
for _, ev := range events {
keys, 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))
if !slices.Equal(got, want) {
t.Errorf("%s: %s declares %v, the Event declares %v", ev.Name, registryPath, got, want)
}
}
for name := range declared {
t.Errorf("%s declares %s, which no Event in this package emits", registryPath, name)
}
}