From 9a46381348531e3aa86ac2963ea6f7ce4fe46d97 Mon Sep 17 00:00:00 2001 From: Ralf Grubenmann Date: Mon, 10 Aug 2026 09:04:39 +0200 Subject: [PATCH 1/4] feat: change openmeter bridge to Lago --- .github/workflows/pipeline.yaml | 4 +- justfile | 2 +- .../Dockerfile | 0 .../go.mod | 13 +- .../go.sum | 31 -- src/otlp-bridge/main.go | 318 ++++++++++++++++++ src/otlp-openmeter-bridge/main.go | 264 --------------- .../configmap.yaml | 4 +- .../deployment.yaml | 31 +- .../secret.yaml | 2 +- .../service.yaml | 7 +- values.yaml | 14 +- 12 files changed, 357 insertions(+), 333 deletions(-) rename src/{otlp-openmeter-bridge => otlp-bridge}/Dockerfile (100%) rename src/{otlp-openmeter-bridge => otlp-bridge}/go.mod (52%) rename src/{otlp-openmeter-bridge => otlp-bridge}/go.sum (59%) create mode 100644 src/otlp-bridge/main.go delete mode 100644 src/otlp-openmeter-bridge/main.go rename templates/{otlp-openmeter-bridge => otlp-bridge}/configmap.yaml (64%) rename templates/{otlp-openmeter-bridge => otlp-bridge}/deployment.yaml (53%) rename templates/{otlp-openmeter-bridge => otlp-bridge}/secret.yaml (86%) rename templates/{otlp-openmeter-bridge => otlp-bridge}/service.yaml (69%) diff --git a/.github/workflows/pipeline.yaml b/.github/workflows/pipeline.yaml index 18a6181..73a171d 100644 --- a/.github/workflows/pipeline.yaml +++ b/.github/workflows/pipeline.yaml @@ -28,8 +28,8 @@ jobs: include: - name: init context: src/init - - name: otlp-openmeter-bridge - context: src/otlp-openmeter-bridge + - name: otlp-bridge + context: src/otlp-bridge uses: ./.github/workflows/docker.yaml with: name: ${{ matrix.name }} diff --git a/justfile b/justfile index 7347ac3..e1d42d9 100644 --- a/justfile +++ b/justfile @@ -5,7 +5,7 @@ root_dir := `git rev-parse --show-toplevel` flake_dir := root_dir / "tools/nix" output_dir := root_dir / ".output" build_dir := output_dir / "build" -go_modules := "src/init src/otlp-openmeter-bridge" +go_modules := "src/init src/otlp-bridge" # Manage nix environment. [group('modules')] diff --git a/src/otlp-openmeter-bridge/Dockerfile b/src/otlp-bridge/Dockerfile similarity index 100% rename from src/otlp-openmeter-bridge/Dockerfile rename to src/otlp-bridge/Dockerfile diff --git a/src/otlp-openmeter-bridge/go.mod b/src/otlp-bridge/go.mod similarity index 52% rename from src/otlp-openmeter-bridge/go.mod rename to src/otlp-bridge/go.mod index 6f5e16c..8f4b66d 100644 --- a/src/otlp-openmeter-bridge/go.mod +++ b/src/otlp-bridge/go.mod @@ -1,25 +1,18 @@ -module github.com/sdsc-vllm/otlp-openmeter-bridge +module github.com/sdsc-vllm/otlp-bridge go 1.24.0 require ( - github.com/cloudevents/sdk-go/v2 v2.16.2 - github.com/google/uuid v1.6.0 go.opentelemetry.io/proto/otlp v1.10.0 google.golang.org/grpc v1.79.2 - google.golang.org/protobuf v1.36.11 ) require ( github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 // indirect - github.com/json-iterator/go v1.1.12 // indirect - github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect - github.com/modern-go/reflect2 v1.0.2 // indirect - go.uber.org/multierr v1.11.0 // indirect - go.uber.org/zap v1.27.0 // indirect golang.org/x/net v0.50.0 // indirect golang.org/x/sys v0.41.0 // indirect golang.org/x/text v0.34.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260209200024-4cfbd4190f57 // indirect -) + google.golang.org/protobuf v1.36.11 // indirect +) \ No newline at end of file diff --git a/src/otlp-openmeter-bridge/go.sum b/src/otlp-bridge/go.sum similarity index 59% rename from src/otlp-openmeter-bridge/go.sum rename to src/otlp-bridge/go.sum index 3a06317..830acec 100644 --- a/src/otlp-openmeter-bridge/go.sum +++ b/src/otlp-bridge/go.sum @@ -1,10 +1,5 @@ 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/cloudevents/sdk-go/v2 v2.16.2 h1:ZYDFrYke4FD+jM8TZTJJO6JhKHzOQl2oqpFK1D+NnQM= -github.com/cloudevents/sdk-go/v2 v2.16.2/go.mod h1:laOcGImm4nVJEU+PHnUrKL56CKmRL65RlQF0kRmW/kg= -github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -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.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= @@ -13,26 +8,10 @@ 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/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= 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/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= -github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= -github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= -github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd h1:TRLaZ9cD/w8PVh93nsPXa1VrQ6jlwL5oN8l14QlcNfg= -github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= -github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= -github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= -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/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= -github.com/stretchr/testify v1.11.0 h1:ib4sjIrwZKxE5u/Japgo/7SJV3PvgjGiRNAvTVGqQl8= -github.com/stretchr/testify v1.11.0/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= -github.com/valyala/bytebufferpool v1.0.0 h1:GqA5TC/0021Y/b9FG4Oi9Mr3q7XYx6KllzawFIhcdPw= -github.com/valyala/bytebufferpool v1.0.0/go.mod h1:6bBcMArwyJ5K/AmCkWv1jt77kVWyCJ6HpOuEn7z0Csc= 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.39.0 h1:8yPrr/S0ND9QEfTfdP9V+SiwT4E0G7Y5MO7p85nis48= @@ -47,20 +26,12 @@ go.opentelemetry.io/otel/trace v1.39.0 h1:2d2vfpEDmCJ5zVYz7ijaJdOF59xLomrvj7bjt6 go.opentelemetry.io/otel/trace v1.39.0/go.mod h1:88w4/PnZSazkGzz/w84VHpQafiU4EtqqlVdxWy+rNOA= 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= -go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= -go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= -go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= -go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= golang.org/x/net v0.50.0 h1:ucWh9eiCGyDR3vtzso0WMQinm2Dnt8cFMuQa9K33J60= golang.org/x/net v0.50.0/go.mod h1:UgoSli3F/pBgdJBHCTc+tp3gmrU4XswgGRgtnwWTfyM= golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk= golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= -golang.org/x/time v0.12.0 h1:ScB/8o8olJvc+CQPWrK3fPZNfh7qgwCrY0zJmoEQLSE= -golang.org/x/time v0.12.0/go.mod h1:CDIdPxbZBQxdj6cxyCIdrNogrJKMJ7pr37NYpMcMDSg= gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk= gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E= google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 h1:JLQynH/LBHfCTSbDWl+py8C+Rg/k1OVH3xfcaiANuF0= @@ -71,5 +42,3 @@ google.golang.org/grpc v1.79.2 h1:fRMD94s2tITpyJGtBBn7MkMseNpOZU8ZxgC3MMBaXRU= google.golang.org/grpc v1.79.2/go.mod h1:KmT0Kjez+0dde/v2j9vzwoAScgEPx/Bw1CYChhHLrHQ= 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/src/otlp-bridge/main.go b/src/otlp-bridge/main.go new file mode 100644 index 0000000..1c57512 --- /dev/null +++ b/src/otlp-bridge/main.go @@ -0,0 +1,318 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "log" + "net" + "net/http" + "os" + "os/signal" + "strconv" + "strings" + "syscall" + "time" + + "google.golang.org/grpc" + + logspb "go.opentelemetry.io/proto/otlp/collector/logs/v1" + commonv1 "go.opentelemetry.io/proto/otlp/common/v1" + logsrc "go.opentelemetry.io/proto/otlp/logs/v1" +) + +type Config struct { + RemoteAPIURL string + BearerToken string + GRPCListenAddr string + MaxRetries int + RetryWaitSeconds int + MetricCode string + CustomerID string + OrganizationSlug string +} + +// LagoEvent maps to the Lago Event schema for ingestion +type LagoEvent struct { + TransactionID string `json:"transaction_id"` + ExternalSubscriptionID string `json:"external_subscription_id"` + Code string `json:"code"` + Timestamp int64 `json:"timestamp"` + Properties map[string]interface{} `json:"properties,omitempty"` +} + +// LagoIngestRequest maps to the Lago Event ingestion request +type LagoIngestRequest struct { + Event LagoEvent `json:"event"` +} + +type server struct { + logspb.UnimplementedLogsServiceServer + config *Config + client *http.Client +} + +func loadConfig() *Config { + return &Config{ + RemoteAPIURL: os.Getenv("REMOTE_API_URL"), + BearerToken: os.Getenv("BEARER_TOKEN"), + GRPCListenAddr: os.Getenv("GRPC_LISTEN_ADDR"), + MaxRetries: 3, + RetryWaitSeconds: 2, + MetricCode: os.Getenv("METRIC_CODE"), + CustomerID: os.Getenv("CUSTOMER_ID"), + OrganizationSlug: os.Getenv("ORGANIZATION_SLUG"), + } +} + +func (s *server) Export(ctx context.Context, req *logspb.ExportLogsServiceRequest) (*logspb.ExportLogsServiceResponse, error) { + for _, rl := range req.GetResourceLogs() { + for _, sl := range rl.GetScopeLogs() { + for _, lr := range sl.GetLogRecords() { + if err := s.processLogRecord(ctx, lr); err != nil { + log.Printf("Error processing log record: %v", err) + } + } + } + } + return &logspb.ExportLogsServiceResponse{}, nil +} + +func (s *server) processLogRecord(ctx context.Context, lr *logsrc.LogRecord) error { + attrs := make(map[string]string) + for _, kv := range lr.Attributes { + attrs[kv.Key] = AnyValueToString(kv.Value) + } + + // Extract fields from OTLP attributes + method := attrs["method"] + path := attrs["request.path"] + duration := attrs["duration"] + responseCode := attrs["response_code"] + + if code, err := strconv.Atoi(responseCode); err != nil || code > 300 { + // only log actual requests + return nil + } + + xForwardedFor := attrs["x-forwarded-for"] + subject := attrs["x-sub"] + xUserName := attrs["x-user-name"] + customerId := s.config.CustomerID + if customerId == "" { + log.Printf("No CUSTOMER_ID configured, skipping event") + return nil + } + xRequestID := attrs["x-request-id"] + genAIRequestModel := attrs["gen_ai.request.model"] + genAIResponseModel := attrs["gen_ai.response.model"] + genAIProviderName := attrs["gen_ai.provider.name"] + genAIUsageInput := attrs["gen_ai.usage.input_tokens"] + genAIUsageOutput := attrs["gen_ai.usage.output_tokens"] + genAIUsageTotal := attrs["gen_ai.usage.total_tokens"] + + // Determine metric code from env or derive from context + metricCode := s.config.MetricCode + if metricCode == "" { + metricCode = "genai.tokens" + } + + // Build properties from OTLP attributes + properties := make(map[string]interface{}) + if method != "" { + properties["method"] = method + } + if path != "" { + properties["request_path"] = path + } + if duration != "" { + properties["duration"] = duration + } + if responseCode != "" { + properties["response_code"] = responseCode + } + if xForwardedFor != "" { + properties["x_forwarded_for"] = xForwardedFor + } + if subject != "" { + properties["x-sub"] = subject + } + if xUserName != "" { + properties["x-user-name"] = xUserName + } + if xRequestID != "" { + properties["x_request_id"] = xRequestID + } + if genAIRequestModel != "" { + properties["gen_ai.request.model"] = genAIRequestModel + } + if genAIResponseModel != "" { + properties["gen_ai.response.model"] = genAIResponseModel + } + if genAIProviderName != "" { + properties["gen_ai.provider.name"] = genAIProviderName + } + if genAIUsageInput != "" { + properties["gen_ai.usage.input_tokens"] = genAIUsageInput + } + if genAIUsageOutput != "" { + properties["gen_ai.usage.output_tokens"] = genAIUsageOutput + } + if genAIUsageTotal != "" { + properties["gen_ai.usage.total_tokens"] = genAIUsageTotal + } + + // Generate structured transaction ID: {type}_{date}_{customer}_{category}_{request_id} + // Example: genai_20240314_cust42_gpt4_x-request-id-123 + dateStr := time.Now().UTC().Format("20060102") + modelSlug := "unknown" + if genAIRequestModel != "" { + modelSlug = genAIRequestModel + } else if genAIResponseModel != "" { + modelSlug = genAIResponseModel + } + txnID := fmt.Sprintf("genai_%s_%s_%s_%s", dateStr, customerId, modelSlug, xRequestID) + + // Use the log record's timestamp if available, otherwise use current time + var timestamp int64 + if lr.GetTimeUnixNano() > 0 { + timestamp = int64(lr.GetTimeUnixNano() / 1000000000) // convert nanoseconds to seconds + } else { + timestamp = time.Now().UTC().Unix() + } + + // Build the Lago event + event := LagoEvent{ + TransactionID: txnID, + ExternalSubscriptionID: customerId, + Code: metricCode, + Timestamp: timestamp, + Properties: properties, + } + + return s.sendWithRetry(ctx, event) +} + +func AnyValueToString(av *commonv1.AnyValue) string { + switch av.Value.(type) { + case *commonv1.AnyValue_StringValue: + return av.GetStringValue() + case *commonv1.AnyValue_IntValue: + return fmt.Sprintf("%d", av.GetIntValue()) + case *commonv1.AnyValue_DoubleValue: + return fmt.Sprintf("%f", av.GetDoubleValue()) + case *commonv1.AnyValue_BoolValue: + return fmt.Sprintf("%t", av.GetBoolValue()) + case *commonv1.AnyValue_ArrayValue: + // Recursively convert array elements + return fmt.Sprintf("%v", av.GetArrayValue()) + case *commonv1.AnyValue_KvlistValue: + // Recursively convert map values + m := make(map[string]interface{}) + for _, kv := range av.GetKvlistValue().GetValues() { + m[kv.Key] = AnyValueToString(kv.Value) + } + return fmt.Sprintf("%v", m) + default: + return "" + } +} + +func (s *server) sendWithRetry(ctx context.Context, event LagoEvent) error { + var err error + for i := 0; i <= s.config.MaxRetries; i++ { + if i > 0 { + log.Printf("Retrying send (%d/%d) after error: %v", i, s.config.MaxRetries, err) + time.Sleep(time.Duration(s.config.RetryWaitSeconds) * time.Second) + } + + err = s.sendEvent(ctx, event) + if err == nil { + return nil + } + + // Check if it's a 429 rate limit error - use backoff + if strings.Contains(err.Error(), "status 429") { + backoff := time.Duration(i+1) * time.Duration(s.config.RetryWaitSeconds) * time.Second + log.Printf("Rate limited, backing off for %v before retry", backoff) + time.Sleep(backoff) + } + } + return fmt.Errorf("failed to send event after %d retries: %w", s.config.MaxRetries, err) +} + +func (s *server) sendEvent(ctx context.Context, event LagoEvent) error { + reqBody := LagoIngestRequest{ + Event: event, + } + + body, err := json.Marshal(reqBody) + if err != nil { + return fmt.Errorf("failed to marshal request: %w", err) + } + + url := fmt.Sprintf("%s/api/v1/events", s.config.RemoteAPIURL) + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(body)) + if err != nil { + return fmt.Errorf("failed to create request: %w", err) + } + + req.Header.Set("Content-Type", "application/json") + if s.config.BearerToken != "" { + req.Header.Set("Authorization", "Bearer "+s.config.BearerToken) + } + + resp, err := s.client.Do(req) + if err != nil { + return fmt.Errorf("failed to send request: %w", err) + } + defer resp.Body.Close() + + respBody, _ := io.ReadAll(resp.Body) + + if resp.StatusCode >= 300 { + return fmt.Errorf("remote API returned status %d: %s", resp.StatusCode, string(respBody)) + } + + log.Printf("Successfully ingested event: transaction_id=%s external_subscription_id=%s", event.TransactionID, event.ExternalSubscriptionID) + return nil +} + +func main() { + cfg := loadConfig() + if cfg.RemoteAPIURL == "" || cfg.GRPCListenAddr == "" || cfg.OrganizationSlug == "" { + log.Fatal("REMOTE_API_URL, GRPC_LISTEN_ADDR, and ORGANIZATION_SLUG must be set") + } + + lis, err := net.Listen("tcp", cfg.GRPCListenAddr) + if err != nil { + log.Fatalf("failed to listen: %v", err) + } + + s := grpc.NewServer() + srv := &server{ + config: cfg, + client: &http.Client{ + Timeout: 10 * time.Second, + }, + } + + logspb.RegisterLogsServiceServer(s, srv) + + stop := make(chan os.Signal, 1) + signal.Notify(stop, os.Interrupt, syscall.SIGTERM) + + go func() { + log.Printf("Starting gRPC server on %s", cfg.GRPCListenAddr) + if err := s.Serve(lis); err != nil { + log.Fatalf("failed to serve: %v", err) + } + }() + + <-stop + log.Println("Shutting down gRPC server...") + s.GracefulStop() + log.Println("Server stopped") +} diff --git a/src/otlp-openmeter-bridge/main.go b/src/otlp-openmeter-bridge/main.go deleted file mode 100644 index ec8a0ea..0000000 --- a/src/otlp-openmeter-bridge/main.go +++ /dev/null @@ -1,264 +0,0 @@ -package main - -import ( - "bytes" - "context" - "encoding/json" - "fmt" - "io" - "log" - "net" - "net/http" - "os" - "os/signal" - "syscall" - "time" - - cloudevents "github.com/cloudevents/sdk-go/v2" - "github.com/google/uuid" - "google.golang.org/grpc" - - logspb "go.opentelemetry.io/proto/otlp/collector/logs/v1" - commonv1 "go.opentelemetry.io/proto/otlp/common/v1" - logsrc "go.opentelemetry.io/proto/otlp/logs/v1" -) - -type Config struct { - RemoteAPIURL string - BearerToken string - GRPCListenAddr string - MaxRetries int - RetryWaitSeconds int - EventType string -} - -type UserData struct { - Key string `json:"key"` - Name string `json:"name"` -} - -type CloudEventData struct { - Method string `json:"method"` - Path string `json:"request_path"` - Duration string `json:"duration"` - ResponseCode string `json:"response_code"` - XForwardedFor string `json:"x_forwarded_for"` - XSub string `json:"x_sub"` - XUserName string `json:"x_user_name"` - XRequestID string `json:"x_request_id"` - GenAIRequestModel string `json:"model"` - GenAIResponseModel string `json:"response_model"` - GenAIProviderName string `json:"provider"` - GenAIUsageInput string `json:"input_tokens"` - GenAIUsageOutput string `json:"output_tokens"` - GenAIUsageTotal string `json:"total_tokens"` -} - -type server struct { - logspb.UnimplementedLogsServiceServer - config *Config - client *http.Client -} - -func loadConfig() *Config { - return &Config{ - RemoteAPIURL: os.Getenv("REMOTE_API_URL"), - BearerToken: os.Getenv("BEARER_TOKEN"), - GRPCListenAddr: os.Getenv("GRPC_LISTEN_ADDR"), - MaxRetries: 3, - RetryWaitSeconds: 2, - EventType: os.Getenv("EVENT_TYPE"), - } -} - -func (s *server) Export(ctx context.Context, req *logspb.ExportLogsServiceRequest) (*logspb.ExportLogsServiceResponse, error) { - for _, rl := range req.GetResourceLogs() { - for _, sl := range rl.GetScopeLogs() { - for _, lr := range sl.GetLogRecords() { - if err := s.processLogRecord(ctx, lr); err != nil { - log.Printf("Error processing log record: %v", err) - } - } - } - } - return &logspb.ExportLogsServiceResponse{}, nil -} - -func (s *server) processLogRecord(ctx context.Context, lr *logsrc.LogRecord) error { - attrs := make(map[string]string) - for _, kv := range lr.Attributes { - attrs[kv.Key] = AnyValueToString(kv.Value) - } - - data := CloudEventData{ - Method: attrs["method"], - Path: attrs["request.path"], - Duration: attrs["duration"], - ResponseCode: attrs["response_code"], - XForwardedFor: attrs["x-forwarded-for"], - XSub: attrs["x-sub"], - XUserName: attrs["x-user-name"], - XRequestID: attrs["x-request-id"], - GenAIRequestModel: attrs["gen_ai.request.model"], - GenAIResponseModel: attrs["gen_ai.response.model"], - GenAIProviderName: attrs["gen_ai.provider.name"], - GenAIUsageInput: attrs["gen_ai.usage.input_tokens"], - GenAIUsageOutput: attrs["gen_ai.usage.output_tokens"], - GenAIUsageTotal: attrs["gen_ai.usage.total_tokens"], - } - - subject := attrs["x-sub"] - - event := cloudevents.NewEvent() - event.SetID(uuid.New().String()) - event.SetSource("otlp-openmeter-bridge") - event.SetType(s.config.EventType) - event.SetSubject(subject) - event.SetTime(time.Now()) - - if err := event.SetData(cloudevents.ApplicationJSON, data); err != nil { - return err - } - - return s.sendWithRetry(ctx, event, data) -} - -func AnyValueToString(av *commonv1.AnyValue) string { - switch av.Value.(type) { - case *commonv1.AnyValue_StringValue: - return av.GetStringValue() - case *commonv1.AnyValue_IntValue: - return fmt.Sprintf("%d", av.GetIntValue()) - case *commonv1.AnyValue_DoubleValue: - return fmt.Sprintf("%f", av.GetDoubleValue()) - case *commonv1.AnyValue_BoolValue: - return fmt.Sprintf("%t", av.GetBoolValue()) - case *commonv1.AnyValue_ArrayValue: - // Recursively convert array elements - return fmt.Sprintf("%v", av.GetArrayValue()) - case *commonv1.AnyValue_KvlistValue: - // Recursively convert map values - m := make(map[string]interface{}) - for _, kv := range av.GetKvlistValue().GetValues() { - m[kv.Key] = AnyValueToString(kv.Value) - } - return fmt.Sprintf("%v", m) - default: - return "" - } -} - -func (s *server) sendWithRetry(ctx context.Context, event cloudevents.Event, eventData CloudEventData) error { - var err error - for i := 0; i <= s.config.MaxRetries; i++ { - if i > 0 { - log.Printf("Retrying send (%d/%d) after error: %v", i, s.config.MaxRetries, err) - time.Sleep(time.Duration(s.config.RetryWaitSeconds) * time.Second) - } - - err = s.sendEvent(ctx, event, eventData) - if err == nil { - return nil - } - } - return fmt.Errorf("failed to send event after %d retries: %w", s.config.MaxRetries, err) -} - -func (s *server) sendEvent(ctx context.Context, event cloudevents.Event, eventData CloudEventData) error { - // create consumer first - userdata := UserData{ - Key: event.Subject(), - Name: eventData.XUserName, - } - b, err := json.Marshal(userdata) - if err != nil { - return err - } - log.Printf("creating customer at %s", fmt.Sprintf("%s/v3/openmeter/customers", s.config.RemoteAPIURL)) - req, err := http.NewRequestWithContext(ctx, http.MethodPost, fmt.Sprintf("%s/v3/openmeter/customers", s.config.RemoteAPIURL), bytes.NewReader(b)) - req.Header.Set("Content-Type", "application/json") - if s.config.BearerToken != "" { - req.Header.Set("Authorization", "Bearer "+s.config.BearerToken) - } - resp, err := s.client.Do(req) - if err != nil { - return err - } - defer resp.Body.Close() - - if resp.StatusCode >= 300 && resp.StatusCode != 409 { - body, _ := io.ReadAll(resp.Body) - return fmt.Errorf("remote API returned status %d: %s", resp.StatusCode, string(body)) - } - - b, err = json.Marshal(event) - if err != nil { - return err - } - log.Printf("creating event at %s", fmt.Sprintf("%s/v3/openmeter/events", s.config.RemoteAPIURL)) - req, err = http.NewRequestWithContext(ctx, http.MethodPost, fmt.Sprintf("%s/v3/openmeter/events", s.config.RemoteAPIURL), bytes.NewReader(b)) - if err != nil { - return err - } - - req.Header.Set("Content-Type", "application/json") - req.Header.Set("ce-specversion", event.SpecVersion()) - req.Header.Set("ce-type", event.Type()) - req.Header.Set("ce-source", event.Source()) - req.Header.Set("ce-subject", event.Subject()) - req.Header.Set("ce-id", event.ID()) - - if s.config.BearerToken != "" { - req.Header.Set("Authorization", "Bearer "+s.config.BearerToken) - } - - resp, err = s.client.Do(req) - if err != nil { - return err - } - defer resp.Body.Close() - - if resp.StatusCode >= 300 { - body, _ := io.ReadAll(resp.Body) - return fmt.Errorf("remote API returned status %d: %s", resp.StatusCode, string(body)) - } - - return nil -} - -func main() { - cfg := loadConfig() - if cfg.RemoteAPIURL == "" || cfg.GRPCListenAddr == "" { - log.Fatal("REMOTE_API_URL and GRPC_LISTEN_ADDR must be set") - } - - lis, err := net.Listen("tcp", cfg.GRPCListenAddr) - if err != nil { - log.Fatalf("failed to listen: %v", err) - } - - s := grpc.NewServer() - srv := &server{ - config: cfg, - client: &http.Client{ - Timeout: 10 * time.Second, - }, - } - - logspb.RegisterLogsServiceServer(s, srv) - - stop := make(chan os.Signal, 1) - signal.Notify(stop, os.Interrupt, syscall.SIGTERM) - - go func() { - log.Printf("Starting gRPC server on %s", cfg.GRPCListenAddr) - if err := s.Serve(lis); err != nil { - log.Fatalf("failed to serve: %v", err) - } - }() - - <-stop - log.Println("Shutting down gRPC server...") - s.GracefulStop() - log.Println("Server stopped") -} diff --git a/templates/otlp-openmeter-bridge/configmap.yaml b/templates/otlp-bridge/configmap.yaml similarity index 64% rename from templates/otlp-openmeter-bridge/configmap.yaml rename to templates/otlp-bridge/configmap.yaml index 62902a3..ea6cd13 100644 --- a/templates/otlp-openmeter-bridge/configmap.yaml +++ b/templates/otlp-bridge/configmap.yaml @@ -2,10 +2,10 @@ apiVersion: v1 kind: ConfigMap metadata: - name: {{ include "telemetry.fullname" . }} + name: otlp-bridge-config namespace: {{ .Release.Namespace }} labels: release: {{ .Release.Name }} data: - remote-api-url: {{ .Values.envoy.telemetry.openmeterUrl | quote }} + remote-api-url: {{ .Values.envoy.telemetry.lagoUrl | quote }} {{- end }} diff --git a/templates/otlp-openmeter-bridge/deployment.yaml b/templates/otlp-bridge/deployment.yaml similarity index 53% rename from templates/otlp-openmeter-bridge/deployment.yaml rename to templates/otlp-bridge/deployment.yaml index c999984..2f16d0a 100644 --- a/templates/otlp-openmeter-bridge/deployment.yaml +++ b/templates/otlp-bridge/deployment.yaml @@ -2,31 +2,30 @@ apiVersion: apps/v1 kind: Deployment metadata: - name: {{ include "telemetry.fullname" . }} + name: otlp-bridge namespace: {{ .Release.Namespace }} labels: - app: otlp-openmeter-bridge + app: otlp-bridge release: {{ .Release.Name }} spec: - replicas: {{ .Values.envoy.telemetry.replicas }} + replicas: {{ .Values.envoy.telemetry.replicas | default 1 }} selector: matchLabels: - app: otlp-openmeter-bridge - release: {{ .Release.Name }} + app: otlp-bridge template: metadata: labels: - app: otlp-openmeter-bridge + app: otlp-bridge release: {{ .Release.Name }} spec: - {{- if .Values.global.imagePullSecrets }} + {{- if .Values.global.imagePullSecrets }} imagePullSecrets: {{ toYaml .Values.global.imagePullSecrets | indent 8 }} - {{- end }} + {{- end }} containers: - name: bridge - image: {{ .Values.envoy.telemetry.image.repository }}:{{ .Values.envoy.telemetry.image.tag | default .Chart.AppVersion }} - imagePullPolicy: {{ .Values.envoy.telemetry.image.pullPolicy }} + image: {{ .Values.envoy.telemetry.image.repository | default "ghcr.io/swissdatasciencecenter/sdsc-llm-deployment/otlp-bridge" }}:{{ .Values.envoy.telemetry.image.tag | default "latest" }} + imagePullPolicy: {{ .Values.envoy.telemetry.image.pullPolicy | default "IfNotPresent" }} ports: - containerPort: 4317 name: grpc @@ -34,16 +33,20 @@ spec: - name: REMOTE_API_URL valueFrom: configMapKeyRef: - name: {{ include "telemetry.fullname" . }} + name: otlp-bridge-config key: remote-api-url - name: GRPC_LISTEN_ADDR value: ":4317" - - name: EVENT_TYPE - value: "llm_call" + - name: METRIC_CODE + value: {{ .Values.envoy.telemetry.metricCode | default "genai.tokens" | quote }} + - name: CUSTOMER_ID + value: {{ .Values.envoy.telemetry.customerId | quote }} + - name: ORGANIZATION_SLUG + value: {{ .Values.envoy.telemetry.organizationSlug | quote }} - name: BEARER_TOKEN valueFrom: secretKeyRef: - name: {{ include "telemetry.fullname" . }} + name: otlp-bridge-secret key: bearer-token resources: {{- toYaml .Values.envoy.telemetry.resources | nindent 12 }} diff --git a/templates/otlp-openmeter-bridge/secret.yaml b/templates/otlp-bridge/secret.yaml similarity index 86% rename from templates/otlp-openmeter-bridge/secret.yaml rename to templates/otlp-bridge/secret.yaml index e2a7a8b..3d70d9f 100644 --- a/templates/otlp-openmeter-bridge/secret.yaml +++ b/templates/otlp-bridge/secret.yaml @@ -2,7 +2,7 @@ apiVersion: v1 kind: Secret metadata: - name: {{ include "telemetry.fullname" . }} + name: otlp-bridge-secret namespace: {{ .Release.Namespace }} labels: release: {{ .Release.Name }} diff --git a/templates/otlp-openmeter-bridge/service.yaml b/templates/otlp-bridge/service.yaml similarity index 69% rename from templates/otlp-openmeter-bridge/service.yaml rename to templates/otlp-bridge/service.yaml index 015c607..3c2d873 100644 --- a/templates/otlp-openmeter-bridge/service.yaml +++ b/templates/otlp-bridge/service.yaml @@ -2,15 +2,14 @@ apiVersion: v1 kind: Service metadata: - name: {{ include "telemetry.fullname" . }} + name: otlp-bridge namespace: {{ .Release.Namespace }} labels: - app: otlp-openmeter-bridge + app: otlp-bridge release: {{ .Release.Name }} spec: selector: - app: otlp-openmeter-bridge - release: {{ .Release.Name }} + app: otlp-bridge ports: - protocol: TCP port: 4317 diff --git a/values.yaml b/values.yaml index 2a464dd..f5d0cc4 100644 --- a/values.yaml +++ b/values.yaml @@ -15,13 +15,19 @@ envoy: resources: telemetry: enabled: false - # OpenMeter API endpoint for billing (required when telemetry is enabled) - openmeterUrl: "" - # Bearer token for authenticating with OpenMeter (always required when telemetry is enabled) + # Lago API endpoint for event ingestion (required when telemetry is enabled) + lagoUrl: "" + # Billable metric code for ingested events (default: genai.tokens) + metricCode: "" + # Fixed Lago customer ID for all ingested events + customerId: "" + # Lago organization slug for event routing + organizationSlug: "" + # Bearer token for authenticating with Lago (always required when telemetry is enabled) bearerToken: "" # OTLP receiver container settings image: - repository: ghcr.io/swissdatasciencecenter/llm-serving/otlp-openmeter-bridge + repository: ghcr.io/swissdatasciencecenter/llm-serving/otlp-bridge tag: "" # defaults to .Chart.AppVersion pullPolicy: IfNotPresent replicas: 1 From d2905152a3e89fbda86036e440e7b03485c551a2 Mon Sep 17 00:00:00 2001 From: cmdoret Date: Mon, 10 Aug 2026 14:28:00 +0200 Subject: [PATCH 2/4] fix(oltp): drop OrganizationSlug --- src/otlp-bridge/main.go | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/src/otlp-bridge/main.go b/src/otlp-bridge/main.go index 1c57512..f39dd81 100644 --- a/src/otlp-bridge/main.go +++ b/src/otlp-bridge/main.go @@ -31,7 +31,6 @@ type Config struct { RetryWaitSeconds int MetricCode string CustomerID string - OrganizationSlug string } // LagoEvent maps to the Lago Event schema for ingestion @@ -63,7 +62,6 @@ func loadConfig() *Config { RetryWaitSeconds: 2, MetricCode: os.Getenv("METRIC_CODE"), CustomerID: os.Getenv("CUSTOMER_ID"), - OrganizationSlug: os.Getenv("ORGANIZATION_SLUG"), } } @@ -282,8 +280,8 @@ func (s *server) sendEvent(ctx context.Context, event LagoEvent) error { func main() { cfg := loadConfig() - if cfg.RemoteAPIURL == "" || cfg.GRPCListenAddr == "" || cfg.OrganizationSlug == "" { - log.Fatal("REMOTE_API_URL, GRPC_LISTEN_ADDR, and ORGANIZATION_SLUG must be set") + if cfg.RemoteAPIURL == "" || cfg.GRPCListenAddr == "" { + log.Fatal("REMOTE_API_URL and GRPC_LISTEN_ADDR must be set") } lis, err := net.Listen("tcp", cfg.GRPCListenAddr) From a4cb4193b0d38db6732bad718ec7f000609a82d5 Mon Sep 17 00:00:00 2001 From: cmdoret Date: Mon, 10 Aug 2026 14:28:55 +0200 Subject: [PATCH 3/4] fix(lago): drop oganizationSlug --- templates/otlp-bridge/deployment.yaml | 2 -- values.yaml | 2 -- 2 files changed, 4 deletions(-) diff --git a/templates/otlp-bridge/deployment.yaml b/templates/otlp-bridge/deployment.yaml index 2f16d0a..8998ea3 100644 --- a/templates/otlp-bridge/deployment.yaml +++ b/templates/otlp-bridge/deployment.yaml @@ -41,8 +41,6 @@ spec: value: {{ .Values.envoy.telemetry.metricCode | default "genai.tokens" | quote }} - name: CUSTOMER_ID value: {{ .Values.envoy.telemetry.customerId | quote }} - - name: ORGANIZATION_SLUG - value: {{ .Values.envoy.telemetry.organizationSlug | quote }} - name: BEARER_TOKEN valueFrom: secretKeyRef: diff --git a/values.yaml b/values.yaml index f5d0cc4..4d701cb 100644 --- a/values.yaml +++ b/values.yaml @@ -21,8 +21,6 @@ envoy: metricCode: "" # Fixed Lago customer ID for all ingested events customerId: "" - # Lago organization slug for event routing - organizationSlug: "" # Bearer token for authenticating with Lago (always required when telemetry is enabled) bearerToken: "" # OTLP receiver container settings From 49cdeef60a40f5f037d995fa904f4618bf45fbe3 Mon Sep 17 00:00:00 2001 From: cmdoret Date: Mon, 10 Aug 2026 14:41:49 +0200 Subject: [PATCH 4/4] fix(envoy): authentik gate block --- templates/envoy/gateway.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/templates/envoy/gateway.yaml b/templates/envoy/gateway.yaml index ed598c3..ed8c530 100644 --- a/templates/envoy/gateway.yaml +++ b/templates/envoy/gateway.yaml @@ -37,10 +37,10 @@ spec: certificateRefs: - kind: Secret name: {{ include "envoy.fullname" . }}-wildcard-https + {{- end }} infrastructure: parametersRef: group: gateway.envoyproxy.io kind: EnvoyProxy name: {{ include "envoy.fullname" . }} - {{- end }} {{- end }}