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
- Variables
- type Batch
- type Committer
- type Cursor
- type DescribeFunc
- type Detacher
- type Final
- type FinalCommitter
- type Manager
- type OpenFunc
- type Options
- type Renewer
- type Run
- func (r *Run) Detach(ctx context.Context) (query.ManagedFinish, error)
- func (r *Run) FinalWindow() Window
- func (r *Run) Finished() <-chan struct{}
- func (r *Run) ProbeStatus() Status
- func (r *Run) Sample(ctx context.Context) (any, error)
- func (r *Run) Status() query.ManagedStatus
- func (r *Run) Stop(ctx context.Context) (query.ManagedFinish, error)
- type Source
- type Status
- type Window
- type Writer
Constants ¶
const DefaultMaxPages = 64
Variables ¶
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 ¶
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 DescribeFunc ¶
type Detacher ¶
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 ¶
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
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 Run ¶
type Run struct {
// contains filtered or unexported fields
}
Run implements query.ManagedRun over one Source.
func (*Run) FinalWindow ¶
FinalWindow is what Stop drained and finalized. It is empty before Stop.
func (*Run) ProbeStatus ¶
ProbeStatus returns the source checkpoint without the session projection.
func (*Run) Status ¶
func (r *Run) Status() query.ManagedStatus
Status returns the state query.SessionRegistry mirrors into its record.
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.