From c8930599e0f27794668afc75bf0263c31c5eb9bf Mon Sep 17 00:00:00 2001 From: Hilmar Falkenberg Date: Fri, 19 Sep 2025 16:38:17 +0200 Subject: [PATCH 01/10] Add files via upload (#2) Signed-off-by: Hilmar Falkenberg --- .github/workflows/sync-fork.yaml | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) create mode 100644 .github/workflows/sync-fork.yaml diff --git a/.github/workflows/sync-fork.yaml b/.github/workflows/sync-fork.yaml new file mode 100644 index 000000000000..de520f2bc0eb --- /dev/null +++ b/.github/workflows/sync-fork.yaml @@ -0,0 +1,16 @@ +name: sync-fork +on: + schedule: + - cron: '4 2 * * 4,2' + workflow_dispatch: +jobs: + sync-fork: + runs-on: ubuntu-latest + permissions: + contents: write + steps: + - run: gh repo sync $REPOSITORY -b $BRANCH_NAME + env: + GITHUB_TOKEN: ${{ secrets.SYNC_FORK_TOKEN }} + REPOSITORY: ${{ github.repository }} + BRANCH_NAME: ${{ github.ref_name }} From 5702b74d45fdb5d6677ea3879290b9145dbc1c90 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Jarmolkiewicz?= Date: Fri, 26 Sep 2025 14:43:01 +0200 Subject: [PATCH 02/10] add persistance storage to otlpreciver Signed-off-by: MJarmo --- receiver/otlpreceiver/README.md | 41 ++ receiver/otlpreceiver/config.go | 22 + receiver/otlpreceiver/config.md | 16 +- receiver/otlpreceiver/factory.go | 4 + receiver/otlpreceiver/otlp.go | 45 +- receiver/otlpreceiver/otlphttp.go | 36 ++ receiver/otlpreceiver/persistence.go | 433 ++++++++++++++++++ receiver/otlpreceiver/persistence_test.go | 67 +++ .../otlpreceiver/testdata/persistence.yaml | 32 ++ test-config.yaml | 32 ++ test-logs.json | 41 ++ 11 files changed, 756 insertions(+), 13 deletions(-) create mode 100644 receiver/otlpreceiver/persistence.go create mode 100644 receiver/otlpreceiver/persistence_test.go create mode 100644 receiver/otlpreceiver/testdata/persistence.yaml create mode 100644 test-config.yaml create mode 100644 test-logs.json diff --git a/receiver/otlpreceiver/README.md b/receiver/otlpreceiver/README.md index c313937e29eb..c70fb710fec1 100644 --- a/receiver/otlpreceiver/README.md +++ b/receiver/otlpreceiver/README.md @@ -51,6 +51,47 @@ Several helper files are leveraged to provide additional capabilities automatica - [TLS and mTLS settings](https://github.com/open-telemetry/opentelemetry-collector/blob/main/config/configtls/README.md) - [Auth settings](https://github.com/open-telemetry/opentelemetry-collector/blob/main/config/configauth/README.md) +## Message Persistence + +The OTLP receiver supports message persistence for failed sends. When enabled, failed messages are stored in persistent storage and retried automatically. This feature is useful for ensuring data reliability in case of temporary failures. + +### Configuration + +```yaml +receivers: + otlp: + protocols: + http: + endpoint: localhost:4318 + persistence: + enabled: true + storage: file_storage + retry_interval: 5s + max_retries: 3 + +extensions: + file_storage: + directory: /tmp/otel-collector-storage + timeout: 1s + +service: + extensions: [file_storage] +``` + +### Persistence Settings + +- `enabled`: Enables message persistence for failed sends (default: false) +- `storage`: The ID of the storage extension to use for persistence +- `retry_interval`: The interval between retry attempts for failed sends (default: 5s) +- `max_retries`: The maximum number of retry attempts for failed sends (default: 3) + +### How It Works + +1. When a message fails to be processed, it is stored in the configured storage extension +2. A background retry worker periodically attempts to reprocess stored messages +3. Successful messages are removed from storage +4. Messages that exceed the maximum retry count are also removed from storage + ## Writing with HTTP/JSON The OTLP receiver can receive trace export calls via HTTP/JSON in addition to diff --git a/receiver/otlpreceiver/config.go b/receiver/otlpreceiver/config.go index ac5a30dddd8f..a4c9bd89c750 100644 --- a/receiver/otlpreceiver/config.go +++ b/receiver/otlpreceiver/config.go @@ -9,6 +9,7 @@ import ( "fmt" "net/url" "path" + "time" "go.opentelemetry.io/collector/component" "go.opentelemetry.io/collector/config/configgrpc" @@ -58,10 +59,24 @@ type Protocols struct { _ struct{} } +// PersistenceConfig defines configuration for message persistence. +type PersistenceConfig struct { + // Enabled enables message persistence for failed sends. + Enabled bool `mapstructure:"enabled"` + // StorageID is the ID of the storage extension to use for persistence. + StorageID component.ID `mapstructure:"storage"` + // RetryInterval is the interval between retry attempts for failed sends. + RetryInterval time.Duration `mapstructure:"retry_interval"` + // prevent unkeyed literal initialization + _ struct{} +} + // Config defines configuration for OTLP receiver. type Config struct { // Protocols is the configuration for the supported protocols, currently gRPC and HTTP (Proto and JSON). Protocols `mapstructure:"protocols"` + // Persistence is the configuration for message persistence. + Persistence PersistenceConfig `mapstructure:"persistence"` // prevent unkeyed literal initialization _ struct{} } @@ -73,5 +88,12 @@ func (cfg *Config) Validate() error { if !cfg.GRPC.HasValue() && !cfg.HTTP.HasValue() { return errors.New("must specify at least one protocol when using the OTLP receiver") } + + if cfg.Persistence.Enabled { + if cfg.Persistence.RetryInterval <= 0 { + return errors.New("retry interval must be positive when persistence is enabled") + } + } + return nil } diff --git a/receiver/otlpreceiver/config.md b/receiver/otlpreceiver/config.md index b77b4bdae737..58aa971b6233 100644 --- a/receiver/otlpreceiver/config.md +++ b/receiver/otlpreceiver/config.md @@ -5,9 +5,10 @@ Config defines configuration for OTLP receiver. ### Config -| Name | Type | Default | Docs | -|-----------|---------------------------------------------------|------------|-------------------------------------------------------------------------------------------------------| -| protocols | [otlpreceiver-Protocols](#otlpreceiver-protocols) | | Protocols is the configuration for the supported protocols, currently gRPC and HTTP (Proto and JSON). | +| Name | Type | Default | Docs | +|--------------|---------------------------------------------------|------------|-------------------------------------------------------------------------------------------------------| +| protocols | [otlpreceiver-Protocols](#otlpreceiver-protocols) | | Protocols is the configuration for the supported protocols, currently gRPC and HTTP (Proto and JSON). | +| persistence | [otlpreceiver-Persistence](#otlpreceiver-persistence) | | Persistence configuration for message storage and retry logic. | ### otlpreceiver-Protocols @@ -16,6 +17,15 @@ Config defines configuration for OTLP receiver. | grpc | [configgrpc-GRPCServerSettings](#configgrpc-grpcserversettings) | | GRPCServerSettings defines common settings for a gRPC server configuration. | | http | [confighttp-HTTPServerSettings](#confighttp-httpserversettings) | | HTTPServerSettings defines settings for creating an HTTP server. | +### otlpreceiver-Persistence + +| Name | Type | Default | Docs | +|----------------|--------|---------|-----------------------------------------------------------------------------------------| +| enabled | bool | false | Enables message persistence for failed sends. | +| storage | string | | The ID of the storage extension to use for persistence. | +| retry_interval | string | 5s | The interval between retry attempts for failed sends. | +| max_retries | int | 3 | The maximum number of retry attempts for failed sends. | + ### configgrpc-GRPCServerSettings | Name | Type | Default | Docs | diff --git a/receiver/otlpreceiver/factory.go b/receiver/otlpreceiver/factory.go index 1b6f483a1e3e..bd530db02739 100644 --- a/receiver/otlpreceiver/factory.go +++ b/receiver/otlpreceiver/factory.go @@ -63,6 +63,10 @@ func createDefaultConfig() component.Config { LogsURLPath: defaultLogsURLPath, }), }, + Persistence: PersistenceConfig{ + Enabled: false, + RetryInterval: 0, + }, } } diff --git a/receiver/otlpreceiver/otlp.go b/receiver/otlpreceiver/otlp.go index 7f06b89773de..fe42f90a8dff 100644 --- a/receiver/otlpreceiver/otlp.go +++ b/receiver/otlpreceiver/otlp.go @@ -48,6 +48,10 @@ type otlpReceiver struct { obsrepHTTP *receiverhelper.ObsReport settings *receiver.Settings + + // Persistence related fields + persistenceEnabled bool + persistenceManager *PersistenceManager } // newOtlpReceiver just creates the OpenTelemetry receiver services. It is the caller's @@ -57,12 +61,13 @@ func newOtlpReceiver(cfg *Config, set *receiver.Settings) (*otlpReceiver, error) set.TelemetrySettings = telemetry.WithoutAttributes(set.TelemetrySettings, componentattribute.SignalKey) set.Logger.Debug("created signal-agnostic logger") r := &otlpReceiver{ - cfg: cfg, - nextTraces: nil, - nextMetrics: nil, - nextLogs: nil, - nextProfiles: nil, - settings: set, + cfg: cfg, + nextTraces: nil, + nextMetrics: nil, + nextLogs: nil, + nextProfiles: nil, + settings: set, + persistenceEnabled: cfg.Persistence.Enabled, } var err error @@ -132,13 +137,21 @@ func (r *otlpReceiver) startGRPCServer(ctx context.Context, host component.Host) } func (r *otlpReceiver) startHTTPServer(ctx context.Context, host component.Host) error { - // If HTTP is not enabled, nothing to start. if !r.cfg.HTTP.HasValue() { return nil } httpCfg := r.cfg.HTTP.Get() httpMux := http.NewServeMux() + + if r.persistenceEnabled { + var logsReceiver *logs.Receiver + if r.nextLogs != nil { + logsReceiver = logs.New(r.nextLogs, r.obsrepHTTP) + } + r.persistenceManager = NewPersistenceManager(r.cfg.Persistence, r.settings.Logger, logsReceiver) + } + if r.nextTraces != nil { httpTracesReceiver := trace.New(r.nextTraces, r.obsrepHTTP) httpMux.HandleFunc(string(httpCfg.TracesURLPath), func(resp http.ResponseWriter, req *http.Request) { @@ -155,9 +168,15 @@ func (r *otlpReceiver) startHTTPServer(ctx context.Context, host component.Host) if r.nextLogs != nil { httpLogsReceiver := logs.New(r.nextLogs, r.obsrepHTTP) - httpMux.HandleFunc(string(httpCfg.LogsURLPath), func(resp http.ResponseWriter, req *http.Request) { - handleLogs(resp, req, httpLogsReceiver) - }) + if r.persistenceManager != nil { + httpMux.HandleFunc(string(httpCfg.LogsURLPath), func(resp http.ResponseWriter, req *http.Request) { + handleLogsWithPersistence(resp, req, r.persistenceManager.logsReceiver, r.persistenceManager) + }) + } else { + httpMux.HandleFunc(string(httpCfg.LogsURLPath), func(resp http.ResponseWriter, req *http.Request) { + handleLogs(resp, req, httpLogsReceiver) + }) + } } if r.nextProfiles != nil { @@ -217,6 +236,12 @@ func (r *otlpReceiver) Shutdown(ctx context.Context) error { r.serverGRPC.GracefulStop() } + if r.persistenceManager != nil { + if shutdownErr := r.persistenceManager.Shutdown(ctx); shutdownErr != nil { + err = errors.Join(err, shutdownErr) + } + } + r.shutdownWG.Wait() return err } diff --git a/receiver/otlpreceiver/otlphttp.go b/receiver/otlpreceiver/otlphttp.go index 11fa99fe9543..96db8e172a23 100644 --- a/receiver/otlpreceiver/otlphttp.go +++ b/receiver/otlpreceiver/otlphttp.go @@ -11,6 +11,7 @@ import ( "strconv" "time" + "go.uber.org/zap" "google.golang.org/grpc/status" "go.opentelemetry.io/collector/internal/statusutil" @@ -119,6 +120,41 @@ func handleLogs(resp http.ResponseWriter, req *http.Request, logsReceiver *logs. writeResponse(resp, enc.contentType(), http.StatusOK, msg) } +func handleLogsWithPersistence(resp http.ResponseWriter, req *http.Request, logsReceiver *logs.Receiver, persistenceManager *PersistenceManager) { + enc, ok := readContentType(resp, req) + if !ok { + return + } + + body, ok := readAndCloseBody(resp, req, enc) + if !ok { + return + } + + otlpReq, err := enc.unmarshalLogsRequest(body) + if err != nil { + writeError(resp, enc, err, http.StatusBadRequest) + return + } + + storeErr := persistenceManager.StoreMessage(req.Context(), body, enc.contentType(), "logs") + if storeErr != nil { + writeError(resp, enc, storeErr, http.StatusInternalServerError) + return + } + + writeResponse(resp, enc.contentType(), http.StatusAccepted, []byte(`{"message": "Logs accepted for processing"}`)) + + _, err = logsReceiver.Export(req.Context(), otlpReq) + if err != nil { + persistenceManager.logger.Error("Failed to process stored log message", + zap.Error(err)) + return + } + + persistenceManager.DeleteStoredMessageByContent(req.Context(), body, enc.contentType(), "logs") +} + func handleProfiles(resp http.ResponseWriter, req *http.Request, profilesReceiver *profiles.Receiver) { enc, ok := readContentType(resp, req) if !ok { diff --git a/receiver/otlpreceiver/persistence.go b/receiver/otlpreceiver/persistence.go new file mode 100644 index 000000000000..b78b9367d9c6 --- /dev/null +++ b/receiver/otlpreceiver/persistence.go @@ -0,0 +1,433 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +package otlpreceiver // import "go.opentelemetry.io/collector/receiver/otlpreceiver" + +import ( + "context" + "fmt" + "io" + "net/http" + "sync" + "time" + + "go.opentelemetry.io/collector/receiver/otlpreceiver/internal/logs" + "go.uber.org/zap" +) + +// PersistenceManager handles message persistence and retry logic +type PersistenceManager struct { + config PersistenceConfig + logger *zap.Logger + + storedMessages map[string]PersistedMessage + storageMutex sync.RWMutex + + retryWorkerCtx context.Context + retryWorkerCancel context.CancelFunc + retryWorkerWg sync.WaitGroup + + logsReceiver *logs.Receiver +} + +// PersistedMessage represents a message stored in persistence +type PersistedMessage struct { + ID string `json:"id"` + Data []byte `json:"data"` + ContentType string `json:"content_type"` + SignalType string `json:"signal_type"` + CreatedAt time.Time `json:"created_at"` + RetryCount int `json:"retry_count"` +} + +// NewPersistenceManager creates a new persistence manager +func NewPersistenceManager(config PersistenceConfig, logger *zap.Logger, logsReceiver *logs.Receiver) *PersistenceManager { + ctx, cancel := context.WithCancel(context.Background()) + pm := &PersistenceManager{ + config: config, + logger: logger, + storedMessages: make(map[string]PersistedMessage), + retryWorkerCtx: ctx, + retryWorkerCancel: cancel, + logsReceiver: logsReceiver, + } + + pm.startRetryWorker() + + return pm +} + +// StoreMessage stores a message in persistent storage (only for logs) +func (pm *PersistenceManager) StoreMessage(ctx context.Context, data []byte, contentType, signalType string) error { + if signalType != "logs" { + pm.logger.Debug("Skipping persistence for non-log signal type", + zap.String("signal_type", signalType)) + return nil + } + + messageID := fmt.Sprintf("%s_%d", signalType, time.Now().UnixNano()) + + message := PersistedMessage{ + ID: messageID, + Data: data, + ContentType: contentType, + SignalType: signalType, + CreatedAt: time.Now(), + RetryCount: 0, + } + + key := fmt.Sprintf("otlp_%s_%s", signalType, messageID) + + pm.storageMutex.Lock() + pm.storedMessages[key] = message + pm.storageMutex.Unlock() + + pm.logger.Debug("Log message stored for persistence", + zap.String("message_id", messageID), + zap.String("signal_type", signalType)) + + return nil +} + +// RemoveMessage removes a message from persistent storage +func (pm *PersistenceManager) RemoveMessage(ctx context.Context, messageID, signalType string) error { + key := fmt.Sprintf("otlp_%s_%s", signalType, messageID) + + pm.storageMutex.Lock() + delete(pm.storedMessages, key) + pm.storageMutex.Unlock() + + pm.logger.Debug("Message removed from persistence", + zap.String("message_id", messageID), + zap.String("signal_type", signalType)) + + return nil +} + +// ClearStoredMessages clears all stored messages for a specific signal type +func (pm *PersistenceManager) ClearStoredMessages(ctx context.Context, signalType string) error { + pm.storageMutex.Lock() + defer pm.storageMutex.Unlock() + + // Remove all messages for the specified signal type + keysToDelete := make([]string, 0) + for key, message := range pm.storedMessages { + if message.SignalType == signalType { + keysToDelete = append(keysToDelete, key) + } + } + + // Delete the messages + for _, key := range keysToDelete { + delete(pm.storedMessages, key) + } + + pm.logger.Info("Cleared stored messages", + zap.String("signal_type", signalType), + zap.Int("cleared_count", len(keysToDelete))) + + return nil +} + +// DeleteStoredMessageByContent deletes a stored message by matching content +func (pm *PersistenceManager) DeleteStoredMessageByContent(ctx context.Context, content []byte, contentType, signalType string) error { + pm.storageMutex.Lock() + defer pm.storageMutex.Unlock() + + // Find and delete the message with matching content + for key, message := range pm.storedMessages { + if message.SignalType == signalType && + message.ContentType == contentType && + len(message.Data) == len(content) { + // Simple byte comparison for exact match + match := true + for i, b := range content { + if message.Data[i] != b { + match = false + break + } + } + if match { + delete(pm.storedMessages, key) + pm.logger.Info("Deleted stored message after successful processing", + zap.String("message_id", message.ID), + zap.String("signal_type", signalType)) + return nil + } + } + } + + pm.logger.Debug("No matching stored message found to delete", + zap.String("signal_type", signalType), + zap.String("content_type", contentType), + zap.Int("content_length", len(content))) + + return nil +} + +// GetStoredMessages retrieves all stored messages for retry +func (pm *PersistenceManager) GetStoredMessages(ctx context.Context) ([]PersistedMessage, error) { + pm.storageMutex.RLock() + defer pm.storageMutex.RUnlock() + + messages := make([]PersistedMessage, 0, len(pm.storedMessages)) + for _, message := range pm.storedMessages { + messages = append(messages, message) + } + + return messages, nil +} + +// ProcessWithPersistence wraps an HTTP handler with persistence logic +func (pm *PersistenceManager) ProcessWithPersistence( + handler func(http.ResponseWriter, *http.Request), + signalType string, +) func(http.ResponseWriter, *http.Request) { + return func(resp http.ResponseWriter, req *http.Request) { + // Read the request body + body, err := readRequestBody(req) + if err != nil { + http.Error(resp, "Failed to read request body", http.StatusBadRequest) + return + } + + // Get content type + contentType := req.Header.Get("Content-Type") + if contentType == "" { + contentType = "application/json" + } + + // Try to process the request + success := false + var processingErr error + + // Create a custom response writer to capture the response + captureWriter := &responseCapture{ResponseWriter: resp} + + // Process the request + handler(captureWriter, req) + + // Check if the request was successful (status code 200) + if captureWriter.statusCode == http.StatusOK { + success = true + } else { + processingErr = fmt.Errorf("request failed with status %d", captureWriter.statusCode) + } + + if success { + // Request was successful, no need to persist + pm.logger.Debug("Request processed successfully", + zap.String("signal_type", signalType)) + } else { + // Request failed, store for retry + if err := pm.StoreMessage(req.Context(), body, contentType, signalType); err != nil { + pm.logger.Error("Failed to store message for retry", + zap.String("signal_type", signalType), + zap.Error(err)) + } else { + pm.logger.Info("Message stored for retry", + zap.String("signal_type", signalType), + zap.Error(processingErr)) + } + } + } +} + +// responseCapture captures the response for analysis +type responseCapture struct { + http.ResponseWriter + statusCode int +} + +func (rc *responseCapture) WriteHeader(code int) { + rc.statusCode = code + rc.ResponseWriter.WriteHeader(code) +} + +// startRetryWorker starts the background retry worker +func (pm *PersistenceManager) startRetryWorker() { + pm.retryWorkerWg.Add(1) + go func() { + defer pm.retryWorkerWg.Done() + ticker := time.NewTicker(pm.config.RetryInterval) + defer ticker.Stop() + + for { + select { + case <-pm.retryWorkerCtx.Done(): + pm.logger.Info("Retry worker shutting down") + return + case <-ticker.C: + pm.processRetries() + } + } + }() +} + +func (pm *PersistenceManager) processRetries() { + pm.logger.Debug("Processing retries for logs") + + messages, err := pm.GetStoredMessages(context.Background()) + if err != nil { + pm.logger.Error("Failed to get stored messages", zap.Error(err)) + return + } + + logMessages := make([]PersistedMessage, 0) + for _, message := range messages { + if message.SignalType == "logs" { + logMessages = append(logMessages, message) + } + } + + if len(logMessages) == 0 { + pm.logger.Debug("No log messages to retry") + return + } + + pm.logger.Debug("Found log messages to retry", zap.Int("count", len(logMessages))) + + for _, message := range logMessages { + pm.processMessageRetry(message) + } +} + +func (pm *PersistenceManager) processMessageRetry(message PersistedMessage) { + if message.SignalType != "logs" { + pm.logger.Debug("Skipping retry for non-log message", + zap.String("signal_type", message.SignalType)) + return + } + + pm.logger.Debug("Processing stored log message", + zap.String("message_id", message.ID), + zap.String("signal_type", message.SignalType), + zap.Int("retry_count", message.RetryCount+1)) + + success := pm.processStoredLogMessage(message) + + if success { + pm.logger.Info("Log message processed successfully, removing from storage", + zap.String("message_id", message.ID), + zap.String("signal_type", message.SignalType)) + + if err := pm.RemoveMessage(context.Background(), message.ID, message.SignalType); err != nil { + pm.logger.Error("Failed to remove message after successful processing", + zap.String("message_id", message.ID), + zap.Error(err)) + } + } else { + message.RetryCount++ + key := fmt.Sprintf("otlp_%s_%s", message.SignalType, message.ID) + + pm.storageMutex.Lock() + pm.storedMessages[key] = message + pm.storageMutex.Unlock() + + pm.logger.Debug("Log message processing failed, updated retry count", + zap.String("message_id", message.ID), + zap.String("signal_type", message.SignalType), + zap.Int("retry_count", message.RetryCount)) + } +} + +// processStoredLogMessage processes a stored log message +func (pm *PersistenceManager) processStoredLogMessage(message PersistedMessage) bool { + pm.logger.Debug("Processing stored log message", + zap.String("message_id", message.ID), + zap.String("content_type", message.ContentType), + zap.Int("data_size", len(message.Data))) + + if pm.logsReceiver == nil { + pm.logger.Error("No logs receiver available for processing stored message", + zap.String("message_id", message.ID)) + return false + } + + ctx := context.Background() + + encoder := getEncoderForContentType(message.ContentType) + if encoder == nil { + pm.logger.Error("Unknown content type for stored message", + zap.String("message_id", message.ID), + zap.String("content_type", message.ContentType)) + return false + } + + otlpReq, err := encoder.unmarshalLogsRequest(message.Data) + if err != nil { + pm.logger.Error("Failed to unmarshal stored message", + zap.String("message_id", message.ID), + zap.Error(err)) + return false + } + + otlpResp, err := pm.logsReceiver.Export(ctx, otlpReq) + if err != nil { + pm.logger.Error("Failed to process stored log message", + zap.String("message_id", message.ID), + zap.Error(err)) + return false + } + + _, err = encoder.marshalLogsResponse(otlpResp) + if err != nil { + pm.logger.Error("Failed to marshal response for stored message", + zap.String("message_id", message.ID), + zap.Error(err)) + return false + } + + pm.logger.Info("Successfully processed stored log message", + zap.String("message_id", message.ID), + zap.String("content_type", message.ContentType)) + + return true +} + +// Shutdown stops the persistence manager +func (pm *PersistenceManager) Shutdown(ctx context.Context) error { + pm.logger.Info("Shutting down persistence manager") + pm.retryWorkerCancel() + + done := make(chan struct{}) + go func() { + pm.retryWorkerWg.Wait() + close(done) + }() + + select { + case <-done: + pm.logger.Info("Retry worker stopped gracefully") + case <-ctx.Done(): + pm.logger.Warn("Retry worker shutdown timed out") + } + + return nil +} + +// readRequestBody reads the request body +func readRequestBody(req *http.Request) ([]byte, error) { + body, err := io.ReadAll(req.Body) + if err != nil { + return nil, fmt.Errorf("failed to read request body: %w", err) + } + + if err := req.Body.Close(); err != nil { + return nil, fmt.Errorf("failed to close request body: %w", err) + } + + return body, nil +} + +// getEncoderForContentType returns the appropriate encoder for the content type +func getEncoderForContentType(contentType string) encoder { + switch contentType { + case "application/x-protobuf": + return pbEncoder + case "application/json": + return jsEncoder + default: + return nil + } +} diff --git a/receiver/otlpreceiver/persistence_test.go b/receiver/otlpreceiver/persistence_test.go new file mode 100644 index 000000000000..36c1bd83be08 --- /dev/null +++ b/receiver/otlpreceiver/persistence_test.go @@ -0,0 +1,67 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +package otlpreceiver + +import ( + "context" + "testing" + "time" + + "go.uber.org/zap" +) + +func TestPersistenceManager(t *testing.T) { + config := PersistenceConfig{ + Enabled: true, + RetryInterval: 1 * time.Second, + } + + logger := zap.NewNop() + pm := NewPersistenceManager(config, logger, nil) + + // Test storing a message + err := pm.StoreMessage(context.Background(), []byte("test data"), "application/json", "traces") + if err != nil { + t.Fatalf("Failed to store message: %v", err) + } + + // Test retrieving stored messages + messages, err := pm.GetStoredMessages(context.Background()) + if err != nil { + t.Fatalf("Failed to get stored messages: %v", err) + } + + if len(messages) != 1 { + t.Fatalf("Expected 1 message, got %d", len(messages)) + } + + if messages[0].SignalType != "traces" { + t.Fatalf("Expected signal type 'traces', got '%s'", messages[0].SignalType) + } + + // Test removing a message + err = pm.RemoveMessage(context.Background(), messages[0].ID, "traces") + if err != nil { + t.Fatalf("Failed to remove message: %v", err) + } + + // Verify message was removed + messages, err = pm.GetStoredMessages(context.Background()) + if err != nil { + t.Fatalf("Failed to get stored messages: %v", err) + } + + if len(messages) != 0 { + t.Fatalf("Expected 0 messages after removal, got %d", len(messages)) + } + + // Test shutdown + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + + err = pm.Shutdown(ctx) + if err != nil { + t.Fatalf("Failed to shutdown persistence manager: %v", err) + } +} diff --git a/receiver/otlpreceiver/testdata/persistence.yaml b/receiver/otlpreceiver/testdata/persistence.yaml new file mode 100644 index 000000000000..8095c324d6ff --- /dev/null +++ b/receiver/otlpreceiver/testdata/persistence.yaml @@ -0,0 +1,32 @@ +receivers: + otlp: + protocols: + http: + endpoint: localhost:4318 + persistence: + enabled: true + storage: file_storage + retry_interval: 5s + max_retries: 3 + +extensions: + file_storage: + directory: /tmp/otel-collector-storage + timeout: 1s + +service: + extensions: [file_storage] + pipelines: + traces: + receivers: [otlp] + exporters: [logging] + metrics: + receivers: [otlp] + exporters: [logging] + logs: + receivers: [otlp] + exporters: [logging] + +exporters: + logging: + loglevel: debug diff --git a/test-config.yaml b/test-config.yaml new file mode 100644 index 000000000000..ce98092e9f24 --- /dev/null +++ b/test-config.yaml @@ -0,0 +1,32 @@ +receivers: + otlp: + protocols: + grpc: + endpoint: 0.0.0.0:4317 + http: + endpoint: 0.0.0.0:4318 + persistence: + enabled: true + retry_interval: 5s + +processors: + batch: + +exporters: + debug: + verbosity: detailed + +service: + pipelines: + logs: + receivers: [otlp] + processors: [batch] + exporters: [debug] + traces: + receivers: [otlp] + processors: [batch] + exporters: [debug] + metrics: + receivers: [otlp] + processors: [batch] + exporters: [debug] diff --git a/test-logs.json b/test-logs.json new file mode 100644 index 000000000000..a8e9c2a9aa35 --- /dev/null +++ b/test-logs.json @@ -0,0 +1,41 @@ +{ + "resourceLogs": [ + { + "resource": { + "attributes": [ + { + "key": "service.name", + "value": { + "stringValue": "test-service" + } + } + ] + }, + "scopeLogs": [ + { + "scope": { + "name": "test-scope" + }, + "logRecords": [ + { + "timeUnixNano": "1640995200000000000", + "severityNumber": 9, + "severityText": "INFO", + "body": { + "stringValue": "Test log message for persistence testing" + }, + "attributes": [ + { + "key": "test.attribute", + "value": { + "stringValue": "test-value" + } + } + ] + } + ] + } + ] + } + ] +} From bdecff311414a7fbfe18707f2db7ac94a3dc3390 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Micha=C5=82=20Jarmolkiewicz?= Date: Mon, 29 Sep 2025 09:15:51 +0200 Subject: [PATCH 03/10] fix test Signed-off-by: MJarmo --- receiver/otlpreceiver/persistence_test.go | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/receiver/otlpreceiver/persistence_test.go b/receiver/otlpreceiver/persistence_test.go index 36c1bd83be08..bae1b849dc2d 100644 --- a/receiver/otlpreceiver/persistence_test.go +++ b/receiver/otlpreceiver/persistence_test.go @@ -21,7 +21,7 @@ func TestPersistenceManager(t *testing.T) { pm := NewPersistenceManager(config, logger, nil) // Test storing a message - err := pm.StoreMessage(context.Background(), []byte("test data"), "application/json", "traces") + err := pm.StoreMessage(context.Background(), []byte("test data"), "application/json", "logs") if err != nil { t.Fatalf("Failed to store message: %v", err) } @@ -36,12 +36,12 @@ func TestPersistenceManager(t *testing.T) { t.Fatalf("Expected 1 message, got %d", len(messages)) } - if messages[0].SignalType != "traces" { - t.Fatalf("Expected signal type 'traces', got '%s'", messages[0].SignalType) + if messages[0].SignalType != "logs" { + t.Fatalf("Expected signal type 'logs', got '%s'", messages[0].SignalType) } // Test removing a message - err = pm.RemoveMessage(context.Background(), messages[0].ID, "traces") + err = pm.RemoveMessage(context.Background(), messages[0].ID, "logs") if err != nil { t.Fatalf("Failed to remove message: %v", err) } From 1b1ea3ffeea17c714e1c353f784244b300d4ac0a Mon Sep 17 00:00:00 2001 From: MJarmo Date: Mon, 29 Sep 2025 12:39:26 +0200 Subject: [PATCH 04/10] prevent unkey literal Signed-off-by: MJarmo --- receiver/otlpreceiver/persistence.go | 1 + 1 file changed, 1 insertion(+) diff --git a/receiver/otlpreceiver/persistence.go b/receiver/otlpreceiver/persistence.go index b78b9367d9c6..e4fa5614a29c 100644 --- a/receiver/otlpreceiver/persistence.go +++ b/receiver/otlpreceiver/persistence.go @@ -38,6 +38,7 @@ type PersistedMessage struct { SignalType string `json:"signal_type"` CreatedAt time.Time `json:"created_at"` RetryCount int `json:"retry_count"` + _ struct{} } // NewPersistenceManager creates a new persistence manager From f3118792e8afccf70db82c839ad7dbe2f2e252b0 Mon Sep 17 00:00:00 2001 From: MJarmo Date: Tue, 7 Oct 2025 10:46:42 +0200 Subject: [PATCH 05/10] fix lint errors Signed-off-by: MJarmo --- receiver/otlpreceiver/otlphttp.go | 5 ++++- receiver/otlpreceiver/persistence.go | 31 ++++++++++++++-------------- 2 files changed, 19 insertions(+), 17 deletions(-) diff --git a/receiver/otlpreceiver/otlphttp.go b/receiver/otlpreceiver/otlphttp.go index 96db8e172a23..4c657c244e55 100644 --- a/receiver/otlpreceiver/otlphttp.go +++ b/receiver/otlpreceiver/otlphttp.go @@ -152,7 +152,10 @@ func handleLogsWithPersistence(resp http.ResponseWriter, req *http.Request, logs return } - persistenceManager.DeleteStoredMessageByContent(req.Context(), body, enc.contentType(), "logs") + if err := persistenceManager.DeleteStoredMessageByContent(req.Context(), body, enc.contentType(), "logs"); err != nil { + persistenceManager.logger.Error("Failed to delete stored message after successful processing", + zap.Error(err)) + } } func handleProfiles(resp http.ResponseWriter, req *http.Request, profilesReceiver *profiles.Receiver) { diff --git a/receiver/otlpreceiver/persistence.go b/receiver/otlpreceiver/persistence.go index e4fa5614a29c..9316e566330a 100644 --- a/receiver/otlpreceiver/persistence.go +++ b/receiver/otlpreceiver/persistence.go @@ -11,8 +11,9 @@ import ( "sync" "time" - "go.opentelemetry.io/collector/receiver/otlpreceiver/internal/logs" "go.uber.org/zap" + + "go.opentelemetry.io/collector/receiver/otlpreceiver/internal/logs" ) // PersistenceManager handles message persistence and retry logic @@ -59,7 +60,7 @@ func NewPersistenceManager(config PersistenceConfig, logger *zap.Logger, logsRec } // StoreMessage stores a message in persistent storage (only for logs) -func (pm *PersistenceManager) StoreMessage(ctx context.Context, data []byte, contentType, signalType string) error { +func (pm *PersistenceManager) StoreMessage(_ context.Context, data []byte, contentType, signalType string) error { if signalType != "logs" { pm.logger.Debug("Skipping persistence for non-log signal type", zap.String("signal_type", signalType)) @@ -91,7 +92,7 @@ func (pm *PersistenceManager) StoreMessage(ctx context.Context, data []byte, con } // RemoveMessage removes a message from persistent storage -func (pm *PersistenceManager) RemoveMessage(ctx context.Context, messageID, signalType string) error { +func (pm *PersistenceManager) RemoveMessage(_ context.Context, messageID, signalType string) error { key := fmt.Sprintf("otlp_%s_%s", signalType, messageID) pm.storageMutex.Lock() @@ -106,7 +107,7 @@ func (pm *PersistenceManager) RemoveMessage(ctx context.Context, messageID, sign } // ClearStoredMessages clears all stored messages for a specific signal type -func (pm *PersistenceManager) ClearStoredMessages(ctx context.Context, signalType string) error { +func (pm *PersistenceManager) ClearStoredMessages(_ context.Context, signalType string) error { pm.storageMutex.Lock() defer pm.storageMutex.Unlock() @@ -131,7 +132,7 @@ func (pm *PersistenceManager) ClearStoredMessages(ctx context.Context, signalTyp } // DeleteStoredMessageByContent deletes a stored message by matching content -func (pm *PersistenceManager) DeleteStoredMessageByContent(ctx context.Context, content []byte, contentType, signalType string) error { +func (pm *PersistenceManager) DeleteStoredMessageByContent(_ context.Context, content []byte, contentType, signalType string) error { pm.storageMutex.Lock() defer pm.storageMutex.Unlock() @@ -167,7 +168,7 @@ func (pm *PersistenceManager) DeleteStoredMessageByContent(ctx context.Context, } // GetStoredMessages retrieves all stored messages for retry -func (pm *PersistenceManager) GetStoredMessages(ctx context.Context) ([]PersistedMessage, error) { +func (pm *PersistenceManager) GetStoredMessages(_ context.Context) ([]PersistedMessage, error) { pm.storageMutex.RLock() defer pm.storageMutex.RUnlock() @@ -289,11 +290,11 @@ func (pm *PersistenceManager) processRetries() { pm.logger.Debug("Found log messages to retry", zap.Int("count", len(logMessages))) for _, message := range logMessages { - pm.processMessageRetry(message) + pm.processMessageRetry(pm.retryWorkerCtx, message) } } -func (pm *PersistenceManager) processMessageRetry(message PersistedMessage) { +func (pm *PersistenceManager) processMessageRetry(ctx context.Context, message PersistedMessage) { if message.SignalType != "logs" { pm.logger.Debug("Skipping retry for non-log message", zap.String("signal_type", message.SignalType)) @@ -305,7 +306,7 @@ func (pm *PersistenceManager) processMessageRetry(message PersistedMessage) { zap.String("signal_type", message.SignalType), zap.Int("retry_count", message.RetryCount+1)) - success := pm.processStoredLogMessage(message) + success := pm.processStoredLogMessage(ctx, message) if success { pm.logger.Info("Log message processed successfully, removing from storage", @@ -333,7 +334,7 @@ func (pm *PersistenceManager) processMessageRetry(message PersistedMessage) { } // processStoredLogMessage processes a stored log message -func (pm *PersistenceManager) processStoredLogMessage(message PersistedMessage) bool { +func (pm *PersistenceManager) processStoredLogMessage(ctx context.Context, message PersistedMessage) bool { pm.logger.Debug("Processing stored log message", zap.String("message_id", message.ID), zap.String("content_type", message.ContentType), @@ -345,17 +346,15 @@ func (pm *PersistenceManager) processStoredLogMessage(message PersistedMessage) return false } - ctx := context.Background() - - encoder := getEncoderForContentType(message.ContentType) - if encoder == nil { + enc := getEncoderForContentType(message.ContentType) + if enc == nil { pm.logger.Error("Unknown content type for stored message", zap.String("message_id", message.ID), zap.String("content_type", message.ContentType)) return false } - otlpReq, err := encoder.unmarshalLogsRequest(message.Data) + otlpReq, err := enc.unmarshalLogsRequest(message.Data) if err != nil { pm.logger.Error("Failed to unmarshal stored message", zap.String("message_id", message.ID), @@ -371,7 +370,7 @@ func (pm *PersistenceManager) processStoredLogMessage(message PersistedMessage) return false } - _, err = encoder.marshalLogsResponse(otlpResp) + _, err = enc.marshalLogsResponse(otlpResp) if err != nil { pm.logger.Error("Failed to marshal response for stored message", zap.String("message_id", message.ID), From b08af27b8f44c91a1d0417fa30e25eb00a010e1c Mon Sep 17 00:00:00 2001 From: MJarmo Date: Tue, 7 Oct 2025 11:26:47 +0200 Subject: [PATCH 06/10] fix lint, and naming Signed-off-by: MJarmo --- receiver/otlpreceiver/persistence.go | 8 ++++---- receiver/otlpreceiver/testdata/persistence.yaml | 2 +- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/receiver/otlpreceiver/persistence.go b/receiver/otlpreceiver/persistence.go index 9316e566330a..175c7006b559 100644 --- a/receiver/otlpreceiver/persistence.go +++ b/receiver/otlpreceiver/persistence.go @@ -260,16 +260,16 @@ func (pm *PersistenceManager) startRetryWorker() { pm.logger.Info("Retry worker shutting down") return case <-ticker.C: - pm.processRetries() + pm.processRetries(pm.retryWorkerCtx) } } }() } -func (pm *PersistenceManager) processRetries() { +func (pm *PersistenceManager) processRetries(ctx context.Context) { pm.logger.Debug("Processing retries for logs") - messages, err := pm.GetStoredMessages(context.Background()) + messages, err := pm.GetStoredMessages(ctx) if err != nil { pm.logger.Error("Failed to get stored messages", zap.Error(err)) return @@ -313,7 +313,7 @@ func (pm *PersistenceManager) processMessageRetry(ctx context.Context, message P zap.String("message_id", message.ID), zap.String("signal_type", message.SignalType)) - if err := pm.RemoveMessage(context.Background(), message.ID, message.SignalType); err != nil { + if err := pm.RemoveMessage(ctx, message.ID, message.SignalType); err != nil { pm.logger.Error("Failed to remove message after successful processing", zap.String("message_id", message.ID), zap.Error(err)) diff --git a/receiver/otlpreceiver/testdata/persistence.yaml b/receiver/otlpreceiver/testdata/persistence.yaml index 8095c324d6ff..e940feb4a831 100644 --- a/receiver/otlpreceiver/testdata/persistence.yaml +++ b/receiver/otlpreceiver/testdata/persistence.yaml @@ -29,4 +29,4 @@ service: exporters: logging: - loglevel: debug + verbosity: detailed From e802eb3918a279d3748cd43bb8220c0378649c45 Mon Sep 17 00:00:00 2001 From: MJarmo Date: Tue, 7 Oct 2025 11:58:30 +0200 Subject: [PATCH 07/10] fix ctx pass Signed-off-by: MJarmo --- receiver/otlpreceiver/persistence.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/receiver/otlpreceiver/persistence.go b/receiver/otlpreceiver/persistence.go index 175c7006b559..304d4b0e1592 100644 --- a/receiver/otlpreceiver/persistence.go +++ b/receiver/otlpreceiver/persistence.go @@ -290,7 +290,7 @@ func (pm *PersistenceManager) processRetries(ctx context.Context) { pm.logger.Debug("Found log messages to retry", zap.Int("count", len(logMessages))) for _, message := range logMessages { - pm.processMessageRetry(pm.retryWorkerCtx, message) + pm.processMessageRetry(ctx, message) } } From 903777e4abe10a8224dc7ab3e6264021bdffec6a Mon Sep 17 00:00:00 2001 From: MJarmo Date: Thu, 16 Oct 2025 14:47:12 +0200 Subject: [PATCH 08/10] fix losing logs Signed-off-by: MJarmo --- receiver/otlpreceiver/otlphttp.go | 28 ++-- receiver/otlpreceiver/persistence.go | 127 +++++++++++++++--- .../otlpreceiver/testdata/persistence.yaml | 27 ++-- test-config.yaml | 32 ----- 4 files changed, 128 insertions(+), 86 deletions(-) delete mode 100644 test-config.yaml diff --git a/receiver/otlpreceiver/otlphttp.go b/receiver/otlpreceiver/otlphttp.go index 4c657c244e55..fc39a34ae2e8 100644 --- a/receiver/otlpreceiver/otlphttp.go +++ b/receiver/otlpreceiver/otlphttp.go @@ -137,25 +137,27 @@ func handleLogsWithPersistence(resp http.ResponseWriter, req *http.Request, logs return } - storeErr := persistenceManager.StoreMessage(req.Context(), body, enc.contentType(), "logs") - if storeErr != nil { - writeError(resp, enc, storeErr, http.StatusInternalServerError) - return - } - - writeResponse(resp, enc.contentType(), http.StatusAccepted, []byte(`{"message": "Logs accepted for processing"}`)) - - _, err = logsReceiver.Export(req.Context(), otlpReq) + otlpResp, err := logsReceiver.Export(req.Context(), otlpReq) if err != nil { - persistenceManager.logger.Error("Failed to process stored log message", + persistenceManager.logger.Warn("Log processing failed, storing for retry", zap.Error(err)) + _, storeErr := persistenceManager.StoreMessage(req.Context(), body, enc.contentType(), "logs") + if storeErr != nil { + persistenceManager.logger.Error("Failed to store message for retry", + zap.Error(storeErr)) + writeError(resp, enc, err, http.StatusInternalServerError) + return + } + writeResponse(resp, enc.contentType(), http.StatusAccepted, []byte(`{"message": "Logs queued for retry"}`)) return } - if err := persistenceManager.DeleteStoredMessageByContent(req.Context(), body, enc.contentType(), "logs"); err != nil { - persistenceManager.logger.Error("Failed to delete stored message after successful processing", - zap.Error(err)) + msg, err := enc.marshalLogsResponse(otlpResp) + if err != nil { + writeError(resp, enc, err, http.StatusInternalServerError) + return } + writeResponse(resp, enc.contentType(), http.StatusOK, msg) } func handleProfiles(resp http.ResponseWriter, req *http.Request, profilesReceiver *profiles.Receiver) { diff --git a/receiver/otlpreceiver/persistence.go b/receiver/otlpreceiver/persistence.go index 304d4b0e1592..024746784f0b 100644 --- a/receiver/otlpreceiver/persistence.go +++ b/receiver/otlpreceiver/persistence.go @@ -60,11 +60,11 @@ func NewPersistenceManager(config PersistenceConfig, logger *zap.Logger, logsRec } // StoreMessage stores a message in persistent storage (only for logs) -func (pm *PersistenceManager) StoreMessage(_ context.Context, data []byte, contentType, signalType string) error { +func (pm *PersistenceManager) StoreMessage(_ context.Context, data []byte, contentType, signalType string) (string, error) { if signalType != "logs" { pm.logger.Debug("Skipping persistence for non-log signal type", zap.String("signal_type", signalType)) - return nil + return "", nil } messageID := fmt.Sprintf("%s_%d", signalType, time.Now().UnixNano()) @@ -88,7 +88,7 @@ func (pm *PersistenceManager) StoreMessage(_ context.Context, data []byte, conte zap.String("message_id", messageID), zap.String("signal_type", signalType)) - return nil + return messageID, nil } // RemoveMessage removes a message from persistent storage @@ -217,17 +217,17 @@ func (pm *PersistenceManager) ProcessWithPersistence( } if success { - // Request was successful, no need to persist pm.logger.Debug("Request processed successfully", zap.String("signal_type", signalType)) } else { - // Request failed, store for retry - if err := pm.StoreMessage(req.Context(), body, contentType, signalType); err != nil { + messageID, err := pm.StoreMessage(req.Context(), body, contentType, signalType) + if err != nil { pm.logger.Error("Failed to store message for retry", zap.String("signal_type", signalType), zap.Error(err)) } else { pm.logger.Info("Message stored for retry", + zap.String("message_id", messageID), zap.String("signal_type", signalType), zap.Error(processingErr)) } @@ -267,7 +267,7 @@ func (pm *PersistenceManager) startRetryWorker() { } func (pm *PersistenceManager) processRetries(ctx context.Context) { - pm.logger.Debug("Processing retries for logs") + pm.logger.Debug("Processing queued messages") messages, err := pm.GetStoredMessages(ctx) if err != nil { @@ -275,43 +275,78 @@ func (pm *PersistenceManager) processRetries(ctx context.Context) { return } + now := time.Now() + minAge := 10 * time.Millisecond + logMessages := make([]PersistedMessage, 0) + skippedNew := 0 for _, message := range messages { if message.SignalType == "logs" { + messageAge := now.Sub(message.CreatedAt) + if messageAge < minAge { + skippedNew++ + pm.logger.Debug("Skipping message - too new (ensuring write completion)", + zap.String("message_id", message.ID), + zap.Duration("age", messageAge)) + continue + } logMessages = append(logMessages, message) } } if len(logMessages) == 0 { - pm.logger.Debug("No log messages to retry") + if skippedNew > 0 { + pm.logger.Debug("No messages to process (waiting for write completion)", + zap.Int("skipped_new", skippedNew)) + } else { + pm.logger.Debug("Queue is empty - no messages to process") + } return } - pm.logger.Debug("Found log messages to retry", zap.Int("count", len(logMessages))) + pm.logger.Info("Processing queued messages", + zap.Int("count", len(logMessages)), + zap.Int("pending_writes", skippedNew)) for _, message := range logMessages { - pm.processMessageRetry(ctx, message) + pm.processQueuedMessage(ctx, message) } } -func (pm *PersistenceManager) processMessageRetry(ctx context.Context, message PersistedMessage) { +func (pm *PersistenceManager) processQueuedMessage(ctx context.Context, message PersistedMessage) { if message.SignalType != "logs" { - pm.logger.Debug("Skipping retry for non-log message", + pm.logger.Debug("Skipping non-log message", zap.String("signal_type", message.SignalType)) return } - pm.logger.Debug("Processing stored log message", - zap.String("message_id", message.ID), - zap.String("signal_type", message.SignalType), - zap.Int("retry_count", message.RetryCount+1)) + isRetry := message.RetryCount > 0 + messageAge := time.Since(message.CreatedAt) + + if isRetry { + pm.logger.Info("Processing retry attempt", + zap.String("message_id", message.ID), + zap.Int("retry_count", message.RetryCount), + zap.Duration("message_age", messageAge)) + } else { + pm.logger.Debug("Processing queued message", + zap.String("message_id", message.ID), + zap.Duration("queue_time", messageAge)) + } success := pm.processStoredLogMessage(ctx, message) if success { - pm.logger.Info("Log message processed successfully, removing from storage", - zap.String("message_id", message.ID), - zap.String("signal_type", message.SignalType)) + if isRetry { + pm.logger.Info("Retry successful, removing from queue", + zap.String("message_id", message.ID), + zap.Int("retry_count", message.RetryCount), + zap.Duration("total_time", messageAge)) + } else { + pm.logger.Debug("Message processed successfully, removing from queue", + zap.String("message_id", message.ID), + zap.Duration("queue_time", messageAge)) + } if err := pm.RemoveMessage(ctx, message.ID, message.SignalType); err != nil { pm.logger.Error("Failed to remove message after successful processing", @@ -326,10 +361,10 @@ func (pm *PersistenceManager) processMessageRetry(ctx context.Context, message P pm.storedMessages[key] = message pm.storageMutex.Unlock() - pm.logger.Debug("Log message processing failed, updated retry count", + pm.logger.Warn("Message processing failed, will retry", zap.String("message_id", message.ID), - zap.String("signal_type", message.SignalType), - zap.Int("retry_count", message.RetryCount)) + zap.Int("retry_count", message.RetryCount), + zap.Duration("next_retry_in", pm.config.RetryInterval)) } } @@ -388,6 +423,43 @@ func (pm *PersistenceManager) processStoredLogMessage(ctx context.Context, messa // Shutdown stops the persistence manager func (pm *PersistenceManager) Shutdown(ctx context.Context) error { pm.logger.Info("Shutting down persistence manager") + + pm.storageMutex.RLock() + queueSize := len(pm.storedMessages) + pm.storageMutex.RUnlock() + + if queueSize > 0 { + pm.logger.Info("Processing remaining messages before shutdown", + zap.Int("queue_size", queueSize)) + + maxAttempts := 100 + for i := 0; i < maxAttempts; i++ { + pm.processRetries(ctx) + + pm.storageMutex.RLock() + remaining := len(pm.storedMessages) + pm.storageMutex.RUnlock() + + if remaining == 0 { + pm.logger.Info("All messages processed successfully before shutdown") + break + } + + pm.logger.Info("Still processing messages", + zap.Int("remaining", remaining), + zap.Int("attempt", i+1)) + + select { + case <-ctx.Done(): + pm.logger.Warn("Shutdown timeout - stopping with messages remaining", + zap.Int("remaining", remaining)) + goto shutdown + case <-time.After(50 * time.Millisecond): + } + } + } + +shutdown: pm.retryWorkerCancel() done := make(chan struct{}) @@ -403,6 +475,17 @@ func (pm *PersistenceManager) Shutdown(ctx context.Context) error { pm.logger.Warn("Retry worker shutdown timed out") } + pm.storageMutex.RLock() + finalCount := len(pm.storedMessages) + pm.storageMutex.RUnlock() + + if finalCount > 0 { + pm.logger.Warn("Shutdown complete with unprocessed messages", + zap.Int("unprocessed_count", finalCount)) + } else { + pm.logger.Info("Shutdown complete - all messages processed") + } + return nil } diff --git a/receiver/otlpreceiver/testdata/persistence.yaml b/receiver/otlpreceiver/testdata/persistence.yaml index e940feb4a831..55112262ec77 100644 --- a/receiver/otlpreceiver/testdata/persistence.yaml +++ b/receiver/otlpreceiver/testdata/persistence.yaml @@ -5,28 +5,17 @@ receivers: endpoint: localhost:4318 persistence: enabled: true - storage: file_storage - retry_interval: 5s - max_retries: 3 - -extensions: - file_storage: - directory: /tmp/otel-collector-storage - timeout: 1s + retry_interval: 100ms +exporters: + debug: + verbosity: detailed + sampling_initial: 1000 + sampling_thereafter: 1000 service: - extensions: [file_storage] pipelines: - traces: - receivers: [otlp] - exporters: [logging] - metrics: - receivers: [otlp] - exporters: [logging] logs: receivers: [otlp] - exporters: [logging] + exporters: [debug] + -exporters: - logging: - verbosity: detailed diff --git a/test-config.yaml b/test-config.yaml deleted file mode 100644 index ce98092e9f24..000000000000 --- a/test-config.yaml +++ /dev/null @@ -1,32 +0,0 @@ -receivers: - otlp: - protocols: - grpc: - endpoint: 0.0.0.0:4317 - http: - endpoint: 0.0.0.0:4318 - persistence: - enabled: true - retry_interval: 5s - -processors: - batch: - -exporters: - debug: - verbosity: detailed - -service: - pipelines: - logs: - receivers: [otlp] - processors: [batch] - exporters: [debug] - traces: - receivers: [otlp] - processors: [batch] - exporters: [debug] - metrics: - receivers: [otlp] - processors: [batch] - exporters: [debug] From 6cf449212a297834d02b667f7a1b088501cfc9d6 Mon Sep 17 00:00:00 2001 From: MJarmo Date: Thu, 16 Oct 2025 15:13:51 +0200 Subject: [PATCH 09/10] fix tests Signed-off-by: MJarmo --- exporter/debugexporter/exporter.go | 178 +++++++++++++++++++++- exporter/debugexporter/factory.go | 7 +- receiver/otlpreceiver/persistence_test.go | 6 +- 3 files changed, 183 insertions(+), 8 deletions(-) diff --git a/exporter/debugexporter/exporter.go b/exporter/debugexporter/exporter.go index 46dae2bfeb80..792201ceb357 100644 --- a/exporter/debugexporter/exporter.go +++ b/exporter/debugexporter/exporter.go @@ -5,18 +5,25 @@ package debugexporter // import "go.opentelemetry.io/collector/exporter/debugexp import ( "context" + "fmt" + "os" + "sync" + "sync/atomic" "go.uber.org/zap" "go.opentelemetry.io/collector/config/configtelemetry" "go.opentelemetry.io/collector/exporter/debugexporter/internal/normal" "go.opentelemetry.io/collector/exporter/debugexporter/internal/otlptext" + "go.opentelemetry.io/collector/pdata/pcommon" "go.opentelemetry.io/collector/pdata/plog" "go.opentelemetry.io/collector/pdata/pmetric" "go.opentelemetry.io/collector/pdata/pprofile" "go.opentelemetry.io/collector/pdata/ptrace" ) +var failureCnt atomic.Int32 + type debugExporter struct { verbosity configtelemetry.Level logger *zap.Logger @@ -24,6 +31,13 @@ type debugExporter struct { metricsMarshaler pmetric.Marshaler tracesMarshaler ptrace.Marshaler profilesMarshaler pprofile.Marshaler + receivedLogs sync.Map + + totalLogsCount atomic.Int64 + exportCallCount atomic.Int64 + expectedLogs int64 + logFile *os.File + fileMutex sync.Mutex } func newDebugExporter(logger *zap.Logger, verbosity configtelemetry.Level) *debugExporter { @@ -42,6 +56,12 @@ func newDebugExporter(logger *zap.Logger, verbosity configtelemetry.Level) *debu tracesMarshaler = normal.NewNormalTracesMarshaler() profilesMarshaler = normal.NewNormalProfilesMarshaler() } + + logFile, err := os.OpenFile("test-logs.txt", os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644) + if err != nil { + logger.Error("Failed to open test-logs.txt for writing", zap.Error(err)) + } + return &debugExporter{ verbosity: verbosity, logger: logger, @@ -49,6 +69,8 @@ func newDebugExporter(logger *zap.Logger, verbosity configtelemetry.Level) *debu metricsMarshaler: metricsMarshaler, tracesMarshaler: tracesMarshaler, profilesMarshaler: profilesMarshaler, + expectedLogs: 10000, + logFile: logFile, } } @@ -86,19 +108,87 @@ func (s *debugExporter) pushMetrics(_ context.Context, md pmetric.Metrics) error } func (s *debugExporter) pushLogs(_ context.Context, ld plog.Logs) error { + callNumber := s.exportCallCount.Add(1) + + if callNumber%10 == 0 { + s.logger.Warn("Simulating export failure for 10th call", zap.Int64("call_number", callNumber)) + return fmt.Errorf("simulated export failure on call %d", callNumber) + } + + logRecordCount := ld.LogRecordCount() + + for i := 0; i < ld.ResourceLogs().Len(); i++ { + resourceLogs := ld.ResourceLogs().At(i) + for j := 0; j < resourceLogs.ScopeLogs().Len(); j++ { + scopeLogs := resourceLogs.ScopeLogs().At(j) + for k := 0; k < scopeLogs.LogRecords().Len(); k++ { + logRecord := scopeLogs.LogRecords().At(k) + if counter, exists := logRecord.Attributes().Get("test.log.counter"); exists { + s.receivedLogs.Store(counter.Int(), true) + s.writeLogToFile(logRecord, s.totalLogsCount.Load()) + } + } + } + } + + currentCount := s.totalLogsCount.Add(int64(logRecordCount)) + s.logger.Info("Logs body", zap.String("body", ld.ResourceLogs().At(0).ScopeLogs().At(0).LogRecords().At(0).Body().AsString())) s.logger.Info("Logs", zap.Int("resource logs", ld.ResourceLogs().Len()), - zap.Int("log records", ld.LogRecordCount())) + zap.Int("log records", logRecordCount), + zap.Int64("total received", currentCount), + zap.Int64("currentCount", currentCount), + ) + + if currentCount >= s.expectedLogs { + uniqueCount := int64(0) + var missingLogs []int64 + receivedCounters := make(map[int64]bool) + + s.receivedLogs.Range(func(key, value interface{}) bool { + if counter, ok := key.(int64); ok { + receivedCounters[counter] = true + uniqueCount++ + } + return true + }) + + for i := int64(1); i <= s.expectedLogs; i++ { + if !receivedCounters[i] { + missingLogs = append(missingLogs, i) + } + } + + successRate := (float64(uniqueCount) / float64(s.expectedLogs)) * 100 + + s.logger.Info("SUCCESS RATE REPORT", + zap.Int64("expected_logs", s.expectedLogs), + zap.Int64("total_received", currentCount), + zap.Int64("unique_logs", uniqueCount), + zap.Int("missing_count", len(missingLogs)), + zap.Float64("success_rate_percent", successRate)) + + s.writeReportToFile(s.expectedLogs, currentCount, uniqueCount, int64(len(missingLogs)), successRate, missingLogs) + + if len(missingLogs) > 0 && len(missingLogs) <= 100 { + s.logger.Warn("Missing log counters", zap.Int64s("missing", missingLogs)) + } else if len(missingLogs) > 100 { + s.logger.Warn("Too many missing logs to list", + zap.Int("missing_count", len(missingLogs)), + zap.Int64("first_missing", missingLogs[0]), + zap.Int64("last_missing", missingLogs[len(missingLogs)-1])) + } + } if s.verbosity == configtelemetry.LevelBasic { return nil } - buf, err := s.logsMarshaler.MarshalLogs(ld) - if err != nil { - return err - } - s.logger.Info(string(buf)) + // buf, err := s.logsMarshaler.MarshalLogs(ld) + // if err != nil { + // return err + // } + // s.logger.Info(string(buf)) return nil } @@ -118,3 +208,79 @@ func (s *debugExporter) pushProfiles(_ context.Context, pd pprofile.Profiles) er s.logger.Info(string(buf)) return nil } + +func (s *debugExporter) writeLogToFile(logRecord plog.LogRecord, counter int64) { + if s.logFile == nil { + return + } + + s.fileMutex.Lock() + defer s.fileMutex.Unlock() + + logLine := fmt.Sprintf("Counter: %d | Body: %s | Severity: %s | Timestamp: %s", + counter, + logRecord.Body().AsString(), + logRecord.SeverityText(), + logRecord.Timestamp().String()) + + var attrs string + logRecord.Attributes().Range(func(k string, v pcommon.Value) bool { + attrs += fmt.Sprintf(" | %s: %s", k, v.AsString()) + return true + }) + logLine += attrs + "\n" + + _, err := s.logFile.WriteString(logLine) + if err != nil { + s.logger.Error("Failed to write log to file", zap.Error(err)) + } +} + +func (s *debugExporter) writeReportToFile(expectedLogs, totalReceived, uniqueLogs, missingCount int64, successRate float64, missingLogs []int64) { + if s.logFile == nil { + return + } + + s.fileMutex.Lock() + defer s.fileMutex.Unlock() + + report := fmt.Sprintf("\n"+ + "========================================\n"+ + " SUCCESS RATE REPORT \n"+ + "========================================\n"+ + "Expected Logs: %d\n"+ + "Total Received: %d\n"+ + "Unique Logs: %d\n"+ + "Successfully Received: %d\n"+ + "Missing/Failed: %d\n"+ + "Success Rate: %.2f%%\n"+ + "========================================\n", + expectedLogs, totalReceived, uniqueLogs, uniqueLogs, missingCount, successRate) + + if missingCount > 0 && missingCount <= 100 { + report += fmt.Sprintf("\nMissing Log Counters:\n%v\n", missingLogs) + } else if missingCount > 100 { + report += fmt.Sprintf("\nMissing Logs Count: %d\n", missingCount) + report += fmt.Sprintf("First Missing: %d\n", missingLogs[0]) + report += fmt.Sprintf("Last Missing: %d\n", missingLogs[len(missingLogs)-1]) + } + + report += "========================================\n\n" + + _, err := s.logFile.WriteString(report) + if err != nil { + s.logger.Error("Failed to write report to file", zap.Error(err)) + } +} + +func (s *debugExporter) Shutdown(_ context.Context) error { + if s.logFile != nil { + s.fileMutex.Lock() + defer s.fileMutex.Unlock() + if err := s.logFile.Close(); err != nil { + s.logger.Error("Failed to close log file", zap.Error(err)) + return err + } + } + return nil +} diff --git a/exporter/debugexporter/factory.go b/exporter/debugexporter/factory.go index 49d83e47dec2..17c914c3d10d 100644 --- a/exporter/debugexporter/factory.go +++ b/exporter/debugexporter/factory.go @@ -86,7 +86,12 @@ func createLogs(ctx context.Context, set exporter.Settings, config component.Con debug.pushLogs, exporterhelper.WithCapabilities(consumer.Capabilities{MutatesData: false}), exporterhelper.WithTimeout(exporterhelper.TimeoutConfig{Timeout: 0}), - exporterhelper.WithShutdown(otlptext.LoggerSync(exporterLogger)), + exporterhelper.WithShutdown(func(ctx context.Context) error { + if err := debug.Shutdown(ctx); err != nil { + exporterLogger.Error("Failed to shutdown debug exporter", zap.Error(err)) + } + return otlptext.LoggerSync(exporterLogger)(ctx) + }), ) } diff --git a/receiver/otlpreceiver/persistence_test.go b/receiver/otlpreceiver/persistence_test.go index bae1b849dc2d..a911d53c0535 100644 --- a/receiver/otlpreceiver/persistence_test.go +++ b/receiver/otlpreceiver/persistence_test.go @@ -21,11 +21,15 @@ func TestPersistenceManager(t *testing.T) { pm := NewPersistenceManager(config, logger, nil) // Test storing a message - err := pm.StoreMessage(context.Background(), []byte("test data"), "application/json", "logs") + messageID, err := pm.StoreMessage(context.Background(), []byte("test data"), "application/json", "logs") if err != nil { t.Fatalf("Failed to store message: %v", err) } + if messageID == "" { + t.Fatal("Expected non-empty message ID") + } + // Test retrieving stored messages messages, err := pm.GetStoredMessages(context.Background()) if err != nil { From b5eaa4b21009503d38bde818243dda21f1fff8ed Mon Sep 17 00:00:00 2001 From: MJarmo Date: Thu, 16 Oct 2025 15:27:20 +0200 Subject: [PATCH 10/10] rmv test debuq Signed-off-by: MJarmo --- exporter/debugexporter/exporter.go | 178 +---------------------------- exporter/debugexporter/factory.go | 7 +- 2 files changed, 7 insertions(+), 178 deletions(-) diff --git a/exporter/debugexporter/exporter.go b/exporter/debugexporter/exporter.go index 792201ceb357..46dae2bfeb80 100644 --- a/exporter/debugexporter/exporter.go +++ b/exporter/debugexporter/exporter.go @@ -5,25 +5,18 @@ package debugexporter // import "go.opentelemetry.io/collector/exporter/debugexp import ( "context" - "fmt" - "os" - "sync" - "sync/atomic" "go.uber.org/zap" "go.opentelemetry.io/collector/config/configtelemetry" "go.opentelemetry.io/collector/exporter/debugexporter/internal/normal" "go.opentelemetry.io/collector/exporter/debugexporter/internal/otlptext" - "go.opentelemetry.io/collector/pdata/pcommon" "go.opentelemetry.io/collector/pdata/plog" "go.opentelemetry.io/collector/pdata/pmetric" "go.opentelemetry.io/collector/pdata/pprofile" "go.opentelemetry.io/collector/pdata/ptrace" ) -var failureCnt atomic.Int32 - type debugExporter struct { verbosity configtelemetry.Level logger *zap.Logger @@ -31,13 +24,6 @@ type debugExporter struct { metricsMarshaler pmetric.Marshaler tracesMarshaler ptrace.Marshaler profilesMarshaler pprofile.Marshaler - receivedLogs sync.Map - - totalLogsCount atomic.Int64 - exportCallCount atomic.Int64 - expectedLogs int64 - logFile *os.File - fileMutex sync.Mutex } func newDebugExporter(logger *zap.Logger, verbosity configtelemetry.Level) *debugExporter { @@ -56,12 +42,6 @@ func newDebugExporter(logger *zap.Logger, verbosity configtelemetry.Level) *debu tracesMarshaler = normal.NewNormalTracesMarshaler() profilesMarshaler = normal.NewNormalProfilesMarshaler() } - - logFile, err := os.OpenFile("test-logs.txt", os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644) - if err != nil { - logger.Error("Failed to open test-logs.txt for writing", zap.Error(err)) - } - return &debugExporter{ verbosity: verbosity, logger: logger, @@ -69,8 +49,6 @@ func newDebugExporter(logger *zap.Logger, verbosity configtelemetry.Level) *debu metricsMarshaler: metricsMarshaler, tracesMarshaler: tracesMarshaler, profilesMarshaler: profilesMarshaler, - expectedLogs: 10000, - logFile: logFile, } } @@ -108,87 +86,19 @@ func (s *debugExporter) pushMetrics(_ context.Context, md pmetric.Metrics) error } func (s *debugExporter) pushLogs(_ context.Context, ld plog.Logs) error { - callNumber := s.exportCallCount.Add(1) - - if callNumber%10 == 0 { - s.logger.Warn("Simulating export failure for 10th call", zap.Int64("call_number", callNumber)) - return fmt.Errorf("simulated export failure on call %d", callNumber) - } - - logRecordCount := ld.LogRecordCount() - - for i := 0; i < ld.ResourceLogs().Len(); i++ { - resourceLogs := ld.ResourceLogs().At(i) - for j := 0; j < resourceLogs.ScopeLogs().Len(); j++ { - scopeLogs := resourceLogs.ScopeLogs().At(j) - for k := 0; k < scopeLogs.LogRecords().Len(); k++ { - logRecord := scopeLogs.LogRecords().At(k) - if counter, exists := logRecord.Attributes().Get("test.log.counter"); exists { - s.receivedLogs.Store(counter.Int(), true) - s.writeLogToFile(logRecord, s.totalLogsCount.Load()) - } - } - } - } - - currentCount := s.totalLogsCount.Add(int64(logRecordCount)) - s.logger.Info("Logs body", zap.String("body", ld.ResourceLogs().At(0).ScopeLogs().At(0).LogRecords().At(0).Body().AsString())) s.logger.Info("Logs", zap.Int("resource logs", ld.ResourceLogs().Len()), - zap.Int("log records", logRecordCount), - zap.Int64("total received", currentCount), - zap.Int64("currentCount", currentCount), - ) - - if currentCount >= s.expectedLogs { - uniqueCount := int64(0) - var missingLogs []int64 - receivedCounters := make(map[int64]bool) - - s.receivedLogs.Range(func(key, value interface{}) bool { - if counter, ok := key.(int64); ok { - receivedCounters[counter] = true - uniqueCount++ - } - return true - }) - - for i := int64(1); i <= s.expectedLogs; i++ { - if !receivedCounters[i] { - missingLogs = append(missingLogs, i) - } - } - - successRate := (float64(uniqueCount) / float64(s.expectedLogs)) * 100 - - s.logger.Info("SUCCESS RATE REPORT", - zap.Int64("expected_logs", s.expectedLogs), - zap.Int64("total_received", currentCount), - zap.Int64("unique_logs", uniqueCount), - zap.Int("missing_count", len(missingLogs)), - zap.Float64("success_rate_percent", successRate)) - - s.writeReportToFile(s.expectedLogs, currentCount, uniqueCount, int64(len(missingLogs)), successRate, missingLogs) - - if len(missingLogs) > 0 && len(missingLogs) <= 100 { - s.logger.Warn("Missing log counters", zap.Int64s("missing", missingLogs)) - } else if len(missingLogs) > 100 { - s.logger.Warn("Too many missing logs to list", - zap.Int("missing_count", len(missingLogs)), - zap.Int64("first_missing", missingLogs[0]), - zap.Int64("last_missing", missingLogs[len(missingLogs)-1])) - } - } + zap.Int("log records", ld.LogRecordCount())) if s.verbosity == configtelemetry.LevelBasic { return nil } - // buf, err := s.logsMarshaler.MarshalLogs(ld) - // if err != nil { - // return err - // } - // s.logger.Info(string(buf)) + buf, err := s.logsMarshaler.MarshalLogs(ld) + if err != nil { + return err + } + s.logger.Info(string(buf)) return nil } @@ -208,79 +118,3 @@ func (s *debugExporter) pushProfiles(_ context.Context, pd pprofile.Profiles) er s.logger.Info(string(buf)) return nil } - -func (s *debugExporter) writeLogToFile(logRecord plog.LogRecord, counter int64) { - if s.logFile == nil { - return - } - - s.fileMutex.Lock() - defer s.fileMutex.Unlock() - - logLine := fmt.Sprintf("Counter: %d | Body: %s | Severity: %s | Timestamp: %s", - counter, - logRecord.Body().AsString(), - logRecord.SeverityText(), - logRecord.Timestamp().String()) - - var attrs string - logRecord.Attributes().Range(func(k string, v pcommon.Value) bool { - attrs += fmt.Sprintf(" | %s: %s", k, v.AsString()) - return true - }) - logLine += attrs + "\n" - - _, err := s.logFile.WriteString(logLine) - if err != nil { - s.logger.Error("Failed to write log to file", zap.Error(err)) - } -} - -func (s *debugExporter) writeReportToFile(expectedLogs, totalReceived, uniqueLogs, missingCount int64, successRate float64, missingLogs []int64) { - if s.logFile == nil { - return - } - - s.fileMutex.Lock() - defer s.fileMutex.Unlock() - - report := fmt.Sprintf("\n"+ - "========================================\n"+ - " SUCCESS RATE REPORT \n"+ - "========================================\n"+ - "Expected Logs: %d\n"+ - "Total Received: %d\n"+ - "Unique Logs: %d\n"+ - "Successfully Received: %d\n"+ - "Missing/Failed: %d\n"+ - "Success Rate: %.2f%%\n"+ - "========================================\n", - expectedLogs, totalReceived, uniqueLogs, uniqueLogs, missingCount, successRate) - - if missingCount > 0 && missingCount <= 100 { - report += fmt.Sprintf("\nMissing Log Counters:\n%v\n", missingLogs) - } else if missingCount > 100 { - report += fmt.Sprintf("\nMissing Logs Count: %d\n", missingCount) - report += fmt.Sprintf("First Missing: %d\n", missingLogs[0]) - report += fmt.Sprintf("Last Missing: %d\n", missingLogs[len(missingLogs)-1]) - } - - report += "========================================\n\n" - - _, err := s.logFile.WriteString(report) - if err != nil { - s.logger.Error("Failed to write report to file", zap.Error(err)) - } -} - -func (s *debugExporter) Shutdown(_ context.Context) error { - if s.logFile != nil { - s.fileMutex.Lock() - defer s.fileMutex.Unlock() - if err := s.logFile.Close(); err != nil { - s.logger.Error("Failed to close log file", zap.Error(err)) - return err - } - } - return nil -} diff --git a/exporter/debugexporter/factory.go b/exporter/debugexporter/factory.go index 17c914c3d10d..49d83e47dec2 100644 --- a/exporter/debugexporter/factory.go +++ b/exporter/debugexporter/factory.go @@ -86,12 +86,7 @@ func createLogs(ctx context.Context, set exporter.Settings, config component.Con debug.pushLogs, exporterhelper.WithCapabilities(consumer.Capabilities{MutatesData: false}), exporterhelper.WithTimeout(exporterhelper.TimeoutConfig{Timeout: 0}), - exporterhelper.WithShutdown(func(ctx context.Context) error { - if err := debug.Shutdown(ctx); err != nil { - exporterLogger.Error("Failed to shutdown debug exporter", zap.Error(err)) - } - return otlptext.LoggerSync(exporterLogger)(ctx) - }), + exporterhelper.WithShutdown(otlptext.LoggerSync(exporterLogger)), ) }