cluster

package
v2.2.1 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: AGPL-3.0 Imports: 22 Imported by: 0

Documentation

Overview

Package cluster lets several Studio replicas share one database: it registers each replica (Node) so others can tell live replicas from dead ones, and carries events every replica must see (Log). See features/ClusterControlPlane.md.

On Postgres it is always on, whether one replica runs or ten: a deployment that forgot to switch it on would fail silently. On SQLite, which only one process can serve, the node registry still runs (one row) and the event log is inert.

Index

Constants

View Source
const DefaultLeaseTTL = 30 * time.Second

DefaultLeaseTTL is how long the leader lease lasts without renewal unless LeadershipOptions.TTL says otherwise.

View Source
const DefaultLogRetention = 15 * time.Minute

DefaultLogRetention is how long the event log keeps rows unless LogOptions.Retention says otherwise: the longest a reader can fall behind and still read every event.

View Source
const LeaderLease = "leader"

LeaderLease is the lease singleton background work runs under: at most one replica holds it, and jobs that must run once for the whole cluster (aggregations, alerts, syncs, cleanups) check Leadership.IsLeader first.

View Source
const MaxLabelLength = 64

MaxLabelLength bounds NodeOptions.Label (the column's size).

View Source
const RelayTopic = "bus.event"

RelayTopic is the cluster log topic bus events travel on.

Variables

View Source
var (
	HeartbeatInterval = 5 * time.Second
	LivenessWindow    = 20 * time.Second
	NodeRetention     = time.Hour
)

Defaults for the node registry. A node counts as live while its row was refreshed within LivenessWindow; that is four refreshes, so one slow database round trip never makes a live node look dead. Rows of nodes that stopped without removing theirs (a crash, a kill) are deleted once they are NodeRetention old, by whichever node prunes first: nothing refers to a node that long dead (its leases, claims and edge streams were taken over long before), and until then the rows show what ran recently.

Functions

func Holder

func Holder(ctx context.Context, db *gorm.DB, name string) (string, *time.Time, error)

Holder returns the replica holding the lease and when it expires, by the database. An expired or missing lease returns "".

func IsLive

func IsLive(db *gorm.DB, nodeID string) (bool, error)

IsLive reports whether nodeID has a fresh registration.

func LiveNodes

func LiveNodes(db *gorm.DB) ([]models.ClusterNode, error)

LiveNodes returns the nodes whose registration is fresh, by the database's clock (never the caller's, so clock skew between replicas cannot make a live node look dead).

func NewNodeID

func NewNodeID() string

NewNodeID returns an ID for this process: hostname, pid and a random suffix. The suffix matters in containers, where every process may be pid 1 and a restarted replica must not inherit its predecessor's claims.

func Owners

func Owners(ctx context.Context, db *gorm.DB, nodeIDs []string) (map[string]Owner, error)

Owners returns, for the node IDs that own edge streams (a page of edges), each node's label and liveness, in one query. IDs with no registration, and empty IDs, are left out.

func RelayedByDefault

func RelayedByDefault(ev eventbridge.Event) bool

RelayedByDefault reports whether a bus event must reach every replica: events for edges (DirDown), which each replica forwards to the edges whose streams it holds, and object change events (system.*), which keep every replica's caches and plugins current. Nothing an edge sent (FromEdge) is relayed by default: object changes come from Studio, and an edge's system.* event must not make other replicas reload plugins or caches.

Types

type Event

type Event struct {
	ID        int64
	Topic     string
	Origin    string // node ID of the publisher
	Payload   []byte
	CreatedAt time.Time
}

Event is one entry of the cluster event log.

type Handler

type Handler func(Event)

Handler handles an event published by another replica. It runs on the log's reader goroutine, one event at a time in id order, so it must be quick (hand long work to a goroutine) and idempotent: delivery is at least once.

type Leadership

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

Leadership keeps this replica's claim on a named lease. The lease is taken and renewed with conditional writes against the database's clock, so at most one replica holds it. A replica believes it leads only until its last successful renewal plus the TTL minus one renewal period, by its own monotonic clock: it stops believing before the database would let another replica take over, even if it cannot reach the database to find out.

With SQLite, which one process serves, that process leads from Start to Stop whatever the lease row says: a restart after a crash must not wait for its dead predecessor's lease to expire. It still records itself as the holder, for the cluster status.

On Postgres a replica restarted after a crash does not wait for its dead predecessor's lease either, when it can tell that the holder is that predecessor (see predecessorGone); any other holder keeps the lease until it expires.

func NewLeadership

func NewLeadership(db *gorm.DB, name, node string, opts LeadershipOptions) *Leadership

NewLeadership returns node's claim on the lease called name. It does nothing until Start.

func (*Leadership) IsLeader

func (l *Leadership) IsLeader() bool

IsLeader reports whether this replica holds the lease now. It is cheap: call it before every run of singleton work.

func (*Leadership) OnChange

func (l *Leadership) OnChange(fn func(leading bool))

OnChange registers fn to be called (on the lease goroutine) when this replica gains (true) or loses (false) the lease.

func (*Leadership) Start

func (l *Leadership) Start()

Start tries to take the lease at once, then keeps renewing or trying.

func (*Leadership) Stop

func (l *Leadership) Stop()

Stop releases the lease if this replica holds it, so another replica takes over at its next attempt instead of waiting for the TTL.

type LeadershipOptions

type LeadershipOptions struct {
	// TTL is how long the lease lasts without renewal (default 30 s);
	// Renew how often the holder renews it and others try to take it
	// (default TTL/3).
	TTL   time.Duration
	Renew time.Duration
	// PredecessorSilence is how long a holder that ran on this machine
	// under this hostname, but in another pid namespace (typically this
	// container's predecessor in its pod), must have been silent in the
	// node registry before this replica takes the lease from it (default
	// two HeartbeatIntervals). See predecessorGone.
	PredecessorSilence time.Duration
}

LeadershipOptions tune a Leadership. The zero value gives the defaults.

type Log

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

Log is the cluster event log: events every other replica must see, kept in a table (cluster_events) that every replica reads. Delivery to each live replica is at least once and in id order; a replica never receives its own events. A replica that starts later does not receive what was published before it started: replicas build their state from the database at start, so they need no history.

func NewLog

func NewLog(db *gorm.DB, node string, opts LogOptions) *Log

NewLog returns the event log for node. It does nothing until Start.

func (*Log) Enabled

func (l *Log) Enabled() bool

Enabled reports whether this database carries the log (Postgres). On other databases Publish does nothing and no event is delivered.

func (*Log) Publish

func (l *Log) Publish(ctx context.Context, topic string, payload []byte) error

Publish appends an event for every other replica and wakes them. It returns once the event is durable. On a database without the log it does nothing.

func (*Log) PublishBatch

func (l *Log) PublishBatch(ctx context.Context, topic string, payloads [][]byte) error

PublishBatch appends one event per payload, all on topic, in one statement with one notification: for a burst of events (an edge's batch of plugin payloads) that would otherwise cost a round trip each. Either every event is appended or none is.

func (*Log) Start

func (l *Log) Start(ctx context.Context) error

Start begins reading. On Postgres it positions the log after the events already published (so none is replayed), listens for notifications, and reads on every notification, every PollInterval and after every listener reconnect. The notifications only make delivery faster: when the listener cannot connect, the log reads every PollInterval and keeps trying to listen in the background. On other databases it does nothing.

func (*Log) Stats

func (l *Log) Stats() LogStats

Stats returns the log's progress.

func (*Log) Stop

func (l *Log) Stop()

Stop ends reading. No handler runs after it returns. It returns promptly whether or not Start ran or succeeded; a Start after Stop does nothing.

func (*Log) Subscribe

func (l *Log) Subscribe(topic string, h Handler) (unsubscribe func())

Subscribe registers h for events on topic ("*" for every topic). It may be called before or after Start. The returned function removes it.

type LogOptions

type LogOptions struct {
	// PollInterval is how often the log is read without a notification
	// (default 5 s). Notifications only make delivery faster: a lost one
	// (listener reconnecting, NOTIFY dropped) delays an event by at most
	// this long. Every notification wakes a read, and so does every
	// listener reconnect, so the poll only covers a listener whose
	// connection died unnoticed. Each poll is two queries on every
	// replica, idle or not.
	PollInterval time.Duration
	// Window is how far back every read looks again for rows it has not
	// seen (default 30 s). An id is allocated before its row commits, so a
	// row can become visible after a higher id was already read; the window
	// is how long such a straggler is still picked up.
	Window time.Duration
	// Retention is how long rows are kept (default 15 min); PruneInterval
	// how often they are pruned (default 1 min). Retention is at least
	// twice Window.
	Retention     time.Duration
	PruneInterval time.Duration
	// BatchSize bounds one read (default 500); a full batch is followed by
	// another read at once.
	BatchSize int
	// ListenerDSN is the connection string for LISTEN; empty means the one
	// the database was opened with. Tests give each replica its own.
	ListenerDSN string
}

LogOptions tune the event log. The zero value gives the defaults.

type LogStats

type LogStats struct {
	Enabled bool
	// Listening reports whether notifications wake the reader; without
	// them (the listener cannot connect) it reads every PollInterval.
	Listening          bool
	Cursor             int64     // highest id handled
	Delivered          uint64    // events handed to handlers
	LastRead           time.Time // last successful read
	ListenerReconnects uint64
	ReadErrors         uint64
	LastError          string
}

LogStats describe a log's progress, for status pages and tests.

type Node

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

Node is this replica's entry in the registry.

func StartNode

func StartNode(db *gorm.DB, id, version string, opts ...NodeOptions) (*Node, error)

StartNode registers the replica and keeps its row fresh until Stop. The first write happens before it returns, so a failure to reach the database is reported rather than discovered later. At most one NodeOptions is used; without one the replica is unlabelled and may lead.

func (*Node) ID

func (n *Node) ID() string

ID is the node's ID.

func (*Node) Stop

func (n *Node) Stop(ctx context.Context)

Stop ends the refreshes and removes the node's row, so other replicas know at once that it is gone rather than after the liveness window.

type NodeOptions

type NodeOptions struct {
	// Label names the replica on the status page and the Edge Gateways
	// page ("studio", "dashboard", "mdcb-eu-1"): at most MaxLabelLength
	// characters, letters, digits, spaces and . _ - : / ( ) only. Empty
	// leaves it unlabelled.
	Label string
	// NeverLeads records that the replica never takes the leader lease (a
	// headless control plane). It is recorded for operators; the replica's
	// own code decides whether it contends for the lease.
	NeverLeads bool
}

NodeOptions describe a replica to the others and to operators.

type NodeStatus

type NodeStatus struct {
	NodeID    string    `json:"node_id"`
	Hostname  string    `json:"hostname"`
	Version   string    `json:"version"`
	StartedAt time.Time `json:"started_at"`
	LastSeen  time.Time `json:"last_seen"`
	// Label names the replica for operators (empty when it has none).
	Label string `json:"label"`
	// LeaderEligible is false for a replica that never takes the leader
	// lease (a headless control plane).
	LeaderEligible bool `json:"leader_eligible"`
	// Edges is how many connected edges hold a stream to this replica.
	Edges int64 `json:"edges"`
	// Leader: this replica holds the leader lease.
	Leader bool `json:"leader"`
	// Self: the replica answering the request.
	Self bool `json:"self"`
}

NodeStatus is one live replica as the cluster status reports it.

type Owner

type Owner struct {
	Label string
	// Live: the replica's registration is fresh. An edge whose owner is not
	// live reconnects to another replica shortly.
	Live bool
}

Owner is the replica holding an edge's stream, as operators see it.

type PushBacklog

type PushBacklog struct {
	Pending  int64 `json:"pending"`   // waiting for the edge's stream
	InFlight int64 `json:"in_flight"` // claimed or sent, waiting for the edge
	// OldestPending is when the oldest pending command was created.
	OldestPending *time.Time `json:"oldest_pending,omitempty"`
	// Unowned counts in-flight commands whose replica is no longer live;
	// a janitor returns them to pending within seconds, so a lasting
	// non-zero value means no replica's janitor runs.
	Unowned int64 `json:"unowned"`
}

PushBacklog counts configuration push commands that are not finished.

type Relay

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

Relay carries bus events between control-plane replicas through the cluster event log. Events published on this replica's bus that the filter selects are written to the log; events other replicas wrote are published on this replica's bus with RelayedFrom set, which also stops them being relayed again. On a database without the log (SQLite, one process) it relays nothing.

Delivery follows the log: at least once to every live replica, in order per publishing replica. Bus events keep their ID, so consumers that persist events (webhooks) store each one once.

func NewRelay

func NewRelay(log *Log, bus eventbridge.Bus, opts RelayOptions) *Relay

NewRelay returns a relay between bus and log. It does nothing until Start.

func (*Relay) Start

func (r *Relay) Start()

Start begins relaying in both directions.

func (*Relay) Stats

func (r *Relay) Stats() RelayStats

Stats returns the relay's counters.

func (*Relay) Stop

func (r *Relay) Stop()

Stop ends relaying. Events still queued are written first, within the publish timeout each.

type RelayOptions

type RelayOptions struct {
	// Filter picks the events to relay (default RelayedByDefault).
	Filter func(eventbridge.Event) bool
	// QueueSize bounds the events waiting to be written to the log
	// (default 4096). When it is full, events are dropped and counted.
	QueueSize int
	// PublishTimeout bounds one write to the log (default 5 s); a failed
	// write is retried with backoff up to MaxRetries times (default 5).
	PublishTimeout time.Duration
	MaxRetries     int
}

RelayOptions tune a Relay. The zero value gives the defaults.

type RelayStats

type RelayStats struct {
	Sent     uint64 // events written to the log
	Received uint64 // events from other replicas republished locally
	Dropped  uint64 // events lost: queue full, or the log write kept failing
}

RelayStats describe a relay's traffic, for the cluster status page.

type Status

type Status struct {
	Self     string       `json:"self"`
	Database string       `json:"database"` // postgres | sqlite
	Nodes    []NodeStatus `json:"nodes"`
	// EdgesWithoutReplica counts connected edges whose owning replica is
	// not live: they reconnect elsewhere shortly.
	EdgesWithoutReplica int64 `json:"edges_without_replica"`

	Leader          string     `json:"leader"`
	LeaderExpiresAt *time.Time `json:"leader_expires_at,omitempty"`

	// EventLog and Relay are this replica's own.
	EventLog LogStats   `json:"event_log"`
	Relay    RelayStats `json:"relay"`

	Pushes PushBacklog `json:"pushes"`

	// Warnings explain anything that needs attention.
	Warnings []string `json:"warnings"`
}

Status is the control plane as one replica sees it.

func Snapshot

func Snapshot(ctx context.Context, db *gorm.DB, self string, log *Log, relay *Relay) (*Status, error)

Snapshot reports the cluster from self's point of view. log and relay may be nil.

Jump to

Keyboard shortcuts

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