fleet

package
v0.2.8 Latest Latest
Warning

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

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

Documentation

Overview

Package fleet is the part of every CandaWS engine that is the same part.

The five services in this directory are five different shapes — a chain, a broker, a write quorum, a dispatcher and a fan-in — and they share almost nothing. What they do share is the seam between an engine and the host that watches it: a set of goroutines started against one context, a token that makes "run once" mean it, a fan-out that hands the latest view to whoever subscribed, and a per-goroutine random stream that a seed fixes.

Those four are here rather than copied five times because they are exactly the kind of thing that gets copied five times and then drifts: a fan-out that drops in four services and blocks in the fifth is a demo whose protocol is hostage to a browser, and nothing about the fifth service would say so.

Everything here is CSP-native, which is the same rule the engines follow. Feed owns its subscriber set in one goroutine and is reached by sending to it; Once is a token in a channel rather than a flag behind a lock; the one sync.WaitGroup in the package is Crew's, counting goroutines that have returned, which is a leaf counter and guards no protocol.

Index

Constants

This section is empty.

Variables

View Source
var ErrAlreadyRunning = errors.New("candaws: this engine is already running")

ErrAlreadyRunning is a second call to an engine's Run.

One engine runs once: its goroutines are started against the context Run was given, so a second call would start a second set against a second context and leave two of everything writing into one set of channels.

Functions

func Jitter

func Jitter(seed int64, index uint64) *rand.Rand

Jitter is one goroutine's own random stream, seeded from an engine's seed and that goroutine's index.

Per-goroutine rather than shared: a package-level source would make two engines in one test binary perturb each other, and a source behind a lock would be the one piece of shared mutable state in a package whose whole claim is that it has none. Equal seeds give equal sequences; they do not give an equal scheduler, so an engine is reproducible in its delays and not in its interleavings.

Types

type Crew

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

Crew is a set of goroutines that live and die with one context.

The sync.WaitGroup is the one in this package and it is a leaf counter: it counts goroutines that have returned and guards no protocol state. Everything those goroutines say to each other, they say on channels.

func (*Crew) Go

func (crew *Crew) Go(ctx context.Context, goroutine func(ctx context.Context))

Go starts one goroutine against the crew's context.

func (*Crew) Wait

func (crew *Crew) Wait()

Wait returns when every goroutine the crew started has returned.

An engine's Run blocks on it, so a caller that waits for Run knows the goroutines are actually gone rather than merely asked to go — which is what makes a specification's cleanup a real join instead of a hope.

type Feed

type Feed[V any] struct {
	// contains filtered or unexported fields
}

Feed is the fan-out at the end of an engine: one goroutine owning the subscriber set and the last view minted, reachable only by sending to it.

V is the engine's own view type and is a type parameter rather than an opaque one for the reason the widget SDK gives about its own state: nothing an author writes has a reason not to know it. A feed never inspects a view, so the constraint is the widest one and the erasure count stays zero.

func NewFeed

func NewFeed[V any](buffer int) *Feed[V]

NewFeed builds a feed and starts nothing. The goroutine belongs to Feed.Run, because it belongs to the context Run is given.

A buffer of zero is legal and means every subscriber sees only the view that arrives while it is reading, which is almost never what a caller wants; the engines here pass a small depth.

func (*Feed[V]) Publish

func (feed *Feed[V]) Publish(ctx context.Context, view V) bool

Publish hands one view to the feed, and reports false when the context ended before the feed took it.

It blocks rather than dropping, because the engine is the one producer and a view it minted and then discarded is a round nothing can observe. Backpressure on subscribers is handled the other way — see [Feed.offer].

func (*Feed[V]) Run

func (feed *Feed[V]) Run(ctx context.Context)

Run owns the subscriber set until the context ends, then closes every subscriber's channel and returns.

The set, the latest view and the "has there been one" flag are locals rather than fields, which is the whole argument: nothing outside this goroutine can name them, so there is nothing for a lock to protect and no way to forget one.

func (*Feed[V]) Subscribe

func (feed *Feed[V]) Subscribe(ctx context.Context) (<-chan V, error)

Subscribe returns a channel carrying every view minted from now on, plus the current one if there is one. The feed closes it when the feed stops.

type Once

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

Once is a start token that can be taken exactly once.

It is a channel holding one value rather than a boolean behind a lock, because a token that can be taken once is the whole of what "run once" means and a channel already is one.

func NewOnce

func NewOnce() Once

NewOnce mints an untaken token.

func (Once) Take

func (once Once) Take() bool

Take reports whether this caller is the one that got the token.

Jump to

Keyboard shortcuts

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