Files
amcs/vendor/github.com/getsentry/sentry-go/internal/telemetry/processor.go
T
Hein 1adf50e3db
CI / build-and-test (push) Failing after 1s
Release / release (push) Failing after 19m26s
fix(go.sum): update ResolveSpec dependency to v1.0.87
2026-06-23 13:17:16 +02:00

56 lines
1.6 KiB
Go

package telemetry
import (
"context"
"time"
"github.com/getsentry/sentry-go/internal/protocol"
"github.com/getsentry/sentry-go/internal/ratelimit"
"github.com/getsentry/sentry-go/report"
)
// Processor is the top-level object that wraps the scheduler and buffers.
type Processor struct {
scheduler *Scheduler
}
// NewProcessor creates a new Processor with the given configuration.
func NewProcessor(
buffers map[ratelimit.Category]Buffer[protocol.TelemetryItem],
transport protocol.TelemetryTransport,
dsn *protocol.Dsn,
sdkInfo func() *protocol.SdkInfo,
recorder report.ClientReportRecorder,
) *Processor {
scheduler := NewScheduler(buffers, transport, dsn, sdkInfo, recorder)
scheduler.Start()
return &Processor{
scheduler: scheduler,
}
}
// Add adds a TelemetryItem to the appropriate buffer based on its category.
//
// The processor should call MakeSerializationSafe to eliminate any race on user mutable fields,
// since the serialization happens on a background goroutine.
func (b *Processor) Add(item protocol.TelemetryItem) bool {
item.MakeSerializationSafe()
return b.scheduler.Add(item)
}
// Flush forces all buffers to flush within the given timeout.
func (b *Processor) Flush(timeout time.Duration) bool {
return b.scheduler.Flush(timeout)
}
// FlushWithContext flushes with a custom context for cancellation.
func (b *Processor) FlushWithContext(ctx context.Context) bool {
return b.scheduler.FlushWithContext(ctx)
}
// Close stops the buffer, flushes remaining data, and releases resources.
func (b *Processor) Close(timeout time.Duration) {
b.scheduler.Stop(timeout)
}