mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
serverboot: give the ateoms a LoggerProvider over the relay (#1857)
Step 2 of #1748, plus the consumer rule from step 1b. `LoggingOptions` gains `ExporterConn` and `RelayCapable`, mirroring `TracingOptions`. Both ateoms initialize logging after metrics with the relay connection and defer a nil-guarded shutdown. `OTEL_LOGS_EXPORTER` still defaults to `none`, so no ateom emits a record; this is the provider the actor events land on when emission moves. A relay-capable log resource carries `ate.otlp.relay` like traces and metrics do. `newLoggerProvider` is split from `InitLogging` so the test asserts that on an emitted record's resource, the way `TestMeterProviderRelayAttribute` does for metrics. `docs/observability.md` states the consumer rule that replaced the relay denylist dropped in #1800: a lifecycle record is authoritative only under ateapi's resource, because the relay admits only ateom resources. Follow-up for the emission step, not here: the controller does not inject `OTEL_LOGS_EXPORTER` into worker pods, so no environment exercises log-over-relay yet.
This commit is contained in:
@@ -24,6 +24,7 @@ import (
|
||||
"go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc"
|
||||
"go.opentelemetry.io/otel/log/global"
|
||||
sdklog "go.opentelemetry.io/otel/sdk/log"
|
||||
"google.golang.org/grpc"
|
||||
)
|
||||
|
||||
const logsExporterEnv = "OTEL_LOGS_EXPORTER"
|
||||
@@ -85,6 +86,11 @@ type LoggingOptions struct {
|
||||
// Exporter is required. Build it with ResolveLogsExporter so
|
||||
// OTEL_LOGS_EXPORTER overrides the component default.
|
||||
Exporter LogsExporter
|
||||
// ExporterConn and RelayCapable are the logs counterpart of the same fields
|
||||
// on TracingOptions: ateom passes its relay connection and marks itself
|
||||
// relay-capable so the resource carries ate.otlp.relay.
|
||||
ExporterConn *grpc.ClientConn
|
||||
RelayCapable bool
|
||||
}
|
||||
|
||||
// InitLogging registers a global LoggerProvider for the records components emit
|
||||
@@ -111,25 +117,43 @@ func InitLogging(ctx context.Context, opts LoggingOptions) (*sdklog.LoggerProvid
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
res, err := newResource(ctx, opts.ServiceName)
|
||||
lp, err := newLoggerProvider(ctx, opts)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create logger resource: %w", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
exporter, err := otlploggrpc.New(ctx,
|
||||
// Matches the trace and metric exporters: GKE managed telemetry does not
|
||||
// support validating the TLS certs of the collector.
|
||||
otlploggrpc.WithInsecure(),
|
||||
)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create OTLP log exporter: %w", err)
|
||||
}
|
||||
|
||||
lp := sdklog.NewLoggerProvider(
|
||||
sdklog.WithResource(res),
|
||||
sdklog.WithProcessor(sdklog.NewBatchProcessor(exporter)),
|
||||
)
|
||||
global.SetLoggerProvider(lp)
|
||||
slog.InfoContext(ctx, "Logging initialized", slog.String("exporter", string(opts.Exporter)))
|
||||
return lp, nil
|
||||
}
|
||||
|
||||
// newLoggerProvider is InitLogging without the global registration. Tests add a
|
||||
// processor to read the resource off an emitted record; the provider does not
|
||||
// expose it otherwise.
|
||||
func newLoggerProvider(ctx context.Context, opts LoggingOptions, extra ...sdklog.Processor) (*sdklog.LoggerProvider, error) {
|
||||
res, err := newResource(ctx, opts.ServiceName, relayAttrs(opts.RelayCapable, opts.ExporterConn)...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create logger resource: %w", err)
|
||||
}
|
||||
|
||||
// Matches the trace and metric exporters: GKE managed telemetry does not
|
||||
// support validating the TLS certs of the collector.
|
||||
expOpts := []otlploggrpc.Option{otlploggrpc.WithInsecure()}
|
||||
if opts.ExporterConn != nil {
|
||||
// WithGRPCConn takes precedence over endpoint/credential options, so
|
||||
// WithInsecure above is inert on this path.
|
||||
expOpts = append(expOpts, otlploggrpc.WithGRPCConn(opts.ExporterConn))
|
||||
}
|
||||
exporter, err := otlploggrpc.New(ctx, expOpts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create OTLP log exporter: %w", err)
|
||||
}
|
||||
|
||||
popts := []sdklog.LoggerProviderOption{
|
||||
sdklog.WithResource(res),
|
||||
sdklog.WithProcessor(sdklog.NewBatchProcessor(exporter)),
|
||||
}
|
||||
for _, p := range extra {
|
||||
popts = append(popts, sdklog.WithProcessor(p))
|
||||
}
|
||||
return sdklog.NewLoggerProvider(popts...), nil
|
||||
}
|
||||
|
||||
@@ -18,7 +18,13 @@ import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
otellog "go.opentelemetry.io/otel/log"
|
||||
"go.opentelemetry.io/otel/log/global"
|
||||
sdklog "go.opentelemetry.io/otel/sdk/log"
|
||||
"go.opentelemetry.io/otel/sdk/resource"
|
||||
"google.golang.org/grpc"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
)
|
||||
|
||||
func TestResolveLogsExporter(t *testing.T) {
|
||||
@@ -140,6 +146,66 @@ func TestInitLoggingDisabled(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// The log path's half of the relay decision, asserted on the resource an
|
||||
// emitted record carries rather than on what relayAttrs returns, so dropping
|
||||
// the attrs in newLoggerProvider fails here. TestRelayAttrs pins the values.
|
||||
func TestLoggerProviderRelayAttribute(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
relayCapable bool
|
||||
conn *grpc.ClientConn
|
||||
want string // "" means absent
|
||||
}{
|
||||
{name: "relay", relayCapable: true, conn: lazyConn(t), want: "relay"},
|
||||
{name: "direct", relayCapable: true, want: "direct"},
|
||||
{name: "not relay capable", want: ""},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
got, ok := emittedResource(t, tc.relayCapable, tc.conn)[string(ateattr.OTLPRelayKey)]
|
||||
if ok != (tc.want != "") || got != tc.want {
|
||||
t.Errorf("%s = %q (present %t), want %q", string(ateattr.OTLPRelayKey), got, ok, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// emittedResource builds the provider the way InitLogging does, emits one
|
||||
// record through a capturing processor, and returns that record's resource.
|
||||
func emittedResource(t *testing.T, relayCapable bool, conn *grpc.ClientConn) map[string]string {
|
||||
t.Helper()
|
||||
proc := &captureProcessor{}
|
||||
lp, err := newLoggerProvider(context.Background(), LoggingOptions{
|
||||
ServiceName: "ateom-gvisor",
|
||||
Exporter: LogsExporterOTLP,
|
||||
ExporterConn: conn,
|
||||
RelayCapable: relayCapable,
|
||||
}, proc)
|
||||
if err != nil {
|
||||
t.Fatalf("newLoggerProvider: %v", err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
// A cancelled context skips the OTLP flush; no collector listens here.
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
_ = lp.Shutdown(ctx)
|
||||
})
|
||||
lp.Logger("test").Emit(context.Background(), otellog.Record{})
|
||||
if proc.res == nil {
|
||||
t.Fatal("no record reached the processor")
|
||||
}
|
||||
return resourceAttrs(proc.res)
|
||||
}
|
||||
|
||||
type captureProcessor struct{ res *resource.Resource }
|
||||
|
||||
func (c *captureProcessor) OnEmit(_ context.Context, r *sdklog.Record) error {
|
||||
c.res = r.Resource()
|
||||
return nil
|
||||
}
|
||||
func (c *captureProcessor) Enabled(context.Context, sdklog.EnabledParameters) bool { return true }
|
||||
func (c *captureProcessor) Shutdown(context.Context) error { return nil }
|
||||
func (c *captureProcessor) ForceFlush(context.Context) error { return nil }
|
||||
|
||||
func TestInitLoggingRequiresOptions(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user