coordinator

package
v1.2.0 Latest Latest
Warning

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

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

Documentation

Overview

Package coordinator is the server side of distributed execution. It registers workers that dial in over gRPC, splits a run across them, starts them at one synchronised instant and merges their per-interval snapshots into a single stream.

Registration

A worker opens one bidirectional Connect stream and sends Hello with the server's join token (compared in constant time) and its protocol version; an incompatible major version is refused with an explanation. The server answers with Welcome, an assigned id and the heartbeat interval. Three missed heartbeats mark a worker lost.

Sharding

Load is split into contiguous slices [lo, hi) of [0, 1) in proportion to each worker's CPUs, optionally first by region. The engine assigns every open-model arrival and every closed-model user to exactly one slice and partitions unique test data by worker index, so the workers together run exactly the single-engine plan.

Clock synchronisation

Before a run the server exchanges several ClockPing/ClockPong round trips with each worker and keeps the offset from the one with the lowest round trip (NTP style). T0 is sent in each worker's own clock, so a worker whose clock is off by seconds still starts within a fraction of the round trip of everyone else.

Worker loss

A lost worker's data up to the loss is kept and the loss is reported as an event and as a window in the result. Its share is deliberately not moved to other workers: the engine fixes a worker's slice of the arrival schedule at T0, and restarting that slice elsewhere mid-run would either replay the missed arrivals as a burst or need the schedule to resume at an offset. The run continues degraded instead and the report says, for which window, how much of the planned load was missing. That keeps the numbers honest, which matters more than reaching the planned rate.

Transport security

Three modes, from strongest:

  • Mutual TLS with the built-in CA (Config.CA, stampede server --worker-mtls). Workers enroll once per certificate lifetime with the Enroll RPC, proving they know the join token without sending it, and then connect with their own certificate. See internal/pki.
  • TLS with a server certificate from your own CA (ServerOptions); workers verify it with the system roots or a CA they are given, and send the join token in Hello.
  • No TLS, for trusted networks only: the join token travels in clear.

Index

Constants

View Source
const (
	WorkerFinished = "finished"
	WorkerFailed   = "failed"
	WorkerLost     = "lost"
)

Worker states in a Result.

View Source
const DefaultCertValidity = 24 * time.Hour

DefaultCertValidity is how long an enrolled worker's certificate lasts. Workers renew at half-life, so rotating the join token locks every worker out within this time.

Variables

This section is empty.

Functions

func ServerOptions

func ServerOptions(tlsCfg *tls.Config) []grpc.ServerOption

ServerOptions returns the gRPC server options the worker protocol expects: keepalives that notice dead connections within seconds, and TLS when tlsCfg is not nil.

Types

type Capacity

type Capacity struct {
	CPUs        int      `json:"cpus"`
	MemoryBytes uint64   `json:"memoryBytes"`
	MaxVUs      int      `json:"maxVUs"`
	Protocols   []string `json:"protocols"`
}

Capacity is what a worker advertised.

type Config

type Config struct {
	// JoinToken is the secret workers present in Hello. Workers are
	// refused while it is empty.
	JoinToken string
	// HeartbeatInterval is how often workers send heartbeats (default
	// 1s). Three missed heartbeats mark a worker lost.
	HeartbeatInterval time.Duration
	Logger            *slog.Logger

	// ClockSamples is the number of ping/pong round trips per clock
	// synchronisation (default 8).
	ClockSamples int
	// CA, when set, turns on mutual TLS: workers enroll for a certificate
	// with Enroll and must present it on Connect, where it replaces the
	// join token as their credential. The gRPC server must use
	// CA.ServerTLS.
	CA *pki.CA
	// CertValidity is the lifetime of enrolled worker certificates
	// (default DefaultCertValidity).
	CertValidity time.Duration

	// MergeGrace is how long after an interval ends the merged snapshot
	// waits for slow workers before it is emitted without them (default
	// two snapshot intervals).
	MergeGrace time.Duration

	// TakeoverLead is how far ahead a spare worker that takes over a lost
	// worker's share is told to start (default 2s): time to synchronise
	// its clock and prepare the scenario.
	TakeoverLead time.Duration
	// NoTakeover keeps lost shares unassigned: the run continues degraded.
	NoTakeover bool
}

Config configures a coordinator.

type Coordinator

type Coordinator struct {
	workerv1.UnimplementedWorkerServiceServer
	// contains filtered or unexported fields
}

Coordinator registers workers and runs distributed tests on them.

func New

func New(cfg Config) *Coordinator

New returns a coordinator. Register it on a gRPC server to accept workers.

func (*Coordinator) Connect

Connect implements WorkerService.

func (*Coordinator) Enroll

Enroll implements WorkerService.

func (*Coordinator) GRPCService

func (c *Coordinator) GRPCService() workerv1.WorkerServiceServer

GRPCService returns the WorkerService implementation, for callers that register services themselves.

func (*Coordinator) Register

func (c *Coordinator) Register(s *grpc.Server)

Register adds the WorkerService to s.

func (*Coordinator) Start

func (c *Coordinator) Start(ctx context.Context, spec RunSpec) (*Run, error)

Start splits spec across idle workers, synchronises their clocks and starts them at a common T0. It returns once every worker has accepted its share, or an error (and no load) when there is not enough capacity or a worker refuses. ctx bounds only this setup; the run itself ends when the plan completes or on Stop or Kill.

func (*Coordinator) Workers

func (c *Coordinator) Workers() []WorkerInfo

Workers lists registered workers, sorted by id. A worker that dropped its connection during a run stays listed, disconnected, until the run ends.

type EventType

type EventType string

EventType names a run event.

const (
	// EventWorkerJoined: the worker accepted its share of the run.
	EventWorkerJoined EventType = "worker-joined"
	// EventWorkerStarted: the worker's first snapshot arrived.
	EventWorkerStarted EventType = "worker-started"
	// EventWorkerSaturated: the worker reported itself saturated for an
	// interval; Reasons says why.
	EventWorkerSaturated EventType = "worker-saturated"
	// EventWorkerLost: three heartbeats missed (or the worker never
	// finished after being stopped). Its data up to the loss is kept.
	EventWorkerLost EventType = "worker-lost"
	// EventDegraded follows a loss: the lost share is not reassigned and
	// the run continues with the remaining workers.
	EventDegraded EventType = "degraded"
	// EventTakeover follows a loss when a spare worker takes over the lost
	// share from a later interval; until then the run is short of it.
	EventTakeover EventType = "takeover"
	// EventWorkerFinished: the worker's part of the run ended.
	EventWorkerFinished EventType = "worker-finished"
	// EventWorkerFailed: the worker reported an error and stopped.
	EventWorkerFailed EventType = "worker-failed"
	// EventIntervalIncomplete: an interval was emitted without some
	// workers' data (listed in Workers); late data still counts in the
	// final result.
	EventIntervalIncomplete EventType = "interval-incomplete"
)

Run events.

type Result

type Result struct {
	RunID      string
	T0         time.Time
	End        time.Time
	Interval   time.Duration
	StopReason string
	// Phases are run-wide per-step phase histograms merged across workers.
	Phases map[int]*[metrics.NumPhases]*metrics.Histogram
	// PeakVUs sums the workers' peaks: exact for closed-model plans,
	// whose shares peak together.
	PeakVUs int
	// Snapshots are the merged intervals in order, including data that
	// arrived too late for the live stream.
	Snapshots []*metrics.Snapshot
	Workers   []WorkerSummary
	// Degraded is true when a worker was lost or failed mid-run.
	Degraded bool
	// LateSnapshots counts worker snapshots that missed the live stream.
	LateSnapshots uint64
}

Result is the outcome of a distributed run.

type Run

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

Run is a distributed run in progress.

func (*Run) Done

func (r *Run) Done() <-chan struct{}

Done is closed when the run has ended and its Result is ready.

func (*Run) Events

func (r *Run) Events() <-chan RunEvent

Events delivers run events and is closed when the run ends. Events that do not fit the buffer are dropped; the Result records the same facts.

func (*Run) ID

func (r *Run) ID() string

ID is the run id.

func (*Run) Interval

func (r *Run) Interval() time.Duration

Interval is the snapshot interval.

func (*Run) Kill

func (r *Run) Kill()

Kill stops all load immediately: the stop goes straight onto every worker's stream, and workers abandon in-flight requests.

func (*Run) Snapshots

func (r *Run) Snapshots() <-chan *metrics.Snapshot

Snapshots delivers one merged snapshot per interval across workers, in order, and is closed when the run ends. Consumers should keep up; if the buffer (sized for the whole plan) ever fills, live snapshots are dropped, and Result.Snapshots still holds every interval.

func (*Run) Stop

func (r *Run) Stop()

Stop ends the load phase gracefully: in-flight iterations get the scenario's graceful stop period.

func (*Run) StopWithReason

func (r *Run) StopWithReason(reason string)

StopWithReason is Stop with a reason recorded as the stop reason, for example "breakpoint reached".

func (*Run) T0

func (r *Run) T0() time.Time

T0 is the synchronised start time in the server's clock.

func (*Run) Wait

func (r *Run) Wait(ctx context.Context) (*Result, error)

Wait blocks until the run ends and returns its result. The error is set when no worker finished its share (every one failed or was lost); the result is still returned with whatever data arrived.

type RunEvent

type RunEvent struct {
	Time       time.Time `json:"time"`
	Type       EventType `json:"type"`
	Worker     string    `json:"worker,omitempty"`
	WorkerName string    `json:"workerName,omitempty"`
	Interval   int64     `json:"interval"`
	Message    string    `json:"message"`
	Reasons    []string  `json:"reasons,omitempty"`
	Workers    []string  `json:"workers,omitempty"`
}

RunEvent is something that happened during a run.

type RunSpec

type RunSpec struct {
	ID       string
	Scenario []byte // scenario YAML
	Env      map[string]string
	Secrets  map[string]string
	// AllowHosts lists the public hosts workers may send requests to,
	// normally the verified target host. Private hosts are always allowed.
	AllowHosts []string
	// Workers is how many workers to use (0 = every idle connected one).
	Workers int
	// Regions optionally splits load by worker region, for example
	// {"mumbai": 0.5, "frankfurt": 0.5}. Fractions are normalised.
	Regions map[string]float64
	// StartDelay is the lead time before the synchronised T0 (default 3s).
	StartDelay time.Duration
}

RunSpec describes a distributed run.

type Window

type Window struct {
	From int64 `json:"from"`
	To   int64 `json:"to"`
}

Window is an inclusive range of intervals.

type WorkerInfo

type WorkerInfo struct {
	ID             string            `json:"id"`
	Name           string            `json:"name"`
	Version        string            `json:"version"`
	Region         string            `json:"region,omitempty"`
	Labels         map[string]string `json:"labels,omitempty"`
	Capacity       Capacity          `json:"capacity"`
	Connected      bool              `json:"connected"`
	ConnectedSince time.Time         `json:"connectedSince"`
	LastHeartbeat  time.Time         `json:"lastHeartbeat"`
	Health         wire.Health       `json:"health"`
	// ClockOffset is the worker's clock minus the server's, from the
	// latest synchronisation.
	ClockOffset time.Duration `json:"clockOffset"`
	CurrentRun  string        `json:"currentRun,omitempty"`
}

WorkerInfo describes a registered worker.

type WorkerSummary

type WorkerSummary struct {
	ID          string        `json:"id"`
	Name        string        `json:"name"`
	Region      string        `json:"region,omitempty"`
	Index       int           `json:"index"`
	ShareLo     float64       `json:"shareLo"`
	ShareHi     float64       `json:"shareHi"`
	ClockOffset time.Duration `json:"clockOffset"`
	ClockRTT    time.Duration `json:"clockRTT"`
	State       string        `json:"state"`
	StopReason  string        `json:"stopReason,omitempty"`
	Error       string        `json:"error,omitempty"`
	PeakVUs     int           `json:"peakVUs"`
	Requests    uint64        `json:"requests"`
	// Iterations counts iterations started.
	Iterations uint64 `json:"iterations"`
	Snapshots  int    `json:"snapshots"`
	// Saturated lists the intervals in which the worker reported itself
	// saturated, and SaturationReasons the distinct reasons.
	Saturated         []Window `json:"saturated,omitempty"`
	SaturationReasons []string `json:"saturationReasons,omitempty"`
	// Lost is the window from the loss until the share was taken over, or
	// to the end of the run.
	Lost *Window `json:"lost,omitempty"`
	// Replaces names the lost worker whose share this one took over, from
	// interval JoinedAt.
	Replaces string `json:"replaces,omitempty"`
	JoinedAt int64  `json:"joinedAt,omitempty"`
	// Missing lists intervals emitted live without this worker's data.
	Missing []int64 `json:"missing,omitempty"`
}

WorkerSummary is one worker's part in a run.

Jump to

Keyboard shortcuts

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