diff --git a/.changeset/bright-otters-observe.md b/.changeset/bright-otters-observe.md new file mode 100644 index 00000000..7b2a83b2 --- /dev/null +++ b/.changeset/bright-otters-observe.md @@ -0,0 +1,5 @@ +--- +"posthog-go": minor +--- + +Add the OpenTelemetry bridge for AI observability as an independently installable nested Go module. diff --git a/.github/dependabot.yml b/.github/dependabot.yml index d230166f..912bb236 100644 --- a/.github/dependabot.yml +++ b/.github/dependabot.yml @@ -10,6 +10,16 @@ updates: go-dependencies: patterns: - "*" + - package-ecosystem: "gomod" + directory: "/otel" + schedule: + interval: weekly + cooldown: + default-days: 7 + groups: + go-dependencies: + patterns: + - "*" - package-ecosystem: "npm" directory: "/" schedule: diff --git a/.github/workflows/fmt.yml b/.github/workflows/fmt.yml index 8ded3e9c..8a21c032 100644 --- a/.github/workflows/fmt.yml +++ b/.github/workflows/fmt.yml @@ -48,3 +48,30 @@ jobs: run: | go mod tidy git diff --exit-code || (echo "go.mod/go.sum are not tidy. Run \`go mod tidy\` locally and commit the changes." && exit 1) + + # The jobs above run at the repository root, and `./...` never descends into a + # nested module, so neither reaches otel/. + otel-fmt-tidy: + name: Check OTel module formatting and go.mod + runs-on: ubuntu-latest + defaults: + run: + working-directory: otel + steps: + - name: Check out source code + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + + - name: Set up Go + uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 + with: + go-version-file: otel/go.mod + + - name: Check formatting + run: | + UNFORMATTED=$(gofmt -l .) + test -z "$UNFORMATTED" || (echo "Files are not formatted:" && echo "$UNFORMATTED" && echo "Run \`gofmt -w .\` in otel/ and commit the changes." && exit 1) + + - name: Check go.mod/go.sum + run: | + go mod tidy + git diff --exit-code . || (echo "otel/go.mod or otel/go.sum are not tidy. Run \`go mod tidy\` in otel/ and commit the changes." && exit 1) diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 52d7df10..977054d7 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -136,14 +136,37 @@ jobs: branch: main env: GITHUB_TOKEN: ${{ steps.releaser.outputs.token }} + + - name: Create OTel module tag + if: steps.commit-version-bump.outputs.commit-hash != '' + env: + GH_TOKEN: ${{ steps.releaser.outputs.token }} + NEW_VERSION: ${{ steps.apply-changesets.outputs.new-version }} + RELEASE_COMMIT: ${{ steps.commit-version-bump.outputs.commit-hash }} + run: | + TAG="otel/v${NEW_VERSION}" + EXISTING_SHA=$(git ls-remote origin "refs/tags/${TAG}" | cut -f1) + if [ -n "$EXISTING_SHA" ]; then + if [ "$EXISTING_SHA" != "$RELEASE_COMMIT" ]; then + echo "Tag ${TAG} already exists at ${EXISTING_SHA}, expected ${RELEASE_COMMIT}" >&2 + exit 1 + fi + echo "Tag ${TAG} already exists at the release commit" + else + gh api --method POST "repos/${GITHUB_REPOSITORY}/git/refs" \ + -f ref="refs/tags/${TAG}" \ + -f sha="$RELEASE_COMMIT" + fi + - name: Create GitHub release if: steps.commit-version-bump.outputs.commit-hash != '' env: GH_TOKEN: ${{ steps.releaser.outputs.token }} NEW_VERSION: ${{ steps.apply-changesets.outputs.new-version }} + RELEASE_COMMIT: ${{ steps.commit-version-bump.outputs.commit-hash }} run: | CHANGELOG_ENTRY=$(awk -v defText="see CHANGELOG.md" '/^## /{if (flag) exit; flag=1} flag && /^##$/{exit} flag; END{if (!flag) print defText}' CHANGELOG.md) - gh release create "v${NEW_VERSION}" --target main --title "${NEW_VERSION}" --notes "$CHANGELOG_ENTRY" + gh release create "v${NEW_VERSION}" --target "$RELEASE_COMMIT" --title "${NEW_VERSION}" --notes "$CHANGELOG_ENTRY" - name: Dispatch posthog upgrade for posthog-go if: steps.commit-version-bump.outputs.commit-hash != '' diff --git a/.github/workflows/unit-tests.yml b/.github/workflows/unit-tests.yml index ce8f6c6a..fc8efb48 100644 --- a/.github/workflows/unit-tests.yml +++ b/.github/workflows/unit-tests.yml @@ -35,6 +35,30 @@ jobs: - name: Build run: go build . + otel-bridge: + name: OTel bridge (Go ${{ matrix.go-version }}) + runs-on: ubuntu-latest + strategy: + fail-fast: false + matrix: + # Test the minimum supported Go version from go.mod and the latest stable Go 1 release. + go-version: ['1.25', '1.x'] + steps: + - name: Check out source code + uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 + + - name: Set up Go ${{ matrix.go-version }} + uses: actions/setup-go@b7ad1dad31e06c5925ef5d2fc7ad053ef454303e # v7.0.0 + with: + go-version: ${{ matrix.go-version }} + + - name: Build, vet, and test + working-directory: otel + run: | + go build ./... + go vet ./... + go test -race -count=1 -timeout=5m ./... + public-api: name: Public API runs-on: ubuntu-latest diff --git a/README.md b/README.md index 371efd6c..f195ae61 100644 --- a/README.md +++ b/README.md @@ -11,6 +11,13 @@ SDK usage examples and code snippets live in the official documentation so they - [Go library docs](https://posthog.com/docs/libraries/go) +## AI observability + +The [`otel`](otel) module is an OpenTelemetry bridge that forwards AI spans +(`gen_ai.*`, `llm.*`, and similar) to PostHog AI observability, including spans +from a Google Agent Development Kit (ADK) for Go agent. It is a separate Go +module, so the core SDK stays free of OpenTelemetry dependencies. + ## Questions? ### [Visit the community forum.](https://posthog.com/questions) diff --git a/otel/LICENSE.md b/otel/LICENSE.md new file mode 100644 index 00000000..922b77cb --- /dev/null +++ b/otel/LICENSE.md @@ -0,0 +1,23 @@ +The MIT License (MIT) + +Copyright (c) 2020 PostHog (part of Hiberly Inc) + +Copyright (c) 2016 Segment, Inc. + +Permission is hereby granted, free of charge, to any person obtaining a copy +of this software and associated documentation files (the "Software"), to deal +in the Software without restriction, including without limitation the rights +to use, copy, modify, merge, publish, distribute, sublicense, and/or sell +copies of the Software, and to permit persons to whom the Software is +furnished to do so, subject to the following conditions: + +The above copyright notice and this permission notice shall be included in all +copies or substantial portions of the Software. + +THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +SOFTWARE. diff --git a/otel/README.md b/otel/README.md new file mode 100644 index 00000000..9f33ae1d --- /dev/null +++ b/otel/README.md @@ -0,0 +1,52 @@ +# PostHog OpenTelemetry bridge for AI observability + +`posthogotel` forwards OpenTelemetry AI spans to [PostHog AI observability](https://posthog.com/docs/ai-engineering/observability). + +It keeps only spans that follow a known AI semantic convention — a span whose +name or any attribute key starts with `gen_ai.`, `llm.`, `ai.`, or +`traceloop.` — and drops every other span. Kept spans go over OTLP/HTTP to the +PostHog `/i/v0/ai/otel` endpoint with the project API key as a bearer token. + +This is a separate Go module, so the core `posthog-go` SDK does not depend on +OpenTelemetry. + +## Usage + +`SpanProcessor` is the recommended integration. Register it on the +`TracerProvider` your application already owns — rather than replacing the +global provider with a new one — so your resource, sampler, and existing +exporters are kept and tracers already handed out (such as ADK Go's) route +through it. Shut down the processor with a fresh context to flush its buffered +spans without shutting down the application-owned provider. + +If you don't already have a `TracerProvider`, create one and register the +processor on it, as [`example/`](example) does. Use `WithHost` for a host other +than PostHog US cloud. For a framework that accepts only a span exporter, use +`NewExporter` instead and pair it with your own batch span processor. + +## Failed generations + +PostHog decides that a generation failed from the OpenTelemetry span status, so set the +status to the error code when a model call fails. Recording the error on the span is not +enough on its own: in OpenTelemetry for Go that only adds an exception event and leaves +the span status unset, so the failed generation reaches PostHog looking successful with an +empty response. The Python and JavaScript instrumentation sets the status for you, which +is why this step is specific to Go. Once the status is set, PostHog fills in the error +message and HTTP status from the recorded exception event. + +## Google Agent Development Kit (ADK) for Go + +[ADK Go](https://google.golang.org/adk) instruments its agents with +OpenTelemetry and emits `gen_ai.*` spans on the global tracer provider. Register +the PostHog span processor on that provider before you run the agent, and the +agent's `gen_ai.*` spans reach PostHog with no further code. + +Those spans carry the generation's model, token counts, latency, and finish +reason. They do **not** carry prompt and response message content: ADK Go emits +message bodies as OpenTelemetry **log records** (event names +`gen_ai.system.message`, `gen_ai.user.message`, and `gen_ai.choice`), gated +behind the `OTEL_INSTRUMENTATION_GENAI_CAPTURE_MESSAGE_CONTENT` environment +variable, not as span attributes. This bridge forwards spans only, so prompt and +response fields stay empty for ADK Go generations. + +See [`example/`](example) for a runnable program. diff --git a/otel/config.go b/otel/config.go new file mode 100644 index 00000000..35da78f5 --- /dev/null +++ b/otel/config.go @@ -0,0 +1,141 @@ +package posthogotel + +import ( + "context" + "errors" + "fmt" + "net/url" + "strings" + + "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp" + sdktrace "go.opentelemetry.io/otel/sdk/trace" +) + +const ( + // DefaultHost is the PostHog US cloud host used when WithHost is not set. + DefaultHost = "https://us.i.posthog.com" + + // ingestPath is the PostHog AI observability OTLP endpoint path. + ingestPath = "/i/v0/ai/otel" + + // maxSpansPerRequest is the maximum number of AI spans the PostHog AI + // observability endpoint accepts in a single OTLP request. Larger requests + // are rejected with a non-retryable HTTP 400 and the whole batch is lost, so + // batches must be capped at this limit. + maxSpansPerRequest = 100 +) + +// errEmptyAPIKey is returned when the project API key is missing. +var errEmptyAPIKey = errors.New("posthogotel: apiKey must not be empty") + +// errInvalidHost is returned when the configured host is not an absolute http +// or https URL with a hostname. +var errInvalidHost = errors.New("posthogotel: host must be an absolute http or https URL, for example https://us.i.posthog.com") + +// config holds the resolved settings for the exporter and the span processor. +type config struct { + host string + // endpoint is the resolved OTLP URL, host joined with ingestPath. + endpoint string +} + +// Option configures the exporter or the span processor. +type Option func(*config) + +// WithHost sets the PostHog host, for example "https://eu.i.posthog.com". +// An empty or blank value keeps DefaultHost. +func WithHost(host string) Option { + return func(c *config) { + if h := strings.TrimSpace(host); h != "" { + c.host = h + } + } +} + +func newConfig(opts ...Option) (config, error) { + c := config{host: DefaultHost} + for _, opt := range opts { + opt(&c) + } + c.host = strings.TrimRight(c.host, "/") + endpoint, err := resolveEndpoint(c.host) + if err != nil { + return config{}, err + } + c.endpoint = endpoint + return c, nil +} + +// resolveEndpoint builds the OTLP URL for host and rejects a host that +// otlptracehttp.WithEndpointURL would silently discard. On a parse failure it +// keeps its localhost defaults, so spans go nowhere while the request still +// carries the API key; a scheme-less host such as "us.i.posthog.com" parses but +// yields an empty endpoint. Requiring an absolute http or https URL with a +// hostname turns both into an upfront error. +// +// ingestPath is joined rather than concatenated. A host that carries a query or +// fragment, such as "https://us.i.posthog.com?region=eu", would concatenate into +// a URL whose path is empty, and the exporter would then fall back to the OTLP +// default "/v1/traces" and send every AI span, with the API key attached, to a +// path PostHog does not serve. +func resolveEndpoint(host string) (string, error) { + u, err := url.Parse(host) + if err != nil { + return "", fmt.Errorf("%w: %v", errInvalidHost, err) + } + if (u.Scheme != "http" && u.Scheme != "https") || u.Hostname() == "" { + return "", errInvalidHost + } + return u.JoinPath(ingestPath).String(), nil +} + +// newOTLPExporter builds an OTLP/HTTP exporter that targets the PostHog AI +// observability endpoint with the project API key as a bearer token. The +// exporter is wrapped so that no single request exceeds maxSpansPerRequest, +// which protects both public entry points regardless of the batch size of the +// span processor that feeds them. +func newOTLPExporter(ctx context.Context, apiKey string, cfg config) (sdktrace.SpanExporter, error) { + apiKey = strings.TrimSpace(apiKey) + exporter, err := otlptracehttp.New(ctx, + otlptracehttp.WithEndpointURL(cfg.endpoint), + otlptracehttp.WithHeaders(map[string]string{ + "Authorization": "Bearer " + apiKey, + }), + ) + if err != nil { + return nil, err + } + return &chunkingExporter{inner: exporter, limit: maxSpansPerRequest}, nil +} + +// chunkingExporter splits each ExportSpans batch into requests of at most limit +// spans. The PostHog AI observability endpoint rejects larger requests with a +// non-retryable HTTP 400 that discards the whole batch, and nothing below this +// module splits a batch: the OTLP exporter turns whatever slice it receives +// into exactly one request. Chunking here caps every request for both the +// SpanProcessor and the caller-supplied Exporter path. +type chunkingExporter struct { + inner sdktrace.SpanExporter + limit int +} + +var _ sdktrace.SpanExporter = (*chunkingExporter)(nil) + +// ExportSpans forwards spans to the inner exporter in slices of at most limit. +func (e *chunkingExporter) ExportSpans(ctx context.Context, spans []sdktrace.ReadOnlySpan) error { + for start := 0; start < len(spans); start += e.limit { + end := start + e.limit + if end > len(spans) { + end = len(spans) + } + if err := e.inner.ExportSpans(ctx, spans[start:end]); err != nil { + return err + } + } + return nil +} + +// Shutdown shuts down the inner exporter. +func (e *chunkingExporter) Shutdown(ctx context.Context) error { + return e.inner.Shutdown(ctx) +} diff --git a/otel/doc.go b/otel/doc.go new file mode 100644 index 00000000..75337ae9 --- /dev/null +++ b/otel/doc.go @@ -0,0 +1,21 @@ +// Package posthogotel forwards OpenTelemetry AI spans to PostHog AI observability. +// +// It keeps only spans that follow a known AI semantic convention. A span +// qualifies when its name or any of its attribute keys starts with one of +// "gen_ai.", "llm.", "ai.", or "traceloop.". Every other span is dropped. +// Kept spans go over OTLP/HTTP to the PostHog "/i/v0/ai/otel" endpoint with the +// project API key in an Authorization: Bearer header. +// +// The package offers two integrations: +// +// - SpanProcessor is the recommended integration. It filters spans, batches +// the AI spans, and exports them. Register it with +// TracerProvider.RegisterSpanProcessor (or the WithSpanProcessor option). +// +// - Exporter is for setups that supply their own span processor, or +// frameworks that accept only a span exporter. It filters spans and +// delegates the AI spans to an OTLP/HTTP exporter. +// +// This is a separate Go module so that the core posthog-go SDK does not depend +// on OpenTelemetry. +package posthogotel diff --git a/otel/example/main.go b/otel/example/main.go new file mode 100644 index 00000000..a8307d88 --- /dev/null +++ b/otel/example/main.go @@ -0,0 +1,117 @@ +// Command example sends AI spans to PostHog AI observability through the +// posthogotel OpenTelemetry bridge. +// +// Google's Agent Development Kit for Go (google.golang.org/adk) instruments +// its agents with OpenTelemetry and emits gen_ai.* spans on the OpenTelemetry +// global tracer provider. Register the PostHog span processor on that provider +// before you run the agent, and the agent's gen_ai.* spans reach PostHog with +// no further code. This example emits one synthetic gen_ai.* generation with +// representative model, message, token, and response attributes in place of a +// live agent so that it runs without model credentials and is easy to inspect +// in the PostHog UI. +// +// Run it with: +// +// POSTHOG_PROJECT_API_KEY=phc_xxx go run . +// +// Set POSTHOG_ENDPOINT to target a host other than PostHog US cloud, for +// example https://eu.i.posthog.com. +package main + +import ( + "context" + "log" + "os" + "time" + + posthogotel "github.com/posthog/posthog-go/otel" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + sdkresource "go.opentelemetry.io/otel/sdk/resource" + sdktrace "go.opentelemetry.io/otel/sdk/trace" +) + +func main() { + apiKey := os.Getenv("POSTHOG_PROJECT_API_KEY") + if apiKey == "" { + log.Fatal("set POSTHOG_PROJECT_API_KEY to your PostHog project API key") + } + + ctx := context.Background() + + var opts []posthogotel.Option + if host := os.Getenv("POSTHOG_ENDPOINT"); host != "" { + opts = append(opts, posthogotel.WithHost(host)) + } + + processor, err := posthogotel.NewSpanProcessor(ctx, apiKey, opts...) + if err != nil { + log.Fatalf("create PostHog span processor: %v", err) + } + + // Register the processor on the global tracer provider that ADK Go uses. + // The resource attributes make the synthetic example easy to identify in + // PostHog without affecting how real instrumented applications are wired. + resource := sdkresource.NewWithAttributes("", + attribute.String("service.name", "posthog-go-otel-example"), + attribute.String("posthog.distinct_id", "posthog-go-otel-example"), + ) + provider := sdktrace.NewTracerProvider( + sdktrace.WithResource(resource), + sdktrace.WithSpanProcessor(processor), + ) + otel.SetTracerProvider(provider) + defer func() { + shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := provider.Shutdown(shutdownCtx); err != nil { + log.Printf("shutdown tracer provider: %v", err) + } + }() + + traceID, runID := runAgentTurn(ctx) + + // ForceFlush blocks until the queued span is exported and surfaces any + // export error. Shutdown alone would not: the batch span processor returns + // only context errors from Shutdown, so a rejected export (for example a bad + // API key or host) never reaches its return value. Return on failure so the + // deferred Shutdown still runs, and report success only once the export is + // actually confirmed. + flushCtx, cancel := context.WithTimeout(ctx, 5*time.Second) + defer cancel() + if err := provider.ForceFlush(flushCtx); err != nil { + log.Printf("send AI span to PostHog: %v", err) + return + } + log.Printf("sent AI span to PostHog (trace_id=%s, example.run_id=%s)", traceID, runID) +} + +// runAgentTurn emits one synthetic gen_ai.* generation with enough attributes +// to validate the conversation, model, provider, token, latency, and trace views. +// ADK Go emits its own gen_ai.* spans, but currently emits message content as log +// records rather than the span attributes used here. +func runAgentTurn(ctx context.Context) (traceID, runID string) { + tracer := otel.Tracer("posthog-go/otel/example") + _, span := tracer.Start(ctx, "chat posthog-go OTel example") + + runID = time.Now().UTC().Format("20060102T150405.000000000Z") + span.SetAttributes( + attribute.String("gen_ai.operation.name", "chat"), + attribute.String("gen_ai.provider.name", "openai"), + attribute.String("gen_ai.request.model", "gpt-4o-mini"), + attribute.String("gen_ai.response.model", "gpt-4o-mini-2024-07-18"), + attribute.String("gen_ai.response.id", "chatcmpl-posthog-go-example"), + attribute.StringSlice("gen_ai.response.finish_reasons", []string{"stop"}), + attribute.String("gen_ai.input.messages", `[{"role":"system","content":"Answer concisely."},{"role":"user","content":"What is PostHog?"}]`), + attribute.String("gen_ai.output.messages", `[{"role":"assistant","content":"PostHog is an open-source product analytics platform."}]`), + attribute.Int("gen_ai.usage.input_tokens", 18), + attribute.Int("gen_ai.usage.output_tokens", 11), + attribute.String("server.address", "api.openai.com"), + attribute.String("example.run_id", runID), + ) + + // Make latency visible in the UI without calling an external model. + time.Sleep(50 * time.Millisecond) + span.End() + return span.SpanContext().TraceID().String(), runID +} diff --git a/otel/exporter.go b/otel/exporter.go new file mode 100644 index 00000000..7eee0599 --- /dev/null +++ b/otel/exporter.go @@ -0,0 +1,51 @@ +package posthogotel + +import ( + "context" + "strings" + + sdktrace "go.opentelemetry.io/otel/sdk/trace" +) + +// Exporter is a span exporter that keeps only AI spans and forwards them to +// PostHog. Use it when you supply your own span processor, or with a framework +// that accepts only a span exporter. For most setups prefer SpanProcessor. +type Exporter struct { + inner sdktrace.SpanExporter +} + +var _ sdktrace.SpanExporter = (*Exporter)(nil) + +// NewExporter builds an Exporter for the given project API key. +func NewExporter(ctx context.Context, apiKey string, opts ...Option) (*Exporter, error) { + if strings.TrimSpace(apiKey) == "" { + return nil, errEmptyAPIKey + } + cfg, err := newConfig(opts...) + if err != nil { + return nil, err + } + inner, err := newOTLPExporter(ctx, apiKey, cfg) + if err != nil { + return nil, err + } + return &Exporter{inner: inner}, nil +} + +// ExportSpans forwards the AI spans in the batch and drops the rest. It returns +// early without a request when the batch has no AI spans. +func (e *Exporter) ExportSpans(ctx context.Context, spans []sdktrace.ReadOnlySpan) error { + aiSpans := filterAISpans(spans) + if len(aiSpans) == 0 { + return nil + } + for _, span := range aiSpans { + warnIfPostHogAIGateway(span) + } + return e.inner.ExportSpans(ctx, aiSpans) +} + +// Shutdown shuts down the underlying OTLP exporter. +func (e *Exporter) Shutdown(ctx context.Context) error { + return e.inner.Shutdown(ctx) +} diff --git a/otel/gateway.go b/otel/gateway.go new file mode 100644 index 00000000..3ded39c2 --- /dev/null +++ b/otel/gateway.go @@ -0,0 +1,88 @@ +package posthogotel + +import ( + "log" + "net/url" + "regexp" + "strings" + + "go.opentelemetry.io/otel/attribute" + sdktrace "go.opentelemetry.io/otel/sdk/trace" +) + +// posthogAIGatewayHosts are the deployed PostHog AI Gateway hosts. The gateway +// captures its own $ai_generation on every routed call, so a service that both +// routes through it and exports spans through this bridge double-counts (and, +// for billable products, double-bills) every generation. Keep in sync with the +// sibling PostHog SDKs. +var posthogAIGatewayHosts = map[string]struct{}{ + "gateway.posthog.com": {}, + "gateway.us.posthog.com": {}, + "gateway.eu.posthog.com": {}, + "ai-gateway.us.posthog.com": {}, + "ai-gateway.eu.posthog.com": {}, +} + +// gatewayURLAttributes are the span attribute keys whose host identifies the +// PostHog AI Gateway. They follow the GenAI/HTTP semantic conventions: +// server.address is a bare host and url.full a full URL. +var gatewayURLAttributes = [...]string{"server.address", "url.full"} + +// gatewayDocsURL points at the AI observability docs referenced by the warning. +const gatewayDocsURL = "https://posthog.com/docs/ai-observability" + +// schemeRE matches a leading URL scheme, so a bare host without one can be +// tolerated (for example "gateway.us.posthog.com/v1"). +var schemeRE = regexp.MustCompile(`(?i)^[a-z][a-z0-9+.-]*://`) + +// isPostHogAIGatewayURL reports whether baseURL points at a known PostHog AI +// Gateway host. +func isPostHogAIGatewayURL(baseURL string) bool { + if baseURL == "" { + return false + } + raw := baseURL + if !schemeRE.MatchString(raw) { + raw = "https://" + raw + } + u, err := url.Parse(raw) + if err != nil { + return false + } + host := strings.ToLower(u.Hostname()) + if host == "" { + return false + } + _, ok := posthogAIGatewayHosts[host] + return ok +} + +// warnIfPostHogAIGateway logs a warning when a span's host/URL attributes point +// at the PostHog AI Gateway, which captures its own $ai_generation. It warns on +// every gateway span by design: the misconfiguration is easy to miss and a +// doubled bill is worse than a noisy log. It never drops the span, because the +// span carries data the gateway never sees. +func warnIfPostHogAIGateway(span sdktrace.ReadOnlySpan) { + for _, attr := range span.Attributes() { + if !isGatewayURLAttribute(string(attr.Key)) { + continue + } + if attr.Value.Type() != attribute.STRING || !isPostHogAIGatewayURL(attr.Value.AsString()) { + continue + } + log.Printf("[PostHog] This OpenTelemetry bridge is exporting spans from a call routed "+ + "through the PostHog AI Gateway, which captures its own $ai_generation. Every such call "+ + "is double-counted and double-billed. Use one or the other — see %s.", gatewayDocsURL) + return + } +} + +// isGatewayURLAttribute reports whether key is one of the gateway URL attributes. +func isGatewayURLAttribute(key string) bool { + for _, k := range gatewayURLAttributes { + if key == k { + return true + } + } + return false +} diff --git a/otel/go.mod b/otel/go.mod new file mode 100644 index 00000000..66263f6a --- /dev/null +++ b/otel/go.mod @@ -0,0 +1,30 @@ +module github.com/posthog/posthog-go/otel + +go 1.25.0 + +require ( + go.opentelemetry.io/otel v1.43.0 + go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.43.0 + go.opentelemetry.io/otel/sdk v1.43.0 + go.opentelemetry.io/proto/otlp v1.10.0 + google.golang.org/protobuf v1.36.11 +) + +require ( + github.com/cenkalti/backoff/v5 v5.0.3 // indirect + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/go-logr/logr v1.4.3 // indirect + github.com/go-logr/stdr v1.2.2 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect + go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.43.0 // indirect + go.opentelemetry.io/otel/metric v1.43.0 // indirect + go.opentelemetry.io/otel/trace v1.43.0 // indirect + golang.org/x/net v0.55.0 // indirect + golang.org/x/sys v0.45.0 // indirect + golang.org/x/text v0.37.0 // indirect + google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 // indirect + google.golang.org/grpc v1.82.1 // indirect +) diff --git a/otel/go.sum b/otel/go.sum new file mode 100644 index 00000000..2380b568 --- /dev/null +++ b/otel/go.sum @@ -0,0 +1,61 @@ +github.com/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1xcsSM= +github.com/cenkalti/backoff/v5 v5.0.3/go.mod h1:rkhZdG3JZukswDf7f0cwqPNk4K0sa+F97BxZthm/crw= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= +github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 h1:HWRh5R2+9EifMyIHV7ZV+MIZqgz+PMpZ14Jynv3O2Zs= +github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0/go.mod h1:JfhWUomR1baixubs02l85lZYYOm7LV6om4ceouMv45c= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +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/otel v1.43.0 h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I= +go.opentelemetry.io/otel v1.43.0/go.mod h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.43.0 h1:88Y4s2C8oTui1LGM6bTWkw0ICGcOLCAI5l6zsD1j20k= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.43.0/go.mod h1:Vl1/iaggsuRlrHf/hfPJPvVag77kKyvrLeD10kpMl+A= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.43.0 h1:3iZJKlCZufyRzPzlQhUIWVmfltrXuGyfjREgGP3UUjc= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.43.0/go.mod h1:/G+nUPfhq2e+qiXMGxMwumDrP5jtzU+mWN7/sjT2rak= +go.opentelemetry.io/otel/metric v1.43.0 h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM= +go.opentelemetry.io/otel/metric v1.43.0/go.mod h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY= +go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg= +go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg= +go.opentelemetry.io/otel/sdk/metric v1.43.0 h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw= +go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A= +go.opentelemetry.io/otel/trace v1.43.0 h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A= +go.opentelemetry.io/otel/trace v1.43.0/go.mod h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0= +go.opentelemetry.io/proto/otlp v1.10.0 h1:IQRWgT5srOCYfiWnpqUYz9CVmbO8bFmKcwYxpuCSL2g= +go.opentelemetry.io/proto/otlp v1.10.0/go.mod h1:/CV4QoCR/S9yaPj8utp3lvQPoqMtxXdzn7ozvvozVqk= +go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= +go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= +golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8= +golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww= +golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY= +golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= +golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= +gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= +gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= +google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 h1:yQugLulqltosq0B/f8l4w9VryjV+N/5gcW0jQ3N8Qec= +google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478/go.mod h1:C6ADNqOxbgdUUeRTU+LCHDPB9ttAMCTff6auwCVa4uc= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 h1:RmoJA1ujG+/lRGNfUnOMfhCy5EipVMyvUE+KNbPbTlw= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= +google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE= +google.golang.org/grpc v1.82.1/go.mod h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA= +google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= +google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/otel/posthogotel_test.go b/otel/posthogotel_test.go new file mode 100644 index 00000000..812fee84 --- /dev/null +++ b/otel/posthogotel_test.go @@ -0,0 +1,518 @@ +package posthogotel + +import ( + "bytes" + "context" + "errors" + "io" + "log" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + "time" + + "go.opentelemetry.io/otel/attribute" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" + coltracepb "go.opentelemetry.io/proto/otlp/collector/trace/v1" + "google.golang.org/protobuf/proto" +) + +// recordSpan starts and ends a span with the given name and attribute keys, +// then returns the resulting ReadOnlySpan. +func recordSpan(t *testing.T, name string, attrKeys ...string) sdktrace.ReadOnlySpan { + t.Helper() + recorder := tracetest.NewSpanRecorder() + provider := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(recorder)) + attrs := make([]attribute.KeyValue, len(attrKeys)) + for i, key := range attrKeys { + attrs[i] = attribute.String(key, "value") + } + _, span := provider.Tracer("test").Start(context.Background(), name) + span.SetAttributes(attrs...) + span.End() + ended := recorder.Ended() + if len(ended) != 1 { + t.Fatalf("expected 1 recorded span, got %d", len(ended)) + } + return ended[0] +} + +func TestIsAISpan(t *testing.T) { + cases := []struct { + name string + spanName string + attrKeys []string + want bool + }{ + {"name gen_ai prefix", "gen_ai.chat", nil, true}, + {"name llm prefix", "llm.request", nil, true}, + {"name ai prefix", "ai.completion", nil, true}, + {"name traceloop prefix", "traceloop.workflow", nil, true}, + {"attribute key prefix", "handler", []string{"gen_ai.system"}, true}, + {"non-ai name and attributes", "http.request", []string{"http.method"}, false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + span := recordSpan(t, tc.spanName, tc.attrKeys...) + if got := IsAISpan(span); got != tc.want { + t.Errorf("IsAISpan(%q) = %v, want %v", tc.spanName, got, tc.want) + } + }) + } +} + +// recordSpanWithAttr records a span carrying a single string attribute. +func recordSpanWithAttr(t *testing.T, name, key, value string) sdktrace.ReadOnlySpan { + t.Helper() + recorder := tracetest.NewSpanRecorder() + provider := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(recorder)) + _, span := provider.Tracer("test").Start(context.Background(), name) + span.SetAttributes(attribute.String(key, value)) + span.End() + ended := recorder.Ended() + if len(ended) != 1 { + t.Fatalf("expected 1 recorded span, got %d", len(ended)) + } + return ended[0] +} + +func TestIsPostHogAIGatewayURL(t *testing.T) { + cases := []struct { + in string + want bool + }{ + {"gateway.us.posthog.com", true}, + {"https://gateway.us.posthog.com/v1", true}, + {"GATEWAY.US.POSTHOG.COM", true}, + {"ai-gateway.eu.posthog.com", true}, + {"gateway.us.posthog.com/v1/chat", true}, + {"api.openai.com", false}, + {"https://us.i.posthog.com", false}, + {"", false}, + } + for _, tc := range cases { + if got := isPostHogAIGatewayURL(tc.in); got != tc.want { + t.Errorf("isPostHogAIGatewayURL(%q) = %v, want %v", tc.in, got, tc.want) + } + } +} + +func TestWarnIfPostHogAIGateway(t *testing.T) { + var buf bytes.Buffer + orig := log.Writer() + log.SetOutput(&buf) + t.Cleanup(func() { log.SetOutput(orig) }) + + cases := []struct { + name string + attrKey string + attrVal string + wantWarn bool + }{ + {"server.address gateway host", "server.address", "gateway.us.posthog.com", true}, + {"url.full gateway host", "url.full", "https://ai-gateway.eu.posthog.com/v1/chat", true}, + {"non-gateway server.address", "server.address", "api.openai.com", false}, + {"gateway host on unrelated attribute", "http.url", "gateway.us.posthog.com", false}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + buf.Reset() + span := recordSpanWithAttr(t, "gen_ai.chat", tc.attrKey, tc.attrVal) + warnIfPostHogAIGateway(span) + warned := strings.Contains(buf.String(), "PostHog AI Gateway") + if warned != tc.wantWarn { + t.Errorf("warned = %v, want %v (log=%q)", warned, tc.wantWarn, buf.String()) + } + }) + } +} + +func TestNewSpanProcessorRejectsEmptyAPIKey(t *testing.T) { + if _, err := NewSpanProcessor(context.Background(), " "); err != errEmptyAPIKey { + t.Errorf("expected errEmptyAPIKey, got %v", err) + } +} + +func TestNewExporterRejectsEmptyAPIKey(t *testing.T) { + if _, err := NewExporter(context.Background(), ""); err != errEmptyAPIKey { + t.Errorf("expected errEmptyAPIKey, got %v", err) + } +} + +func TestExporterTrimsAPIKeyWhitespace(t *testing.T) { + server := newOTLPServer(t) + exporter, err := NewExporter(context.Background(), " phc_test\n", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewExporter: %v", err) + } + defer exporter.Shutdown(context.Background()) + + span := recordSpan(t, "gen_ai.chat") + if err := exporter.ExportSpans(context.Background(), []sdktrace.ReadOnlySpan{span}); err != nil { + t.Fatalf("ExportSpans: %v", err) + } + + _, auth, _, _ := server.snapshot() + if want := "Bearer phc_test"; auth != want { + t.Errorf("Authorization = %q, want %q", auth, want) + } +} + +// otlpServer records the export requests it receives. +type otlpServer struct { + server *httptest.Server + mu sync.Mutex + auth string + path string + names []string + calls int + perCall []int +} + +func newOTLPServer(t *testing.T) *otlpServer { + t.Helper() + s := &otlpServer{} + s.server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + req := &coltracepb.ExportTraceServiceRequest{} + if err := proto.Unmarshal(body, req); err != nil { + t.Errorf("failed to decode OTLP request: %v", err) + } + s.mu.Lock() + s.calls++ + s.auth = r.Header.Get("Authorization") + s.path = r.URL.Path + count := 0 + for _, rs := range req.GetResourceSpans() { + for _, ss := range rs.GetScopeSpans() { + for _, span := range ss.GetSpans() { + s.names = append(s.names, span.GetName()) + count++ + } + } + } + s.perCall = append(s.perCall, count) + s.mu.Unlock() + + resp, _ := proto.Marshal(&coltracepb.ExportTraceServiceResponse{}) + w.Header().Set("Content-Type", "application/x-protobuf") + _, _ = w.Write(resp) + })) + t.Cleanup(s.server.Close) + return s +} + +func (s *otlpServer) snapshot() (calls int, auth, path string, names []string) { + s.mu.Lock() + defer s.mu.Unlock() + return s.calls, s.auth, s.path, append([]string(nil), s.names...) +} + +// batchSizes returns the number of spans carried by each export request, in +// the order the requests arrived. +func (s *otlpServer) batchSizes() []int { + s.mu.Lock() + defer s.mu.Unlock() + return append([]int(nil), s.perCall...) +} + +// emitSpans sends one AI span and one non-AI span through the provider. +func emitSpans(provider *sdktrace.TracerProvider) { + tracer := provider.Tracer("test") + _, ai := tracer.Start(context.Background(), "gen_ai.chat") + ai.End() + _, other := tracer.Start(context.Background(), "http.request") + other.End() +} + +func TestSpanProcessorExportsOnlyAISpans(t *testing.T) { + server := newOTLPServer(t) + + processor, err := NewSpanProcessor(context.Background(), "phc_test", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewSpanProcessor: %v", err) + } + provider := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(processor)) + emitSpans(provider) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := provider.ForceFlush(ctx); err != nil { + t.Fatalf("ForceFlush: %v", err) + } + if err := provider.Shutdown(ctx); err != nil { + t.Fatalf("Shutdown: %v", err) + } + + _, auth, path, names := server.snapshot() + if want := "Bearer phc_test"; auth != want { + t.Errorf("Authorization = %q, want %q", auth, want) + } + if want := ingestPath; path != want { + t.Errorf("path = %q, want %q", path, want) + } + if len(names) != 1 || names[0] != "gen_ai.chat" { + t.Errorf("exported span names = %v, want [gen_ai.chat]", names) + } +} + +func TestSpanProcessorDropsNonAISpansWithoutRequest(t *testing.T) { + server := newOTLPServer(t) + + processor, err := NewSpanProcessor(context.Background(), "phc_test", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewSpanProcessor: %v", err) + } + provider := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(processor)) + + _, span := provider.Tracer("test").Start(context.Background(), "http.request") + span.End() + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := provider.ForceFlush(ctx); err != nil { + t.Fatalf("ForceFlush: %v", err) + } + if err := provider.Shutdown(ctx); err != nil { + t.Fatalf("Shutdown: %v", err) + } + + if calls, _, _, _ := server.snapshot(); calls != 0 { + t.Errorf("expected no export request, got %d", calls) + } +} + +func TestSpanProcessorKeepsBatchesWithinEndpointLimit(t *testing.T) { + server := newOTLPServer(t) + + processor, err := NewSpanProcessor(context.Background(), "phc_test", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewSpanProcessor: %v", err) + } + provider := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(processor)) + + // Emit more AI spans than the endpoint accepts in a single request. With the + // SDK default batch size (512) they would be sent as one oversized request + // that the endpoint rejects with a non-retryable 400. + const total = 2*maxSpansPerRequest + 5 + tracer := provider.Tracer("test") + for i := 0; i < total; i++ { + _, span := tracer.Start(context.Background(), "gen_ai.chat") + span.End() + } + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := provider.ForceFlush(ctx); err != nil { + t.Fatalf("ForceFlush: %v", err) + } + if err := provider.Shutdown(ctx); err != nil { + t.Fatalf("Shutdown: %v", err) + } + + _, _, _, names := server.snapshot() + if len(names) != total { + t.Errorf("exported %d spans, want %d", len(names), total) + } + for i, n := range server.batchSizes() { + if n > maxSpansPerRequest { + t.Errorf("request %d carried %d spans, exceeds endpoint limit %d", i, n, maxSpansPerRequest) + } + } +} + +func TestExporterWarnsIfPostHogAIGateway(t *testing.T) { + server := newOTLPServer(t) + exporter, err := NewExporter(context.Background(), "phc_test", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewExporter: %v", err) + } + defer exporter.Shutdown(context.Background()) + + var buf bytes.Buffer + orig := log.Writer() + log.SetOutput(&buf) + t.Cleanup(func() { log.SetOutput(orig) }) + + span := recordSpanWithAttr(t, "gen_ai.chat", "server.address", "gateway.us.posthog.com") + if err := exporter.ExportSpans(context.Background(), []sdktrace.ReadOnlySpan{span}); err != nil { + t.Fatalf("ExportSpans: %v", err) + } + if !strings.Contains(buf.String(), "PostHog AI Gateway") { + t.Errorf("expected gateway warning, got %q", buf.String()) + } +} + +func TestExporterExportsOnlyAISpans(t *testing.T) { + server := newOTLPServer(t) + + exporter, err := NewExporter(context.Background(), "phc_test", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewExporter: %v", err) + } + provider := sdktrace.NewTracerProvider(sdktrace.WithBatcher(exporter)) + emitSpans(provider) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := provider.Shutdown(ctx); err != nil { + t.Fatalf("Shutdown: %v", err) + } + + _, _, path, names := server.snapshot() + if want := ingestPath; path != want { + t.Errorf("path = %q, want %q", path, want) + } + if len(names) != 1 || names[0] != "gen_ai.chat" { + t.Errorf("exported span names = %v, want [gen_ai.chat]", names) + } +} + +func TestExporterKeepsBatchesWithinEndpointLimit(t *testing.T) { + server := newOTLPServer(t) + + exporter, err := NewExporter(context.Background(), "phc_test", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewExporter: %v", err) + } + // WithBatcher uses the SDK default batch size (512), so without chunking the + // exporter would hand more than the endpoint's limit to a single request. + provider := sdktrace.NewTracerProvider(sdktrace.WithBatcher(exporter)) + + const total = 2*maxSpansPerRequest + 5 + tracer := provider.Tracer("test") + for i := 0; i < total; i++ { + _, span := tracer.Start(context.Background(), "gen_ai.chat") + span.End() + } + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := provider.Shutdown(ctx); err != nil { + t.Fatalf("Shutdown: %v", err) + } + + _, _, _, names := server.snapshot() + if len(names) != total { + t.Errorf("exported %d spans, want %d", len(names), total) + } + for i, n := range server.batchSizes() { + if n > maxSpansPerRequest { + t.Errorf("request %d carried %d spans, exceeds endpoint limit %d", i, n, maxSpansPerRequest) + } + } +} + +func TestExporterExportSpansSkipsRequestWhenNoAISpans(t *testing.T) { + server := newOTLPServer(t) + + exporter, err := NewExporter(context.Background(), "phc_test", WithHost(server.server.URL)) + if err != nil { + t.Fatalf("NewExporter: %v", err) + } + defer exporter.Shutdown(context.Background()) + + span := recordSpan(t, "http.request") + if err := exporter.ExportSpans(context.Background(), []sdktrace.ReadOnlySpan{span}); err != nil { + t.Fatalf("ExportSpans: %v", err) + } + + if calls, _, _, _ := server.snapshot(); calls != 0 { + t.Errorf("expected no export request, got %d", calls) + } +} + +func TestWithHostFallsBackToDefault(t *testing.T) { + cfg, err := newConfig(WithHost(" ")) + if err != nil { + t.Fatalf("newConfig: %v", err) + } + if cfg.host != DefaultHost { + t.Errorf("host = %q, want default %q", cfg.host, DefaultHost) + } + cfg, err = newConfig(WithHost("https://eu.i.posthog.com/")) + if err != nil { + t.Fatalf("newConfig: %v", err) + } + if cfg.host != "https://eu.i.posthog.com" { + t.Errorf("host = %q, want trailing slash trimmed", cfg.host) + } +} + +func TestNewConfigRejectsInvalidHost(t *testing.T) { + invalid := []string{ + "us.i.posthog.com", // missing scheme + "https://a b.example.com", // space in host + "https://ex.com:port", // invalid port + "http://[::1", // malformed host + "ftp://example.com", // wrong scheme + "https://", // no hostname + } + for _, host := range invalid { + if _, err := newConfig(WithHost(host)); !errors.Is(err, errInvalidHost) { + t.Errorf("newConfig(WithHost(%q)) err = %v, want errInvalidHost", host, err) + } + } + + valid := []struct { + host string + wantEndpoint string + }{ + {"https://us.i.posthog.com", "https://us.i.posthog.com" + ingestPath}, + {"https://eu.i.posthog.com/", "https://eu.i.posthog.com" + ingestPath}, + {"http://localhost:8000", "http://localhost:8000" + ingestPath}, + // A host carrying a query or fragment must still resolve to ingestPath. + // Concatenating would leave the path empty, and the OTLP exporter would + // fall back to its "/v1/traces" default. + {"https://us.i.posthog.com?region=eu", "https://us.i.posthog.com" + ingestPath + "?region=eu"}, + {"https://us.i.posthog.com#frag", "https://us.i.posthog.com" + ingestPath + "#frag"}, + {"https://proxy.example.com/posthog", "https://proxy.example.com/posthog" + ingestPath}, + } + for _, tc := range valid { + cfg, err := newConfig(WithHost(tc.host)) + if err != nil { + t.Errorf("newConfig(WithHost(%q)) err = %v, want nil", tc.host, err) + continue + } + if cfg.endpoint != tc.wantEndpoint { + t.Errorf("newConfig(WithHost(%q)).endpoint = %q, want %q", tc.host, cfg.endpoint, tc.wantEndpoint) + } + } +} + +func TestConstructorsRejectInvalidHost(t *testing.T) { + if _, err := NewSpanProcessor(context.Background(), "phc_test", WithHost("us.i.posthog.com")); !errors.Is(err, errInvalidHost) { + t.Errorf("NewSpanProcessor err = %v, want errInvalidHost", err) + } + if _, err := NewExporter(context.Background(), "phc_test", WithHost("us.i.posthog.com")); !errors.Is(err, errInvalidHost) { + t.Errorf("NewExporter err = %v, want errInvalidHost", err) + } +} + +func TestExporterTargetsIngestPathForHostWithQuery(t *testing.T) { + server := newOTLPServer(t) + + // The OTLP exporter silently falls back to "/v1/traces" when the endpoint + // URL has no path, so assert the request the server actually receives. + exporter, err := NewExporter(context.Background(), "phc_test", WithHost(server.server.URL+"?region=eu")) + if err != nil { + t.Fatalf("NewExporter: %v", err) + } + defer exporter.Shutdown(context.Background()) + + span := recordSpan(t, "gen_ai.chat") + if err := exporter.ExportSpans(context.Background(), []sdktrace.ReadOnlySpan{span}); err != nil { + t.Fatalf("ExportSpans: %v", err) + } + + calls, _, path, _ := server.snapshot() + if calls != 1 { + t.Fatalf("export requests = %d, want 1", calls) + } + if path != ingestPath { + t.Errorf("path = %q, want %q", path, ingestPath) + } +} diff --git a/otel/processor.go b/otel/processor.go new file mode 100644 index 00000000..b56fda78 --- /dev/null +++ b/otel/processor.go @@ -0,0 +1,59 @@ +package posthogotel + +import ( + "context" + "strings" + + sdktrace "go.opentelemetry.io/otel/sdk/trace" +) + +// SpanProcessor is a self-contained span processor that keeps only AI spans, +// batches them, and exports them to PostHog. Register it with +// TracerProvider.RegisterSpanProcessor or the sdktrace.WithSpanProcessor option. +type SpanProcessor struct { + inner sdktrace.SpanProcessor +} + +var _ sdktrace.SpanProcessor = (*SpanProcessor)(nil) + +// NewSpanProcessor builds a SpanProcessor for the given project API key. It +// wraps a batch span processor around the PostHog OTLP exporter. +func NewSpanProcessor(ctx context.Context, apiKey string, opts ...Option) (*SpanProcessor, error) { + if strings.TrimSpace(apiKey) == "" { + return nil, errEmptyAPIKey + } + cfg, err := newConfig(opts...) + if err != nil { + return nil, err + } + exporter, err := newOTLPExporter(ctx, apiKey, cfg) + if err != nil { + return nil, err + } + return &SpanProcessor{inner: sdktrace.NewBatchSpanProcessor( + exporter, + sdktrace.WithMaxExportBatchSize(maxSpansPerRequest), + )}, nil +} + +// OnStart does no work. Filtering happens in OnEnd, once the span is complete. +func (p *SpanProcessor) OnStart(context.Context, sdktrace.ReadWriteSpan) {} + +// OnEnd forwards AI spans to the batch processor and drops the rest. +func (p *SpanProcessor) OnEnd(s sdktrace.ReadOnlySpan) { + if !IsAISpan(s) { + return + } + warnIfPostHogAIGateway(s) + p.inner.OnEnd(s) +} + +// Shutdown shuts down the underlying batch span processor. +func (p *SpanProcessor) Shutdown(ctx context.Context) error { + return p.inner.Shutdown(ctx) +} + +// ForceFlush flushes the pending AI spans through the batch span processor. +func (p *SpanProcessor) ForceFlush(ctx context.Context) error { + return p.inner.ForceFlush(ctx) +} diff --git a/otel/spans.go b/otel/spans.go new file mode 100644 index 00000000..8f2bfeeb --- /dev/null +++ b/otel/spans.go @@ -0,0 +1,45 @@ +package posthogotel + +import ( + "strings" + + sdktrace "go.opentelemetry.io/otel/sdk/trace" +) + +// aiSpanPrefixes are the known AI semantic convention prefixes. +var aiSpanPrefixes = []string{"gen_ai.", "llm.", "ai.", "traceloop."} + +// IsAISpan reports whether a span follows a known AI semantic convention. +// It returns true when the span name, or any of its attribute keys, starts +// with one of the aiSpanPrefixes. +func IsAISpan(span sdktrace.ReadOnlySpan) bool { + if hasAIPrefix(span.Name()) { + return true + } + for _, attr := range span.Attributes() { + if hasAIPrefix(string(attr.Key)) { + return true + } + } + return false +} + +func hasAIPrefix(s string) bool { + for _, prefix := range aiSpanPrefixes { + if strings.HasPrefix(s, prefix) { + return true + } + } + return false +} + +// filterAISpans returns the subset of spans that IsAISpan accepts. +func filterAISpans(spans []sdktrace.ReadOnlySpan) []sdktrace.ReadOnlySpan { + aiSpans := make([]sdktrace.ReadOnlySpan, 0, len(spans)) + for _, span := range spans { + if IsAISpan(span) { + aiSpans = append(aiSpans, span) + } + } + return aiSpans +}