probe

package
v0.1.39 Latest Latest
Warning

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

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

Documentation

Overview

Package probe coordinates cursor-based external sources with durable record streams. Source adapters own transport and decoding; this package owns exclusive writers, cursor fences, append ordering and finalization.

Index

Constants

View Source
const DefaultMaxPages = 64

Variables

View Source
var ErrAlreadyManaged = errors.New("probe is already managed")

Functions

This section is empty.

Types

type Batch

type Batch struct {
	Generation string
	Next       int64
	Write      int64
	More       bool
	Active     bool
	Rows       []recordstore.Row
	Summary    any
}

Batch is one source page. Next is the cursor after the page, Write is the source's current exclusive write cursor, and More requests another page in the same sample.

type Committer

type Committer interface {
	Commit(context.Context, Batch) error
}

Committer advances source-local decoding state after a batch is durable. A retried Commit must be safe after the record store already accepted the batch, because a process can fail between those two operations.

type Cursor

type Cursor struct {
	Generation string `json:"generation,omitempty"`
	Next       int64  `json:"next"`
}

Cursor is the next position to request from one source generation.

type DescribeFunc

type DescribeFunc func(ctx context.Context, stream string) (*query.EventsRef, error)

type Detacher

type Detacher interface {
	Detach(context.Context) error
}

Detacher releases local resources without freezing, finalizing, sealing or releasing the external source.

type Final

type Final struct {
	Rows    []recordstore.Row
	Summary any
	Warning string
	Result  any
}

Final holds rows materialized only after the source is frozen, such as calls that were still open when a trace ended.

type FinalCommitter

type FinalCommitter interface {
	CommitFinal(context.Context, Final) error
}

FinalCommitter publishes source-local finalization state after its rows are durable and before the stream is sealed.

type Manager

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

Manager admits at most one local writer for an external source identity.

func NewManager

func NewManager() *Manager

func (*Manager) Arm

func (m *Manager) Arm(ctx context.Context, options Options, open OpenFunc) (_ *Run, resultErr error)

Arm reserves the identity before opening the source, creates its empty stream, and returns only after its first durable reference is available.

func (*Manager) Get

func (m *Manager) Get(identity string) (*Run, bool)

Get returns the locally managed run for identity. A source still arming is reported as absent until Arm returns it.

type OpenFunc

type OpenFunc func(context.Context) (Source, error)

type Options

type Options struct {
	Identity string
	Stream   string
	Kind     string
	Backend  Writer
	Describe DescribeFunc
	Cursor   Cursor
	MaxPages int

	// RenewEvery enables background lease renewal. The opened Source must
	// implement Renewer when it is positive.
	RenewEvery time.Duration
}

Options configure one exclusively managed probe generation.

type Renewer

type Renewer interface {
	Renew(context.Context) error
}

Renewer extends a source ownership lease.

type Run

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

Run implements query.ManagedRun over one Source.

func (*Run) Detach

func (r *Run) Detach(ctx context.Context) (query.ManagedFinish, error)

Detach ends local ownership without touching the source or its stream.

func (*Run) FinalWindow

func (r *Run) FinalWindow() Window

FinalWindow is what Stop drained and finalized. It is empty before Stop.

func (*Run) Finished

func (r *Run) Finished() <-chan struct{}

func (*Run) ProbeStatus

func (r *Run) ProbeStatus() Status

ProbeStatus returns the source checkpoint without the session projection.

func (*Run) Sample

func (r *Run) Sample(ctx context.Context) (any, error)

Sample reads pages until the source cursor catches its write cursor.

func (*Run) Status

func (r *Run) Status() query.ManagedStatus

Status returns the state query.SessionRegistry mirrors into its record.

func (*Run) Stop

func (r *Run) Stop(ctx context.Context) (query.ManagedFinish, error)

Stop freezes the source, drains it, materializes final rows, seals the stream, and only then releases the source claim.

type Source

type Source interface {
	Read(context.Context, Cursor) (Batch, error)
	Freeze(context.Context) error
	Finalize(context.Context) (Final, error)
	Release(context.Context) error
}

Source reads and closes one external cursor source. Freeze must stop new source writes; Release gives up its remote ownership claim.

type Status

type Status struct {
	Cursor       Cursor `json:"cursor"`
	SourceActive bool   `json:"sourceActive"`
	Appended     int64  `json:"appended"`
	Skipped      int64  `json:"skipped"`
	Source       any    `json:"source,omitempty"`
}

Status is the durable restart checkpoint a managed session records in its summary. A successor may pass Cursor back through Options.

type Window

type Window struct {
	Cursor  Cursor             `json:"cursor"`
	Stored  recordstore.Window `json:"stored"`
	Rows    int64              `json:"rows"`
	Skipped int64              `json:"skipped"`
	Active  bool               `json:"active"`
	Source  any                `json:"source,omitempty"`
	Events  *query.EventsRef   `json:"events,omitempty"`
}

Window describes what one Sample appended and the cursor it committed.

type Writer

type Writer interface {
	Append(ctx context.Context, stream, kind string, rows []recordstore.Row) (recordstore.AppendResult, error)
	Seal(ctx context.Context, stream string) error
}

Writer is the record-store surface a probe needs. recordstore.Backend and cmd/query record result stores both satisfy it.

Jump to

Keyboard shortcuts

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