archive

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

Documentation

Index

Constants

View Source
const (
	DefaultTTL   = 7 * 24 * time.Hour
	DefaultQuota = int64(10 << 30)
)
View Source
const AgentBootstrapResource = "agent-bootstrap"

AgentBootstrapResource names the agent-restore bootstrap within an archive. It is agent-specific: final agent/subscription control records plus the recipe closure of their snapshots. The OTLP signal streams alongside it are generic telemetry resources.

View Source
const BootstrapContentType = "application/vnd.dagger.telemetry.bootstrap"
View Source
const ManifestVersion = 4
View Source
const MaxBootstrapPayloadSize = 64 << 20

MaxBootstrapPayloadSize bounds the encoded payload of each bootstrap frame.

View Source
const MaxTitleRunes = 120

MaxTitleRunes bounds an archive title. Titles come from trace records any session producer can emit, and listings print them to a terminal.

Variables

View Source
var (
	ErrCleanMiss = errors.New("engine archive clean miss")
	ErrState     = errors.New("engine archive state failure")
	ErrCorrupt   = errors.New("engine archive corruption")
	ErrTransient = errors.New("engine archive transient failure")
)
View Source
var ErrBootstrapIncomplete = errors.New("bootstrap stream ended before terminal frame")

ErrBootstrapIncomplete identifies a finite bootstrap response that ended before its terminal frame.

Functions

func BuildBootstrap

func BuildBootstrap(header BootstrapHeader, signals []BootstrapSignal) ([]byte, int64, error)

func DecodeBootstrap

func DecodeBootstrap(r io.Reader, onHeader func(BootstrapHeader) error, consume func([]byte) error) (BootstrapHeader, BootstrapTerminal, error)

DecodeBootstrap requires a header first and a terminal last, rejects trailing bytes, verifies the terminal checksum, calls onHeader before reading any logs frame, and calls consume with each logs frame's payload. EOF before the terminal is an interruption, never a successful finite response.

func IsCleanMiss

func IsCleanMiss(err error) bool

IsCleanMiss reports whether the connected engine definitively has no usable archive and a caller may try another source.

func SanitizeTitle

func SanitizeTitle(title string) string

SanitizeTitle reduces a trace-derived title to one bounded, printable line: control and format characters (including bidi overrides) are dropped or turned into spaces, whitespace runs collapse, and long titles are cut with an ellipsis.

func VerifyClosure

func VerifyClosure(roots []string, load func(string) (*callpbv1.Call, error)) (map[string]*callpbv1.Call, error)

VerifyClosure loads the recipe closure of the anchors without evaluating any recipe, failing on a missing or cyclic dependency. The returned payloads are exactly the dependency closure of the anchors.

func WriteBootstrapFrame

func WriteBootstrapFrame(w io.Writer, kind BootstrapFrameKind, payload []byte) error

Types

type AgentRevision

type AgentRevision struct {
	Key      agentcontrol.Key
	Revision int64
}

type Bootstrap

type Bootstrap struct {
	File    string `json:"file"`
	Records int64  `json:"records"`
}

type BootstrapBatch

type BootstrapBatch struct {
	Logs *collogspb.ExportLogsServiceRequest
}

BootstrapBatch is one decoded OTLP logs batch from a bootstrap response.

type BootstrapFrameKind

type BootstrapFrameKind byte
const (
	BootstrapFrameHeader   BootstrapFrameKind = 1
	BootstrapFrameLogs     BootstrapFrameKind = 3
	BootstrapFrameTerminal BootstrapFrameKind = 4
)

func ReadBootstrapFrame

func ReadBootstrapFrame(r io.Reader) (BootstrapFrameKind, []byte, error)

type BootstrapHeader

type BootstrapHeader struct {
	Version    int        `json:"version"`
	TraceID    string     `json:"traceID"`
	SealAt     string     `json:"sealAt"`
	HighWater  HighWater  `json:"highWater"`
	Completion Completion `json:"completion"`
}

type BootstrapResult

type BootstrapResult struct {
	Header   BootstrapHeader
	Terminal BootstrapTerminal
}

BootstrapResult is the immutable cut described by a bootstrap.

type BootstrapSignal

type BootstrapSignal struct {
	Payload []byte
	Records int64
}

BootstrapSignal is one encoded OTLP logs batch of Records log records.

type BootstrapTerminal

type BootstrapTerminal struct {
	LogRecords int64  `json:"logRecords"`
	SHA256     string `json:"sha256"`
}

type Client

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

Client reads telemetry archives through a connected engine's HTTP transport.

func NewClient

func NewClient(httpClient HTTPDoer) *Client

NewClient creates an archive client for the connected engine transport.

func NewClientWithURL

func NewClientWithURL(httpClient HTTPDoer, baseURL string) (*Client, error)

NewClientWithURL creates an archive client with an explicit base URL. It is useful for HTTP proxies and tests; connected engine clients should use NewClient.

func (*Client) Bootstrap

func (c *Client) Bootstrap(ctx context.Context, traceID string, consume func(BootstrapHeader, BootstrapBatch) error) (BootstrapResult, error)

Bootstrap buffers the complete raw bootstrap and checks its header and terminal checksum before invoking consume, so a truncated or corrupt bootstrap never applies part of itself. Consumers may then apply the batches and wait for their frontend barrier, without loading unrelated historical telemetry. The engine verified the roster and recipe closure when it built the bootstrap.

func (*Client) List

func (c *Client) List(ctx context.Context, opts ListOptions) (Page, error)

List returns one page of archives.

func (*Client) ListAll

func (c *Client) ListAll(ctx context.Context, opts ListOptions) ([]Manifest, error)

ListAll follows list cursors until the engine returns the final page.

func (*Client) Logs

func (c *Client) Logs(ctx context.Context, traceID string, opts StreamOptions, consume func(int64, *collogspb.ExportLogsServiceRequest) error) (int64, error)

Logs reads a finite framed log stream. It returns the last safe resume cursor.

func (*Client) Metrics

func (c *Client) Metrics(ctx context.Context, traceID string, opts StreamOptions, consume func(int64, *colmetricspb.ExportMetricsServiceRequest) error) (int64, error)

Metrics reads a finite framed metric stream. It returns the last safe resume cursor.

func (*Client) Traces

func (c *Client) Traces(ctx context.Context, traceID string, opts StreamOptions, consume func(int64, *coltracepb.ExportTraceServiceRequest) error) (int64, error)

Traces reads a finite framed trace stream. It returns the last safe resume cursor.

func (*Client) Unsealed

func (c *Client) Unsealed(ctx context.Context, traceID string) (UnsealedArchive, error)

Unsealed looks up an interrupted or incomplete archive for a best-effort read of what it recorded; Cut is the current end of its store. It fails with an ErrState request error for any other state, including a sealed (closed) archive.

type Completion

type Completion struct {
	Agents        []AgentRevision        `json:"agents"`
	Subscriptions []SubscriptionRevision `json:"subscriptions"`
}

Completion is the final roster the producer witnessed when the archive was sealed: the agents and subscriptions a restore installs from the bootstrap.

func Witness

func Witness(want agentcontrol.Expectation) Completion

type Config

type Config struct {
	Root        string
	TTL         time.Duration
	QuotaBytes  int64
	Now         func() time.Time
	RemoveStore func(string) (bool, error)
	// StoreSize reports the on-disk size of a main client's telemetry store.
	// It sizes archives that end unsealed, so they count toward the quota.
	StoreSize func(clientID string) (int64, error)
}

type ErrorKind

type ErrorKind string

ErrorKind groups archive failures by the recovery decision a caller should make.

const (
	ErrorCleanMiss ErrorKind = "clean_miss"
	ErrorState     ErrorKind = "state"
	ErrorCorrupt   ErrorKind = "corrupt"
	ErrorTransient ErrorKind = "transient"
)

type Failure

type Failure struct {
	Kind  FailureKind `json:"kind"`
	State State       `json:"state,omitempty"`
	Err   error       `json:"-"`
}

func (*Failure) Error

func (e *Failure) Error() string

func (*Failure) Unwrap

func (e *Failure) Unwrap() error

type FailureKind

type FailureKind string
const (
	FailureNotFound FailureKind = "not_found"
	FailureState    FailureKind = "state"
	FailureCorrupt  FailureKind = "corrupt"
	FailureIO       FailureKind = "io"
)

type FinalizeInput

type FinalizeInput struct {
	HighWater        HighWater
	SealAt           time.Time
	StoreSizeBytes   int64
	BootstrapBytes   []byte
	BootstrapRecords int64
}

FinalizeInput is the sealed cut. The caller built and verified the bootstrap for exactly this cut.

type HTTPDoer

type HTTPDoer interface {
	Do(*http.Request) (*http.Response, error)
}

HTTPDoer is the transport required by the archive client. engine/client's DirectConn implements this interface and carries the connected session's authentication and routing metadata.

type HighWater

type HighWater struct {
	Spans   int64 `json:"spans"`
	Logs    int64 `json:"logs"`
	Metrics int64 `json:"metrics"`
}

type ListOptions

type ListOptions struct {
	After          string
	Limit          int
	ExcludeTraceID string
}

ListOptions controls one archive list page. The engine also excludes the archive currently being captured for this connected session. ExcludeTraceID can omit one additional archive client-side.

type Manager

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

func NewManager

func NewManager(cfg Config) (*Manager, error)

func (*Manager) BeginFinalizing

func (m *Manager) BeginFinalizing(traceID string) error

func (*Manager) BootstrapPath

func (m *Manager) BootstrapPath(traceID string) string

BootstrapPath is where a sealed archive's agent bootstrap is stored.

func (*Manager) Finalize

func (m *Manager) Finalize(traceID string, in FinalizeInput) (Manifest, error)

func (*Manager) GC

func (m *Manager) GC() (overage int64, err error)

GC deletes expired archives, then the oldest ended archives (closed or unsealed) until the retained ones fit the quota. The newest ended archive is kept even when it alone exceeds the quota; overage reports by how much. An archive whose store is still open is left in place and retried by the next GC.

func (*Manager) KeepSet

func (m *Manager) KeepSet() map[string]bool

func (*Manager) List

func (m *Manager) List(after, excludeTraceID string, limit int) Page

func (*Manager) Manifest

func (m *Manager) Manifest(traceID string) (Manifest, error)

func (*Manager) MarkIncomplete

func (m *Manager) MarkIncomplete(traceID string, cause error) error

func (*Manager) Register

func (m *Manager) Register(traceID, mainClientID string) (Manifest, error)

Register creates the active archive for a trace. The first session to register a trace owns its archive; a later registration for the same trace (e.g. a nested session that inherited TRACEPARENT) fails.

func (*Manager) SetTitle

func (m *Manager) SetTitle(traceID, title string) error

SetTitle records the session title the main client published into an active archive's trace. It is persisted immediately, so an archive recovered after an engine crash keeps it, and the sealed manifest inherits it. The title is sanitized first; an unchanged or empty title writes nothing.

type Manifest

type Manifest struct {
	Version      int        `json:"version"`
	TraceID      string     `json:"traceID"`
	MainClientID string     `json:"mainClientID"`
	State        State      `json:"state"`
	Title        string     `json:"title,omitempty"`
	StartedAt    time.Time  `json:"startedAt"`
	ClosedAt     *time.Time `json:"closedAt,omitempty"`
	ExpiresAt    time.Time  `json:"expiresAt"`
	SealAt       *time.Time `json:"sealAt,omitempty"`
	SizeBytes    int64      `json:"sizeBytes"`
	Bootstrap    Bootstrap  `json:"bootstrap,omitempty"`
	HighWater    HighWater  `json:"highWater"`
	Failure      string     `json:"failure,omitempty"`
}

Manifest describes one archive. An archive is identified by its trace ID alone: the first session to register a trace owns its archive.

type Page

type Page struct {
	Archives []Manifest `json:"archives"`
	Next     string     `json:"next,omitempty"`
}

type RequestError

type RequestError struct {
	Kind       ErrorKind
	Failure    FailureKind
	State      State
	StatusCode int
	Err        error
}

RequestError is a typed archive transport or protocol failure.

func (*RequestError) Error

func (e *RequestError) Error() string

func (*RequestError) Is

func (e *RequestError) Is(target error) bool

func (*RequestError) Unwrap

func (e *RequestError) Unwrap() error

type State

type State string
const (
	StateActive      State = "active"
	StateFinalizing  State = "finalizing"
	StateClosed      State = "closed"
	StateInterrupted State = "interrupted"
	StateIncomplete  State = "incomplete"
)

func (State) Unsealed

func (s State) Unsealed() bool

Unsealed reports whether an archive's session ended without a verified seal: the engine stopped before finalizing it, or finalization failed.

type StreamOptions

type StreamOptions struct {
	Cursor    int64
	HighWater int64
	Unsealed  bool
}

StreamOptions fixes one finite signal read to a high-water cursor. Cursor is the last batch successfully acknowledged by the caller. Unsealed reads an interrupted or incomplete archive at the cut returned by Unsealed; its log stream then includes agent control records.

type SubscriptionRevision

type SubscriptionRevision struct {
	Key      agentcontrol.EdgeKey
	Revision int64
}

type UnsealedArchive

type UnsealedArchive struct {
	Cut HighWater
}

UnsealedArchive is an archive whose session ended without a seal. It has no bootstrap; stream it with StreamOptions.Unsealed at Cut.

Jump to

Keyboard shortcuts

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