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
- Variables
- func RPCErrorCode(err error) string
- type CallToolParams
- type CallToolResult
- type CaptureTx
- type Client
- func (c *Client) Audit(ctx context.Context, params any) error
- func (c *Client) Deenroll(ctx context.Context, reason string) error
- func (c *Client) FetchBundle(ctx context.Context, have string) (*protocol.BundleFetchResult, error)
- func (c *Client) Prepare() error
- func (c *Client) Run(ctx context.Context) error
- func (c *Client) Status() Status
- type Clock
- type Conn
- type Content
- type DivergenceError
- type Entry
- type EpochState
- type Handler
- type Hooks
- type Identity
- type InitializeResult
- type MemOptions
- type MemStore
- func (s *MemStore) Coalesce() (int, error)
- func (s *MemStore) Commit(epoch protocol.EpochID, seq uint64, chainHash protocol.Hash) error
- func (s *MemStore) DiscardAbove(seq uint64) error
- func (s *MemStore) Do(fn func(tx Tx) error) error
- func (s *MemStore) Entries(fromSeq uint64) []*Entry
- func (s *MemStore) Epoch() (EpochState, bool)
- func (s *MemStore) Halted() (string, bool)
- func (s *MemStore) Identity() Identity
- func (s *MemStore) Incarnation() uint64
- func (s *MemStore) LastCommitted() (protocol.ChainPoint, bool)
- func (s *MemStore) MarkRegistered() error
- func (s *MemStore) MarkTransmitted(seqs ...uint64) error
- func (s *MemStore) Notify() <-chan struct{}
- func (s *MemStore) OpenEpoch(reason string, prev *protocol.EpochID, prevHead *uint64) (protocol.EpochID, error)
- func (s *MemStore) Reopen() *MemStore
- func (s *MemStore) SetHalted(code string) error
- func (s *MemStore) SetIdentity(id Identity) error
- func (s *MemStore) Transmitted() map[protocol.RecordID]protocol.Hash
- func (s *MemStore) WriterID() protocol.WriterID
- type Options
- type RecordState
- type RejectError
- type Request
- type ServerInfo
- type Status
- type StopError
- type Store
- type Tool
- type ToolsListResult
- type Transport
- type Tx
Constants ¶
const ( HaltSuperseded = "superseded" HaltDeenrolled = "deenrolled" )
Halt codes in addition to the hello rejection codes of SPEC 8.3.
const MCPProtocolVersion = "2025-06-18"
MCP message shapes served over the reverse channel (SPEC 9.4, 9.5).
const MaxBackoff = 5 * time.Minute
MaxBackoff caps the delay between connection attempts.
Variables ¶
var ( ErrDisconnected = errors.New("connection closed") ErrNotConnected = errors.New("no active session") )
Session errors.
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 ¶
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 (*Client) Deenroll ¶
Deenroll requests explicit de-enrollment; on success the credential is deleted and Run stops.
func (*Client) FetchBundle ¶
FetchBundle requests the assigned rule bundle over the active session.
func (*Client) Prepare ¶
Prepare opens an initial epoch offline and emits its first checkpoint, so records can be appended before connecting.
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 DivergenceError ¶
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.
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 ¶
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 ¶
Coalesce folds the oldest run of two or more never-transmitted non-checkpoint records into a range (SPEC 8.7).
func (*MemStore) DiscardAbove ¶
func (*MemStore) Epoch ¶
func (s *MemStore) Epoch() (EpochState, bool)
func (*MemStore) Incarnation ¶
func (*MemStore) LastCommitted ¶
func (s *MemStore) LastCommitted() (protocol.ChainPoint, bool)
func (*MemStore) MarkRegistered ¶
func (*MemStore) MarkTransmitted ¶
func (*MemStore) Reopen ¶
Reopen copies the durable state into a new handle with the next incarnation, as a restart does; halt is not copied.
func (*MemStore) SetIdentity ¶
func (*MemStore) Transmitted ¶
Transmitted returns every record ID this spool ever marked transmitted, with its record hash.
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 and maximum 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 ¶
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 ¶
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).
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.
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. |