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
- func ServerOptions(tlsCfg *tls.Config) []grpc.ServerOption
- type Capacity
- type Config
- type Coordinator
- func (c *Coordinator) Connect(stream workerv1.WorkerService_ConnectServer) error
- func (c *Coordinator) Enroll(ctx context.Context, req *workerv1.EnrollRequest) (*workerv1.EnrollResponse, error)
- func (c *Coordinator) GRPCService() workerv1.WorkerServiceServer
- func (c *Coordinator) Register(s *grpc.Server)
- func (c *Coordinator) Start(ctx context.Context, spec RunSpec) (*Run, error)
- func (c *Coordinator) Workers() []WorkerInfo
- type EventType
- type Result
- type Run
- func (r *Run) Done() <-chan struct{}
- func (r *Run) Events() <-chan RunEvent
- func (r *Run) ID() string
- func (r *Run) Interval() time.Duration
- func (r *Run) Kill()
- func (r *Run) Snapshots() <-chan *metrics.Snapshot
- func (r *Run) Stop()
- func (r *Run) StopWithReason(reason string)
- func (r *Run) T0() time.Time
- func (r *Run) Wait(ctx context.Context) (*Result, error)
- type RunEvent
- type RunSpec
- type Window
- type WorkerInfo
- type WorkerSummary
Constants ¶
const ( WorkerFinished = "finished" WorkerFailed = "failed" WorkerLost = "lost" )
Worker states in a Result.
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 ¶
func (c *Coordinator) Connect(stream workerv1.WorkerService_ConnectServer) error
Connect implements WorkerService.
func (*Coordinator) Enroll ¶
func (c *Coordinator) Enroll(ctx context.Context, req *workerv1.EnrollRequest) (*workerv1.EnrollResponse, error)
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 ¶
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 ¶
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) 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 ¶
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 ¶
StopWithReason is Stop with a reason recorded as the stop reason, for example "breakpoint reached".
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 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"`
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.