Reference

Domain Addon Modules

Domain addon modules: scheduler, settings, storage, notifications, verification, links, metrics, search, i18n, apilog, content, workers, coordination, websockets, and flowexec.

Domain Addon Modules

github.com/redelay/go-modules provides domain-level feature modules that extend the Redelay framework. Each module is added via a blank import and works with zero configuration defaults.

Looking for raw infrastructure nodes (Redis/Mongo/Kafka/HTTP/file/SMS/FTP/SFTP/Postgres/ClickHouse/RabbitMQ) to compose directly in FlowDSL flows? See the Infrastructure Modules reference and the Node Catalog.

Quick start

shell
go get github.com/redelay/go-modules
go
import (
    _ "github.com/redelay/go-modules/scheduler"
    _ "github.com/redelay/go-modules/metrics"
    // add more as needed
)

Module overview

ModulePackagePurpose
schedulergo-modules/schedulerDistributed cron with leader election
settingsgo-framework/modules/settingsKey-value settings API
storagego-modules/storageObject storage (S3/MinIO) with S3 list / head / copy / delete-many nodes
notificationsgo-modules/notificationsReal-time SSE notifications
verificationgo-modules/verificationEmail address verification
linksgo-modules/linksURL shortener
metricsgo-modules/metricsPrometheus metrics
flowexecgo-flowdsl/flowexec/moduleFlow CRUD, versioning, run launch, live SSE events
searchgo-modules/searchPluggable vector search (Qdrant today, OpenSearch planned) with index CRUD + FlowDSL upsert/query nodes
i18ngo-modules/i18nLocale-resolution middleware, Mongo message catalog with fallback chain, admin CRUD, translate/resolve-locale nodes. Also drives customer-app UI-label i18n: write-once POST /i18n/report-keys (frontend auto-registers the $t keys it uses; I18N_AUTOREGISTER), admin POST /admin/i18n/translations/translate-missing (bulk AI-translate every active locale via the ai module), and a per-key Translation.Context hint. Reusable frontend $t plumbing — catalog + ICU plurals via Intl.PluralRules + interpolation + build-time extractor — ships in redelay-starter-kit/templates/nuxt-i18n. The admin Translations catalog adds a Default-text column, per-row Re-translate, and per-key context; localized admin fields (LocalizedField) take a per-field context.
flowtemplatesgo-modules/flowtemplatesCurated starter flows surfaced via flowexec.TemplateProvider
coordinationgo-modules/coordinationCross-container Coordinator — pub/sub, streams, KV (NATS or Redis backend)
websocketsgo-modules/wsWebSocket hub — rooms + cross-replica fan-out, health/info/broadcast control plane
securitygo-modules/securityPer-user security — activity log, trusted devices, rate limits, auto-restriction; security/admin for IP blocks + user restriction
eventadmingo-modules/eventadminEvent-bus admin (admin-api only) — event definitions, manual publish, event log, webhook subscriptions
metrics/admingo-modules/metrics/adminDatabase observability (admin-api only) — live ClickHouse/Mongo/Redis stats
apiloggo-modules/apilogTime-series log of outbound external API calls, queryable in admin; Mask() redaction + apilog/record debug node
aigo-modules/aiProvider-agnostic, LLM-assisted content authoring on go-ai (stateless — no DB). ai/admin (admin-api only): POST /admin/ai/translate ({text,from,to,context?}) + POST /admin/ai/generate. Model selection via the ai.models settings group — profile (a named llm-chat profile: provider+model+creds+tuning, overrides the flat fields) / provider (options_ref llm.providers) / model / temperature — admin-editable, no restart; per-consumer profiles let each feature run its own model. translate validates format (placeholders / ICU-plural structure / HTML tags preserved — corrupt locales re-requested then omitted) and translates ICU-plural sources one locale per call; context disambiguates ambiguous terms.
contentgo-modules/contentCMS-lite scope-scoped localized pages (/pages/{slug}) + reusable fragments (/fragments/{key}). Prose may embed [[ module.key ]] settings tokens, resolved per request against the PUBLIC settings gate. Fragment leaves are localized with an explicit _i18n marker; both pages and fragments carry crud.TranslationState (ready / stale / retranslate). See Localized content.
workersgo-modules/workersCross-instance worker registry with heartbeat tracking

audit (go-modules/audit) is a library package, not a module — it has no factory and is not blank-imported. Entities embed History []audit.Event and append via audit.Push.

FlowDSL node subpackages

Each module (except metrics and flowexec) ships a flowdsl/ subpackage that registers FlowDSL nodes as a companion module. Enable them with a blank import:

go
import (
    _ "github.com/redelay/go-modules/scheduler"
    _ "github.com/redelay/go-modules/scheduler/flowdsl"   // 4 FlowDSL nodes
)
SubpackageFactoryNodesKinds
scheduler/flowdslscheduler-flowdsl42 source (with settings_schema), 2 action
settings/flowdslsettings-flowdsl31 transform (with NotFound), 2 action
storage/flowdslstorage-flowdsl84 upload/download/delete/presign + 4 S3 list / head / copy / delete-many
notifications/flowdslnotifications-flowdsl32 action, 1 source
verification/flowdslverification-flowdsl21 action, 1 router (Valid/Invalid)
links/flowdsllinks-flowdsl32 action, 1 transform (with NotFound)
search/flowdslsearch-flowdsl22 action (search-upsert, search-query)

For the full node catalog — including core, events, redisops, mongoops, and every infrastructure module — browse the Node Catalog.


scheduler

Distributed cron job scheduler with MongoDB persistence and Redis-based leader election.

Only the elected leader runs jobs at any given time — all replicas coordinate automatically.

A schedule fires an event — event_type + event_action on a cron expression or once at a fixed time. The event goes onto the bus, where any consumer or FlowDSL flow can react to it. The scheduler is a timer, not an executor: it decides when, the event bus decides what happens next.

Environment variables

VariableDefaultDescription
SCHEDULER_COLLECTIONscheduler_schedulesMongoDB collection for schedule definitions
SCHEDULER_LOG_COLLECTIONscheduler_logsMongoDB collection for execution logs
SCHEDULER_POLL_INTERVAL_SECONDS10How often to check for due schedules
SCHEDULER_LOG_TTL_DAYS30How long execution history is kept. 0 keeps it forever
SCHEDULER_LOCK_TTL_SECONDS30Redis leader-election lock TTL
SCHEDULER_LOCK_KEYredelay:<app-name>:scheduler:leaderRedis key for the leader lock. Namespaced by APP_NAME

Leader election — read this before sharing a Redis

Every replica polls; the Redis lock decides which one dispatches. Three things about that lock exist because each was got wrong once, and the failure they produce is the same every time: schedules fire on time, every execution logs success, and nothing is delivered.

The lock key is namespaced per app. It used to be the constant redelay:scheduler:leader in every deployment, so any process reaching the same Redis competed for one lock — a second app, a staging environment pointed at the same instance, or a developer's go run left over from yesterday. Whichever won dispatched everyone's schedules, into whatever transport it happened to have. A 27-hour-old local build holding no Kafka connection at all once owned a deployment's cron for hours this way.

The lock records its holder (host/pid-N), not 1. A lock that says nothing about its owner cannot answer "which process is running my cron", and answering it otherwise means stopping containers one at a time.

The lock is only renewed while still owned. A blind EXPIRE extends whatever is there, so a replica that lost the lock would keep refreshing the new leader's claim — two dispatchers, and a lock that never changes hands again.

Freezing schedules (before a migration, during an incident)

POST /schedules/{id}/pause and /resume set the state explicitly and are idempotent — 200 when it is already there. Use these to script a freeze:

shell
# Pause everything. Safe to re-run.
for id in $(curl -s -H "Authorization: Bearer $TOKEN" \
      "$API/api/v1/schedules/?limit=200" | jq -r '.items[].id'); do
  curl -s -X POST -H "Authorization: Bearer $TOKEN" \
      "$API/api/v1/schedules/$id/pause" > /dev/null
done

/toggle exists too but flips, which is right for a switch in a UI and wrong for a script: re-running a toggle-based freeze turns everything back on, and two operators running it at once can leave the fleet in either state.

resume deliberately leaves a completed one-time schedule completed — re-arming it would fire it a second time for reasons nobody chose.

An update leaves status alone unless enabled is sent, so editing a cron expression or a description during a freeze does not lift it.

Observability

GET /health includes a scheduler check. It is deliberately not unhealthy for a follower — that is the normal state for every replica but one — but it does fail when the lock is held by an owner that does not look like a peer, which is the stale-build case above and is invisible everywhere else.

Each execution log row records host and topic. topic matters because a schedule whose event_type resolves to nothing still logs success; without it there is no way to tell a real publish from a no-op. Every dispatch is also logged at INFO on success as well as failure — a scheduler whose whole job is to publish used to be silent when it worked, so a deployment where every run published into a void looked exactly like a healthy one.

If scheduled events stop arriving, check in this order:

  1. GET /health → the scheduler check names a foreign lock holder.
  2. redis-cli get redelay:<app>:scheduler:leader → who thinks it is leader.
  3. The host field on recent scheduler_logs rows → who actually dispatched.
  4. A duration_ms of 0 on every row is a strong hint that publishing is a no-op: a real broker round-trip is not free.

HTTP surface — scheduler/admin (admin-api only)

The management API is a separate submodule, go-modules/scheduler/admin, blank-imported from the admin-api binary only. The core scheduler module runs the engine and mounts nothing; importing it on the public api gives you reliable schedule firing without exposing schedule management there.

go
// cmd/admin-api only
import _ "github.com/redelay/go-modules/scheduler/admin"

The routes sit at /schedules/ (matching the API contract, not under the admin prefix) and are gated on RequirePermission("admin:access"). A schedule publishes an arbitrary event as an arbitrary actor on a timer, so managing them is an admin capability — the permission check and the admin-api-only import keep the surface off the public api.

MethodPathPurpose
GET/schedules/List schedules (filter by enabled, tag, event_type, event_action)
POST/schedules/Create a schedule
GET/schedules/{schedule_id}Get a schedule
PUT/schedules/{schedule_id}Update (partial; changing the timing re-arms the next run)
DELETE/schedules/{schedule_id}Delete the schedule and its execution history
POST/schedules/{schedule_id}/executeFire now, without shifting the regular timing
GET/schedules/{schedule_id}/executionsList execution log entries (newest first)
POST/schedules/{schedule_id}/toggleEnable / pause

A schedule is either recurring (a 5-field cron_expression, read in its own timezone) or one-time (one_time: true + an execution_time); a one-time schedule moves to a terminal completed state after it fires. A failed publish retries up to max_retries with retry_delay seconds between attempts before waiting for the next occurrence. Six-field cron expressions are rejected — the parser does not enable seconds, so a sixth field would silently shift every position.

shell
curl -X POST /api/v1/schedules/ -H 'Authorization: Bearer <admin>' -d '{
  "name": "nightly-report",
  "event_type": "report", "event_action": "generate",
  "cron_expression": "0 3 * * *", "timezone": "Europe/Warsaw",
  "event_payload": {"kind": "daily"}
}'

FlowDSL alternative

The same effect is available declaratively — a flow with a cron or interval source node, resolved natively by the flow engine:

yaml
nodes:
  - id: nightly
    nodeType: redelay/scheduler-cron-trigger
    settings:
      cron: "0 2 * * *"
  - id: cleanup
    nodeType: mongoops/delete

SchedulesProvider (Schedules() []*ir.Schedule) is metadata only — it populates the module browser's IR and is not wired to the scheduler service.


settings

Key-value settings store backed by MongoDB with in-memory caching and cluster-wide invalidation.

Implements modules.SettingsReader — other modules resolve settings via dependency injection with automatic 3-layer precedence: DB override → environment variable → YAML default. Updates propagate to all workers instantly via the event bus without a restart.

Data model

go
type Setting struct {
    ModuleID  string    // e.g. "ai-llm"
    Key       string    // e.g. "openai_api_key"
    Value     any       // JSON-serializable
    UpdatedAt time.Time
}

Unique index on (module_id, key). Default collection: _settings.

3-layer resolution precedence

PrioritySourceNotes
1 (highest)MongoDB (admin-saved)Written via admin API
2Environment variableExplicit EnvVar field or implicit {MODULE}_{KEY}
3YAML defaultDeclared in module's SettingsDefinition

Implicit env var naming: module ai-llm + key openai_api_key → AI_LLM_OPENAI_API_KEY (hyphens → underscores, uppercased).

Caching and propagation

Settings are cached in memory per worker behind a sync.RWMutex. The cache is designed for correctness across three deployment shapes:

  1. Single binary (one cmd/api with admin imports) — local Service.Set invalidates the local cache; events are unnecessary.
  2. Multi-binary local dev (cmd/api + cmd/admin-api with TRANSPORT=memory) — each binary runs an isolated in-process event bus, so a settings.updated event published by admin-api never reaches cmd/api. The TTL (below) is the only invalidation mechanism in this mode.
  3. Production cluster (N replicas with TRANSPORT=kafka/nats/redis) — every replica subscribes with a unique settings-invalidator-{hostname}-{8-char-uuid} GroupID, so events fan out to all replicas in milliseconds. The TTL acts as a safety net for missed events (consumer restart mid-event, network blip, hostname collision).

Cache rules

Resolution sourceCached?Invalidated by
db — admin-edited row in MongoDBYes, with TTLlocal Service.Set, settings.updated event, TTL expiry
env — environment variableYes, with TTLTTL expiry (env doesn't change at runtime)
default — YAML-declared defaultNever cachedn/a — re-resolved on every miss

Why defaults are not cached: a default returned before any DB row exists would shadow a future admin write whose invalidation event a peer binary failed to receive. Skipping the cache for defaults guarantees the next read after a peer save sees the new DB value immediately.

Save flow

  1. The value is written to MongoDB.
  2. The local cache entry is invalidated on the writing worker.
  3. A settings.updated event is published to the event bus.
  4. Every worker subscribed with a unique GroupID receives the event (fan-out, not round-robin) and invalidates its local cache entry.
  5. On next read, the cache miss triggers a MongoDB re-query.

Cache TTL

VariableDefaultNotes
SETTINGS_CACHE_TTL_SECONDS30Maximum staleness for a DB-sourced cache entry. Set to 0 to disable expiry — only safe in single-binary deployments OR a real-transport cluster where every replica has a unique hostname.

Tuning guidance:

  • Single-binary dev: SETTINGS_CACHE_TTL_SECONDS=0 is fine — local invalidation is authoritative.
  • Multi-binary dev (TRANSPORT=memory, separate cmd/api + cmd/admin-api): leave at default 30. Admin saves become visible to the public binary within 30 seconds even though no event crosses binaries.
  • Production with kafka/nats/redis: leave at default 30 (or raise to 300 for less DB load). Events handle ~ms invalidation; the TTL is just a backstop for missed events.

Access control: public and sensitive flags

Settings default to admin-only. The public API exposes a setting only when its declared schema field opts in with public: true. Sensitive secrets opt in with sensitive: true and are masked on admin reads.

Field flagPublic API (/settings/...)Admin API (/admin/settings/...)
(none)404raw value
public: trueresolved value (DB > ENV > default)raw value
sensitive: true404masked sentinel ••••••••
public: true AND sensitive: truerejected at startup by modval (MV054)—

Orphan DB rows (rows whose (module, key) is not declared in any module's SettingsDefinition) are invisible on the public API and visible verbatim on the admin API.

Mask passthrough: when an admin write submits the •••••••• sentinel as the value for a sensitive: true field, the existing stored value is preserved — admins can re-save a form without retyping every secret. Submitting any other string overwrites normally. Non-sensitive fields treat the sentinel as a literal string.

Example schema:

yaml
settings:
  groups:
    - id: ai-llm.openai
      label: OpenAI
      settings:
        - key: openai_api_key
          label: API key
          type: string
          env_var: OPENAI_API_KEY
          sensitive: true     # masked on admin reads, never readable on public API
        - key: openai_default_model
          label: Default model
          type: string
          public: true        # readable by /settings/ai-llm/openai_default_model

Read-only public routes (core module)

Unauthenticated. Only fields declared with public: true AND sensitive: false are returned.

MethodPathDescription
GET/settings/{module}/{key}Get a public setting value (404 if not declared public: true)
GET/settings/{module}List all public settings for a module (filtered to public: true)

Admin routes (settings/admin submodule)

Require admin:access permission. Sensitive values are masked with •••••••• in responses.

MethodPathDescription
GET/admin/settings/schemasList settings schemas from all registered modules
POST/admin/settings/getBulk-get multiple settings by (module, key) pairs (masks sensitive)
PUT/admin/settings/{module}/{key}Set or update a setting (mask passthrough)
POST/admin/settings/saveAlternative POST endpoint for save (mask passthrough)
DELETE/admin/settings/{module}/{key}Delete a setting (scoped to module, 204)

Internal (service-to-service) routes

Secret-gated, not user-authenticated: the callers are other services, which have no user to act as. Gated on the x-api-secret header (constant-time compare); disabled unless SETTINGS_INTERNAL_SECRET is set — an unset secret refuses everything, so the internal API is opt-in.

MethodPathDescription
GET/internal/settings/schemasEvery module's settings schema, as {settings, groups, subgroups}
POST/internal/settings/refresh-cacheDrop this instance's settings cache

These are the service-to-service counterparts of the admin schemas / refresh-cache endpoints — same effect, different caller. subgroups is always empty: Redelay's settings IR has one level of grouping, and the key is present only for wire compatibility with the legacy shape.

Events published

TopicTriggerPayload
settings.updatedAny save or delete{ moduleId, key, value, deleted }

Environment variables

VariableDefaultDescription
SETTINGS_COLLECTION_settingsMongoDB collection name
SETTINGS_CACHE_TTL_SECONDS30Resolution-cache staleness bound for DB-sourced values; 0 disables expiry. See Caching and propagation.
SETTINGS_INTERNAL_SECRET(unset)Shared secret for the /internal/settings/* service-to-service routes. Unset disables them.

Blank-import

go
import (
    _ "github.com/redelay/go-framework/modules/settings"        // core read-only routes
    _ "github.com/redelay/go-framework/modules/settings/admin"  // admin CRUD routes
    _ "github.com/redelay/go-framework/modules/settings/flowdsl" // FlowDSL nodes
)

FlowDSL nodes (settings/flowdsl)

NodeKindDescription
redelay/settings-gettransformRead a setting by module + key; has NotFound output port
redelay/settings-setactionWrite a setting value
redelay/settings-deleteactionRemove a setting

storage

Object storage module backed by MinIO (S3-compatible).

Implements modules.ObjectStorage — inject into other modules for file uploads and downloads.

Methods

go
PutObject(ctx, bucket, key string, r io.Reader, size int64) error
GetObject(ctx, bucket, key string) (io.ReadCloser, error)
DeleteObject(ctx, bucket, key string) error
PresignURL(ctx, bucket, key string, ttl time.Duration) (string, error)

Environment variables

VariableDefaultDescription
MINIO_ENDPOINTlocalhost:9000MinIO server endpoint
MINIO_ACCESS_KEYminioadminAccess key ID
MINIO_SECRET_KEYminioadminSecret access key
MINIO_BUCKETredelayDefault bucket name
MINIO_USE_SSLfalseUse TLS

notifications

Real-time notifications via Server-Sent Events with MongoDB persistence.

Clients subscribe to GET /notifications/stream (under ROUTE_PREFIX, so /api/v1/notifications/stream by default) and receive events in real time. The module consumes notification.created events from the event bus and fans them out to all connected clients for that user.

Routes

MethodPathDescription
GET/notifications/streamSSE stream for authenticated user
GET/notifications/List notifications (paginated)
POST/notifications/{id}/readMark notification as read

Environment variables

VariableDefaultDescription
NOTIFICATIONS_COLLECTIONnotificationsMongoDB collection name

Event schema (notification.created)

json
{
  "user_id": "string",
  "type": "string",
  "title": "string",
  "body": "string",
  "data": {}
}

verification

Single-use token-based email verification.

Issues verification tokens, stores them in MongoDB, and optionally emits email.send events (consumed by go-module-email) when a verification request is made via the HTTP API.

Routes

MethodPathDescription
POST/verification/requestCreate token and emit verification email
POST/verification/verifyConsume token and return verification result

Store API (for other modules)

The verification module exposes its Store directly so other modules can manage the token lifecycle themselves without going through the HTTP routes. Discover it via Configure():

go
func (m *MyModule) Configure(registry *modules.Registry) error {
    for _, mod := range registry.All() {
        if vm, ok := mod.(*verification.Module); ok {
            m.verifyStore = vm.Store
        }
    }
    return nil
}

// Create a token (returns *verification.VerificationToken):
vt, err := m.verifyStore.Create(ctx, email, verification.KindEmail, 32)

// Consume a token (returns nil if invalid/expired):
vt, err := m.verifyStore.Consume(ctx, token)

Environment variables

VariableDefaultDescription
VERIFICATION_COLLECTIONverification_tokensMongoDB collection name
VERIFICATION_TOKEN_TTL_MINUTES30Token expiry in minutes
VERIFICATION_TOKEN_BYTES32Random bytes for token generation

URL shortener with click tracking.

Short links are stored in MongoDB. The redirect handler is mounted at /{code} (outside the API prefix) via MountProvider.

Routes

MethodPathDescription
POST/links/Create a short link
GET/links/List short links (last 100)
DELETE/links/{code}Delete a short link
GET/{code}Redirect to target URL

Environment variables

VariableDefaultDescription
LINKS_COLLECTIONshort_linksMongoDB collection name
LINKS_CODE_LENGTH6Generated code length in bytes
LINKS_BASE_URL``Base URL prepended to short codes

metrics

Prometheus-compatible metrics endpoint and HTTP request duration middleware.

Mounts /metrics (configurable) via MountProvider and wraps all routes with a http_request_duration_seconds histogram labeled by method, path, and status code.

Environment variables

VariableDefaultDescription
METRICS_PATH/metricsHTTP path for the metrics endpoint

Exposed metrics

MetricTypeLabels
http_request_duration_secondsHistogrammethod, path, status

metrics/admin — database observability (admin-api only)

The metrics/admin submodule adds GET /admin/metrics/database (under the admin prefix, admin:access), returning a live summary of ClickHouse, MongoDB, and Redis activity. Import it from the admin-api binary only.

Every figure is read from the database's own statistics — Redelay does not instrument queries app-side — so a field is populated only where the backend natively reports it:

BackendSourceRealNot available
ClickHousesystem.query_logtotal, avg/max ms, error rate, slow %, top query kinds—
RedisINFO commandstats / statstotal ops, avg ms, top commandsmax, error rate, slow %
MongoDBserverStatustotal ops, ops/minavg/max ms, error rate, slow % (need the profiler)

A nil backend yields a zero-valued section rather than an error, so the endpoint works on a deployment running only some of the three.


security

Per-user security: an activity log, trusted devices, rate limiting, and automatic restriction of suspicious accounts.

The core module mounts under /user (relative to ROUTE_PREFIX) and records request activity through a middleware, restricting a user once their suspicious -activity count crosses a threshold.

Routes (core, authenticated)

MethodPathDescription
GET/user/activitiesThe caller's recent activity
GET/user/rate-limitsThe caller's current rate-limit status
GET/user/security/profileSuspicious count, restriction state, known devices
GET/user/security/suspicious-activitiesThe caller's flagged activity
GET/user/security/trusted-devicesList trusted devices
POST/user/security/trusted-devicesTrust the current device
DELETE/user/security/trusted-devices/{device_id}Untrust a device

security/admin (admin-api only)

Superuser operations under the admin prefix: per-user activity reports (/admin/security/activities/{user_id} + summary), error/path reports, security analytics, and enforcement — POST /admin/security/block-ip/{ip} and POST /admin/security/restrict-user/{user_id}.

Environment variables

VariableDefaultDescription
SECURITY_ACTIVITY_COLLECTIONuser_activitiesActivity log collection
SECURITY_DEVICE_COLLECTIONtrusted_devicesTrusted-device collection
SECURITY_SUSPICIOUS_COLLECTIONsuspicious_activitiesFlagged-activity collection
SECURITY_RESTRICTION_COLLECTIONsecurity_restrictionsActive restrictions
SECURITY_ACTIVITY_RETENTION_DAYS30Activity-log TTL
SECURITY_RESTRICT_THRESHOLD10Suspicious events before auto-restriction
SECURITY_USER_RATE_LIMIT1000Per-user request budget
SECURITY_IP_RATE_LIMIT200Per-IP request budget

eventadmin

Event-bus admin (admin-api only): browse event definitions, publish events by hand, browse the event log, and manage webhook subscriptions.

Mounts under /admin/events (admin:access). Import from the admin-api binary only.

Routes

MethodPathDescription
GET/admin/events/definitionsKnown event types — live registry + runtime-registered
POST/admin/events/definitionRegister a runtime event definition (?definition=<json>)
POST/admin/events/publishPublish an event → bus + log + webhook dispatch
GET/admin/events/logsBrowse the event log (entity_type/action/limit filters)
GET/admin/events/subscriptionsList webhook subscriptions
POST/admin/events/subscriptionRegister a webhook (entity_type/action match; nil = any)

Definitions are read from the live module registry at request time — after every module has registered — so the list is current without depending on module startup ordering.

The event log captures events published through this module (and any module that calls its service). It does not passively sniff every event on the bus: Redelay consumers are per-topic and a catch-all sink would depend on fragile Configure ordering. For a complete trail, add per-topic consumers or use the audit module.

Environment variables

VariableDefaultDescription
EVENTADMIN_DEFINITION_COLLECTIONevent_definitionsRuntime definitions
EVENTADMIN_SUBSCRIPTION_COLLECTIONevent_subscriptionsWebhook subscriptions
EVENTADMIN_LOG_COLLECTIONevent_logsEvent log
EVENTADMIN_LOG_TTL_HOURS168Event-log retention (7 days)
EVENTADMIN_WEBHOOK_TIMEOUT_SECONDS5Per-callback timeout on dispatch

flowexec

Flow CRUD, immutable versioning, run launch, and live SSE event streaming for FlowDSL flows.

go-flowdsl/flowexec/module is the framework-side HTTP surface over go-flowdsl/flowexec. It wires the executor, store, and sink into your app and exposes them as REST + Server-Sent Events. Register it with a blank import:

go
import (
    _ "github.com/redelay/go-flowdsl/runtime"       // your handlers register here
    _ "github.com/redelay/go-flowdsl/flowexec/module"
)

At bootstrap it inspects deps.DB:

deps.DBStoreSink
non-nilmongostore (indexes created automatically)mongots time-series sink (falls back to noop on init error)
nilmemstorenoop

App code registers node handlers on the engine exposed by the module:

go
m := modules.Registry().MustGet("flowexec").(*flowexec.Module)
m.Engine().RegisterHandler("charge_card", chargeCard)
m.Engine().RegisterKindHandler(ir.NodeKindAction, defaultAction)

Live SSE broadcaster

Between the executor and the downstream sink sits a broadcaster that tees every FlowEvent to per-runID subscriber channels. The downstream sink still receives the canonical record. Slow subscribers are dropped (non-blocking send) rather than backpressuring the run — the durable event stream in the sink remains complete, subscribers see a live best-effort feed.

Routes

All routes are mounted relative to ROUTE_PREFIX.

MethodPathDescription
GET/flowsList flows (cursor-paginated)
POST/flowsCreate a flow — {name, description, labels}
GET/flows/{id}Read a flow (head and published pointers)
PATCH/flows/{id}Update flow metadata
DELETE/flows/{id}Delete flow and all versions
POST/flows/{id}/versionsSave a new immutable version — {document, note, labels}
GET/flows/{id}/versionsList versions (newest first, documents omitted)
GET/flows/{id}/versions/{vid}Read a version, full document included

?format=spec query parameter

The version endpoints accept an optional ?format=spec query parameter. When present, the API transparently converts between the Studio spec document format and the internal ir.Workflow using go-flowdsl/spec.

Save a version in spec format:

shell
curl -sX POST "$API/flows/$FLOW/versions?format=spec" \
  -H 'Content-Type: application/json' \
  -d @my-flow.spec.json

The request body is a spec.Document (with flowdsl, info, flows, components). The server calls spec.ToWorkflow() and stores the resulting ir.Workflow.

Read a version in spec format:

shell
curl -s "$API/flows/$FLOW/versions/$VID?format=spec"

The server calls spec.FromWorkflow() on the stored ir.Workflow and returns a spec.Document in the response. This is the format the FlowDSL Studio editor expects.

Without ?format=spec, both endpoints use the raw ir.Workflow format. | POST | /flows/{id}/publish | Flip published pointer — {versionID} | | POST | /flows/{id}/runs | Launch run on the currently published version — {input, labels} | | GET | /runs/{runID}/events | SSE live stream of FlowEvent payloads for a run |

POST /flows/{id}/runs resolves the published version, calls Executor.RunAsync, and returns 202 Accepted with {runID, versionID, versionHash}. If the flow has no published version the response is 409 Conflict.

SSE stream

GET /runs/{runID}/events returns a text/event-stream that writes a : connected prelude, then one SSE frame per FlowEvent:

text
event: node.started
data: {"runID":"run.abc","flowID":"flow.xyz","nodeID":"n1","timestamp":"2026-04-14T10:00:00Z",...}

Heartbeat : keep-alive lines are written every 15 seconds so intermediaries do not drop the connection. The stream closes automatically when a run.completed or run.failed event arrives.

Response shapes

POST /flows/{id}/runs:

json
{
  "runID": "run.018f...",
  "versionID": "v.018f...",
  "versionHash": "f3e1b..."
}

Every event emitted on the SSE stream carries the same versionHash, so clients can tell at a glance whether they are watching the same version they launched.

Usage with historical replay

The module does not serve paginated historical events — use the sink's Reader interface (implemented by mongots) to back a replay endpoint in your own module when needed. The SSE stream is strictly live.