telemetry

package
v1.0.0-beta.15 Latest Latest
Warning

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

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

Documentation

Index

Constants

View Source
const (
	// CallPayloadExportDelay is how long the first payload of a burst waits
	// for the rest of its recipe closure before the batch is exported.
	CallPayloadExportDelay = 5 * time.Millisecond

	// CallPayloadRetryBaseDelay is the backoff after the first failed export
	// of a payload batch; it doubles per consecutive failure up to
	// CallPayloadRetryMaxDelay.
	CallPayloadRetryBaseDelay = 50 * time.Millisecond
	CallPayloadRetryMaxDelay  = 2 * time.Second
	// CallPayloadMaxExportAttempts bounds how often one batch is retried
	// before it is dropped, so a dead client DB cannot wedge the queue
	// forever. Dropping is lossy: the session exporter releases the failed
	// targets, but a later closure walk only re-emits a dropped record if it
	// reaches it through a root that is itself still undelivered. When the
	// root already landed, every walk from it short-circuits at the root's
	// claim and the dropped dependency stays missing for that client unless
	// some other chain happens to include it. The drop is logged with the
	// records' digests so that gap is at least diagnosable.
	CallPayloadMaxExportAttempts = 8
)
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 (
	// LiveContentType identifies the framed binary OTLP subscription protocol.
	LiveContentType = "application/vnd.dagger.otlp.stream"
	// LegacyLiveContentType identifies the original SSE OTLP subscription protocol.
	LegacyLiveContentType = "text/event-stream"
	// LiveCursorHeader resumes a binary subscription after the last consumed row ID.
	LiveCursorHeader = "X-Dagger-Telemetry-Cursor"
	// LegacyLiveCursorHeader resumes an SSE subscription after the last event ID.
	LegacyLiveCursorHeader = "X-Last-Event-ID"

	// MaxLivePayloadSize bounds each protobuf frame while allowing a stream to
	// carry any number of frames.
	MaxLivePayloadSize = 64 << 20
)
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 CloudExportTimeout = 10 * time.Second

CloudExportTimeout bounds one request of a Cloud exporter, the OTLP HTTP exporters' own default.

View Source
const CloudTokenRefreshTimeout = 5 * time.Second

CloudTokenRefreshTimeout bounds one refresh of an expired OAuth token.

View Source
const HeartbeatInterval = 30 * time.Second

Variables

View Source
var (
	ErrArchiveSignalClosed  = errors.New("archive import signal is closed")
	ErrArchiveSpanAbandoned = errors.New("archive span import was permanently abandoned")
)
View Source
var (
	// ErrInvalidLiveFrame identifies malformed or out-of-bounds frame metadata.
	ErrInvalidLiveFrame = errors.New("invalid live telemetry frame")
	// ErrLiveStream identifies an error reported by the stream producer. It is
	// not a transport interruption and must not be retried at the same cursor.
	ErrLiveStream = errors.New("live telemetry stream error")
)
View Source
var ErrInvalidCloudURL = errors.New("DAGGER_CLOUD_URL is not a valid URL")

ErrInvalidCloudURL is the CloudEmitStatus error when DAGGER_CLOUD_URL is not a valid URL.

Functions

func BoundedTokenRefresh

func BoundedTokenRefresh(refresh func(context.Context) (*oauth2.Token, error)) func(context.Context) (*oauth2.Token, error)

BoundedTokenRefresh bounds each call of a token refresh callback by CloudTokenRefreshTimeout, on a context of its own, since exports refresh from background goroutines long after any request.

func CallSpanDigest

func CallSpanDigest(span sdktrace.ReadOnlySpan) (digest string, ok bool)

CallSpanDigest returns the digest a call span delivers the frame of: its dagger.io/dag.digest, the key the producer claims the payload under. ok is false for spans that carry no frame.

func ConfiguredCloudExporters

func ConfiguredCloudExporters(ctx context.Context) (sdktrace.SpanExporter, sdklog.Exporter, sdkmetric.Exporter, bool)

ConfiguredCloudExporters returns this process's Dagger Cloud exporters, built once by NewCloudExporters from the process's Cloud credential: DAGGER_CLOUD_TOKEN, or the `dagger login` token with its current organization, refreshed and persisted as it expires.

func IsCallPayloadRecord

func IsCallPayloadRecord(record sdklog.Record) bool

IsCallPayloadRecord reports whether a record is on the call payload channel: a bytes body whose content type attribute names an encoded call.

func IsCallSpan

func IsCallSpan(span sdktrace.ReadOnlySpan) bool

IsCallSpan reports whether a span carries its call's encoded frame (dagger.io/dag.call), which makes it the frame's delivery.

func MeasuringStreamClientInterceptor

func MeasuringStreamClientInterceptor() grpc.StreamClientInterceptor

func MeasuringUnaryClientInterceptor

func MeasuringUnaryClientInterceptor() grpc.UnaryClientInterceptor

func MeasuringUnaryServerInterceptor

func MeasuringUnaryServerInterceptor() grpc.UnaryServerInterceptor

func NewCloudExporters

func NewCloudExporters(ctx context.Context, cloudAuth *auth.Cloud, tokenRefreshFn func(context.Context) (*oauth2.Token, error), cloudURL string) (sdktrace.SpanExporter, sdklog.Exporter, sdkmetric.Exporter, error)

NewCloudExporters builds OTLP exporters for Dagger Cloud from a Cloud credential. cloudURL is the Cloud API URL; when empty it falls back to DAGGER_CLOUD_URL and then the default Cloud API URL.

Basic (engine token) and OIDC credentials are sent as a static Authorization header. OAuth access tokens expire, so the exporters instead share an HTTP client that sets the current bearer token on each request, refreshed through tokenRefreshFn when it is non-nil.

Each exporter stamps the X-Dagger-Export sequence of its own writer, and uploads use the private Cloud export transport. Every request is bounded by CloudExportTimeout: the OTLP exporters apply their own default timeout only to a client they build themselves.

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. The engine wraps it in WithoutCallSpans: call spans, which carry frames nothing else delivers, take the lossless CallSpanProcessor instead.

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 ProbeCloudURL

func ProbeCloudURL(ctx context.Context, cloudURL string) error

ProbeCloudURL checks that the Cloud API at cloudURL answers this process: one unauthenticated HEAD of its traces endpoint through the Cloud exporters' transport, so it goes through the same proxy and trusts the same certificates as they do. Any HTTP response counts; only failing to get one is an error. ctx bounds the probe.

func ReexportMetricsFromPB added in v0.13.6

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

func ResolveCloudURL

func ResolveCloudURL(cloudURL string) string

ResolveCloudURL returns the Cloud API URL the Cloud exporters use for cloudURL: cloudURL itself, else DAGGER_CLOUD_URL, else the default.

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)

func WithoutCallPayloads

func WithoutCallPayloads(next sdklog.Processor) sdklog.Processor

WithoutCallPayloads wraps an ordinary log processor so neither call payloads nor revisioned agent controls reach its bounded queue. Both have their own protected transport; ordinary processing would duplicate them and allow a recipe burst to evict exec output.

func WithoutCallSpans

func WithoutCallSpans(next sdktrace.SpanProcessor) sdktrace.SpanProcessor

WithoutCallSpans wraps an ordinary span processor so call spans never reach it: CallSpanProcessor carries them, and the ordinary bounded queue would both duplicate them and let them be dropped on overflow.

func WriteLiveError

func WriteLiveError(w io.Writer, cursor int64, streamErr error) error

WriteLiveError reports a producer error that a consumer must not retry at the same cursor. Error text is truncated to keep error frames tightly bounded.

func WriteLiveFrame

func WriteLiveFrame(w io.Writer, cursor int64, payload []byte) error

WriteLiveFrame writes one cursor-addressed protobuf payload.

func WriteLiveHello

func WriteLiveHello(w io.Writer, cursor int64) error

WriteLiveHello writes the marker sent as soon as a subscription is attached, before any data is available. Sending body bytes (rather than relying on a header-only flush) lets the stream announce itself through intermediaries that only flush on body writes.

func WriteLiveTerminal

func WriteLiveTerminal(w io.Writer, cursor int64) error

WriteLiveTerminal writes the marker sent after the server has drained the subscription.

Types

type ArchiveCut

type ArchiveCut struct {
	HighWater ArchiveHighWater
	SealAt    time.Time
}

ArchiveCut identifies one immutable version of a telemetry archive. The bootstrap and all three remainder streams must present this same cut.

type ArchiveHighWater

type ArchiveHighWater struct {
	Spans   int64
	Logs    int64
	Metrics int64
}

ArchiveHighWater is the immutable final cursor of each archive signal.

type ArchiveImportBatch

type ArchiveImportBatch struct {
	// Cursor is the source remainder cursor, or zero for bootstrap. A retry
	// after exporter enqueue only repeats the application barrier.
	Cursor  int64
	Spans   *coltracepb.ExportTraceServiceRequest
	Logs    *collogspb.ExportLogsServiceRequest
	Metrics *colmetricspb.ExportMetricsServiceRequest
}

ArchiveImportBatch carries exactly one OTLP export request. A bootstrap and its remainder intentionally use the same batch type and importer.

type ArchiveSignal

type ArchiveSignal string

ArchiveSignal identifies one independently streamed OTLP signal.

const (
	ArchiveSpans   ArchiveSignal = "spans"
	ArchiveLogs    ArchiveSignal = "logs"
	ArchiveMetrics ArchiveSignal = "metrics"
)

type ArchiveTraceImporter

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

ArchiveTraceImporter imports a bootstrap and its bounded remainder against one fixed source cut. All exporter and barrier calls are serialized so a successful ImportAndWait is an acknowledgment that the frontend has applied that batch, not merely accepted it into an asynchronous dispatch queue.

func NewArchiveTraceImporter

func NewArchiveTraceImporter(sinks TraceImportSinks, cut ArchiveCut) (*ArchiveTraceImporter, error)

func (*ArchiveTraceImporter) AbandonRemainder

func (imp *ArchiveTraceImporter) AbandonRemainder(ctx context.Context, signal ArchiveSignal) error

AbandonRemainder permanently stops one signal after its retry policy is exhausted. Abandoning spans seals the unfinished set at the manifest time and irrevocably rejects later span batches. Other signals have no bearing on the span seal.

func (*ArchiveTraceImporter) CompleteRemainder

func (imp *ArchiveTraceImporter) CompleteRemainder(ctx context.Context, signal ArchiveSignal, cursor int64) error

CompleteRemainder records a bounded signal's terminal cursor. Spans seal at the manifest timestamp only when their remainder reaches its exact high-water mark. Log and metric completion cannot delay span sealing.

func (*ArchiveTraceImporter) ImportAndWait

func (imp *ArchiveTraceImporter) ImportAndWait(ctx context.Context, batch ArchiveImportBatch) error

ImportAndWait imports one batch and waits until the frontend event loop has applied it. Success is the caller's acknowledgment that its archive cursor may advance. Completing bootstrap requires no special call and never seals; the same importer remains open for remainder batches.

func (*ArchiveTraceImporter) Wait

func (imp *ArchiveTraceImporter) Wait(ctx context.Context) error

Wait is an event-loop barrier. It is useful for a terminal bootstrap frame that carries no records. It deliberately does not seal unfinished spans.

type BlockingLogProcessor

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

BlockingLogProcessor batches log records to an exporter without ever dropping one. Unlike the SDK BatchProcessor, which overwrites its oldest record when its queue is full, OnEmit blocks while the queue is full, and a batch whose export fails is retried with backoff until it lands or the processor is shut down. It is for producers that must learn about loss themselves: they emit from a goroutine that can wait, behind their own bounded, counted queue.

func NewBlockingLogProcessor

func NewBlockingLogProcessor(exporter sdklog.Exporter, queueSize, batchSize int, interval time.Duration) *BlockingLogProcessor

NewBlockingLogProcessor exports to exporter in batches of at most batchSize records, at most interval after the first record of a batch arrives, with at most queueSize records waiting.

func (*BlockingLogProcessor) Enabled

func (*BlockingLogProcessor) ForceFlush

func (p *BlockingLogProcessor) ForceFlush(ctx context.Context) error

ForceFlush waits until every record queued before the call is exported.

func (*BlockingLogProcessor) OnEmit

func (p *BlockingLogProcessor) OnEmit(ctx context.Context, record *sdklog.Record) error

OnEmit queues a copy of the record, waiting while the queue is full. It returns ctx's error if ctx ends first, and drops nothing otherwise; after Shutdown it discards the record.

func (*BlockingLogProcessor) Shutdown

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

Shutdown stops accepting records, exports the queued ones within ctx, and shuts the exporter down. Records still queued when ctx ends are lost.

type CallPayloadBatchProcessor

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

CallPayloadBatchProcessor gives immutable call payloads a short, on-demand path to the session exporter while the ordinary log processor retains its lower-frequency batching. OnEmit only clones and enqueues matching records; one worker wakes on the first payload, briefly coalesces its recipe closure, then drains it in bounded exporter batches. The ingress queue is deliberately lossless and therefore unbounded: dropping any one record can leave a client with a call whose nested recipe dependency can never be rebuilt, while the root's per-target claim suppresses every later closure walk that could repair it. Recipe bursts are finite and deduplicated per delivery target; bounding exporter batches, rather than ingress, keeps each persistence operation bounded.

This is the ONLY log transport for payload records (WithoutCallPayloads keeps them out of the ordinary processor), so it owns their retries too: a batch whose export fails goes back to the head of the queue, in order, and is retried with exponential backoff until it lands or CallPayloadMaxExportAttempts is spent. Shutdown keeps retrying within its context; when that ends, it cancels the export in flight and returns only once the worker has stopped, so the caller may shut the exporter down.

func NewCallPayloadBatchProcessor

func NewCallPayloadBatchProcessor(exporter sdklog.Exporter) *CallPayloadBatchProcessor

func NewControlBatchProcessor

func NewControlBatchProcessor(exporter sdklog.Exporter) *CallPayloadBatchProcessor

NewControlBatchProcessor protects revisioned agent and subscription records from the ordinary bounded queue. Capture success is not a delivery receipt: ForceFlush returns the persistence failures of its own pass, and Shutdown returns every terminal persistence failure of the processor's lifetime.

func (*CallPayloadBatchProcessor) Enabled

func (CallPayloadBatchProcessor) ForceFlush

func (batcher CallPayloadBatchProcessor) ForceFlush(ctx context.Context) error

func (*CallPayloadBatchProcessor) OnEmit

func (processor *CallPayloadBatchProcessor) OnEmit(_ context.Context, record *sdklog.Record) error

func (CallPayloadBatchProcessor) Shutdown

func (batcher CallPayloadBatchProcessor) Shutdown(ctx context.Context) error

type CallSpanProcessor

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

CallSpanProcessor is the protected lane for call spans: spans carrying their own encoded call frame as dagger.io/dag.call (core/telemetry.go). That span copy is the only carrier of a spanned call's frame — the producer claims the digest and skips its payload log — so losing the span loses a frame a client, an archive seal or a Cloud restore may need to rebuild a recipe.

It exports both the live start snapshot (as otel.LiveSpanProcessor does, so clients render a call's arguments while it runs) and the end snapshot, through a CoalescingSpanExporter. Delivery mirrors CallPayloadBatchProcessor: an unbounded queue that never blocks or drops on overflow, on-demand coalesced export, in-order retry with backoff up to CallPayloadMaxExportAttempts, ForceFlush reporting only its own pass's drops and Shutdown every terminal loss of the processor's lifetime. Shutdown does not shut the exporter down.

Pair it with WithoutCallSpans around the ordinary live processor, so no call span is exported twice.

func NewCallSpanProcessor

func NewCallSpanProcessor(exporter sdktrace.SpanExporter) *CallSpanProcessor

func (CallSpanProcessor) ForceFlush

func (batcher CallSpanProcessor) ForceFlush(ctx context.Context) error

func (*CallSpanProcessor) OnEnd

func (processor *CallSpanProcessor) OnEnd(span sdktrace.ReadOnlySpan)

func (*CallSpanProcessor) OnStart

func (processor *CallSpanProcessor) OnStart(_ context.Context, span sdktrace.ReadWriteSpan)

func (CallSpanProcessor) Shutdown

func (batcher CallSpanProcessor) Shutdown(ctx context.Context) error

type CloudEmitStatus

type CloudEmitStatus struct {
	Emitting bool
	// Credential is the credential in use: "DAGGER_CLOUD_TOKEN" or
	// "dagger login". Empty when there is none.
	Credential string
	// Org is the org the telemetry goes to, when the credential names one.
	Org string
	// Err is set when the credential cannot be read, or the Cloud URL is not
	// valid (ErrInvalidCloudURL).
	Err error
}

CloudEmitStatus says whether this process sends telemetry to Dagger Cloud, and why. It applies the same rules as ConfiguredCloudExporters.

func CloudEmitStatusFor

func CloudEmitStatusFor(ctx context.Context) CloudEmitStatus

CloudEmitStatusFor returns this process's CloudEmitStatus. It does not ask Dagger Cloud if the credential is valid, but it can refresh an expired login token.

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 LiveFrameKind

type LiveFrameKind uint8

LiveFrameKind distinguishes the frames a live telemetry stream carries.

const (
	// LiveFrameData carries one cursor-addressed protobuf payload.
	LiveFrameData LiveFrameKind = iota
	// LiveFrameHello opens a subscription. Its cursor is the position the
	// server resumed from; it carries no payload and does not advance the
	// consumer's cursor.
	LiveFrameHello
	// LiveFrameTerminal marks a cleanly drained subscription. Its cursor is
	// the last delivered cursor and it carries no payload.
	LiveFrameTerminal
)

func ReadLiveFrame

func ReadLiveFrame(r io.Reader) (kind LiveFrameKind, cursor int64, payload []byte, err error)

ReadLiveFrame reads one frame. Data frames carry a cursor-addressed protobuf payload; hello and terminal frames carry only a cursor. A producer-reported error frame is returned as an error wrapping ErrLiveStream.

func (LiveFrameKind) String

func (k LiveFrameKind) String() string

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 SharedMetricExporter

type SharedMetricExporter struct {
	sdkmetric.Exporter
}

SharedMetricExporter wraps a metric exporter owned elsewhere, for example by a session while each client owns a periodic reader over it: a reader's shutdown shuts its exporter down, so Shutdown is a no-op here and the owner shuts the underlying exporter down itself.

func (SharedMetricExporter) Shutdown

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 TraceImportBarrier

type TraceImportBarrier interface {
	WaitForEventLoop(context.Context) error
}

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. TraceImportBarrier acknowledges application, not exporter enqueue completion.

type TraceImportSinks

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

type TraceImporter

type TraceImporter struct {

	// KeepRoots leaves the imported trace's parentless spans as real roots
	// instead of stamping them passthrough (§5.1.1). Set it when the imported
	// trace IS the session: `dagger trace` renders nothing else, and its root
	// becomes the primary span the whole view hangs from -- a passthrough
	// zoomed span renders only its revealed spans (dagui.DB.RowsView), not
	// its children. A resume leaves it unset: there the live root is primary
	// and the imported one has to get out of the way.
	KeepRoots bool
	// 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.

func (*TraceImporter) SealAt

func (imp *TraceImporter) SealAt(ctx context.Context, at time.Time) error

SealAt ends historical spans at the verified archive close time.

Directories

Path Synopsis
Package cgroupmetrics reports resource use for the cgroup that contains the Dagger engine process.
Package cgroupmetrics reports resource use for the cgroup that contains the Dagger engine process.

Jump to

Keyboard shortcuts

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