Documentation
¶
Index ¶
- Constants
- func ConfiguredCloudExporters(ctx context.Context) (sdktrace.SpanExporter, sdklog.Exporter, sdkmetric.Exporter, bool)
- func MeasuringStreamClientInterceptor() grpc.StreamClientInterceptor
- func MeasuringUnaryClientInterceptor() grpc.UnaryClientInterceptor
- func MeasuringUnaryServerInterceptor() grpc.UnaryServerInterceptor
- func NewLargeQueueLiveSpanProcessor(exp sdktrace.SpanExporter) *telemetry.LiveSpanProcessor
- func NewLogBatchProcessor(exp sdklog.Exporter) *sdklog.BatchProcessor
- func ReexportMetricsFromPB(ctx context.Context, exps []sdkmetric.Exporter, ...) error
- func Task(ctx context.Context, name string, fn func(context.Context) error) (rerr error)
- func TaskRet[T any](ctx context.Context, name string, fn func(context.Context) (T, error)) (ret T, rerr error)
- func URLForTrace(ctx context.Context) (url string, msg string, ok bool)
- type EnvGetter
- type LabelFlag
- type Labels
- func (labels Labels) AsMap() map[string]string
- func (labels Labels) Get(key string) (string, bool)
- func (labels *Labels) UnmarshalJSON(dt []byte) error
- func (labels Labels) UserAgent() string
- func (labels Labels) WithCILabels() Labels
- func (labels Labels) WithCircleCILabels() Labels
- func (labels Labels) WithClientLabels(engineVersion string) Labels
- func (labels Labels) WithEngineLabel(engineName string) Labels
- func (labels Labels) WithGitHubLabels() Labels
- func (labels Labels) WithGitLabLabels() Labels
- func (labels Labels) WithGitLabels(workdir string) Labels
- func (labels Labels) WithHarnessLabels() Labels
- func (labels Labels) WithJenkinsLabels() Labels
- func (labels Labels) WithServerLabels(engineVersion, os, arch string, cacheEnabled bool) Labels
- func (labels Labels) WithVCSLabels(workdir string) Labels
- type LogFanOutExporter
- type OSEnvGetter
- type SpanFanOutExporter
- type SpanHeartbeater
- type TraceImportSinks
- type TraceImporter
- func (imp *TraceImporter) ImportLogs(ctx context.Context, req *collogspb.ExportLogsServiceRequest) error
- func (imp *TraceImporter) ImportMetrics(ctx context.Context, req *colmetricspb.ExportMetricsServiceRequest) error
- func (imp *TraceImporter) ImportSpans(ctx context.Context, req *coltracepb.ExportTraceServiceRequest) error
- func (imp *TraceImporter) Seal(ctx context.Context) error
Constants ¶
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.
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.
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
Types ¶
type Labels ¶
type Labels struct {
// contains filtered or unexported fields
}
func LoadDefaultLabels ¶
func (Labels) AsMap ¶ added in v0.19.7
AsMap returns a reference to the internal labels map. It's not intended to be moodified by the caller.
func (*Labels) UnmarshalJSON ¶
func (Labels) WithCILabels ¶
func (Labels) WithCircleCILabels ¶
func (Labels) WithClientLabels ¶
func (Labels) WithEngineLabel ¶
func (Labels) WithGitHubLabels ¶
func (Labels) WithGitLabLabels ¶
func (Labels) WithGitLabels ¶
func (Labels) WithHarnessLabels ¶ added in v0.14.0
func (Labels) WithJenkinsLabels ¶ added in v0.12.0
func (Labels) WithServerLabels ¶
func (Labels) WithVCSLabels ¶
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 ¶
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.
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.
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
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 ¶
func (imp *TraceImporter) ImportLogs(ctx context.Context, req *collogspb.ExportLogsServiceRequest) error
ImportLogs folds one OTLP log export request of the foreign trace in.
func (*TraceImporter) ImportMetrics ¶
func (imp *TraceImporter) ImportMetrics(ctx context.Context, req *colmetricspb.ExportMetricsServiceRequest) error
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 ¶
func (imp *TraceImporter) ImportSpans(ctx context.Context, req *coltracepb.ExportTraceServiceRequest) error
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.