clientdb

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: 40 Imported by: 0

Documentation

Index

Constants

View Source
const (
	DumpAll     = "all"
	DumpSpans   = "spans"
	DumpLogs    = "logs"
	DumpMetrics = "metrics"
)
View Source
const CollectGarbageAfter = time.Hour

CollectGarbageAfter is the time after which a store is considered garbage.

Variables

This section is empty.

Functions

func CallPayloadBody

func CallPayloadBody(row Log) (*callpbv1.Call, error)

CallPayloadBody decodes the call frame a call-payload log record carries.

func DecodeLogRecord

func DecodeLogRecord(row Log) (sdklog.Record, error)

DecodeLogRecord is strict for restore-critical data; the display conversion intentionally skips malformed rows and must not be used for verification.

func Dump

func Dump(ctx context.Context, root, clientID, selection string, out io.Writer) error

Dump writes the selected persisted stream rows as JSON lines. It opens files read-only and is intended for a closed store, whose tails have been flushed.

func LogsToPB added in v0.13.1

func LogsToPB(dbLog []Log) []*otlplogsv1.ResourceLogs

func MarshalProtoJSONs added in v0.13.1

func MarshalProtoJSONs[T proto.Message](protos []T) ([]byte, error)

func MetricsToPB added in v0.13.6

func MetricsToPB(dbMetrics []Metric) []*otlpmetricsv1.ResourceMetrics

func UnmarshalProtoJSONs added in v0.13.1

func UnmarshalProtoJSONs[T proto.Message](pb []byte, base T, out *[]T) error

Types

type AppendStats

type AppendStats struct {
	LastID          int64
	CapWaitDuration time.Duration
	SpillLagRows    int
	SpillLagBytes   int64
}

AppendStats describes the append-side work and the stream backlog visible when an append returns. CapWaitDuration is non-zero only when the hard memory cap applied backpressure to the caller.

type CallFrame

type CallFrame struct {
	Digest string
	Span   *Span
	Log    *Log
}

CallFrame is where the store holds one call's frame: the newest snapshot of the span that carries it as dagger.io/dag.call, and/or the call-payload log record that carries it. A frame delivered by both transports has both; the log record is the cheaper one to decode (the span's attributes hold everything else about the span too).

type DB added in v0.19.7

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

DB is one client's standalone append-only telemetry store.

func (*DB) AncestorClosure

func (s *DB) AncestorClosure(ids map[string]struct{}) map[string]struct{}

AncestorClosure returns ids plus every member's ancestor chain up to its trace root. A scoped span load includes it so that no loaded span's parent pointer resolves to an unreceived placeholder — dagui would otherwise mistake a placeholder for a root and lose the chain's UI flags.

func (*DB) AppendLogs

func (s *DB) AppendLogs(rows []Log) (AppendStats, error)

func (*DB) AppendMetrics

func (s *DB) AppendMetrics(rows []Metric) (AppendStats, error)

func (*DB) AppendSpans

func (s *DB) AppendSpans(rows []Span) (AppendStats, error)

func (*DB) CallDigests

func (s *DB) CallDigests() map[string]struct{}

CallDigests returns every call digest the store holds a frame for, on either transport -- the seed for a session-wide call search, sized by the number of distinct calls. The returned map is a fresh snapshot the caller owns.

func (*DB) CallPayload

func (s *DB) CallPayload(ctx context.Context, traceID, digest string, cut HighWater) (Log, error)

CallPayload returns a call-payload log row carrying digest's frame, as long as it lies within the cut and belongs to traceID.

A frame reaches the store on one of two carriers (see callLookup): a call-payload log record, or — for a call with a recording span of its own — the span's dagger.io/dag.call attribute, which is then the frame's only delivery. The log row is preferred. A frame that only rode its span is returned as a payload row synthesized from the span row, in exactly the shape the producer emits (core/dag_call_telemetry.go), so bootstrap packing and every reader stay oblivious to the carrier. A synthesized row has no log row ID (0).

The span side is indexed by its newest snapshot. Every snapshot of a call span carries the same frame, and a sealed cut is taken after the session's telemetry shut down, so the newest snapshot lies within it.

func (*DB) CausalChildren

func (s *DB) CausalChildren(spanID string) []string

CausalChildren returns the span IDs that cause-link to the given span — e.g. a service's long-lived exec span cause-links to the API spans that installed the Service value (Container.asService and friends).

func (*DB) CheckTestSpanIDs

func (s *DB) CheckTestSpanIDs() (checks, tests map[string]struct{})

CheckTestSpanIDs returns the spans marked as named checks and as test cases or suites, answered from the span index. This is the seed for resolving a check/test NAME to a span: load these spans (plus ancestors, for the test view's containment walks) into a throwaway dagui.DB and match names there, instead of retaining a whole-session DB just to answer name lookups.

func (*DB) Checkpoint

func (s *DB) Checkpoint(ctx context.Context) (HighWater, error)

Checkpoint is a persistence barrier, not merely an in-memory high water. Producers must be quiesced to use the returned cut for finalization.

func (*DB) ChildSpanIDs

func (s *DB) ChildSpanIDs(spanID string) map[string]struct{}

ChildSpanIDs returns the spans one edge beneath spanID, over the same edges as the log queries: parent→child plus cause-purpose links. Answered from the index alone, so a caller can load a span's immediate surroundings without materializing its whole subtree.

func (*DB) Close added in v0.19.7

func (s *DB) Close() error

func (*DB) ControlRows

func (s *DB) ControlRows(ctx context.Context, traceID string, through int64, want agentcontrol.Expectation) ([]Log, error)

ControlRows selects only final received projections. Verification compares them with a separately supplied, quiesced producer witness.

func (*DB) FlushLogs

func (s *DB) FlushLogs(ctx context.Context) error

FlushLogs writes the in-memory log tail to its file without fsyncing, so restore-critical rows survive the engine process being killed. Session end still Checkpoints for real durability.

func (*DB) FlushSpans

func (s *DB) FlushSpans(ctx context.Context) error

FlushSpans is FlushLogs for the span tail: a spanned call's frame rides its call span alone, so an unsealed restore needs those rows on file too.

func (*DB) HasDescendants

func (s *DB) HasDescendants(spanID string) bool

HasDescendants reports whether any span is nested beneath spanID, following the same edges as the log queries: parent→child plus cause-purpose links.

It answers purely from the in-memory span index -- no stream reads, no subtree materialization -- so it is cheap enough to use as a pre-filter (e.g. "did this tool call produce any child telemetry worth rendering?").

func (*DB) HasSpan

func (s *DB) HasSpan(spanID string) bool

HasSpan reports whether the store has seen any snapshot of spanID.

func (*DB) HighWater

func (s *DB) HighWater() HighWater

HighWater reports the current end of each stream without a persistence barrier. Use it only for a store no producer writes to anymore, such as an archive whose session ended without sealing.

func (*DB) ImportTrace

func (s *DB) ImportTrace(ctx context.Context, id string, load func(*DB) error) (*DB, error)

ImportTrace publishes a complete snapshot exactly once. Failed/canceled downloads remain invisible and may be retried. The caller must retain s while using the returned store, and must not close the returned store.

func (*DB) InspectionStores

func (s *DB) InspectionStores() []*DB

InspectionStores returns the live store followed by imported snapshots in trace ID order. Retain s while using them; do not close the borrowed stores.

func (*DB) Read

func (s *DB) Read() *DB

Read mirrors the current DB handle seam. The append-only store needs no separate read pool, so selectors are bound to the same immutable streams.

func (*DB) SelectCallFrames

func (s *DB) SelectCallFrames(ctx context.Context, digests map[string]struct{}) ([]CallFrame, error)

SelectCallFrames returns the frames the store holds for digests, in ascending digest order, skipping digests it holds no frame for. Answered through the call index, so a recipe rebuild that follows a chain's references loads a row per frame rather than scanning either stream.

func (*DB) SelectLogsBeneathSpan

func (s *DB) SelectLogsBeneathSpan(ctx context.Context, arg SelectLogsBeneathSpanParams) ([]Log, error)

SelectLogsBeneathSpan returns the log rows of the capture rooted at arg.SpanID — the rows attributed to any span in its log scope (the span itself, its cause-link targets, and both sets' subtrees; see SpanLogScope) — in append order, starting after cursor arg.ID, up to arg.Limit rows. The root's own rows are part of the capture: a span's directly-attributed output (e.g. the service stdio records routed to an install span) is exactly what a reader asking about that span wants.

Resolved through the per-span log index: the row IDs come straight from the scope's index entries, so the cost scales with the capture, not with the session — the previous implementation scanned the whole log stream past the cursor, linear in session size, on a path that runs per LLM tool call.

func (*DB) SelectLogsForSpans

func (s *DB) SelectLogsForSpans(ctx context.Context, ids map[string]struct{}, perSpanTail int) ([]Log, error)

SelectLogsForSpans returns the log rows attributed to any span in ids, in append order — the log half of a scoped load. perSpanTail > 0 bounds each span to its newest perSpanTail rows; renderers bound every span's log output to a tail anyway, and the cap keeps one pathological span (e.g. a service that streamed millions of lines) from ballooning the load.

func (*DB) SelectLogsRange

func (s *DB) SelectLogsRange(ctx context.Context, p SelectLogsRangeParams) ([]Log, error)

func (*DB) SelectLogsSince

func (s *DB) SelectLogsSince(ctx context.Context, arg SelectLogsSinceParams) ([]Log, error)

func (*DB) SelectMetricsRange

func (s *DB) SelectMetricsRange(ctx context.Context, p SelectMetricsRangeParams) ([]Metric, error)

func (*DB) SelectMetricsSince

func (s *DB) SelectMetricsSince(ctx context.Context, arg SelectMetricsSinceParams) ([]Metric, error)

func (*DB) SelectSpan

func (s *DB) SelectSpan(ctx context.Context, arg SelectSpanParams) (Span, error)

func (*DB) SelectSpansLatest

func (s *DB) SelectSpansLatest(ctx context.Context, ids map[string]struct{}) ([]Span, error)

SelectSpansLatest returns the newest snapshot row of every span in ids, in append order. A span's snapshots are cumulative, so its newest row alone reconstructs the state that full sequential processing would end with — this is the span half of a scoped load, sized by the scope instead of the session.

func (*DB) SelectSpansRange

func (s *DB) SelectSpansRange(ctx context.Context, p SelectSpansRangeParams) ([]Span, error)

func (*DB) SelectSpansSince

func (s *DB) SelectSpansSince(ctx context.Context, arg SelectSpansSinceParams) ([]Span, error)

func (*DB) SizeBytes

func (s *DB) SizeBytes() (int64, error)

func (*DB) SpanIDs

func (s *DB) SpanIDs() map[string]struct{}

SpanIDs returns every span ID the store has seen a snapshot of -- the seed for a session-wide scoped load (e.g. a name search), sized by the span count rather than by the snapshot stream. The returned map is a fresh snapshot the caller owns.

func (*DB) SpanLogScope

func (s *DB) SpanLogScope(spanID string) map[string]struct{}

SpanLogScope returns the set of span IDs whose log rows belong to a capture rooted at spanID — the same walk SelectLogsBeneathSpan scopes its logs to: the span itself, the spans it cause-links to (e.g. the install spans a service exec span links, which carry the service's stdio records), and everything beneath either over child edges and cause-purpose link edges. The returned map is a fresh snapshot the caller owns; the set only ever grows as spans arrive.

type DBs

type DBs struct {
	Root string
	// contains filtered or unexported fields
}

DBs owns the refcounted set of open per-client telemetry stores. Live client runtimes and readers retain references; the final Close releases the streams.

func NewDBs

func NewDBs(root string) *DBs

func (*DBs) Close

func (r *DBs) Close() error

Close prevents new stores from opening and waits for in-flight opens. Actively referenced stores remain usable and close their streams when their final handle is released.

func (*DBs) GC

func (r *DBs) GC(keep map[string]bool) error

GC removes complete client stores whose newest stream (or transitional SQLite sidecar) is older than CollectGarbageAfter. Grouping files by client keeps a recently active stream from being separated from an older sibling.

func (*DBs) Open

func (r *DBs) Open(ctx context.Context, clientID string) (*DB, error)

func (*DBs) OpenStats

func (r *DBs) OpenStats() OpenStats

func (*DBs) Remove

func (r *DBs) Remove(clientID string) (bool, error)

Remove deletes a store only when no live writer or reader holds a reference. The same per-store lock serializes reopening, so GC cannot unlink a newly reopened archive stream.

func (*DBs) StoreSize

func (r *DBs) StoreSize(clientID string) (int64, error)

StoreSize is the on-disk size of a client's store streams, measured without opening (and replaying) the store. Missing streams count as empty.

type HighWater

type HighWater struct{ Spans, Logs, Metrics int64 }

HighWater is a fixed inclusive cursor; zero denotes an empty stream.

type Log

type Log struct {
	ID                   int64
	TraceID              sql.NullString
	SpanID               sql.NullString
	Timestamp            int64
	SeverityNumber       int64
	SeverityText         string
	Body                 []byte
	Attributes           []byte
	InstrumentationScope []byte
	Resource             []byte
	ResourceSchemaURL    string
}

func SpanCallPayloadLog

func SpanCallPayloadLog(row Span, digest string) (Log, error)

SpanCallPayloadLog synthesizes the call-payload log row for the frame a call span row carries as dagger.io/dag.call: a bytes body holding the encoded call, tagged with the call payload content type and digest, attributed to the span. The frame must carry digest.

type Metric added in v0.13.6

type Metric struct {
	ID   int64
	Data []byte
}

type OpenStats

type OpenStats struct {
	Stores  int
	Streams int
	Refs    int
}

OpenStats is a measured snapshot of currently open telemetry stores. Each referenced store owns exactly three stream handles.

type SelectLogsBeneathSpanParams added in v0.19.0

type SelectLogsBeneathSpanParams struct {
	SpanID sql.NullString
	ID     int64
	Limit  int64
}

type SelectLogsRangeParams

type SelectLogsRangeParams struct{ AfterID, ThroughID, Limit int64 }

type SelectLogsSinceParams

type SelectLogsSinceParams struct {
	ID    int64
	Limit int64
}

type SelectMetricsRangeParams

type SelectMetricsRangeParams struct{ AfterID, ThroughID, Limit int64 }

type SelectMetricsSinceParams added in v0.13.6

type SelectMetricsSinceParams struct {
	ID    int64
	Limit int64
}

type SelectSpanParams added in v0.19.0

type SelectSpanParams struct {
	TraceID string
	SpanID  string
}

type SelectSpansRangeParams

type SelectSpansRangeParams struct{ AfterID, ThroughID, Limit int64 }

type SelectSpansSinceParams

type SelectSpansSinceParams struct {
	ID    int64
	Limit int64
}

type Span

type Span struct {
	ID                     int64
	TraceID                string
	SpanID                 string
	TraceState             string
	ParentSpanID           sql.NullString
	Flags                  int64
	Name                   string
	Kind                   string
	StartTime              int64
	EndTime                sql.NullInt64
	Attributes             []byte
	DroppedAttributesCount int64
	Events                 []byte
	DroppedEventsCount     int64
	Links                  []byte
	DroppedLinksCount      int64
	StatusCode             int64
	StatusMessage          string
	InstrumentationScope   []byte
	Resource               []byte
	ResourceSchemaURL      string
}

func (*Span) ReadOnly

func (span *Span) ReadOnly() sdktrace.ReadOnlySpan

Jump to

Keyboard shortcuts

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