telemetry

package
v1.0.0-beta.14 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 18, 2026 License: Apache-2.0 Imports: 45 Imported by: 0

Documentation

Index

Constants

View Source
const (
	LargeSpanQueueSize       = 16384 // 16Ki — 8x the SDK default; bounds retained snapshots per hop
	LargeSpanExportBatchSize = 2048  // drains a full queue in 8 batches; bounds per-batch references
)

Enlarged, BOUNDED BatchSpanProcessor sizes for the span hops that carry a trace toward Dagger Cloud. The OTel BSP queue is non-blocking and silently DROPS spans on overflow; the SDK default of 2048 slots is too small for a burst like a cold engine build (~15k spans, live-double-emitted into ~30k records), which is why these hops use a larger queue at all.

The queue must stay MODEST as well as bounded, because its worst case is paid per-processor worst case used to multiply by client and ancestor counts. The engine server now owns one processor per session and routes snapshots by their stamped origin ID; the bounded queue still protects against burst retention.

Kept BOUNDED — never BlockOnQueueFull — so telemetry can NEVER stall the build. If a burst still overflows, spans are dropped rather than retained: the wcprof completeness checksum (received < declared, see engine/server/wcprofcount.go) catches the loss loudly and the offline analyzer refuses the trace instead of silently ranking on partial data.

View Source
const (
	LogQueueSize          = 2048
	LogExportInterval     = 250 * time.Millisecond
	LogExportMaxBatchSize = 128
)

Log-record batch settings for client DB log routes (engine/server).

The SDK BatchProcessor's poll loop clones its entire batchSize-long []Record buffer on EVERY export tick while the exporter is ready — even when zero records were dequeued (sdk/log batch.go: TryDequeue's write callback does buf = slices.Clone(buf), and EnqueueExport returns true for an empty slice). A Record is a large struct (~0.5 KiB inline), so the resulting allocation churn is processors × ticks/sec × batchSize × sizeof(Record), INDEPENDENT of actual log volume. The engine therefore uses one processor per session and routes each batch by stamped origin ID.

Interval and batch size each scale that churn linearly. Crucially, OnEmit self-flushes as soon as a full batch accumulates (pollTrigger), so a longer interval delays only SPARSE log records — bursts still export immediately. 2048 queue slots keep the SDK's previous default ceiling explicit for debug accounting. 250ms/128 cuts the volume-independent churn ~10x while keeping trickle latency well below human-noticeable for log output.

View Source
const HeartbeatInterval = 30 * time.Second

Variables

This section is empty.

Functions

func MeasuringStreamClientInterceptor

func MeasuringStreamClientInterceptor() grpc.StreamClientInterceptor

func MeasuringUnaryClientInterceptor

func MeasuringUnaryClientInterceptor() grpc.UnaryClientInterceptor

func MeasuringUnaryServerInterceptor

func MeasuringUnaryServerInterceptor() grpc.UnaryServerInterceptor

func NewLargeQueueLiveSpanProcessor

func NewLargeQueueLiveSpanProcessor(exp sdktrace.SpanExporter) *telemetry.LiveSpanProcessor

NewLargeQueueLiveSpanProcessor is otel.NewLiveSpanProcessor with the enlarged bounded queue above in place of the default 2048-slot one. Used on the CLI→Cloud exporter (internal/cmd/dagger) and the engine's per-client store exporters so a big-burst trace arrives complete.

func NewLogBatchProcessor

func NewLogBatchProcessor(exp sdklog.Exporter) *sdklog.BatchProcessor

NewLogBatchProcessor is the log analog of NewLargeQueueLiveSpanProcessor: the session-owned batch processor for client DB log routes, with the bounded churn settings above in place of the SDK defaults.

func ReexportMetricsFromPB added in v0.13.6

func ReexportMetricsFromPB(ctx context.Context, exps []sdkmetric.Exporter, req *colmetricspb.ExportMetricsServiceRequest) error

func Task

func Task(ctx context.Context, name string, fn func(context.Context) error) (rerr error)

func TaskRet

func TaskRet[T any](ctx context.Context, name string, fn func(context.Context) (T, error)) (ret T, rerr error)

func URLForTrace

func URLForTrace(ctx context.Context) (url string, msg string, ok bool)

Types

type EnvGetter added in v0.19.7

type EnvGetter interface {
	Getenv(key string) string
}

type LabelFlag

type LabelFlag struct {
	Labels
}

func NewLabelFlag

func NewLabelFlag() LabelFlag

func (LabelFlag) Set

func (flag LabelFlag) Set(s string) error

func (LabelFlag) String

func (flag LabelFlag) String() string

func (LabelFlag) Type

func (flag LabelFlag) Type() string

type Labels

type Labels struct {
	// contains filtered or unexported fields
}

func LoadDefaultLabels

func LoadDefaultLabels(workdir, clientEngineVersion string) Labels

func NewLabels added in v0.19.7

func NewLabels(labels map[string]string, eg EnvGetter, ghEventPayload []byte) Labels

func (Labels) AsMap added in v0.19.7

func (labels Labels) AsMap() map[string]string

AsMap returns a reference to the internal labels map. It's not intended to be moodified by the caller.

func (Labels) Get added in v0.19.7

func (labels Labels) Get(key string) (string, bool)

func (*Labels) UnmarshalJSON

func (labels *Labels) UnmarshalJSON(dt []byte) error

func (Labels) UserAgent

func (labels Labels) UserAgent() string

func (Labels) WithCILabels

func (labels Labels) WithCILabels() Labels

func (Labels) WithCircleCILabels

func (labels Labels) WithCircleCILabels() Labels

func (Labels) WithClientLabels

func (labels Labels) WithClientLabels(engineVersion string) Labels

func (Labels) WithEngineLabel

func (labels Labels) WithEngineLabel(engineName string) Labels

func (Labels) WithGitHubLabels

func (labels Labels) WithGitHubLabels() Labels

func (Labels) WithGitLabLabels

func (labels Labels) WithGitLabLabels() Labels

func (Labels) WithGitLabels

func (labels Labels) WithGitLabels(workdir string) Labels

func (Labels) WithHarnessLabels added in v0.14.0

func (labels Labels) WithHarnessLabels() Labels

func (Labels) WithJenkinsLabels added in v0.12.0

func (labels Labels) WithJenkinsLabels() Labels

func (Labels) WithServerLabels

func (labels Labels) WithServerLabels(engineVersion, os, arch string, cacheEnabled bool) Labels

func (Labels) WithVCSLabels

func (labels Labels) WithVCSLabels(workdir string) Labels

type LogFanOutExporter

type LogFanOutExporter struct {
	// contains filtered or unexported fields
}

LogFanOutExporter is the log analog of SpanFanOutExporter: every Export call is delivered to each destination in order.

func NewLogFanOutExporter

func NewLogFanOutExporter(dests ...sdklog.Exporter) *LogFanOutExporter

NewLogFanOutExporter fans log records out to every destination exporter.

func (*LogFanOutExporter) Export

func (e *LogFanOutExporter) Export(ctx context.Context, records []sdklog.Record) error

Export exports to each destination in order. A destination's error does not prevent export to the others; all errors are returned joined.

func (*LogFanOutExporter) ForceFlush

func (e *LogFanOutExporter) ForceFlush(ctx context.Context) error

ForceFlush flushes every destination, joining any errors.

func (*LogFanOutExporter) Shutdown

func (e *LogFanOutExporter) Shutdown(ctx context.Context) error

Shutdown shuts down every destination, joining any errors.

type OSEnvGetter added in v0.19.7

type OSEnvGetter struct{}

func (OSEnvGetter) Getenv added in v0.19.7

func (e OSEnvGetter) Getenv(key string) string

type SpanFanOutExporter

type SpanFanOutExporter struct {
	// contains filtered or unexported fields
}

SpanFanOutExporter delivers every ExportSpans call to each destination in order. See the package comment above for why the engine fans out at the exporter rather than registering per-destination processors.

func NewSpanFanOutExporter

func NewSpanFanOutExporter(dests ...sdktrace.SpanExporter) *SpanFanOutExporter

NewSpanFanOutExporter fans spans out to every destination exporter.

func (*SpanFanOutExporter) ExportSpans

func (e *SpanFanOutExporter) ExportSpans(ctx context.Context, spans []sdktrace.ReadOnlySpan) error

ExportSpans exports to each destination in order. A destination's error does not prevent export to the others; all errors are returned joined.

func (*SpanFanOutExporter) Shutdown

func (e *SpanFanOutExporter) Shutdown(ctx context.Context) error

Shutdown shuts down every destination, joining any errors.

type SpanHeartbeater

type SpanHeartbeater struct {
	sdktrace.SpanExporter
	// contains filtered or unexported fields
}

SpanHeartbeater is a SpanExporter that keeps track of live spans and re-exports them periodically to the underlying SpanExporter to indicate that they are indeed still live.

func NewSpanHeartbeater

func NewSpanHeartbeater(exp sdktrace.SpanExporter) *SpanHeartbeater

func (*SpanHeartbeater) ExportSpans

func (p *SpanHeartbeater) ExportSpans(ctx context.Context, spans []sdktrace.ReadOnlySpan) error

func (*SpanHeartbeater) Shutdown

func (p *SpanHeartbeater) Shutdown(ctx context.Context) error

type TraceImportSinks

type TraceImportSinks struct {
	Spans   sdktrace.SpanExporter
	Logs    sdklog.Exporter
	Metrics sdkmetric.Exporter
}

TraceImportSinks are the exporters an imported trace lands in — the live frontend's own (Frontend.SpanExporter, LogExporter, MetricExporter), which is the whole point: one DB, both sessions. A nil sink drops its stream.

type TraceImporter

type TraceImporter struct {
	// contains filtered or unexported fields
}

TraceImporter folds a foreign trace into a live client's sinks, applying §5.1's two fixes on the way through. Call ImportSpans/ImportLogs/ ImportMetrics per export request, in any order, then Seal exactly once when the stream ends.

Safe for concurrent use: the three streams arrive on separate connections.

func NewTraceImporter

func NewTraceImporter(sinks TraceImportSinks) *TraceImporter

NewTraceImporter returns an importer feeding sinks.

func (*TraceImporter) ImportLogs

ImportLogs folds one OTLP log export request of the foreign trace in.

func (*TraceImporter) ImportMetrics

ImportMetrics folds one OTLP metric export request of the foreign trace in. Imported metrics COUNT (§12): cost and token totals accumulate across a resume rather than restarting at zero.

func (*TraceImporter) ImportSpans

ImportSpans folds one OTLP span export request of the foreign trace in, stamping its roots passthrough and noting which of its spans never ended.

func (*TraceImporter) Seal

func (imp *TraceImporter) Seal(ctx context.Context) error

Seal ends the import: every span the capture left running is re-exported with an end time and the Canceled/LeftRunning marks, the same shape dagui.DB produces for its own root's leftovers (§5.1.2).

The end time is the imported root's, or — when the source session crashed hard enough that even its root never ended — the newest timestamp the capture carried. Neither is the truth about when that work stopped, which nothing recorded; both say "no later than this", which is the honest reading and the one the DB's own sweep already takes.

Idempotent: a second call has nothing left to seal.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL