Documentation
¶
Overview ¶
Package agent implements agent-level runtime services.
Index ¶
- Constants
- type AgentConfig
- type DeviceKeyUpdater
- type HeartbeatClient
- type HeartbeatConfig
- type HeartbeatService
- func (s *HeartbeatService) Run(ctx context.Context) error
- func (s *HeartbeatService) SetBuildRequest(fn func() api.HeartbeatRequest)
- func (s *HeartbeatService) SetOnAuthFailure(fn func())
- func (s *HeartbeatService) SetOnRotateKeys(fn func())
- func (s *HeartbeatService) SetReconcileTrigger(rt ReconcileTrigger)
- type KeyRotator
- type ReconcileTrigger
- type RotateClient
Constants ¶
const ( // DefaultMode is the default operating mode. DefaultMode = "node" // DefaultLogLevel is the default log level. DefaultLogLevel = "info" // DefaultDataDir is the default data directory. DefaultDataDir = "/var/lib/plexd" )
const DefaultHeartbeatInterval = 30 * time.Second
DefaultHeartbeatInterval is the default heartbeat interval.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type AgentConfig ¶
type AgentConfig struct {
// Mode is the operating mode: "node" or "bridge".
// Default: "node"
Mode string `yaml:"mode"`
// LogLevel is the log level: "debug", "info", "warn", "error".
// Default: "info"
LogLevel string `yaml:"log_level"`
// DataDir is the directory for persistent agent data.
// Default: /var/lib/plexd
DataDir string `yaml:"data_dir"`
API api.Config `yaml:"api"`
Registration registration.Config `yaml:"registration"`
Reconcile reconcile.Config `yaml:"reconcile"`
NodeAPI nodeapi.Config `yaml:"node_api"`
Actions actions.Config `yaml:"actions"`
Policy policy.Config `yaml:"policy"`
WireGuard wireguard.Config `yaml:"wireguard"`
Metrics metrics.Config `yaml:"metrics"`
LogFwd logfwd.Config `yaml:"log_fwd"`
AuditFwd auditfwd.Config `yaml:"audit_fwd"`
Integrity integrity.Config `yaml:"integrity"`
Upgrade upgrade.Config `yaml:"upgrade"`
Tunnel tunnel.Config `yaml:"tunnel"`
NAT nat.Config `yaml:"nat"`
PeerExchange peerexchange.Config `yaml:"peer_exchange"`
Bridge bridge.Config `yaml:"bridge"`
Heartbeat HeartbeatConfig `yaml:"heartbeat"`
}
AgentConfig is the top-level configuration for the plexd agent. It aggregates all subsystem configurations and is populated from a YAML configuration file via ParseConfig.
func ParseConfig ¶
func ParseConfig(path string) (*AgentConfig, error)
ParseConfig reads a YAML configuration file and returns an AgentConfig. It applies defaults and validates the configuration.
func (*AgentConfig) ApplyDefaults ¶
func (c *AgentConfig) ApplyDefaults()
ApplyDefaults sets default values for zero-valued fields.
func (*AgentConfig) Validate ¶
func (c *AgentConfig) Validate() error
Validate checks that required fields are set and values are acceptable.
type DeviceKeyUpdater ¶ added in v0.2.0
DeviceKeyUpdater installs a rotated private key on the WireGuard device.
type HeartbeatClient ¶
type HeartbeatClient interface {
Heartbeat(ctx context.Context, nodeID string, req api.HeartbeatRequest) (*api.HeartbeatResponse, error)
}
HeartbeatClient sends heartbeat requests to the control plane.
type HeartbeatConfig ¶
type HeartbeatConfig struct {
// Interval is the heartbeat send interval.
// Default: 30s
Interval time.Duration `yaml:"interval"`
// NodeID is the node identifier (required).
NodeID string `yaml:"node_id"`
}
HeartbeatConfig holds the configuration for the heartbeat service.
func (*HeartbeatConfig) ApplyDefaults ¶
func (c *HeartbeatConfig) ApplyDefaults()
ApplyDefaults sets default values for zero-valued fields.
func (*HeartbeatConfig) Validate ¶
func (c *HeartbeatConfig) Validate() error
Validate checks that required fields are set.
type HeartbeatService ¶
type HeartbeatService struct {
// contains filtered or unexported fields
}
HeartbeatService sends periodic heartbeats to the control plane and dispatches directive flags from the response.
func NewHeartbeatService ¶
func NewHeartbeatService(cfg HeartbeatConfig, client HeartbeatClient, logger *slog.Logger) *HeartbeatService
NewHeartbeatService creates a new HeartbeatService with the given configuration and client. The config is validated; defaults are applied for any zero-valued optional fields.
func (*HeartbeatService) Run ¶
func (s *HeartbeatService) Run(ctx context.Context) error
Run starts the heartbeat loop. It sends one heartbeat immediately and then continues at the configured interval until ctx is cancelled. Run always returns nil.
func (*HeartbeatService) SetBuildRequest ¶
func (s *HeartbeatService) SetBuildRequest(fn func() api.HeartbeatRequest)
SetBuildRequest sets a custom heartbeat request builder. If not set, the service sends a zero-valued HeartbeatRequest.
func (*HeartbeatService) SetOnAuthFailure ¶
func (s *HeartbeatService) SetOnAuthFailure(fn func())
SetOnAuthFailure sets a callback invoked when a heartbeat fails with a 401 Unauthorized error.
func (*HeartbeatService) SetOnRotateKeys ¶
func (s *HeartbeatService) SetOnRotateKeys(fn func())
SetOnRotateKeys sets a callback invoked when the control plane signals that keys should be rotated.
func (*HeartbeatService) SetReconcileTrigger ¶
func (s *HeartbeatService) SetReconcileTrigger(rt ReconcileTrigger)
SetReconcileTrigger sets the reconcile trigger invoked when the control plane requests reconciliation.
type KeyRotator ¶ added in v0.2.0
type KeyRotator struct {
// contains filtered or unexported fields
}
KeyRotator runs the mesh-key rotation state machine. Both rotation signals drive the same single-flight cycle: the repeated heartbeat rotate_keys flag through Rotate, which rate-limits fresh starts, and the one-shot SSE rotate_keys event through RotateNow, which does not. The control plane keeps a rotation pending until the node confirms it, so repeated signals, retries, and the RecoverPending sweep all converge on the same crash-safe outcome.
func NewKeyRotator ¶ added in v0.2.0
func NewKeyRotator(client RotateClient, id *registration.NodeIdentity, dataDir string, device DeviceKeyUpdater, rt ReconcileTrigger, logger *slog.Logger) *KeyRotator
NewKeyRotator builds a KeyRotator that stages fresh keys under dataDir, submits their public half via client, installs the confirmed private key on device (may be nil), and triggers a reconcile through rt once a rotation commits.
func (*KeyRotator) RecoverPending ¶ added in v0.2.0
func (r *KeyRotator) RecoverPending(ctx context.Context)
RecoverPending drives staged-but-uncommitted rotations for the lifetime of ctx. A staged key that the control plane already accepted strands the node unless something resubmits it: the rotation completed server-side, so no further signal arrives, every peer holds the new public key, and the node still runs the old one. That state outlives startup — a submit whose response is lost (proxy 5xx, read timeout, reset connection) leaves exactly the same staged key at any point during the process lifetime — so this loop keeps sweeping instead of returning once nothing is staged.
While a key is staged it resubmits with a capped exponential backoff; with nothing staged it re-checks at the cap and makes no HTTP call, so an idle node never starts a rotation on its own.
func (*KeyRotator) Rotate ¶ added in v0.2.0
func (r *KeyRotator) Rotate(ctx context.Context) error
Rotate runs one single-flight rotation cycle. It stages a fresh key (or reuses one staged by an earlier attempt), submits its public half, and commits the durable identity only after the control plane confirms. A concurrent Rotate that loses the TryLock returns nil without touching the staging file, and starting a fresh rotation within minInterval of the last commit is skipped; a key already staged is always resubmitted. The outcome taxonomy over the submit result is:
- success with a complete receipt: commit with the returned receipt;
- success without a receipt: a control plane that predates this contract, so keep the staged key for a resubmit and return the error;
- 422 keys_rotate_public_key_unchanged: a previous submission already landed (a crash between submit and persist), so commit with a nil receipt to keep the existing LastRotation;
- 409 keys_rotate_no_pending_rotation: a stale or cancelled signal whose staged key never landed, so discard it and return nil;
- a permanent rejection (400, 403, 413, 404 keys_rotate_peer_not_found, 422 keys_rotate_public_key_invalid): no retry changes the outcome, so discard the staged key and return the error;
- anything else (auth, 5xx, transport errors): keep the staged key so RecoverPending or the next signal retries, and return the error.
func (*KeyRotator) RotateNow ¶ added in v0.2.0
func (r *KeyRotator) RotateNow(ctx context.Context) error
RotateNow runs the same rotation cycle as Rotate but starts a fresh rotation even within minInterval of the last commit. It serves the SSE rotate_keys event: the control plane sends it once per rotation decision and never repeats it, so a cooldown skip here would strand the rotation until the heartbeat flag catches up — or forever, for a control plane that signals via SSE alone. A redelivered event is harmless: with no rotation pending the submit is answered 409 keys_rotate_no_pending_rotation and the staged key is discarded without touching the device.
type ReconcileTrigger ¶
type ReconcileTrigger interface {
TriggerReconcile()
}
ReconcileTrigger triggers an immediate reconciliation.
type RotateClient ¶ added in v0.2.0
type RotateClient interface {
RotateKeys(ctx context.Context, req api.KeyRotateRequest) (*api.KeyRotateResponse, error)
}
RotateClient submits a completed rotation's public key to the control plane.