refcp

package
v0.2.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: 16 Imported by: 0

Documentation

Overview

Package refcp is an in-memory reference control plane for contract and end-to-end tests.

Index

Constants

View Source
const (
	SourceSummary = "summary"
	SourceRecord  = "record"
)

Notification sources.

Variables

View Source
var (
	ErrUnknownTarget = errors.New("refcp: unknown target")
	ErrUnknownEpoch  = errors.New("refcp: unknown epoch")
	ErrNotCommitted  = errors.New("refcp: sequence not committed")
	ErrNoSession     = errors.New("refcp: target has no active session")
)

Errors returned by management and query methods.

Functions

This section is empty.

Types

type AuditEntry

type AuditEntry struct {
	Time        time.Time
	Target      string
	Session     string
	Event       string
	Row         int
	Code        string
	Writer      protocol.WriterID
	Incarnation uint64
	Epoch       protocol.EpochID
	Seq         uint64
	Credential  string
	Alarm       bool
	Detail      string
}

AuditEntry is one audited event.

type Boundary

type Boundary struct {
	Epoch    protocol.EpochID
	Seq      uint64
	Expected protocol.Hash
	Got      protocol.Hash
}

Boundary is a reconstruction boundary: a checkpoint that did not match replayed state.

type EdgeChange

type EdgeChange struct {
	Epoch     protocol.EpochID
	Seq       uint64
	Time      time.Time
	Op        protocol.OpKind
	From      string
	Type      string
	To        string
	Attrs     map[string]any
	PrevAttrs map[string]any
	// Coalesced marks a folded range operation; Span is its unavailable interior.
	Coalesced bool
	Span      protocol.Span
}

EdgeChange is one committed edge operation.

type EpochView

type EpochView struct {
	ID         protocol.EpochID
	Owner      protocol.WriterID
	Open       bool
	Head       protocol.ChainPoint
	ClosedAt   uint64
	OpenReason string
	PrevEpoch  *protocol.EpochID
	PrevHead   *uint64
	Watermark  uint64
	Pending    []uint64
	Registered time.Time
}

EpochView describes one epoch in registration order.

type FindingState

type FindingState struct {
	FindingID     string
	DedupKey      string
	State         string
	Transition    string
	Severity      string
	RuleID        string
	BundleVersion string
	FirstSeen     uint64
	LastSeen      uint64
	EvalTime      uint64
	Count         uint64
	Epoch         protocol.EpochID
	Seq           uint64
}

FindingState is the committed lifecycle of one finding episode.

type Notification

type Notification struct {
	Target    string
	FindingID string
	DedupKey  string
	Severity  string
	EvalTime  uint64
	// Late marks a backlog notification decided from the lifecycle summary.
	Late bool
	// LateDelivered marks an evaluation time older than the late-delivery threshold on arrival.
	LateDelivered bool
	Source        string
	Epoch         protocol.EpochID
	Seq           uint64
	Time          time.Time
}

Notification is one notification decision (SPEC 8.5).

type Options

type Options struct {
	Now           func() time.Time
	LateThreshold time.Duration // default 15m
	WindowBytes   int64         // default protocol.DefaultWindowSize
	MaxFrameBytes int64         // default protocol.MaxFramePayload
	Conn          tunnel.ConnOptions
	Logger        *slog.Logger
}

Options configures a Server.

type Server

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

Server is the reference control plane. It is safe for concurrent use.

func New

func New(opts Options) *Server

New returns an empty control plane.

func (*Server) Audit

func (s *Server) Audit(targetID string) []AuditEntry

Audit returns the audit log, for one target or for all when targetID is empty.

func (*Server) BindHost

func (s *Server) BindHost(targetID string) error

BindHost records a pending administrator binding: the next new writer takes over the host (row 14).

func (*Server) Boundaries

func (s *Server) Boundaries(targetID string) ([]Boundary, error)

Boundaries returns the reconstruction boundaries of a target.

func (*Server) CallTool

func (s *Server) CallTool(ctx context.Context, targetID, name string, args any) (*client.CallToolResult, error)

CallTool forwards an MCP tools/call to the target's attached writer.

func (*Server) Close

func (s *Server) Close()

Close ends every session.

func (*Server) CreateHostGroup

func (s *Server) CreateHostGroup() (groupID, token string, err error)

CreateHostGroup returns a host group enrollment token (emx1_g_); each enrollment creates a host target.

func (*Server) CreateTarget

func (s *Server) CreateTarget(targetType string) (targetID, token string, err error)

CreateTarget registers a target and returns its enrollment token (emx1_c_ or emx1_h_).

func (*Server) Deenroll

func (s *Server) Deenroll(targetID string) error

Deenroll revokes every credential of the target and tells an attached writer to stop.

func (*Server) DelayAcks

func (s *Server) DelayAcks(d time.Duration)

DelayAcks delays every acknowledgement by d.

func (*Server) Disconnect

func (s *Server) Disconnect(targetID string) error

Disconnect closes the target's attached session, as a network failure would.

func (*Server) DisconnectAfterRecords

func (s *Server) DisconnectAfterRecords(n int)

DisconnectAfterRecords closes the session right after the n-th next committed record, before acknowledging it.

func (*Server) DropAcks

func (s *Server) DropAcks(n int)

DropAcks commits the next n acknowledged frames without sending their acknowledgement.

func (*Server) EdgeHistory

func (s *Server) EdgeHistory(targetID, uid, edgeType string, from, to time.Time) ([]EdgeChange, error)

EdgeHistory returns committed changes of edges touching uid (optionally one edge type) with record time in [from, to].

func (*Server) Epochs

func (s *Server) Epochs(targetID string) ([]EpochView, error)

Epochs lists a target's epochs in registration order.

func (*Server) Findings

func (s *Server) Findings(targetID string) ([]FindingState, error)

Findings returns the committed lifecycle of every finding of a target, ordered by finding ID.

func (*Server) Heads

func (s *Server) Heads(targetID string) (map[protocol.EpochID]protocol.ChainPoint, error)

Heads returns the committed head of every epoch.

func (*Server) HealthReports

func (s *Server) HealthReports(targetID string) ([]json.RawMessage, error)

HealthReports returns the agent.health payloads received for a target.

func (*Server) ListTools

func (s *Server) ListTools(ctx context.Context, targetID string) ([]client.Tool, error)

ListTools returns the tools advertised by the target's attached writer.

func (*Server) Notifications

func (s *Server) Notifications(targetID string) ([]Notification, error)

Notifications returns the notification log of a target in decision order.

func (*Server) PublishBundle

func (s *Server) PublishBundle(targetType, version string, archive, signature, keyManifest []byte) error

PublishBundle stores the latest bundle for a target type and announces it to attached writers.

func (*Server) PublishKeyManifestChain

func (s *Server) PublishKeyManifestChain(targetType string, chain [][]byte) error

PublishKeyManifestChain sets the earlier key manifests served with bundles for targetType.

func (*Server) Records

func (s *Server) Records(targetID string, id protocol.EpochID) ([]*protocol.Record, error)

Records returns the committed records of an epoch in chain order. Callers must not mutate them.

func (*Server) RejectNext

func (s *Server) RejectNext(code string)

RejectNext rejects the next received record with code instead of applying it.

func (*Server) RejectNextHello

func (s *Server) RejectNextHello(code string)

RejectNextHello rejects the next hello with code.

func (*Server) ResolveConflict

func (s *Server) ResolveConflict(targetID, machineID string) error

ResolveConflict clears an identity conflict; a non-empty machineID becomes the enrolled machine ID.

func (*Server) RotateCredential

func (s *Server) RotateCredential(ctx context.Context, targetID string) (string, error)

RotateCredential sends a new credential to the attached writer and revokes the old one once it is persisted.

func (*Server) ServeHTTP

func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request)

ServeHTTP serves the enrollment endpoint and the tunnel.

func (*Server) SetCompat

func (s *Server) SetCompat(targetID string, c protocol.Compat) error

SetCompat sets the compatibility verdict returned in hello results.

func (*Server) SetUnavailable

func (s *Server) SetUnavailable(v bool)

SetUnavailable makes enrollment and new tunnel connections fail with 503; attached sessions continue.

func (*Server) StateAt

func (s *Server) StateAt(targetID string, id protocol.EpochID, seq uint64) (*protocol.State, error)

StateAt reconstructs the state at seq; range interiors return protocol.ErrUnavailable.

func (*Server) StateHashAt

func (s *Server) StateHashAt(targetID string, id protocol.EpochID, seq uint64) (protocol.Hash, error)

StateHashAt returns the state hash at seq from the verified replay of the epoch.

func (*Server) Stats

func (s *Server) Stats(targetID string) (Stats, error)

Stats returns record handling counters of a target.

func (*Server) Summaries

func (s *Server) Summaries(targetID string) ([]protocol.SummaryParams, error)

Summaries returns every lifecycle summary received for a target.

func (*Server) Unavailable

func (s *Server) Unavailable(targetID string, id protocol.EpochID) ([]protocol.Span, error)

Unavailable returns the coalesced range interiors of an epoch.

type Stats

type Stats struct {
	Committed   int
	Duplicates  int
	Rejected    int
	Divergences int
}

Stats counts record handling per target.

Jump to

Keyboard shortcuts

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