Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 37 additions & 6 deletions cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import (
"crypto/tls"
"errors"
"flag"
"log/slog"
"net/http"
"os"
"time"
Expand All @@ -40,7 +41,6 @@ import (
"sigs.k8s.io/controller-runtime/pkg/cache"
"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/healthz"
"sigs.k8s.io/controller-runtime/pkg/log/zap"
"sigs.k8s.io/controller-runtime/pkg/metrics/filters"
metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
"sigs.k8s.io/controller-runtime/pkg/webhook"
Expand All @@ -51,6 +51,7 @@ import (
agentraxv1alpha1 "github.com/gitcommitankit/agentrax/api/v1alpha1"
"github.com/gitcommitankit/agentrax/internal/controller"
"github.com/gitcommitankit/agentrax/internal/metrics"
"github.com/gitcommitankit/agentrax/internal/observability"
"github.com/gitcommitankit/agentrax/internal/quota"
"github.com/gitcommitankit/agentrax/internal/registry"
"github.com/gitcommitankit/agentrax/internal/rollout"
Expand Down Expand Up @@ -124,6 +125,9 @@ func main() {
var gatewayName string
var gatewayNamespace string
var registryAddr string
var otlpEndpoint string
var otlpInsecure bool
var logLevel string
var tlsOpts []func(*tls.Config)
flag.StringVar(&metricsAddr, "metrics-bind-address", "0", "The address the metrics endpoint binds to. "+
"Use :8443 for HTTPS or :8080 for HTTP, or leave as 0 to disable the metrics service.")
Expand All @@ -146,13 +150,40 @@ func main() {
"Namespace of the Gateway API Gateway object used for canary traffic splitting.")
flag.StringVar(&registryAddr, "registry-bind-address", ":9090",
"The address the MCP discovery registry HTTP endpoint binds to.")
opts := zap.Options{
Development: false,
}
opts.BindFlags(flag.CommandLine)
flag.StringVar(&otlpEndpoint, "otlp-endpoint", "",
"gRPC endpoint for the OpenTelemetry trace exporter (e.g. localhost:4317). "+
"Leave empty to disable tracing.")
flag.BoolVar(&otlpInsecure, "otlp-insecure", false,
"Use plaintext (insecure) gRPC connection for OpenTelemetry trace exporter. "+
"Only use for local development; remote collectors use TLS by default.")
flag.StringVar(&logLevel, "log-level", "info",
"Minimum log level to emit. One of: debug, info, warn, error.")
flag.Parse()

ctrl.SetLogger(zap.New(zap.UseFlagOptions(&opts)))
// Configure structured JSON logging via log/slog, bridged to controller-runtime
// through the logr interface. This replaces the default Zap logger.
var slogLevel slog.Level
if err := slogLevel.UnmarshalText([]byte(logLevel)); err != nil {
slogLevel = slog.LevelInfo
}
slogLogger := slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: slogLevel}))
ctrl.SetLogger(observability.NewLogr(slogLogger))

// Initialize OTel TracerProvider. The returned Shutdown must be deferred so
// buffered spans are flushed before the process exits.
startCtx := context.Background()
tpShutdown, err := observability.InitTracerProvider(startCtx, otlpEndpoint, otlpInsecure)
if err != nil {
setupLog.Error(err, "unable to initialize OpenTelemetry TracerProvider")
os.Exit(1)
}
defer func() {
shutCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if shutErr := tpShutdown(shutCtx); shutErr != nil {
setupLog.Error(shutErr, "error shutting down OTel TracerProvider")
}
}()

// if the enable-http2 flag is false (the default), http/2 should be disabled
// due to its vulnerabilities. More specifically, disabling http/2 will
Expand Down
10 changes: 5 additions & 5 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,11 @@ require (
github.com/onsi/gomega v1.33.1
github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring v0.75.0
github.com/prometheus/client_golang v1.19.1
github.com/prometheus/client_model v0.6.1
github.com/stretchr/testify v1.10.0
go.opentelemetry.io/otel v1.34.0
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.27.0
go.opentelemetry.io/otel/sdk v1.34.0
go.opentelemetry.io/otel/trace v1.34.0
k8s.io/api v0.31.0
k8s.io/apiextensions-apiserver v0.31.0
k8s.io/apimachinery v0.31.0
Expand Down Expand Up @@ -57,6 +60,7 @@ require (
github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 // indirect
github.com/prometheus/client_model v0.6.1 // indirect
github.com/prometheus/common v0.55.0 // indirect
github.com/prometheus/procfs v0.15.1 // indirect
github.com/spf13/cobra v1.8.1 // indirect
Expand All @@ -65,12 +69,8 @@ require (
github.com/x448/float16 v0.8.4 // indirect
go.opentelemetry.io/auto/sdk v1.1.0 // indirect
go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.53.0 // indirect
go.opentelemetry.io/otel v1.34.0 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.28.0 // indirect
go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.27.0 // indirect
go.opentelemetry.io/otel/metric v1.34.0 // indirect
go.opentelemetry.io/otel/sdk v1.34.0 // indirect
go.opentelemetry.io/otel/trace v1.34.0 // indirect
go.opentelemetry.io/proto/otlp v1.3.1 // indirect
go.uber.org/multierr v1.11.0 // indirect
go.uber.org/zap v1.26.0 // indirect
Expand Down
61 changes: 55 additions & 6 deletions internal/controller/agentdeployment_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,9 @@ import (

"github.com/go-logr/logr"
monitoringv1 "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace"
apimeta "k8s.io/apimachinery/pkg/api/meta"

agentraxv1alpha1 "github.com/gitcommitankit/agentrax/api/v1alpha1"
Expand Down Expand Up @@ -125,7 +128,23 @@ type AgentDeploymentReconciler struct {
// It creates and self-heals a Deployment, Service, and (when Prometheus Operator is present)
// a ServiceMonitor as owned child resources, then updates status conditions.
func (r *AgentDeploymentReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
// Start a root OTel span for the entire reconcile loop. The span is ended via
// defer so all exit paths (early return, error, requeue) are covered.
ctx, span := observability.Tracer.Start(ctx, "reconcile",
trace.WithAttributes(
attribute.String("tenant", req.Namespace),
attribute.String("name", req.Name),
),
)
defer span.End()

// Enrich the controller-runtime logger with trace_id/span_id so every log
// line emitted inside this reconcile cycle is correlated to the active OTel trace.
logger := log.FromContext(ctx)
if sc := span.SpanContext(); sc.IsValid() {
logger = logger.WithValues("trace_id", sc.TraceID().String(), "span_id", sc.SpanID().String())
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
ctx = log.IntoContext(ctx, logger)

// Observe reconcile wall-clock duration on every exit path, including early
// returns, errors, and requeues. The tenant label uses the request namespace
Expand All @@ -138,13 +157,19 @@ func (r *AgentDeploymentReconciler) Reconcile(ctx context.Context, req ctrl.Requ
}()

// 1. Fetch the AgentDeployment; return immediately if it has been deleted.
fetchCtx, fetchSpan := observability.Tracer.Start(ctx, "fetch_crd")
ad := &agentraxv1alpha1.AgentDeployment{}
if err := r.Get(ctx, req.NamespacedName, ad); err != nil {
if err := r.Get(fetchCtx, req.NamespacedName, ad); err != nil {
if apierrors.IsNotFound(err) {
fetchSpan.End()
return ctrl.Result{}, nil
}
fetchSpan.RecordError(err)
fetchSpan.SetStatus(codes.Error, "fetch_crd failed")
fetchSpan.End()
return ctrl.Result{}, fmt.Errorf("fetching AgentDeployment: %w", err)
}
fetchSpan.End()

// 2. Handle finalizer lifecycle.
if ad.DeletionTimestamp.IsZero() {
Expand Down Expand Up @@ -202,37 +227,61 @@ func (r *AgentDeploymentReconciler) Reconcile(ctx context.Context, req ctrl.Requ
}
}

// 4–7. Reconcile child resources (Deployment, Service, ServiceMonitor, HPA)
// under a single span. Each helper propagates ctx so sub-operations can be
// correlated if they are instrumented in future phases.
childrenCtx, childrenSpan := observability.Tracer.Start(ctx, "reconcile_children")

// 4. Reconcile child Deployment.
if err := r.reconcileDeployment(ctx, ad); err != nil {
if err := r.reconcileDeployment(childrenCtx, ad); err != nil {
childrenSpan.RecordError(err)
childrenSpan.SetStatus(codes.Error, "reconcile_children failed")
childrenSpan.End()
return ctrl.Result{}, fmt.Errorf("reconciling deployment: %w", err)
}

// 5. Reconcile child Service.
if err := r.reconcileService(ctx, ad); err != nil {
if err := r.reconcileService(childrenCtx, ad); err != nil {
childrenSpan.RecordError(err)
childrenSpan.SetStatus(codes.Error, "reconcile_children failed")
childrenSpan.End()
return ctrl.Result{}, fmt.Errorf("reconciling service: %w", err)
}

// 6. Reconcile ServiceMonitor when Prometheus Operator is present.
if err := r.reconcileServiceMonitor(ctx, ad); err != nil {
if err := r.reconcileServiceMonitor(childrenCtx, ad); err != nil {
childrenSpan.RecordError(err)
childrenSpan.SetStatus(codes.Error, "reconcile_children failed")
childrenSpan.End()
return ctrl.Result{}, fmt.Errorf("reconciling servicemonitor: %w", err)
}

// 7. Reconcile the managed HPA (skip during active canary — Phase 4 owns it).
// reconcileHPA also returns the quota evaluation state so updateStatus can
// write the correct QuotaLimited condition onto the freshly re-fetched object.
hpaResult, qs, err := r.reconcileHPA(ctx, ad)
hpaResult, qs, err := r.reconcileHPA(childrenCtx, ad)
if err != nil {
childrenSpan.RecordError(err)
childrenSpan.SetStatus(codes.Error, "reconcile_children failed")
childrenSpan.End()
return ctrl.Result{}, fmt.Errorf("reconciling hpa: %w", err)
}
childrenSpan.End()

// 8. Derive status from the live Deployment and update it — always last.
// We continue into updateStatus even when hpaResult requests a requeue so
// that the QuotaLimited condition is written in the same reconcile cycle.
// Return the shorter of the two requeue intervals.
statusResult, err := r.updateStatus(ctx, ad, logger, qs)
statusCtx, statusSpan := observability.Tracer.Start(ctx, "update_status")
statusResult, err := r.updateStatus(statusCtx, ad, logger, qs)
if err != nil {
statusSpan.RecordError(err)
statusSpan.SetStatus(codes.Error, "update_status failed")
statusSpan.End()
return statusResult, fmt.Errorf("updating status: %w", err)
}
statusSpan.End()

if hpaResult.RequeueAfter > 0 {
if statusResult.RequeueAfter == 0 || hpaResult.RequeueAfter < statusResult.RequeueAfter {
return hpaResult, nil
Expand Down
40 changes: 40 additions & 0 deletions internal/observability/logging.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
package observability

import (
"context"
"io"
"log/slog"

"github.com/go-logr/logr"
"go.opentelemetry.io/otel/trace"
)

// NewJSONLogger returns a [*slog.Logger] that writes structured JSON records to
// w. Pass os.Stdout in production; pass a [*bytes.Buffer] in tests.
func NewJSONLogger(w io.Writer) *slog.Logger {
return slog.New(slog.NewJSONHandler(w, &slog.HandlerOptions{
Level: slog.LevelInfo,
}))
}

// WithTraceContext returns a child logger pre-populated with "trace_id" and
// "span_id" fields extracted from the active OpenTelemetry span in ctx. If ctx
// carries no valid span the original logger is returned unchanged.
func WithTraceContext(ctx context.Context, logger *slog.Logger) *slog.Logger {
span := trace.SpanFromContext(ctx)
sc := span.SpanContext()
if !sc.IsValid() {
return logger
}
return logger.With(
"trace_id", sc.TraceID().String(),
"span_id", sc.SpanID().String(),
)
}

// NewLogr wraps a [*slog.Logger] as a [logr.Logger] so it can be passed to
// controller-runtime's [ctrl.SetLogger]. All controller-runtime log output
// (including reconcile errors) then flows through the slog JSON handler.
func NewLogr(logger *slog.Logger) logr.Logger {
return logr.FromSlogHandler(logger.Handler())
}
81 changes: 81 additions & 0 deletions internal/observability/tracing.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
// Package observability provides OpenTelemetry tracing initialization for the Agentrax operator.
//
// Use [InitTracerProvider] in main() to configure the global OTel tracer. When
// endpoint is empty the function installs a no-op provider so callers need not
// guard on whether tracing is enabled. The returned Shutdown function must be
// deferred in main() to flush buffered spans before process exit.
package observability

import (
"context"
"fmt"

"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc"
"go.opentelemetry.io/otel/propagation"
"go.opentelemetry.io/otel/sdk/resource"
sdktrace "go.opentelemetry.io/otel/sdk/trace"
semconv "go.opentelemetry.io/otel/semconv/v1.26.0"
"go.opentelemetry.io/otel/trace"
"go.opentelemetry.io/otel/trace/noop"
)

// TracerName is the instrumentation library name used for all Agentrax OTel spans.
const TracerName = "agentrax.io/controller"

// Tracer is the package-level tracer used by controller code. It is set by
// [InitTracerProvider] and defaults to the global no-op tracer so it is always
// safe to call, even when tracing is disabled.
var Tracer trace.Tracer = noop.NewTracerProvider().Tracer(TracerName)

// InitTracerProvider configures the global OpenTelemetry TracerProvider and
// installs a W3C TraceContext propagator. When endpoint is empty it installs a
// no-op provider (tracing disabled). TLS is used by default for exporter connections;
// set insecure to true only for local development with plaintext collectors.
// The returned Shutdown function flushes and stops the exporter; it must be
// deferred in main().
func InitTracerProvider(ctx context.Context, endpoint string, insecure bool) (func(context.Context) error, error) {
if endpoint == "" {
// No-op: tracing disabled. Global tracer stays as the default no-op.
return func(context.Context) error { return nil }, nil
}

opts := []otlptracegrpc.Option{
otlptracegrpc.WithEndpoint(endpoint),
}
if insecure {
opts = append(opts, otlptracegrpc.WithInsecure())
}

exp, err := otlptracegrpc.New(ctx, opts...)
if err != nil {
return nil, fmt.Errorf("creating OTLP gRPC exporter: %w", err)
}

res, err := resource.Merge(
resource.Default(),
resource.NewSchemaless(
semconv.ServiceName("agentrax-operator"),
),
)
if err != nil {
return nil, fmt.Errorf("creating OTel resource: %w", err)
}

tp := sdktrace.NewTracerProvider(
sdktrace.WithBatcher(exp),
sdktrace.WithResource(res),
)

otel.SetTracerProvider(tp)
// W3C TraceContext + Baggage propagation — required for cross-service correlation.
otel.SetTextMapPropagator(propagation.NewCompositeTextMapPropagator(
propagation.TraceContext{},
propagation.Baggage{},
))

// Update the package-level tracer to use the real provider.
Tracer = tp.Tracer(TracerName)

return tp.Shutdown, nil
}
Loading
Loading