nodeapi

package
v0.16.0 Latest Latest
Warning

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

Go to latest
Published: Oct 4, 2026 License: Apache-2.0 Imports: 28 Imported by: 0

Documentation

Overview

Package nodeapi implements the in-cluster HTTPS API between node agents and the coordinator (docs/architecture.md).

Index

Constants

View Source
const (
	DefaultClientTimeout   = 30 * time.Second
	DefaultLongPollTimeout = DefaultLongPollMax + 20*time.Second
)

Client defaults.

View Source
const (
	DefaultLongPollMax  = 55 * time.Second
	DefaultAuthLogEvery = 10 * time.Second
)

Server defaults.

View Source
const (
	MaxRequestBytes  = 16 << 20
	MaxResponseBytes = 20 << 20
)

Body size caps.

View Source
const (
	MaxRuleStatuses = 1024
	MaxRuleReason   = 512
)

Bounds of the rule states a registration carries.

View Source
const AuditImpersonation = "impersonation"

AuditImpersonation is the AuditEvent kind for a submission naming another node.

View Source
const ContentType = "application/cbor"

ContentType is the media type of every request and response body.

View Source
const DefaultAudience = "exitmesh-coordinator"

DefaultAudience is the projected token audience the coordinator accepts.

View Source
const DefaultMaxTokenLifetime = time.Hour

DefaultMaxTokenLifetime bounds exp minus iat; the chart projects 600 s tokens.

Variables

View Source
var (
	// ErrUnauthenticated reports a missing or invalid bearer token.
	ErrUnauthenticated = errors.New("nodeapi: unauthenticated")
	// ErrAuthUnavailable reports that issuer discovery or key retrieval failed.
	ErrAuthUnavailable = errors.New("nodeapi: token verification unavailable")
)
View Source
var ErrInvalid = errors.New("nodeapi: invalid message")

ErrInvalid reports a message that fails decoding or validation.

View Source
var ErrUnavailable = errors.New("coordinator unavailable")

ErrUnavailable marks a transient backend condition; the node is told to retry without an error being logged.

View Source
var ErrUnknownTask = errors.New("nodeapi: unknown task")

ErrUnknownTask is returned by Backend.TaskResult for a task not assigned to the node.

Functions

func IsMalformed

func IsMalformed(err error) bool

IsMalformed reports whether the request content was refused; authentication, authorization, and transport failures never are.

func IsRetryable

func IsRetryable(err error) bool

IsRetryable reports whether err is an Error worth retrying.

func Marshal

func Marshal(v any) ([]byte, error)

Marshal validates v when it is an API message and encodes it with deterministic CBOR.

func NewServer

func NewServer(o Options) http.Handler

NewServer returns the node API handler; the http.Server write timeout must exceed LongPollMax.

func RuleStateRank

func RuleStateRank(state string) int

RuleStateRank orders rule states from the most to the least severe for reporting.

func Unmarshal

func Unmarshal(b []byte, v any) error

Unmarshal strictly decodes b into v and validates it, rejecting bodies above MaxResponseBytes.

Types

type AuditEvent

type AuditEvent struct {
	Kind              string
	Time              time.Time
	AuthenticatedNode string
	ClaimedNode       string
	Namespace         string
	Pod               string
	PodUID            string
	Method            string
	Path              string
	RemoteAddr        string
}

AuditEvent records a security-relevant rejection.

type AuthOptions

type AuthOptions struct {
	// APIServer is the API server base URL serving issuer discovery and /openid/v1/jwks.
	APIServer string
	// HTTPClient carries the API server transport and credentials.
	HTTPClient     *http.Client
	Audience       string
	Namespace      string
	ServiceAccount string
	ResolvePod     NodeResolver
	MaxLifetime    time.Duration
	Clock          func() time.Time
}

AuthOptions configures an Authenticator.

type Authenticator

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

Authenticator verifies projected tokens offline, so a token outlives its pod until expiry.

func NewAuthenticator

func NewAuthenticator(o AuthOptions) (*Authenticator, error)

NewAuthenticator validates options; discovery runs on first use and is retried after failures.

func (*Authenticator) Authenticate

func (a *Authenticator) Authenticate(ctx context.Context, token string) (Identity, error)

Authenticate verifies a bearer token and returns the node agent identity it proves.

type Backend

type Backend interface {
	Register(ctx context.Context, node string, req RegisterRequest) (RegisterResponse, error)
	// Submit returns only after the items are durably stored; queue names the sequence space of the item sequences.
	Submit(ctx context.Context, node, queue string, items []Item) (ackedThrough uint64, err error)
	// Bundle returns nil when have is current; changed is closed once the answer may differ.
	Bundle(ctx context.Context, node, have string) (*BundlePayload, <-chan struct{}, error)
	// Kube reports nothing new while the revision is at most since.
	Kube(ctx context.Context, node string, since uint64) (KubeUpdate, <-chan struct{}, error)
	Tasks(ctx context.Context, node string) ([]Task, <-chan struct{}, error)
	TaskResult(ctx context.Context, node string, res TaskResult) error
}

Backend is the coordinator side of the API; node is always the authenticated node.

type BundlePayload

type BundlePayload struct {
	Version     string `cbor:"1,keyasint"`
	Archive     []byte `cbor:"2,keyasint"`
	Signature   []byte `cbor:"3,keyasint,omitempty"`
	KeyManifest []byte `cbor:"4,keyasint,omitempty"`
	// KeyManifestChain is forwarded unchanged from the control plane (ascending sequence).
	KeyManifestChain [][]byte `cbor:"5,keyasint,omitempty"`
}

BundlePayload is a signed rule bundle for distribution to node agents.

type Client

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

Client is the node agent side of the API. Retries are the caller's job.

func NewClient

func NewClient(o ClientOptions) (*Client, error)

NewClient builds a client for the coordinator Service.

func (*Client) PostResult

func (c *Client) PostResult(ctx context.Context, res TaskResult) error

PostResult reports a task result.

func (*Client) Register

func (c *Client) Register(ctx context.Context, req RegisterRequest) (RegisterResponse, error)

Register announces the node and returns the target bundle version.

func (*Client) Submit

func (c *Client) Submit(ctx context.Context, queue string, items []Item) (uint64, error)

Submit sends items of the queue in ascending sequence order and returns the highest durably acknowledged sequence.

func (*Client) WaitBundle

func (c *Client) WaitBundle(ctx context.Context, have string) (*BundlePayload, error)

WaitBundle long-polls for a target bundle other than have; nil means no change before the poll ended.

func (*Client) WaitKube

func (c *Client) WaitKube(ctx context.Context, since uint64) (*KubeUpdate, error)

WaitKube long-polls for a kube series revision above since; nil means no change.

func (*Client) WaitTasks

func (c *Client) WaitTasks(ctx context.Context) ([]Task, error)

WaitTasks long-polls for pending tasks; nil means none arrived.

type ClientOptions

type ClientOptions struct {
	BaseURL string
	// CAFile verifies the coordinator certificate; empty uses the system roots.
	CAFile string
	// TokenFile is reread on every request because the kubelet rotates it.
	TokenFile       string
	Node            string
	Timeout         time.Duration
	LongPollTimeout time.Duration
}

ClientOptions configures NewClient.

type Error

type Error struct {
	Op         string
	StatusCode int
	Message    string
	Retryable  bool
	// Malformed is set when the request content itself is refused (HTTP 400 or 413) or cannot be encoded.
	Malformed bool
	Err       error
}

Error is a failed API call; Retryable is set for network errors, 5xx, 408, and 429.

func (*Error) Error

func (e *Error) Error() string

func (*Error) Unwrap

func (e *Error) Unwrap() error

type Identity

type Identity struct {
	Node           string
	Namespace      string
	ServiceAccount string
	Pod            string
	PodUID         string
	Expiry         time.Time
}

Identity is an authenticated node agent.

type Item

type Item struct {
	Seq     uint64             `cbor:"1,keyasint"`
	Kind    ItemKind           `cbor:"2,keyasint"`
	Finding []byte             `cbor:"3,keyasint,omitempty"`
	Facts   []metricfacts.Fact `cbor:"4,keyasint,omitempty"`
	Part    *Part              `cbor:"5,keyasint,omitempty"`
}

Item is one node queue entry.

type ItemKind

type ItemKind string

ItemKind is the payload type of a queue item.

const (
	KindFinding     ItemKind = "finding"
	KindMetricFacts ItemKind = "metric_facts"
	KindSeries      ItemKind = "series"
)

Item kinds.

type KubeSeries

type KubeSeries struct {
	Labels map[string]string `cbor:"1,keyasint,omitempty"`
	Value  float64           `cbor:"2,keyasint"`
}

KubeSeries is one kube_* series value.

type KubeUpdate

type KubeUpdate struct {
	Revision uint64       `cbor:"1,keyasint"`
	Series   []KubeSeries `cbor:"2,keyasint,omitempty"`
}

KubeUpdate is the kube_* series set for one node at a revision.

type NodeResolver

type NodeResolver func(namespace, name, uid string) (nodeName string, ok bool)

NodeResolver maps a pod claim to the pod's node; it must return false when the pod's UID differs.

func PodNodeResolver

func PodNodeResolver(lookup func(namespace, name string) (nodeName, uid string, ok bool)) NodeResolver

PodNodeResolver builds a NodeResolver from a pod lookup, enforcing the UID match and a scheduled pod.

type Options

type Options struct {
	Auth         *Authenticator
	Backend      Backend
	Audit        func(AuditEvent)
	Clock        func() time.Time
	MaxBody      int64
	LongPollMax  time.Duration
	Logger       *slog.Logger
	AuthLogEvery time.Duration
}

Options configures NewServer.

type Part

type Part struct {
	RuleID        string   `cbor:"1,keyasint"`
	RuleVersion   int      `cbor:"2,keyasint,omitempty"`
	BundleVersion string   `cbor:"3,keyasint,omitempty"`
	EvalTimeMs    int64    `cbor:"4,keyasint"`
	Samples       []Sample `cbor:"5,keyasint,omitempty"`
}

Part is engine.Part flattened for the wire.

func PartFromEngine

func PartFromEngine(p engine.Part) (Part, error)

PartFromEngine flattens an engine part, rejecting histogram samples and non-finite values.

func (Part) Engine

func (p Part) Engine() engine.Part

Engine converts the wire part back to an engine part.

type Process added in v0.4.0

type Process struct {
	UID          int      `cbor:"1,keyasint" json:"uid"`
	GID          int      `cbor:"2,keyasint" json:"gid"`
	Capabilities []string `cbor:"3,keyasint,omitempty" json:"capabilities"`
}

Process is a process identity.

type QueueUsage

type QueueUsage struct {
	Bytes    int64  `cbor:"1,keyasint,omitempty"`
	Capacity int64  `cbor:"2,keyasint,omitempty"`
	Items    uint64 `cbor:"3,keyasint,omitempty"`
	Acked    uint64 `cbor:"4,keyasint,omitempty"`
	Next     uint64 `cbor:"5,keyasint,omitempty"`
}

QueueUsage mirrors the node agent queue occupancy.

type RegisterRequest

type RegisterRequest struct {
	Node          string            `cbor:"1,keyasint"`
	AgentVersion  string            `cbor:"2,keyasint,omitempty"`
	BundleVersion string            `cbor:"3,keyasint,omitempty"`
	Capabilities  []string          `cbor:"4,keyasint,omitempty"`
	Coverage      map[string]string `cbor:"5,keyasint,omitempty"`
	Warming       bool              `cbor:"6,keyasint,omitempty"`
	QueueUsage    QueueUsage        `cbor:"7,keyasint"`
	// Rules are the node-local rule states, at most MaxRuleStatuses, most severe first.
	Rules []RuleStatus `cbor:"8,keyasint,omitempty"`
	// Process is the node agent's user and effective capabilities, so a root agent is visible on the connector.
	Process *Process `cbor:"9,keyasint,omitempty"`
}

RegisterRequest is the body of POST /v1/node/register.

type RegisterResponse

type RegisterResponse struct {
	TargetBundle string `cbor:"1,keyasint,omitempty"`
	ServerTimeMs int64  `cbor:"2,keyasint"`
}

RegisterResponse answers a registration.

type RuleStatus

type RuleStatus struct {
	RuleID          string `cbor:"1,keyasint"`
	Version         int    `cbor:"2,keyasint,omitempty"`
	State           string `cbor:"3,keyasint"`
	Reason          string `cbor:"4,keyasint,omitempty"`
	LastEvalMs      int64  `cbor:"5,keyasint,omitempty"`
	BudgetLimited   bool   `cbor:"6,keyasint,omitempty"`
	EvidenceLimited bool   `cbor:"7,keyasint,omitempty"`
}

RuleStatus is one node-local rule state.

type Sample

type Sample struct {
	Labels      map[string]string `cbor:"1,keyasint,omitempty"`
	Value       float64           `cbor:"2,keyasint"`
	TimestampMs int64             `cbor:"3,keyasint"`
}

Sample is one pre-aggregated series point.

type SubmitRequest

type SubmitRequest struct {
	Node  string `cbor:"1,keyasint"`
	Items []Item `cbor:"2,keyasint,omitempty"`
	// Queue identifies the node queue whose sequence space Items belong to; empty from agents that predate it.
	Queue string `cbor:"3,keyasint,omitempty"`
}

SubmitRequest is the body of POST /v1/node/records.

type SubmitResponse

type SubmitResponse struct {
	AckedThrough uint64 `cbor:"1,keyasint"`
}

SubmitResponse acknowledges items durably stored by the coordinator.

type Task

type Task struct {
	ID         string   `cbor:"1,keyasint"`
	Kind       TaskKind `cbor:"2,keyasint"`
	Payload    []byte   `cbor:"3,keyasint,omitempty"`
	DeadlineMs int64    `cbor:"4,keyasint,omitempty"`
}

Task is an investigation or aggregation task assigned to a node.

type TaskKind

type TaskKind string

TaskKind is the type of a coordinator task.

const (
	TaskPromQLQuery TaskKind = "promql_query"
	TaskLogQLQuery  TaskKind = "logql_query"
	TaskLogRead     TaskKind = "log_read"
	TaskEvidence    TaskKind = "evidence_read"
)

Task kinds.

type TaskList

type TaskList struct {
	Tasks []Task `cbor:"1,keyasint,omitempty"`
}

TaskList is the body of a GET /v1/node/tasks response.

type TaskResult

type TaskResult struct {
	ID      string `cbor:"1,keyasint"`
	Payload []byte `cbor:"2,keyasint,omitempty"`
	Error   string `cbor:"3,keyasint,omitempty"`
}

TaskResult is the body of POST /v1/node/tasks/{id}.

Jump to

Keyboard shortcuts

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