Reference

Coordination

Cross-container coordination layer — pub/sub invalidation, fan-out streams, and watchable KV for multi-container Redelay deployments.

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:

ConcernPrimitiveTypical use
Pub/sub invalidationPublish / Subscribe"Flow X was re-published — every container drop your cached route." Fire-and-forget, at-most-once.
Fan-out streamsAppendStream / ReadStreamLive-mode metrics, run events, ephemeral telemetry. Bounded retention; replay within the window.
Watchable KVKVSet / KVGet / KVDelete / KVWatchPublished 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 recordMongo / ClickHouse
To invalidate a cache across every containerCoordinator.Publish or KVWatch
To fan out live metrics from one process to manyCoordinator.AppendStream + ReadStream
To share "the currently-published route for flow X"Coordinator.KVSet + KVWatch

The interface

go
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 entry
  • coordination.StreamEnd — tail only new entries
  • StreamPosition(id) — resume strictly after the given entry

Errors are minimal:

  • coordination.ErrNotFound — returned by KVGet for missing keys
  • coordination.ErrClosed — returned by operations on a closed Coordinator

Backends

Three interchangeable implementations, all conformance-tested against the same test suite:

BackendPackageBest for
noop (in-process)go-framework/runtime/coordination/noopDev, tests, single-container deployments. Stdlib-only, zero external deps.
NATSgo-modules/coordination/natsDefault recommendation. Core NATS pub/sub + JetStream streams + JetStream KV in one system. Clean cluster story (3 nodes → HA).
Redisgo-modules/coordination/redisWhen 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 Replicas setting.
  • Redis Cluster — Each KV bucket is a single hash coord:kv:{bucket} and its pub/sub notification channel is coord:kvn:{bucket}. The {bucket} hash tag co-locates both keys on the same slot so TxPipeline writes + notifications succeed without CROSSSLOT errors. Streams are single-key, also slot-local.

Auto-wiring

Blank-import the module and pick a backend via the COORD env var:

go
import _ "github.com/redelay/go-modules/coordination/module"
COORDBackendExtra env vars
noop (default)In-process—
natsNATSNATS_URL (default nats://localhost:4222)
redisRedisREDIS_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.

go
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:

go
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) via coordFrameSink wired into inproc.Options.Downstream.
  • Read — handleFlowLiveSSE (in flowexec/module/admin) calls coord.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 versionHash are 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

VariableDescriptionDefault
NATS_URLServer URL or comma-separated cluster URLsnats://localhost:4222

coordnats.Config fields (for direct wiring):

FieldDescriptionDefault
URLNATS URL(s)—
StreamPrefixPrefix for JetStream stream namesCOORD_
SubjectPrefixPrefix for JetStream subjects backing streamscoord.stream.
StreamMaxAgeRetention age per stream24h
StreamMaxBytesPer-stream storage cap10 MiB
KVReplicasKV bucket replica count1 (raise to 3 in prod)

Requires JetStream — start NATS with -js.

Redis backend

VariableDescriptionDefault
REDIS_URLhost:port (or redis://host:port)localhost:6379

coordredis.Config fields:

FieldDescriptionDefault
Addr / Password / DBConnection settingslocalhost:6379 / — / 0
ClientPre-built *goredis.Client (overrides Addr/...)—
StreamMaxLenXADD MAXLEN ~ cap per stream10_000
StreamBlockDurationXREAD BLOCK duration200 ms
StreamBatchSizeXREAD COUNT64
KeyPrefixPrefix for every key this coordinator createscoord:

noop backend

FieldDescriptionDefault
streamMaxLen (via NewWithLimit)Per-stream ring buffer size10_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:

go
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

shell
# 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's KVSet via KVWatch.
  • TestRedisE2E_LiveMetricsFanout — producer container writes stream entries; consumer container (simulating admin-api) reads every one.

Deployment shape

text
[ 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)
  • api containers run flows and publish frames to the coord stream.
  • admin-api containers read from the stream to serve Studio Live mode.
  • No direct network path required between api and admin-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

text
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.