client

package
v0.4.0 Latest Latest
Warning

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

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

Documentation

Overview

Package client implements the transport-agnostic writer session of the History Protocol (protocol/SPEC.md sections 8 and 9): hello, resume and drain, windowed replay, acknowledgements, divergence handling with rebaseline, and the reverse MCP channel.

Store contract

The client drives a durable spool through the Store interface. It mirrors the subset of internal/spool used by the writer session; *spool.Spool is adapted to it by a thin wrapper because Store.Do passes the Tx interface rather than a concrete transaction type.

  • WriterID and Incarnation are fixed for the lifetime of the handle. The incarnation is incremented and made durable before the handle is returned.
  • Identity returns the enrolled target, credential, and machine ID; SetIdentity persists a replacement durably before returning (credential rotation and de-enrollment use it).
  • Epoch returns the current epoch and its chain position (the highest sequence ever assigned is Chain.Head). OpenEpoch starts a new epoch: entries above prevHead must have been discarded first, the remaining entries of the previous epoch are dropped, and LastCommitted resets. MarkRegistered records that the control plane accepted the epoch.
  • Do holds the sequence lock for the duration of fn. Every append happens inside Do through Tx.Append, which assigns the next envelope, encodes the record, and links its chain hash. Appends become durable when fn returns nil; if fn fails, its appends are discarded.
  • Entries(from) returns copies of the spooled entries with Seq >= from in chain order, or a bounded prefix of them holding at least one entry. A range record produced by coalescing appears once, at the end of its span.
  • MarkTransmitted moves entries to TransmittedUnconfirmed durably and fails with ErrNotSpooled, marking nothing, if any sequence is no longer spooled. Transmitted entries never change.
  • Commit is cumulative: it verifies the chain hash at seq, deletes every entry at or below it, and records the committed head. A chain hash that differs from the writer's, a sequence above the highest assigned, or a sequence inside a coalesced span returns ErrDivergence.
  • DiscardAbove deletes every entry above seq; it precedes OpenEpoch during a rebaseline.
  • Notify is signaled after appends and commits.
  • SetHalted records a stop code for this writer process; Do refuses to run while halted.

MemStore is a goroutine-safe in-memory Store with the full semantics above, for tests.

Index

Constants

View Source
const (
	HaltSuperseded = "superseded"
	HaltDeenrolled = "deenrolled"
)

Halt codes in addition to the hello rejection codes of SPEC 8.3.

View Source
const MCPProtocolVersion = "2025-06-18"

MCP message shapes served over the reverse channel (SPEC 9.4, 9.5).

Variables

View Source
var (
	ErrDisconnected = errors.New("connection closed")
	ErrNotConnected = errors.New("no active session")
)

Session errors.

View Source
var (
	ErrDivergence    = errors.New("committed chain diverges from the spool")
	ErrHalted        = errors.New("writer halted")
	ErrNoEpoch       = errors.New("no epoch open")
	ErrNotEnrolled   = errors.New("target identity not enrolled")
	ErrNotSpooled    = errors.New("sequence not spooled")
	ErrWrongEpoch    = errors.New("epoch is not the current epoch")
	ErrEntriesRemain = errors.New("entries above the previous head remain")
)

Store errors.

Functions

func RPCErrorCode

func RPCErrorCode(err error) string

RPCErrorCode returns the protocol code carried by a JSON-RPC error, or "".

Types

type CallToolParams

type CallToolParams struct {
	Name      string          `json:"name"`
	Arguments json.RawMessage `json:"arguments,omitempty"`
}

CallToolParams is the tools/call request.

type CallToolResult

type CallToolResult struct {
	Content           []Content       `json:"content"`
	StructuredContent json.RawMessage `json:"structuredContent,omitempty"`
	IsError           bool            `json:"isError"`
}

CallToolResult is the MCP tools/call result.

func ToolResult

func ToolResult(v any, err error) CallToolResult

ToolResult converts a tool outcome into an MCP CallToolResult.

type CaptureTx

type CaptureTx interface {
	Tx
	// Reason is the checkpoint reason the hook must emit.
	Reason() protocol.CheckpointReason
	// Epoch is the epoch being appended to.
	Epoch() EpochState
	// Committed is the committed head H; a replay summary covers findings touched above it.
	Committed() uint64
}

CaptureTx is the transaction handed to Hooks inside Store.Do.

type Client

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

Client is the writer session client.

func New

func New(opts Options) (*Client, error)

New validates options and applies defaults.

func (*Client) Audit

func (c *Client) Audit(ctx context.Context, params any) error

Audit sends an investigation.audit notification.

func (*Client) Deenroll

func (c *Client) Deenroll(ctx context.Context, reason string) error

Deenroll requests explicit de-enrollment; on success the credential is deleted and Run stops.

func (*Client) FetchBundle

func (c *Client) FetchBundle(ctx context.Context, have string) (*protocol.BundleFetchResult, error)

FetchBundle requests the assigned rule bundle over the active session.

func (*Client) Prepare

func (c *Client) Prepare() error

Prepare opens an initial epoch offline and emits its first checkpoint, so records can be appended before connecting.

func (*Client) Run

func (c *Client) Run(ctx context.Context) error

Run connects, resumes, and streams records until ctx ends or the writer must stop.

func (*Client) Status

func (c *Client) Status() Status

Status returns a snapshot of the session and spool backlog.

type Clock

type Clock interface {
	Now() time.Time
	After(d time.Duration) <-chan time.Time
}

Clock abstracts time for backoff, rate limiting, and health reports.

type Conn

type Conn interface {
	Call(ctx context.Context, method string, params, result any) error
	Notify(ctx context.Context, method string, params any) error
	SendBinary(ctx context.Context, frame []byte) error
	Handle(h Handler)
	// Done is closed after the session ends and no Handler invocation is running.
	Done() <-chan struct{}
	Close() error
}

Conn is one bidirectional JSON-RPC session with binary record frames.

type Content

type Content struct {
	Type string `json:"type"`
	Text string `json:"text"`
}

Content is one MCP content block.

type DivergenceError

type DivergenceError struct {
	Head   uint64
	Reason string
}

DivergenceError triggers a rebaseline at the committed head Head (SPEC 8.6).

func (*DivergenceError) Error

func (e *DivergenceError) Error() string

type Entry

type Entry struct {
	Seq       uint64
	Type      protocol.RecordType
	State     RecordState
	Bytes     []byte
	Hash      protocol.Hash
	ChainHash protocol.Hash
}

Entry is one spooled record.

func AppendCheckpoint

func AppendCheckpoint(tx CaptureTx, st *protocol.State, iv protocol.Interval, capabilities []string) (*Entry, error)

AppendCheckpoint appends a checkpoint of st with the reason of tx and, at sequence 1, the previous epoch and head.

type EpochState

type EpochState struct {
	ID         protocol.EpochID
	OpenReason string
	PrevEpoch  *protocol.EpochID
	PrevHead   *uint64
	Registered bool
	Chain      protocol.Chain
}

EpochState is the writer's current epoch.

type Handler

type Handler func(ctx context.Context, req *Request) (any, error)

Handler serves incoming requests and notifications; a *protocol.RPCError is sent as-is, other errors as internal errors.

type Hooks

type Hooks interface {
	// CaptureReplay appends only the replay anchor (reason 5) and returns the lifecycle summary as of it (SPEC 8.4 step 3).
	CaptureReplay(tx CaptureTx) (protocol.SummaryParams, error)
	// Rebaseline appends the full checkpoint that starts a new epoch.
	Rebaseline(tx CaptureTx) error
	BundleAvailable(protocol.BundleAvailableParams)
	Tools() []Tool
	Tool(ctx context.Context, name string, args json.RawMessage) (any, error)
	Health() any
}

Hooks connects the session to the writer's state, findings, rules, and investigation tools.

type Identity

type Identity struct {
	TargetID     string
	TargetType   string
	Credential   string
	CredentialID string
	MachineID    string
}

Identity is the enrolled identity persisted with the spool.

type InitializeResult

type InitializeResult struct {
	ProtocolVersion string         `json:"protocolVersion"`
	Capabilities    map[string]any `json:"capabilities"`
	ServerInfo      ServerInfo     `json:"serverInfo"`
}

InitializeResult is the MCP initialize result.

type MemOptions

type MemOptions struct {
	WriterID protocol.WriterID // zero generates a random writer ID
	Identity Identity
	Now      func() time.Time
}

MemOptions configures a MemStore.

type MemStore

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

MemStore is a goroutine-safe in-memory Store with full spool semantics, for tests.

func NewMemStore

func NewMemStore(opts MemOptions) (*MemStore, error)

NewMemStore creates a spool with incarnation 1.

func (*MemStore) Coalesce

func (s *MemStore) Coalesce() (int, error)

Coalesce folds the oldest run of two or more never-transmitted non-checkpoint records into a range (SPEC 8.7).

func (*MemStore) Commit

func (s *MemStore) Commit(epoch protocol.EpochID, seq uint64, chainHash protocol.Hash) error

func (*MemStore) DiscardAbove

func (s *MemStore) DiscardAbove(seq uint64) error

func (*MemStore) Do

func (s *MemStore) Do(fn func(tx Tx) error) error

Do runs fn under the sequence lock; if fn fails, its appends are discarded.

func (*MemStore) Entries

func (s *MemStore) Entries(fromSeq uint64) []*Entry

func (*MemStore) Epoch

func (s *MemStore) Epoch() (EpochState, bool)

func (*MemStore) Halted

func (s *MemStore) Halted() (string, bool)

func (*MemStore) Identity

func (s *MemStore) Identity() Identity

func (*MemStore) Incarnation

func (s *MemStore) Incarnation() uint64

func (*MemStore) LastCommitted

func (s *MemStore) LastCommitted() (protocol.ChainPoint, bool)

func (*MemStore) MarkRegistered

func (s *MemStore) MarkRegistered() error

func (*MemStore) MarkTransmitted

func (s *MemStore) MarkTransmitted(seqs ...uint64) error

func (*MemStore) Notify

func (s *MemStore) Notify() <-chan struct{}

func (*MemStore) OpenEpoch

func (s *MemStore) OpenEpoch(reason string, prev *protocol.EpochID, prevHead *uint64) (protocol.EpochID, error)

func (*MemStore) Reopen

func (s *MemStore) Reopen() *MemStore

Reopen copies the durable state into a new handle with the next incarnation, as a restart does; halt is not copied.

func (*MemStore) SetHalted

func (s *MemStore) SetHalted(code string) error

func (*MemStore) SetIdentity

func (s *MemStore) SetIdentity(id Identity) error

func (*MemStore) Transmitted

func (s *MemStore) Transmitted() map[protocol.RecordID]protocol.Hash

Transmitted returns every record ID this spool ever marked transmitted, with its record hash.

func (*MemStore) WriterID

func (s *MemStore) WriterID() protocol.WriterID

type Options

type Options struct {
	Store     Store
	Transport Transport
	Hooks     Hooks
	Agent     protocol.AgentInfo
	Logger    *slog.Logger
	Clock     Clock
	// Jitter returns a random duration in [0, ceiling]; the default is uniform (full jitter).
	Jitter               func(ceiling time.Duration) time.Duration
	BackoffBase          time.Duration // default 1s
	BackoffMax           time.Duration // default 5m
	ReplayBytesPerSecond int64         // backlog replay rate limit; 0 is unlimited
	HealthInterval       time.Duration // default 60s; negative disables health reports
	CallTimeout          time.Duration // default 30s
	MaxBatchRecords      int           // default 512
	DisableCompression   bool
}

Options configures a Client.

type RecordState

type RecordState int

RecordState is the delivery state of a spooled record (SPEC 8.1). Committed records are deleted.

const (
	NeverTransmitted RecordState = iota
	TransmittedUnconfirmed
)

func (RecordState) String

func (s RecordState) String() string

type RejectError

type RejectError struct {
	Code string
	Err  error
}

RejectError is a hello rejection after which the writer keeps spooling and retries.

func (*RejectError) Error

func (e *RejectError) Error() string

func (*RejectError) Unwrap

func (e *RejectError) Unwrap() error

type Request

type Request struct {
	Method       string
	Params       json.RawMessage
	Notification bool
}

Request is an incoming JSON-RPC request or notification.

type ServerInfo

type ServerInfo struct {
	Name    string `json:"name"`
	Version string `json:"version"`
}

ServerInfo names the MCP server.

type Status

type Status struct {
	Connected  bool
	SessionID  string
	Epoch      protocol.EpochID
	Registered bool
	Head       uint64
	Watermark  uint64
	// BacklogRecords and BacklogBytes count what Store.Entries(0) returns.
	BacklogRecords int
	BacklogBytes   int64
	Sessions       int
	Rebaselines    int
	LastError      string
	LastErrorCode  string
	Halted         string
	Compat         protocol.Compat
}

Status is a snapshot of the writer session.

type StopError

type StopError struct{ Code string }

StopError ends Run: the writer must stop writing (SPEC 8.3 rows 1, 4 to 7, supersede, de-enrollment).

func (*StopError) Error

func (e *StopError) Error() string

type Store

type Store interface {
	WriterID() protocol.WriterID
	Incarnation() uint64
	Identity() Identity
	SetIdentity(Identity) error
	Epoch() (EpochState, bool)
	OpenEpoch(reason string, prev *protocol.EpochID, prevHead *uint64) (protocol.EpochID, error)
	MarkRegistered() error
	Do(fn func(tx Tx) error) error
	Entries(fromSeq uint64) []*Entry
	MarkTransmitted(seqs ...uint64) error
	Commit(epoch protocol.EpochID, seq uint64, chainHash protocol.Hash) error
	LastCommitted() (protocol.ChainPoint, bool)
	DiscardAbove(seq uint64) error
	Notify() <-chan struct{}
	SetHalted(code string) error
	Halted() (string, bool)
}

Store is the durable spool used by the writer session. See the package documentation.

type Tool

type Tool struct {
	Name        string          `json:"name"`
	Description string          `json:"description,omitempty"`
	InputSchema json.RawMessage `json:"inputSchema"`
}

Tool describes one investigation tool.

type ToolsListResult

type ToolsListResult struct {
	Tools []Tool `json:"tools"`
}

ToolsListResult is the tools/list result.

type Transport

type Transport interface {
	Dial(ctx context.Context) (Conn, error)
}

Transport dials a control plane session.

type Tx

type Tx interface {
	Append(t protocol.RecordType, build func(env protocol.Envelope) (*protocol.Record, error)) (*Entry, error)
}

Tx appends records under the sequence lock.

Directories

Path Synopsis
Package clienttest provides an in-memory writer model implementing client.Hooks over a client.MemStore, for tests.
Package clienttest provides an in-memory writer model implementing client.Hooks over a client.MemStore, for tests.

Jump to

Keyboard shortcuts

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