From 25470fbf0b31f71ef1595d32fb96287b8de80f57 Mon Sep 17 00:00:00 2001 From: MartinForReal Date: Thu, 3 Sep 2026 14:54:21 +0800 Subject: [PATCH] feat(go): add gated OTLP metrics, reconcile events, and log bridge MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implements the Go observability gaps in #2148 as three independent, default-OFF additions so the Go runtime/controller match the OpenTelemetry capabilities of the Python runtime: 1. go/adk — GenAI token usage metrics. Configure a MeterProvider with a periodic OTLP metric exporter and emit gen_ai.client.token.usage (Int64 histogram, unit {token}) with the same attribute set the Python runtime records (gen_ai.token.type input/output, gen_ai.request.model, gen_ai.response.model, gen_ai.provider.name, gen_ai.agent.name). The A2A executor records one input + one output observation per LLM call (prompt tokens; output = candidate + reasoning tokens), skipping streamed Partial events. Gated behind OTEL_METRICS_ENABLED. 2. go/core — Kubernetes Events on reconcile. The RemoteMCPServer and MCPServer (MCPServerTool) discovery reconcilers emit Normal ToolsDiscovered and Warning ValidationFailed / Warning ReconcileFailed via an optional EventRecorder. Emission is skipped when no recorder is wired (the default), so behavior is byte-identical to today. 3. go/core — controller log -> OTLP bridge. InitLoggerProvider builds an OTLP LoggerProvider via autoexport (same resource as traces) and ControllerZapOpts tees the controller zap logger with an otelzap bridge core while preserving stdout. Gated behind OTEL_LOGGING_ENABLED; when unset both functions are no-ops. Promotes the OTLP metric and log SDK dependencies to direct requires (same versions already in the graph) and adds the otelzap bridge dependency. Refs: #2148 Signed-off-by: MartinForReal --- go/adk/cmd/main.go | 15 ++ go/adk/pkg/a2a/executor.go | 41 ++++ go/adk/pkg/a2a/executor_metrics_test.go | 67 +++++++ go/adk/pkg/config/config_usage.go | 11 ++ go/adk/pkg/telemetry/metrics.go | 152 +++++++++++++++ go/adk/pkg/telemetry/metrics_test.go | 176 ++++++++++++++++++ go/adk/pkg/telemetry/tracing.go | 44 ++++- go/core/cmd/controller-v2/main.go | 25 ++- .../controller/mcpserver/reconciler.go | 30 +++ .../controller/mcpserver/reconciler_test.go | 64 +++++++ .../controller/remotemcpserver/reconciler.go | 35 ++++ go/core/internal/telemetry/logging.go | 75 ++++++++ go/core/internal/telemetry/logging_test.go | 50 +++++ go/core/internal/telemetry/tracing.go | 41 ++-- go/go.mod | 11 +- go/go.sum | 4 + 16 files changed, 807 insertions(+), 34 deletions(-) create mode 100644 go/adk/pkg/a2a/executor_metrics_test.go create mode 100644 go/adk/pkg/telemetry/metrics.go create mode 100644 go/adk/pkg/telemetry/metrics_test.go create mode 100644 go/core/internal/telemetry/logging.go create mode 100644 go/core/internal/telemetry/logging_test.go diff --git a/go/adk/cmd/main.go b/go/adk/cmd/main.go index c165bf67f8..bedbfefbc4 100644 --- a/go/adk/cmd/main.go +++ b/go/adk/cmd/main.go @@ -20,6 +20,7 @@ import ( runnerpkg "github.com/kagent-dev/kagent/go/adk/pkg/runner" "github.com/kagent-dev/kagent/go/adk/pkg/session" "github.com/kagent-dev/kagent/go/adk/pkg/telemetry" + "github.com/kagent-dev/kagent/go/api/adk" "go.uber.org/zap" "go.uber.org/zap/zapcore" ) @@ -218,12 +219,15 @@ func main() { } stream := agentConfig.GetStream() + modelName, providerName := resolveModelLabels(agentConfig) executor := a2a.NewKAgentExecutor(a2a.KAgentExecutorConfig{ RunnerConfig: runnerConfig, SessionService: sessionService, Stream: stream, AppName: appName, Logger: logger, + ModelName: modelName, + ProviderName: providerName, }) // Build the agent card. @@ -262,6 +266,17 @@ func main() { } } +// resolveModelLabels derives the gen_ai.request.model / gen_ai.provider.name +// attributes for token-usage metrics from the agent config. The provider name +// is mapped to its OpenTelemetry GenAI semconv value; both are empty strings +// when no model is configured, so the metric simply omits those attributes. +func resolveModelLabels(agentConfig *adk.AgentConfig) (model, provider string) { + if agentConfig == nil || agentConfig.Model == nil { + return "", "" + } + return config.ModelName(agentConfig.Model), telemetry.SemconvProviderName(agentConfig.Model.GetType()) +} + func deriveAppName(kagentName, kagentNamespace string, agentCard *a2atype.AgentCard, logger logr.Logger) string { if kagentNamespace != "" && kagentName != "" { namespace := strings.ReplaceAll(kagentNamespace, "-", "_") diff --git a/go/adk/pkg/a2a/executor.go b/go/adk/pkg/a2a/executor.go index de46a9b181..8390ccad70 100644 --- a/go/adk/pkg/a2a/executor.go +++ b/go/adk/pkg/a2a/executor.go @@ -33,6 +33,11 @@ type KAgentExecutorConfig struct { Stream bool AppName string Logger logr.Logger + // ModelName and ProviderName label the GenAI token-usage metric + // (gen_ai.request.model / gen_ai.provider.name). Both may be empty, in + // which case the corresponding metric attributes are omitted. + ModelName string + ProviderName string } // KAgentExecutor keeps kagent's request/session glue around the upstream ADK @@ -56,6 +61,8 @@ func NewKAgentExecutor(cfg KAgentExecutorConfig) *KAgentExecutor { if cfg.SessionService != nil { runnerConfig.SessionService = cfg.SessionService } + modelName := cfg.ModelName + providerName := cfg.ProviderName builtin := adka2a.NewExecutor(adka2a.ExecutorConfig{ RunnerConfig: runnerConfig, RunConfig: runConfig, @@ -65,6 +72,7 @@ func NewKAgentExecutor(cfg KAgentExecutorConfig) *KAgentExecutor { if event.InvocationID != "" { trace.SpanFromContext(ctx).SetAttributes(attribute.String("gcp.vertex.agent.invocation_id", event.InvocationID)) } + recordTokenUsage(ctx, event, modelName, providerName, cfg.AppName) // Preserve the artifact's protocol type while giving current A2A clients a // common ordering key. A2A #2129 will replace this with native artifact // start/end generations and a task timeline. @@ -88,6 +96,39 @@ func NewKAgentExecutor(cfg KAgentExecutorConfig) *KAgentExecutor { } } +// recordTokenUsage records GenAI token usage for a single agent event on the +// gen_ai.client.token.usage histogram. Partial (streaming) events are skipped: +// a streamed LLM call emits many Partial chunks but usage is reported once on +// the aggregated non-partial event, so this records one input + one output +// observation per LLM call, not per stream chunk. Output combines candidate + +// reasoning tokens, matching the Python runtime's accounting. +func recordTokenUsage(ctx context.Context, adkEvent *adksession.Event, modelName, providerName, agentName string) { + if usage, ok := tokenUsageFromEvent(adkEvent, modelName, providerName, agentName); ok { + telemetry.RecordTokenUsage(ctx, usage) + } +} + +// tokenUsageFromEvent derives GenAI token usage from a single ADK session event. +// Partial (streaming) events are skipped: a streamed LLM call emits many Partial +// chunks but usage is reported once on the aggregated non-partial event, so this +// yields one input + one output observation per LLM call, not per stream chunk. +// Output combines candidate + reasoning tokens, matching the Python runtime's +// accounting. ok=false when there is nothing to record. +func tokenUsageFromEvent(adkEvent *adksession.Event, modelName, providerName, agentName string) (telemetry.TokenUsage, bool) { + usage := adkEvent.UsageMetadata + if usage == nil || adkEvent.Partial { + return telemetry.TokenUsage{}, false + } + return telemetry.TokenUsage{ + RequestModel: modelName, + ResponseModel: adkEvent.ModelVersion, + ProviderName: providerName, + AgentName: agentName, + InputTokens: int64(usage.PromptTokenCount), + OutputTokens: int64(usage.CandidatesTokenCount) + int64(usage.ThoughtsTokenCount), + }, true +} + // UserIDCallInterceptor returns an a2asrv.CallInterceptor that extracts the // x-user-id HTTP header from the incoming request metadata and sets it as the // authenticated user on the CallContext. diff --git a/go/adk/pkg/a2a/executor_metrics_test.go b/go/adk/pkg/a2a/executor_metrics_test.go new file mode 100644 index 0000000000..072c7d7a86 --- /dev/null +++ b/go/adk/pkg/a2a/executor_metrics_test.go @@ -0,0 +1,67 @@ +package a2a + +import ( + "testing" + + adkmodel "google.golang.org/adk/v2/model" + adksession "google.golang.org/adk/v2/session" + "google.golang.org/genai" +) + +// TestTokenUsageFromEvent_Values verifies the input/output token accounting: +// output combines candidate + reasoning tokens, and the model/provider/agent +// labels flow through. +func TestTokenUsageFromEvent_Values(t *testing.T) { + usage := &genai.GenerateContentResponseUsageMetadata{ + PromptTokenCount: 10, + CandidatesTokenCount: 5, + ThoughtsTokenCount: 3, + } + event := &adksession.Event{ + LLMResponse: adkmodel.LLMResponse{ + ModelVersion: "gemini-2.5-flash", + UsageMetadata: usage, + }, + } + + got, ok := tokenUsageFromEvent(event, "gemini-2.5-flash", "gcp.gemini", "my-agent") + if !ok { + t.Fatal("expected usage to be recorded") + } + if got.InputTokens != 10 { + t.Errorf("input = %d, want 10", got.InputTokens) + } + if got.OutputTokens != 8 { + t.Errorf("output = %d, want 8 (candidates 5 + thoughts 3)", got.OutputTokens) + } + if got.RequestModel != "gemini-2.5-flash" || got.ProviderName != "gcp.gemini" || got.AgentName != "my-agent" { + t.Errorf("labels not propagated: %+v", got) + } +} + +// TestTokenUsageFromEvent_SkipsPartialAndNil verifies streamed chunks that carry +// usage metadata are skipped (one observation per LLM call), and events without +// usage produce nothing. +func TestTokenUsageFromEvent_SkipsPartialAndNil(t *testing.T) { + usage := &genai.GenerateContentResponseUsageMetadata{PromptTokenCount: 10, CandidatesTokenCount: 5} + + if _, ok := tokenUsageFromEvent( + &adksession.Event{LLMResponse: adkmodel.LLMResponse{Partial: true, UsageMetadata: usage}}, + "m", "p", "a", + ); ok { + t.Error("partial event must not be recorded") + } + + if _, ok := tokenUsageFromEvent(&adksession.Event{}, "m", "p", "a"); ok { + t.Error("event without usage metadata must not be recorded") + } +} + +// TestRecordTokenUsage_NoopWithUninitializedRecorder verifies the executor +// recording path is a no-op before telemetry metrics are initialized (metrics +// disabled), so it cannot panic. +func TestRecordTokenUsage_NoopWithUninitializedRecorder(t *testing.T) { + usage := &genai.GenerateContentResponseUsageMetadata{PromptTokenCount: 10, CandidatesTokenCount: 5} + event := &adksession.Event{LLMResponse: adkmodel.LLMResponse{UsageMetadata: usage}} + recordTokenUsage(t.Context(), event, "gemini-2.5-flash", "gcp.gemini", "my-agent") +} diff --git a/go/adk/pkg/config/config_usage.go b/go/adk/pkg/config/config_usage.go index 01cea9ae1c..9cc44c3909 100644 --- a/go/adk/pkg/config/config_usage.go +++ b/go/adk/pkg/config/config_usage.go @@ -139,3 +139,14 @@ func getModelName(m adk.Model) string { return "unknown" } } + +// ModelName returns the configured model's identifier (e.g. "gpt-4o"), or "" +// when no model is configured. This labels the gen_ai.request.model token-usage +// metric attribute. +func ModelName(m adk.Model) string { + name := getModelName(m) + if name == "unknown" { + return "" + } + return name +} diff --git a/go/adk/pkg/telemetry/metrics.go b/go/adk/pkg/telemetry/metrics.go new file mode 100644 index 0000000000..9d483aa9ff --- /dev/null +++ b/go/adk/pkg/telemetry/metrics.go @@ -0,0 +1,152 @@ +package telemetry + +import ( + "context" + "net/url" + + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc" + "go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp" + "go.opentelemetry.io/otel/metric" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/resource" +) + +// GenAI token-usage instrumentation. The metric and attribute names follow +// the OpenTelemetry GenAI semantic conventions, and the attribute set matches +// the one upstream google-adk records for gen_ai.client.token.usage, so a +// single dashboard works across the Go and Python runtimes. +const ( + metricGenAIClientTokenUsage = "gen_ai.client.token.usage" + genAIMeterScope = "gcp.vertex.agent" + + attrGenAITokenType = "gen_ai.token.type" + attrGenAIRequestModel = "gen_ai.request.model" + attrGenAIResponseModel = "gen_ai.response.model" + attrGenAIProviderName = "gen_ai.provider.name" + attrGenAIAgentName = "gen_ai.agent.name" + + tokenTypeInput = "input" + tokenTypeOutput = "output" +) + +// tokenUsageHistogram records gen_ai.client.token.usage per LLM call. It is +// set when the meter provider is initialized (metrics enabled); otherwise it +// stays nil and recording is a cheap no-op, keeping the gate default-OFF. +var tokenUsageHistogram metric.Int64Histogram + +// TokenUsage carries the per-LLM-call labels and token counts for one +// recording on the gen_ai.client.token.usage histogram. RequestModel and +// ProviderName are resolved at startup from the agent config; ResponseModel +// falls back to RequestModel when a specific response model is unavailable. +type TokenUsage struct { + RequestModel string + ResponseModel string + ProviderName string + AgentName string + InputTokens int64 + OutputTokens int64 +} + +// RecordTokenUsage records input + output token counts on the +// gen_ai.client.token.usage histogram, one observation per token type. +// Zero/negative counts are skipped, and nothing is recorded when metrics are +// disabled or initialization failed (the instrument is nil). +func RecordTokenUsage(ctx context.Context, usage TokenUsage) { + h := tokenUsageHistogram + if h == nil { + return + } + responseModel := usage.ResponseModel + if responseModel == "" { + responseModel = usage.RequestModel + } + base := []attribute.KeyValue{ + attribute.String(attrGenAIRequestModel, usage.RequestModel), + attribute.String(attrGenAIResponseModel, responseModel), + attribute.String(attrGenAIProviderName, usage.ProviderName), + attribute.String(attrGenAIAgentName, usage.AgentName), + } + recordToken := func(tokenType string, count int64) { + if count <= 0 { + return + } + opts := append([]attribute.KeyValue{attribute.String(attrGenAITokenType, tokenType)}, base...) + h.Record(ctx, count, metric.WithAttributes(opts...)) + } + recordToken(tokenTypeInput, usage.InputTokens) + recordToken(tokenTypeOutput, usage.OutputTokens) +} + +// SemconvProviderName maps a kagent model type to its OpenTelemetry GenAI +// gen_ai.provider.name value. Unknown types pass through unchanged so custom +// providers keep their configured identity. +func SemconvProviderName(modelType string) string { + switch modelType { + case "openai": + return "openai" + case "azure_openai": + return "azure.ai.openai" + case "anthropic": + return "anthropic" + case "gemini", "gemini_vertex_ai", "gemini_anthropic": + return "gcp.gemini" + case "bedrock": + return "aws.bedrock" + default: + return modelType + } +} + +// newMeterProvider builds a MeterProvider with a periodic OTLP metric exporter, +// sharing the endpoint/protocol resolution used by traces and logs. +func newMeterProvider(ctx context.Context, res *resource.Resource) (*sdkmetric.MeterProvider, error) { + protocol := resolveOTLPProtocol("METRICS") + endpoint := resolveEndpoint("METRICS") + + var exporter sdkmetric.Exporter + var err error + switch protocol { + case "http/protobuf": + var opts []otlpmetrichttp.Option + if endpoint != "" { + opts = append(opts, otlpmetrichttp.WithEndpointURL(endpoint)) + } + exporter, err = otlpmetrichttp.New(ctx, opts...) + default: + var opts []otlpmetricgrpc.Option + if endpoint != "" { + if u, parseErr := url.Parse(endpoint); parseErr == nil && u.Scheme != "" && u.Host != "" { + opts = append(opts, otlpmetricgrpc.WithEndpointURL(u.String())) + } else { + opts = append(opts, otlpmetricgrpc.WithEndpoint(endpoint)) + } + } + exporter, err = otlpmetricgrpc.New(ctx, opts...) + } + if err != nil { + return nil, err + } + + return sdkmetric.NewMeterProvider( + sdkmetric.WithReader(sdkmetric.NewPeriodicReader(exporter)), + sdkmetric.WithResource(res), + ), nil +} + +// initTokenUsageRecorder binds the gen_ai.client.token.usage histogram to the +// given meter scope. It is called after setting the global meter provider. +func initTokenUsageRecorder(mp *sdkmetric.MeterProvider) { + meter := mp.Meter(genAIMeterScope) + var err error + tokenUsageHistogram, err = meter.Int64Histogram( + metricGenAIClientTokenUsage, + metric.WithUnit("{token}"), + metric.WithDescription("Number of input and output tokens used by GenAI requests."), + ) + if err != nil { + otel.Handle(err) + tokenUsageHistogram = nil + } +} diff --git a/go/adk/pkg/telemetry/metrics_test.go b/go/adk/pkg/telemetry/metrics_test.go new file mode 100644 index 0000000000..0fb37fbd0d --- /dev/null +++ b/go/adk/pkg/telemetry/metrics_test.go @@ -0,0 +1,176 @@ +package telemetry + +import ( + "context" + "testing" + + "go.opentelemetry.io/otel/attribute" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" +) + +// withTokenRecorder installs a manual meter reader so tests can inspect the +// gen_ai.client.token.usage histogram after RecordTokenUsage calls. +func withTokenRecorder(t *testing.T) *sdkmetric.ManualReader { + t.Helper() + prev := tokenUsageHistogram + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + initTokenUsageRecorder(mp) + t.Cleanup(func() { + tokenUsageHistogram = prev + _ = mp.Shutdown(context.Background()) + }) + return reader +} + +// tokenUsagePoint returns the unique data point on gen_ai.client.token.usage +// whose attributes carry the given token type, failing the test otherwise. +func tokenUsagePoint(t *testing.T, reader *sdkmetric.ManualReader, tokenType string) metricdata.HistogramDataPoint[int64] { + t.Helper() + var rm metricdata.ResourceMetrics + if err := reader.Collect(context.Background(), &rm); err != nil { + t.Fatalf("collect: %v", err) + } + for _, sm := range rm.ScopeMetrics { + for _, m := range sm.Metrics { + if m.Name != metricGenAIClientTokenUsage { + continue + } + hist, ok := m.Data.(metricdata.Histogram[int64]) + if !ok { + t.Fatalf("expected histogram data, got %T", m.Data) + } + for _, dp := range hist.DataPoints { + if v, _ := dp.Attributes.Value(attribute.Key(attrGenAITokenType)); v.AsString() == tokenType { + return dp + } + } + } + } + t.Fatalf("no data point with gen_ai.token.type=%q on %s", tokenType, metricGenAIClientTokenUsage) + return metricdata.HistogramDataPoint[int64]{} +} + +func attrString(t *testing.T, dp metricdata.HistogramDataPoint[int64], key string) string { + t.Helper() + v, ok := dp.Attributes.Value(attribute.Key(key)) + if !ok { + t.Fatalf("attribute %q not set on data point", key) + } + s := v.AsString() + if s == "" { + t.Fatalf("attribute %q empty on data point", key) + } + return s +} + +func TestRecordTokenUsage_RecordsInputAndOutput(t *testing.T) { + reader := withTokenRecorder(t) + + RecordTokenUsage(context.Background(), TokenUsage{ + RequestModel: "gpt-4o", + ResponseModel: "gpt-4o-2024-05-13", + ProviderName: "openai", + AgentName: "my-agent", + InputTokens: 10, + OutputTokens: 5, + }) + + input := tokenUsagePoint(t, reader, tokenTypeInput) + if got := input.Count; got != 1 { + t.Fatalf("input count = %d, want 1", got) + } + if got := input.Sum; got != 10 { + t.Fatalf("input sum = %d, want 10", got) + } + if got := attrString(t, input, attrGenAIRequestModel); got != "gpt-4o" { + t.Errorf("input request model = %q", got) + } + if got := attrString(t, input, attrGenAIResponseModel); got != "gpt-4o-2024-05-13" { + t.Errorf("input response model = %q", got) + } + if got := attrString(t, input, attrGenAIProviderName); got != "openai" { + t.Errorf("input provider = %q", got) + } + if got := attrString(t, input, attrGenAIAgentName); got != "my-agent" { + t.Errorf("input agent = %q", got) + } + + output := tokenUsagePoint(t, reader, tokenTypeOutput) + if got := output.Sum; got != 5 { + t.Fatalf("output sum = %d, want 5", got) + } + if got := output.Count; got != 1 { + t.Fatalf("output count = %d, want 1", got) + } +} + +func TestRecordTokenUsage_ResponseModelFallsBackToRequest(t *testing.T) { + reader := withTokenRecorder(t) + + RecordTokenUsage(context.Background(), TokenUsage{ + RequestModel: "gemini-2.5-flash", + ProviderName: "gcp.gemini", + AgentName: "a", + InputTokens: 3, + }) + + pt := tokenUsagePoint(t, reader, tokenTypeInput) + if got := attrString(t, pt, attrGenAIResponseModel); got != "gemini-2.5-flash" { + t.Errorf("response model = %q, want fallback to request model", got) + } +} + +func TestRecordTokenUsage_SkipsZeroAndNegative(t *testing.T) { + reader := withTokenRecorder(t) + + // All-zero, and negative/zero mixes, must not create data points. + RecordTokenUsage(context.Background(), TokenUsage{ + RequestModel: "gpt-4o", ProviderName: "openai", AgentName: "my-agent", + }) + RecordTokenUsage(context.Background(), TokenUsage{ + RequestModel: "gpt-4o", ProviderName: "openai", InputTokens: -1, + }) + + var rm metricdata.ResourceMetrics + if err := reader.Collect(context.Background(), &rm); err != nil { + t.Fatalf("collect: %v", err) + } + for _, sm := range rm.ScopeMetrics { + for _, m := range sm.Metrics { + if m.Name == metricGenAIClientTokenUsage { + t.Fatalf("expected no data on %s for zero/negative counts", metricGenAIClientTokenUsage) + } + } + } +} + +func TestRecordTokenUsage_NoopWhenNotInitialized(t *testing.T) { + prev := tokenUsageHistogram + tokenUsageHistogram = nil + t.Cleanup(func() { tokenUsageHistogram = prev }) + + RecordTokenUsage(context.Background(), TokenUsage{ + RequestModel: "gpt-4o", ProviderName: "openai", InputTokens: 10, OutputTokens: 5, + }) // must not panic when metrics are disabled +} + +func TestSemconvProviderName(t *testing.T) { + cases := map[string]string{ + "openai": "openai", + "azure_openai": "azure.ai.openai", + "anthropic": "anthropic", + "gemini": "gcp.gemini", + "gemini_vertex_ai": "gcp.gemini", + "gemini_anthropic": "gcp.gemini", + "bedrock": "aws.bedrock", + "ollama": "ollama", + "some-custom": "some-custom", + } + for in, want := range cases { + if got := SemconvProviderName(in); got != want { + t.Errorf("SemconvProviderName(%q) = %q, want %q", in, got, want) + } + } +} diff --git a/go/adk/pkg/telemetry/tracing.go b/go/adk/pkg/telemetry/tracing.go index 3a1d22a936..7aa1c1b64e 100644 --- a/go/adk/pkg/telemetry/tracing.go +++ b/go/adk/pkg/telemetry/tracing.go @@ -49,14 +49,19 @@ func StartInvocationSpan(ctx context.Context) (context.Context, trace.Span) { // 3s and is configurable via KAGENT_TRACE_FLUSH_TIMEOUT_MS. func ForceFlush(ctx context.Context) { type flusher interface{ ForceFlush(context.Context) error } - fp, ok := otel.GetTracerProvider().(flusher) - if !ok { - return - } flushCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), flushTimeout()) defer cancel() - if err := fp.ForceFlush(flushCtx); err != nil { - otel.Handle(err) + if fp, ok := otel.GetTracerProvider().(flusher); ok { + if err := fp.ForceFlush(flushCtx); err != nil { + otel.Handle(err) + } + } + // A periodic metric reader may never fire before the actor is suspended, + // so drain it alongside spans (see newMeterProvider). + if mp, ok := otel.GetMeterProvider().(flusher); ok { + if err := mp.ForceFlush(flushCtx); err != nil { + otel.Handle(err) + } } } @@ -96,6 +101,7 @@ func Init(ctx context.Context, serviceName string, serviceNamespace string) (shu tracingEnabled := strings.EqualFold(strings.TrimSpace(os.Getenv("OTEL_TRACING_ENABLED")), "true") loggingEnabled := strings.EqualFold(strings.TrimSpace(os.Getenv("OTEL_LOGGING_ENABLED")), "true") + metricsEnabled := strings.EqualFold(strings.TrimSpace(os.Getenv("OTEL_METRICS_ENABLED")), "true") otelOpts := []adktelemetry.Option{adktelemetry.WithResource(telemetryResource)} if tracingEnabled { tracerProvider, tpErr := newTracerProvider(ctx, telemetryResource) @@ -116,19 +122,39 @@ func Init(ctx context.Context, serviceName string, serviceNamespace string) (shu if telErr != nil { return nil, true, telErr } - telemetryProviders.SetGlobalOtelProviders() + + var meterShutdown func(context.Context) error + if metricsEnabled { + meterProvider, mpErr := newMeterProvider(ctx, telemetryResource) + if mpErr != nil { + return nil, true, mpErr + } + otel.SetMeterProvider(meterProvider) + initTokenUsageRecorder(meterProvider) + meterShutdown = meterProvider.Shutdown + } + otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator( propagation.TraceContext{}, propagation.Baggage{}, )) - return telemetryProviders.Shutdown, true, nil + if meterShutdown == nil { + return telemetryProviders.Shutdown, true, nil + } + return func(shutdownCtx context.Context) error { + if err := telemetryProviders.Shutdown(shutdownCtx); err != nil { + return err + } + return meterShutdown(shutdownCtx) + }, true, nil } func isTelemetryEnabled() bool { return strings.EqualFold(strings.TrimSpace(os.Getenv("OTEL_TRACING_ENABLED")), "true") || - strings.EqualFold(strings.TrimSpace(os.Getenv("OTEL_LOGGING_ENABLED")), "true") + strings.EqualFold(strings.TrimSpace(os.Getenv("OTEL_LOGGING_ENABLED")), "true") || + strings.EqualFold(strings.TrimSpace(os.Getenv("OTEL_METRICS_ENABLED")), "true") } // resolveOTLPProtocol returns the OTLP protocol for the given signal, diff --git a/go/core/cmd/controller-v2/main.go b/go/core/cmd/controller-v2/main.go index 2c78c2c561..26163c7bff 100644 --- a/go/core/cmd/controller-v2/main.go +++ b/go/core/cmd/controller-v2/main.go @@ -76,7 +76,6 @@ func main() { log.Fatalf("parse ZAP_LOG_LEVEL: %v", err) } } - ctrl.SetLogger(zap.New(zap.Level(logLevel))) // otelgrpc snapshots the global TracerProvider and propagator when its handler // is constructed, so tracing has to be registered before any server is built. shutdownTracing, err := telemetry.InitTracerProvider(ctx, version.Version) @@ -91,6 +90,24 @@ func main() { } }() + // Initialize the OTLP logger provider before building the controller logger + // so the otelzap bridge below binds to the configured global provider. + shutdownLogging, err := telemetry.InitLoggerProvider(ctx, version.Version) + if err != nil { + log.Fatalf("initialize logging: %v", err) + } + defer func() { + shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := shutdownLogging(shutdownCtx); err != nil { + log.Printf("shutdown logging: %v", err) + } + }() + + // OTEL_LOGGING_ENABLED additively tees the controller's logs over OTLP; + // otherwise the tee option is a no-op and this matches upstream logging. + ctrl.SetLogger(zap.New(append([]zap.Opts{zap.Level(logLevel)}, telemetry.ControllerZapOpts()...)...)) + dbURL, err := database.ResolveURL(env("POSTGRES_DATABASE_URL", "postgres://postgres:kagent@kagent-postgresql.kagent.svc.cluster.local:5432/postgres"), os.Getenv("POSTGRES_DATABASE_URL_FILE")) if err != nil { log.Fatal(err) @@ -163,11 +180,13 @@ func main() { log.Fatalf("add reconciler to controller manager: %v", err) } mcpClient := toolservice.NewRuntimeMCPClient(manager.GetClient()) - remoteMCPDiscovery := remotemcpcontroller.New(manager.GetClient(), mcpClient, store) + remoteMCPDiscovery := remotemcpcontroller.New(manager.GetClient(), mcpClient, store). + WithRecorder(manager.GetEventRecorder("remotemcpserver")) if err := remoteMCPDiscovery.SetupWithManager(manager); err != nil { log.Fatalf("set up RemoteMCPServer discovery: %v", err) } - mcpServerDiscovery := mcpservercontroller.New(manager.GetClient(), mcpClient, store) + mcpServerDiscovery := mcpservercontroller.New(manager.GetClient(), mcpClient, store). + WithRecorder(manager.GetEventRecorder("mcpserver-catalog")) if err := mcpServerDiscovery.SetupWithManager(manager); err != nil { log.Fatalf("set up MCPServer discovery: %v", err) } diff --git a/go/core/internal/controller/mcpserver/reconciler.go b/go/core/internal/controller/mcpserver/reconciler.go index d91bfe02f8..80e8ffe6aa 100644 --- a/go/core/internal/controller/mcpserver/reconciler.go +++ b/go/core/internal/controller/mcpserver/reconciler.go @@ -33,6 +33,7 @@ import ( apiMeta "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/tools/events" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/controller" @@ -66,12 +67,23 @@ type Reconciler struct { client client.Client discoverer ToolDiscoverer catalog CatalogStore + // Recorder emits Kubernetes Events on reconcile transitions. Optional; + // event emission is skipped when nil. + Recorder events.EventRecorder } func New(client client.Client, discoverer ToolDiscoverer, catalog CatalogStore) *Reconciler { return &Reconciler{client: client, discoverer: discoverer, catalog: catalog} } +// WithRecorder wires an optional EventRecorder used to surface reconcile +// outcomes via kubectl describe. It returns the receiver so it composes with +// New. +func (r *Reconciler) WithRecorder(recorder events.EventRecorder) *Reconciler { + r.Recorder = recorder + return r +} + func (r *Reconciler) SetupWithManager(manager ctrl.Manager) error { installed, err := controllerEnabled(manager.GetRESTMapper()) if err != nil { @@ -120,6 +132,8 @@ func (r *Reconciler) Reconcile(ctx context.Context, request reconcile.Request) ( Ref: request.NamespacedName, GroupKind: mcpServerGroupKind, }) if err != nil { + r.recordEvent(ctx, server, "Warning", "ReconcileFailed", "Reconcile", + "failed to discover MCPServer tools: %v", err) catalogErr := r.updateCatalog(ctx, server, nil, false) return reconcile.Result{}, errors.Join( fmt.Errorf("discover MCPServer tools: %w", err), @@ -129,6 +143,8 @@ func (r *Reconciler) Reconcile(ctx context.Context, request reconcile.Request) ( discovered, err := toolcatalog.NormalizeTools(tools) if err != nil { + r.recordEvent(ctx, server, "Warning", "ValidationFailed", "Reconcile", + "invalid MCPServer tool discovery: %v", err) catalogErr := r.updateCatalog(ctx, server, nil, false) return reconcile.Result{}, errors.Join( err, @@ -136,11 +152,25 @@ func (r *Reconciler) Reconcile(ctx context.Context, request reconcile.Request) ( ) } if err := r.updateCatalog(ctx, server, discovered, true); err != nil { + r.recordEvent(ctx, server, "Warning", "ReconcileFailed", "Reconcile", + "failed to update MCPServer tool catalog: %v", err) return reconcile.Result{}, fmt.Errorf("update MCPServer tool catalog: %w", err) } + if len(discovered) > 0 { + r.recordEvent(ctx, server, "Normal", "ToolsDiscovered", "Discover", "Discovered %d MCP tools", len(discovered)) + } return reconcile.Result{RequeueAfter: refreshInterval}, nil } +// recordEvent emits a Kubernetes Event against the reconciled object. It is a +// no-op when no Recorder is wired on this reconciler. +func (r *Reconciler) recordEvent(ctx context.Context, object client.Object, eventType, reason, action, messageFmt string, args ...any) { + if r.Recorder == nil { + return + } + r.Recorder.Eventf(object, nil, eventType, reason, action, messageFmt, args...) +} + func isReady(server *kmcp.MCPServer) bool { condition := apiMeta.FindStatusCondition(server.Status.Conditions, string(kmcp.MCPServerConditionReady)) return condition != nil && condition.Status == metav1.ConditionTrue && condition.ObservedGeneration == server.Generation diff --git a/go/core/internal/controller/mcpserver/reconciler_test.go b/go/core/internal/controller/mcpserver/reconciler_test.go index e69dbfb98b..3f6ed9b19c 100644 --- a/go/core/internal/controller/mcpserver/reconciler_test.go +++ b/go/core/internal/controller/mcpserver/reconciler_test.go @@ -19,6 +19,7 @@ package mcpserver import ( "context" "errors" + "strings" "testing" "time" @@ -31,6 +32,7 @@ import ( "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/tools/events" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" "sigs.k8s.io/controller-runtime/pkg/client/fake" @@ -209,3 +211,65 @@ func testClient(t *testing.T, objects ...client.Object) client.Client { func apiMetaTestMapper() *apiMeta.DefaultRESTMapper { return apiMeta.NewDefaultRESTMapper([]schema.GroupVersion{kmcp.GroupVersion}) } + +// TestReconcileEmitsToolsDiscovered verifies a Normal ToolsDiscovered event is +// emitted when a ready MCPServer is successfully discovered. +func TestReconcileEmitsToolsDiscovered(t *testing.T) { + server := readyServer() + discoverer := &fakeDiscoverer{tools: []toolservice.MCPAppTool{{Name: "zeta"}}} + recorder := events.NewFakeRecorder(1) + + result, err := New(testClient(t, server), discoverer, &fakeCatalog{}). + WithRecorder(recorder).Reconcile(t.Context(), ctrl.Request{ + NamespacedName: client.ObjectKeyFromObject(server), + }) + if err != nil { + t.Fatalf("Reconcile() error = %v", err) + } + if result.RequeueAfter != 5*time.Minute { + t.Fatalf("Reconcile() requeue = %s, want 5m", result.RequeueAfter) + } + + select { + case ev := <-recorder.Events: + if !strings.Contains(ev, "Normal ToolsDiscovered ") { + t.Fatalf("unexpected event: %q", ev) + } + default: + t.Fatal("expected a ToolsDiscovered event, got none") + } +} + +// TestReconcileEmitsValidationFailed verifies a malformed discovery triggers a +// Warning ValidationFailed event. +func TestReconcileEmitsValidationFailed(t *testing.T) { + server := readyServer() + discoverer := &fakeDiscoverer{tools: []toolservice.MCPAppTool{{Name: " "}}} + recorder := events.NewFakeRecorder(1) + + if _, err := New(testClient(t, server), discoverer, &fakeCatalog{}). + WithRecorder(recorder).Reconcile(t.Context(), ctrl.Request{ + NamespacedName: client.ObjectKeyFromObject(server), + }); err == nil { + t.Fatal("expected validation error") + } + + ev := <-recorder.Events + if !strings.Contains(ev, "Warning ValidationFailed ") { + t.Fatalf("unexpected event: %q", ev) + } +} + +// TestReconcileNoEventsWithoutRecorder verifies event emission is skipped when +// no recorder is wired, keeping zero behavioral change by default. +func TestReconcileNoEventsWithoutRecorder(t *testing.T) { + server := readyServer() + discoverer := &fakeDiscoverer{tools: []toolservice.MCPAppTool{{Name: "zeta"}}} + + if _, err := New(testClient(t, server), discoverer, &fakeCatalog{}). + Reconcile(t.Context(), ctrl.Request{ + NamespacedName: client.ObjectKeyFromObject(server), + }); err != nil { + t.Fatalf("Reconcile() error = %v", err) + } +} diff --git a/go/core/internal/controller/remotemcpserver/reconciler.go b/go/core/internal/controller/remotemcpserver/reconciler.go index 10215d9b3d..f61e73c7fc 100644 --- a/go/core/internal/controller/remotemcpserver/reconciler.go +++ b/go/core/internal/controller/remotemcpserver/reconciler.go @@ -34,6 +34,7 @@ import ( apiMeta "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/tools/events" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/builder" "sigs.k8s.io/controller-runtime/pkg/client" @@ -68,12 +69,23 @@ type Reconciler struct { client client.Client discoverer ToolDiscoverer catalog CatalogStore + // Recorder emits Kubernetes Events on reconcile transitions. Optional; + // event emission is skipped when nil. + Recorder events.EventRecorder } func New(client client.Client, discoverer ToolDiscoverer, catalog CatalogStore) *Reconciler { return &Reconciler{client: client, discoverer: discoverer, catalog: catalog} } +// WithRecorder wires an optional EventRecorder used to surface reconcile +// outcomes via kubectl describe. It returns the receiver so it composes with +// New. +func (r *Reconciler) WithRecorder(recorder events.EventRecorder) *Reconciler { + r.Recorder = recorder + return r +} + func (r *Reconciler) SetupWithManager(manager ctrl.Manager) error { return ctrl.NewControllerManagedBy(manager). For(&v1alpha3.RemoteMCPServer{}, builder.WithPredicates(predicate.GenerationChangedPredicate{})). @@ -95,6 +107,8 @@ func (r *Reconciler) Reconcile(ctx context.Context, request reconcile.Request) ( Ref: request.NamespacedName, GroupKind: remoteGroupKind, }) if err != nil { + r.recordEvent(ctx, server, "Warning", "ReconcileFailed", "Reconcile", + "failed to discover RemoteMCPServer tools: %v", err) statusErr := r.updateStatus(ctx, server, nil, metav1.ConditionFalse, "DiscoveryFailed", err.Error()) catalogErr := r.updateCatalog(ctx, server, nil, false) return reconcile.Result{}, errors.Join( @@ -106,6 +120,8 @@ func (r *Reconciler) Reconcile(ctx context.Context, request reconcile.Request) ( discovered, err := toolcatalog.NormalizeTools(tools) if err != nil { + r.recordEvent(ctx, server, "Warning", "ValidationFailed", "Reconcile", + "invalid RemoteMCPServer tool discovery: %v", err) statusErr := r.updateStatus(ctx, server, nil, metav1.ConditionFalse, "InvalidDiscovery", err.Error()) catalogErr := r.updateCatalog(ctx, server, nil, false) return reconcile.Result{}, errors.Join( @@ -114,16 +130,35 @@ func (r *Reconciler) Reconcile(ctx context.Context, request reconcile.Request) ( wrapError("clear invalid RemoteMCPServer tool catalog", catalogErr), ) } + // Only surface a Normal event on the transition to a successful discovery, + // so the periodic refresh timer does not re-emit it forever. + firstDiscovery := len(server.Status.DiscoveredTools) == 0 && len(discovered) > 0 message := fmt.Sprintf("Discovered %d MCP tools", len(discovered)) if err := r.updateStatus(ctx, server, discovered, metav1.ConditionTrue, "DiscoverySucceeded", message); err != nil { + r.recordEvent(ctx, server, "Warning", "ReconcileFailed", "Reconcile", + "failed to update RemoteMCPServer discovery status: %v", err) return reconcile.Result{}, fmt.Errorf("update RemoteMCPServer discovery status: %w", err) } if err := r.updateCatalog(ctx, server, discovered, true); err != nil { + r.recordEvent(ctx, server, "Warning", "ReconcileFailed", "Reconcile", + "failed to update RemoteMCPServer tool catalog: %v", err) return reconcile.Result{}, fmt.Errorf("update RemoteMCPServer tool catalog: %w", err) } + if firstDiscovery { + r.recordEvent(ctx, server, "Normal", "ToolsDiscovered", "Discover", "Discovered %d MCP tools", len(discovered)) + } return reconcile.Result{RequeueAfter: refreshInterval}, nil } +// recordEvent emits a Kubernetes Event against the reconciled object. It is a +// no-op when no Recorder is wired on this reconciler. +func (r *Reconciler) recordEvent(ctx context.Context, object client.Object, eventType, reason, action, messageFmt string, args ...any) { + if r.Recorder == nil { + return + } + r.Recorder.Eventf(object, nil, eventType, reason, action, messageFmt, args...) +} + func (r *Reconciler) updateCatalog(ctx context.Context, server *v1alpha3.RemoteMCPServer, tools []*v1alpha3.MCPTool, connected bool) error { name := client.ObjectKeyFromObject(server).String() var lastConnected *time.Time diff --git a/go/core/internal/telemetry/logging.go b/go/core/internal/telemetry/logging.go new file mode 100644 index 0000000000..f90c351f4b --- /dev/null +++ b/go/core/internal/telemetry/logging.go @@ -0,0 +1,75 @@ +package telemetry + +import ( + "context" + "fmt" + + otelzap "go.opentelemetry.io/contrib/bridges/otelzap" + "go.opentelemetry.io/contrib/exporters/autoexport" + logglobal "go.opentelemetry.io/otel/log/global" + sdklog "go.opentelemetry.io/otel/sdk/log" + "go.uber.org/zap" + "go.uber.org/zap/zapcore" + crzap "sigs.k8s.io/controller-runtime/pkg/log/zap" + + "github.com/kagent-dev/kagent/go/core/pkg/env" +) + +// loggerBridgeName is the instrumentation scope under which the controller's +// own logs are emitted by the otelzap bridge, matching the module that owns +// the controller. +const loggerBridgeName = "github.com/kagent-dev/kagent/go/core" + +// InitLoggerProvider configures an OTLP LoggerProvider and registers it as the +// global OTel logger provider. The exporter type and endpoint are read from the +// standard OTEL environment variables via autoexport, mirroring +// InitTracerProvider, so log and trace pipelines share the same OTLP config. +// The returned shutdown function must be called on process exit to flush +// in-flight log records. When OTEL_LOGGING_ENABLED is unset (the default) the +// pipeline is not created and a no-op shutdown is returned. +func InitLoggerProvider(ctx context.Context, serviceVersion string) (func(context.Context) error, error) { + if !env.OtelLoggingEnabled.Get() { + return func(context.Context) error { return nil }, nil + } + + exporter, err := autoexport.NewLogExporter(ctx) + if err != nil { + return nil, fmt.Errorf("create log exporter: %w", err) + } + + res, err := newTelemetryResource(ctx, serviceVersion) + if err != nil { + return nil, err + } + + lp := sdklog.NewLoggerProvider( + sdklog.WithProcessor(sdklog.NewBatchProcessor(exporter)), + sdklog.WithResource(res), + ) + + logglobal.SetLoggerProvider(lp) + + return lp.Shutdown, nil +} + +// ControllerZapOpts returns controller-runtime zap options. When +// OTEL_LOGGING_ENABLED is set it additively tees the controller's stdout zap +// core with an otelzap bridge core, routing the controller's own logs through +// the global OTLP LoggerProvider while preserving stdout logging. When disabled +// it returns no options, leaving the logger byte-identical to upstream. +// +// InitLoggerProvider must be called first so the bridge core binds to the +// configured global LoggerProvider. +func ControllerZapOpts() []crzap.Opts { + if !env.OtelLoggingEnabled.Get() { + return nil + } + bridgeCore := otelzap.NewCore(loggerBridgeName, + otelzap.WithLoggerProvider(logglobal.GetLoggerProvider()), + ) + return []crzap.Opts{ + crzap.RawZapOpts(zap.WrapCore(func(core zapcore.Core) zapcore.Core { + return zapcore.NewTee(core, bridgeCore) + })), + } +} diff --git a/go/core/internal/telemetry/logging_test.go b/go/core/internal/telemetry/logging_test.go new file mode 100644 index 0000000000..45ba624684 --- /dev/null +++ b/go/core/internal/telemetry/logging_test.go @@ -0,0 +1,50 @@ +package telemetry + +import ( + "context" + "testing" + + crzap "sigs.k8s.io/controller-runtime/pkg/log/zap" +) + +// TestControllerZapOpts_DisabledByDefault verifies no zap options are added when +// OTEL_LOGGING_ENABLED is off, keeping the controller logger byte-identical. +func TestControllerZapOpts_DisabledByDefault(t *testing.T) { + t.Setenv("OTEL_LOGGING_ENABLED", "false") + if opts := ControllerZapOpts(); opts != nil { + t.Fatalf("expected nil opts when logging disabled, got %d", len(opts)) + } +} + +// TestControllerZapOpts_EnabledTeesLogger verifies that when +// OTEL_LOGGING_ENABLED is set an otelzap bridge option is returned and the +// resulting logger is usable (the bridge tees onto the stdout core without +// panicking even with the default no-op global LoggerProvider). +func TestControllerZapOpts_EnabledTeesLogger(t *testing.T) { + t.Setenv("OTEL_LOGGING_ENABLED", "true") + + opts := ControllerZapOpts() + if len(opts) != 1 { + t.Fatalf("expected 1 zap opt when logging enabled, got %d", len(opts)) + } + + logger := crzap.New(append([]crzap.Opts{crzap.UseDevMode(true)}, opts...)...) + logger.Info("tee smoke test") +} + +// TestInitLoggerProvider_DisabledNoop verifies the returned shutdown is a safe +// no-op when logging is disabled. +func TestInitLoggerProvider_DisabledNoop(t *testing.T) { + t.Setenv("OTEL_LOGGING_ENABLED", "false") + + shutdown, err := InitLoggerProvider(context.Background(), "test") + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if shutdown == nil { + t.Fatal("expected non-nil shutdown") + } + if err := shutdown(context.Background()); err != nil { + t.Fatalf("no-op shutdown returned error: %v", err) + } +} diff --git a/go/core/internal/telemetry/tracing.go b/go/core/internal/telemetry/tracing.go index 71528adf6e..cdc8be6928 100644 --- a/go/core/internal/telemetry/tracing.go +++ b/go/core/internal/telemetry/tracing.go @@ -38,6 +38,29 @@ func InitTracerProvider(ctx context.Context, serviceVersion string) (func(contex return nil, fmt.Errorf("create span exporter: %w", err) } + res, err := newTelemetryResource(ctx, serviceVersion) + if err != nil { + return nil, err + } + + tp := sdktrace.NewTracerProvider( + sdktrace.WithBatcher(exporter), + sdktrace.WithResource(res), + ) + + otel.SetTracerProvider(tp) + otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator( + propagation.TraceContext{}, + propagation.Baggage{}, + )) + + return tp.Shutdown, nil +} + +// newTelemetryResource builds the OTel resource shared by the controller's +// signal pipelines (traces and logs), so every signal carries the same +// resource attributes. +func newTelemetryResource(ctx context.Context, serviceVersion string) (*resource.Resource, error) { instanceID, err := os.Hostname() if err != nil || instanceID == "" { instanceID = uuid.New().String() @@ -59,25 +82,9 @@ func InitTracerProvider(ctx context.Context, serviceVersion string) (func(contex attrs = append(attrs, semconv.K8SNodeName(node)) } - res, err := resource.New(ctx, + return resource.New(ctx, resource.WithTelemetrySDK(), resource.WithAttributes(attrs...), resource.WithFromEnv(), ) - if err != nil { - return nil, fmt.Errorf("create OTEL resource: %w", err) - } - - tp := sdktrace.NewTracerProvider( - sdktrace.WithBatcher(exporter), - sdktrace.WithResource(res), - ) - - otel.SetTracerProvider(tp) - otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator( - propagation.TraceContext{}, - propagation.Baggage{}, - )) - - return tp.Shutdown, nil } diff --git a/go/go.mod b/go/go.mod index c581c5af4f..444495bedd 100644 --- a/go/go.mod +++ b/go/go.mod @@ -51,16 +51,22 @@ require ( github.com/stretchr/testify v1.12.1 github.com/testcontainers/testcontainers-go v0.44.0 github.com/testcontainers/testcontainers-go/modules/postgres v0.44.0 + go.opentelemetry.io/contrib/bridges/otelzap v0.19.0 go.opentelemetry.io/contrib/exporters/autoexport v0.69.0 go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.69.0 go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.69.0 go.opentelemetry.io/otel v1.44.0 go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc v0.20.0 go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploghttp v0.20.0 + go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.44.0 + go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp v1.44.0 go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.44.0 go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.44.0 + go.opentelemetry.io/otel/log v0.20.0 + go.opentelemetry.io/otel/metric v1.44.0 go.opentelemetry.io/otel/sdk v1.44.0 go.opentelemetry.io/otel/sdk/log v0.20.0 + go.opentelemetry.io/otel/sdk/metric v1.44.0 go.opentelemetry.io/otel/trace v1.44.0 go.uber.org/zap v1.28.0 golang.org/x/sync v0.22.0 @@ -423,16 +429,11 @@ require ( go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/contrib/bridges/prometheus v0.69.0 // indirect go.opentelemetry.io/contrib/detectors/gcp v1.44.0 // indirect - go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.44.0 // indirect - go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp v1.44.0 // indirect go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.44.0 // indirect go.opentelemetry.io/otel/exporters/prometheus v0.66.0 // indirect go.opentelemetry.io/otel/exporters/stdout/stdoutlog v0.20.0 // indirect go.opentelemetry.io/otel/exporters/stdout/stdoutmetric v1.44.0 // indirect go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.44.0 // indirect - go.opentelemetry.io/otel/log v0.20.0 // indirect - go.opentelemetry.io/otel/metric v1.44.0 // indirect - go.opentelemetry.io/otel/sdk/metric v1.44.0 // indirect go.opentelemetry.io/proto/otlp v1.11.0 // indirect go.uber.org/atomic v1.11.0 // indirect go.uber.org/multierr v1.11.0 // indirect diff --git a/go/go.sum b/go/go.sum index 29ab0e602a..3af306f0d0 100644 --- a/go/go.sum +++ b/go/go.sum @@ -1218,6 +1218,8 @@ go.opencensus.io v0.20.2/go.mod h1:6WKK9ahsWS3RSO+PY9ZHZUfv2irvY6gN279GOPZjmmk= go.opencensus.io v0.22.2/go.mod h1:yxeiOL68Rb0Xd1ddK5vPZ/oVn4vY4Ynel7k9FzqtOIw= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= +go.opentelemetry.io/contrib/bridges/otelzap v0.19.0 h1:48Eq3xxFx2KlL/tF7lnl42kKJBDlhNTLRzv0h154JnM= +go.opentelemetry.io/contrib/bridges/otelzap v0.19.0/go.mod h1:cQbV77F0u6HmtZPiQD9oxp2esaOEb4uLqIta6OFIKOk= go.opentelemetry.io/contrib/bridges/prometheus v0.69.0 h1:saQoWg5845Q8TojpqeVStS7zGwVZ6bc5W2PJavTPiBM= go.opentelemetry.io/contrib/bridges/prometheus v0.69.0/go.mod h1:AAaS6xs5AyqMdR3Ir0nSWK+QudL2XM8Vbw5INzUxNc8= go.opentelemetry.io/contrib/detectors/gcp v1.44.0 h1:NmLfL734pJhM0JKaYd2Y28+nY9dPRWYAAbxhRCrKXPw= @@ -1254,6 +1256,8 @@ go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.44.0 h1:bl2S7Ubua0Nms+D go.opentelemetry.io/otel/exporters/stdout/stdouttrace v1.44.0/go.mod h1:L0hRV50XdVIODHUfWEqGRCXQvj2rV82STVo12FMFBU0= go.opentelemetry.io/otel/log v0.20.0 h1:/5i0vuHxCLWUfChWG41K9wkM0jafruPw9NU1/RCJirs= go.opentelemetry.io/otel/log v0.20.0/go.mod h1:wOcMcjsZpG8x7Bak7IhSi/lg8wscV2C1VdrKCLPlt0E= +go.opentelemetry.io/otel/log/logtest v0.20.0 h1:+tsZVE15N+RWyN9lUzsRyw7hMZXNMepGu105Eim82/k= +go.opentelemetry.io/otel/log/logtest v0.20.0/go.mod h1:zS9Ryx9RrEAG2tgapMBSvacwhVSSOGSaSiWWgW3NPlQ= go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc= go.opentelemetry.io/otel/metric v1.44.0/go.mod h1:8O7hanEPBNgEMmybD3s2VBKcgWOCsA6tzHBPODAiquo= go.opentelemetry.io/otel/metric/x v0.66.0 h1:YkCrx1zLOChi9ZcZ6euupOcsgzbVlec7D/xoEU1+cTA=