Coordination
Package: github.com/redelay/go-framework/runtime/coordination
Overview
The coordination layer provides transport-agnostic primitives for cross-container state synchronisation in multi-container Redelay deployments (Docker Swarm, Kubernetes, multi-VM).
It is a deliberately small surface — three concerns, nothing more:
| Concern | Primitive | Typical use |
|---|---|---|
| Pub/sub invalidation | Publish / Subscribe | "Flow X was re-published — every container drop your cached route." Fire-and-forget, at-most-once. |
| Fan-out streams | AppendStream / ReadStream | Live-mode metrics, run events, ephemeral telemetry. Bounded retention; replay within the window. |
| Watchable KV | KVSet / KVGet / KVDelete / KVWatch | Published routes, feature flags, per-tenant config. Small value count, high read rate, low write rate. |
This is not an event bus (use go-events for business events) and not a persistence store (Mongo/ClickHouse remain authoritative). Coordination is the bus for temporarily-synced state that every container must see within tens of milliseconds.
When to use which
| You need... | Use this |
|---|---|
To publish a domain event (user.created) | go-events EventBus |
| To persist a record | Mongo / ClickHouse |
| To invalidate a cache across every container | Coordinator.Publish or KVWatch |
| To fan out live metrics from one process to many | Coordinator.AppendStream + ReadStream |
| To share "the currently-published route for flow X" | Coordinator.KVSet + KVWatch |
The interface
package coordination
type Coordinator interface {
Publish(ctx context.Context, topic string, payload []byte) error
Subscribe(ctx context.Context, topic string, handler Handler) (Subscription, error)
AppendStream(ctx context.Context, stream string, payload []byte) (id string, err error)
ReadStream(ctx context.Context, stream string, from StreamPosition) (StreamReader, error)
KVSet(ctx context.Context, bucket, key string, value []byte) error
KVGet(ctx context.Context, bucket, key string) ([]byte, error)
KVDelete(ctx context.Context, bucket, key string) error
KVWatch(ctx context.Context, bucket, keyPattern string) (KVWatcher, error)
Ping(ctx context.Context) error
Close() error
}
StreamPosition has two canonical values plus any opaque ID returned by AppendStream:
coordination.StreamStart— replay from the oldest retained entrycoordination.StreamEnd— tail only new entriesStreamPosition(id)— resume strictly after the given entry
Errors are minimal:
coordination.ErrNotFound— returned byKVGetfor missing keyscoordination.ErrClosed— returned by operations on a closed Coordinator
Backends
Three interchangeable implementations, all conformance-tested against the same test suite:
| Backend | Package | Best for |
|---|---|---|
| noop (in-process) | go-framework/runtime/coordination/noop | Dev, tests, single-container deployments. Stdlib-only, zero external deps. |
| NATS | go-modules/coordination/nats | Default recommendation. Core NATS pub/sub + JetStream streams + JetStream KV in one system. Clean cluster story (3 nodes → HA). |
| Redis | go-modules/coordination/redis | When you already run Redis for cache/sessions. PUBLISH/PSUBSCRIBE + XADD/XREAD + hash-KV with pub/sub-backed watch. Redis Cluster–safe via hash-tag key layout. |
Every backend passes the same conformance suite so they're behaviourally interchangeable.
Cluster safety
- NATS — JetStream is cluster-native. Subjects flow across the cluster; KV buckets and streams replicate per
Replicassetting. - Redis Cluster — Each KV bucket is a single hash
coord:kv:{bucket}and its pub/sub notification channel iscoord:kvn:{bucket}. The{bucket}hash tag co-locates both keys on the same slot soTxPipelinewrites + notifications succeed withoutCROSSSLOTerrors. Streams are single-key, also slot-local.
Auto-wiring
Blank-import the module and pick a backend via the COORD env var:
import _ "github.com/redelay/go-modules/coordination/module"
COORD | Backend | Extra env vars |
|---|---|---|
noop (default) | In-process | — |
nats | NATS | NATS_URL (default nats://localhost:4222) |
redis | Redis | REDIS_URL (default localhost:6379) |
app.Bootstrap() injects the Coordinator into ModuleDeps.Coordinator. Explicit wiring via app.WithCoordinator(c) always wins over env-based auto-wiring.
import (
"github.com/redelay/go-framework/app"
coordnats "github.com/redelay/go-modules/coordination/nats"
)
coord, _ := coordnats.New(coordnats.Config{
URL: "nats://a.internal:4222,nats://b.internal:4222,nats://c.internal:4222",
KVReplicas: 3,
})
inst, _ := app.Bootstrap(app.WithCoordinator(coord))
Unknown COORD values and connection failures fall back to noop with a warning — coordination is never on the critical path, so a broker outage must not stall request handling.
Using it in a module
ModuleDeps.Coordinator is nil when no backend is wired, so always nil-check:
func (m *Module) Startup(ctx context.Context) error {
if m.deps.Coordinator == nil {
return nil // single-process mode — nothing to do
}
go m.watchRoutes(ctx)
return nil
}
func (m *Module) watchRoutes(ctx context.Context) {
w, err := m.deps.Coordinator.KVWatch(ctx, "published-routes", "*")
if err != nil {
m.logger.Warn("coord watch failed; falling back to polling", zap.Error(err))
return
}
defer w.Close()
for e := range w.Events() {
if e.Deleted {
m.cache.Delete(e.Key)
continue
}
m.cache.Set(e.Key, e.Value)
}
}
Built-in uses
Live-mode metrics fan-out (flowexec)
When a Coordinator is wired, flowexec's inproc.Sink forwards every flushed metrics.Frame to the coordination stream live:<flowID>. The Studio's Live-mode SSE handler reads from that stream, so runs executed on any container surface in the admin-api's canvas within one tick (~1 s).
See Running flows for the operational story. Internally:
- Emit —
coord.AppendStream(ctx, "live:"+flowID, frameJSON)viacoordFrameSinkwired intoinproc.Options.Downstream. - Read —
handleFlowLiveSSE(inflowexec/module/admin) callscoord.ReadStream(ctx, "live:"+flowID, StreamEnd)when a Coordinator is available; falls back to the local in-proc aggregator otherwise. - Filter — frames with no node activity, no edge traffic, and no
versionHashare dropped at the producer (frameHasSignal) so idle aggregators in containers that aren't running the flow don't pollute the stream.
See go-flowdsl → Cross-container live-metrics fan-out for the component-level breakdown and go-flowdsl → Module HTTP surface for the admin-only route table.
Future — published-route invalidation
flowexec will switch PublishedRoute lookups to KVWatch("routes", "*") so every container picks up POST /flows/:id/publish within ~20 ms instead of waiting for the Mongo poll interval. Same pattern applies to any module that caches "which version is live" per flow.
Configuration reference
NATS backend
| Variable | Description | Default |
|---|---|---|
NATS_URL | Server URL or comma-separated cluster URLs | nats://localhost:4222 |
coordnats.Config fields (for direct wiring):
| Field | Description | Default |
|---|---|---|
URL | NATS URL(s) | — |
StreamPrefix | Prefix for JetStream stream names | COORD_ |
SubjectPrefix | Prefix for JetStream subjects backing streams | coord.stream. |
StreamMaxAge | Retention age per stream | 24h |
StreamMaxBytes | Per-stream storage cap | 10 MiB |
KVReplicas | KV bucket replica count | 1 (raise to 3 in prod) |
Requires JetStream — start NATS with -js.
Redis backend
| Variable | Description | Default |
|---|---|---|
REDIS_URL | host:port (or redis://host:port) | localhost:6379 |
coordredis.Config fields:
| Field | Description | Default |
|---|---|---|
Addr / Password / DB | Connection settings | localhost:6379 / — / 0 |
Client | Pre-built *goredis.Client (overrides Addr/...) | — |
StreamMaxLen | XADD MAXLEN ~ cap per stream | 10_000 |
StreamBlockDuration | XREAD BLOCK duration | 200 ms |
StreamBatchSize | XREAD COUNT | 64 |
KeyPrefix | Prefix for every key this coordinator creates | coord: |
noop backend
| Field | Description | Default |
|---|---|---|
streamMaxLen (via NewWithLimit) | Per-stream ring buffer size | 10_000 |
Conformance suite
Every backend must pass the shared suite in runtime/coordination/conformance. Twelve cases cover:
- Pub/sub — exact topic, multiple subscribers, unsubscribe cleanup
- Streams — append + tail, replay from start, resume from an ID
- KV — round-trip set/get/delete, missing-key error, watch snapshot + updates, watch deletes
- Lifecycle — ping, idempotent close
Run from a backend's own test file:
func TestMyBackend(t *testing.T) {
conformance.Run(t, func(t *testing.T) coordination.Coordinator {
c, err := mybackend.New(...)
if err != nil {
t.Fatalf("new: %v", err)
}
t.Cleanup(func() { _ = c.Close() })
return c
})
}
Build and test
# Unit tests (no brokers; noop only)
cd go-framework && go test -race ./runtime/coordination/...
cd go-modules && go test -race ./coordination/...
# Integration tests (Docker required — spins up real NATS + Redis)
cd go-modules && go test -tags integration -timeout 120s ./coordination/integration/...
The integration suite runs the full conformance against testcontainer-hosted NATS 2 (JetStream on) and Redis 7 containers, plus two end-to-end scenarios:
TestNATSE2E_FlowPublishInvalidation— two Coordinator instances (two simulated containers) share a KV bucket; container B receives container A'sKVSetviaKVWatch.TestRedisE2E_LiveMetricsFanout— producer container writes stream entries; consumer container (simulating admin-api) reads every one.
Deployment shape
[ api × N ] [ admin-api × M (Tailscale-gated) ]
\ /
\ /
[ NATS cluster 3+ nodes ] ← COORD backbone
|
[ Mongo replica set ] ← authoritative store
[ Kafka / NATS / Redis ] ← event transport (go-events)
apicontainers run flows and publish frames to the coord stream.admin-apicontainers read from the stream to serve Studio Live mode.- No direct network path required between
apiandadmin-api— they share state through the coord broker only, so admin-api can live on a private network (Tailscale, WireGuard, VPC peering) without opening any surface to the api tier.
Dependency graph
go-framework/runtime/coordination ← interface + noop (stdlib only)
↑
go-modules/coordination/{nats,redis} ← backend impls
↑
go-modules/coordination/module ← COORD env switch
↑
cmd/{api,admin-api} main.go ← blank import
The interface has no external dependencies. Blank-importing go-modules/coordination/module is the single line a host needs to opt in.
Assistant (chat + handoff)
Project-local AI assistant module — FlowDSL-backed chat, persisted conversations with TTL, LLM cost tracking, and a human-handoff path that fans out via the event bus.
Search (vector)
Pluggable vector-search module — Qdrant today, OpenSearch planned. Index CRUD, semantic query, embedding-provider switching, FlowDSL nodes.