Domain Addon Modules
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.
Quick start
go get github.com/redelay/go-modules
import (
_ "github.com/redelay/go-modules/scheduler"
_ "github.com/redelay/go-modules/metrics"
// add more as needed
)
Module overview
| Module | Package | Purpose |
|---|---|---|
| scheduler | go-modules/scheduler | Distributed cron with leader election |
| settings | go-framework/modules/settings | Key-value settings API |
| storage | go-modules/storage | Object storage (S3/MinIO) with S3 list / head / copy / delete-many nodes |
| notifications | go-modules/notifications | Real-time SSE notifications |
| verification | go-modules/verification | Email address verification |
| links | go-modules/links | URL shortener |
| metrics | go-modules/metrics | Prometheus metrics |
| flowexec | go-flowdsl/flowexec/module | Flow CRUD, versioning, run launch, live SSE events |
| search | go-modules/search | Pluggable vector search (Qdrant today, OpenSearch planned) with index CRUD + FlowDSL upsert/query nodes |
| i18n | go-modules/i18n | Locale-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. |
| flowtemplates | go-modules/flowtemplates | Curated starter flows surfaced via flowexec.TemplateProvider |
| coordination | go-modules/coordination | Cross-container Coordinator — pub/sub, streams, KV (NATS or Redis backend) |
| websockets | go-modules/ws | WebSocket hub — rooms + cross-replica fan-out, health/info/broadcast control plane |
| security | go-modules/security | Per-user security — activity log, trusted devices, rate limits, auto-restriction; security/admin for IP blocks + user restriction |
| eventadmin | go-modules/eventadmin | Event-bus admin (admin-api only) — event definitions, manual publish, event log, webhook subscriptions |
| metrics/admin | go-modules/metrics/admin | Database observability (admin-api only) — live ClickHouse/Mongo/Redis stats |
| apilog | go-modules/apilog | Time-series log of outbound external API calls, queryable in admin; Mask() redaction + apilog/record debug node |
| ai | go-modules/ai | Provider-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. |
| content | go-modules/content | CMS-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. |
| workers | go-modules/workers | Cross-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:
import (
_ "github.com/redelay/go-modules/scheduler"
_ "github.com/redelay/go-modules/scheduler/flowdsl" // 4 FlowDSL nodes
)
| Subpackage | Factory | Nodes | Kinds |
|---|---|---|---|
scheduler/flowdsl | scheduler-flowdsl | 4 | 2 source (with settings_schema), 2 action |
settings/flowdsl | settings-flowdsl | 3 | 1 transform (with NotFound), 2 action |
storage/flowdsl | storage-flowdsl | 8 | 4 upload/download/delete/presign + 4 S3 list / head / copy / delete-many |
notifications/flowdsl | notifications-flowdsl | 3 | 2 action, 1 source |
verification/flowdsl | verification-flowdsl | 2 | 1 action, 1 router (Valid/Invalid) |
links/flowdsl | links-flowdsl | 3 | 2 action, 1 transform (with NotFound) |
search/flowdsl | search-flowdsl | 2 | 2 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
| Variable | Default | Description |
|---|---|---|
SCHEDULER_COLLECTION | scheduler_schedules | MongoDB collection for schedule definitions |
SCHEDULER_LOG_COLLECTION | scheduler_logs | MongoDB collection for execution logs |
SCHEDULER_POLL_INTERVAL_SECONDS | 10 | How often to check for due schedules |
SCHEDULER_LOG_TTL_DAYS | 30 | How long execution history is kept. 0 keeps it forever |
SCHEDULER_LOCK_TTL_SECONDS | 30 | Redis leader-election lock TTL |
SCHEDULER_LOCK_KEY | redelay:<app-name>:scheduler:leader | Redis 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:
# 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:
GET /health→ theschedulercheck names a foreign lock holder.redis-cli get redelay:<app>:scheduler:leader→ who thinks it is leader.- The
hostfield on recentscheduler_logsrows → who actually dispatched. - A
duration_msof0on 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.
// 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.
| Method | Path | Purpose |
|---|---|---|
| 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}/execute | Fire now, without shifting the regular timing |
| GET | /schedules/{schedule_id}/executions | List execution log entries (newest first) |
| POST | /schedules/{schedule_id}/toggle | Enable / 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.
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:
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
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
| Priority | Source | Notes |
|---|---|---|
| 1 (highest) | MongoDB (admin-saved) | Written via admin API |
| 2 | Environment variable | Explicit EnvVar field or implicit {MODULE}_{KEY} |
| 3 | YAML default | Declared 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:
- Single binary (one cmd/api with admin imports) — local
Service.Setinvalidates the local cache; events are unnecessary. - Multi-binary local dev (cmd/api + cmd/admin-api with
TRANSPORT=memory) — each binary runs an isolated in-process event bus, so asettings.updatedevent published by admin-api never reaches cmd/api. The TTL (below) is the only invalidation mechanism in this mode. - Production cluster (N replicas with
TRANSPORT=kafka/nats/redis) — every replica subscribes with a uniquesettings-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 source | Cached? | Invalidated by |
|---|---|---|
db — admin-edited row in MongoDB | Yes, with TTL | local Service.Set, settings.updated event, TTL expiry |
env — environment variable | Yes, with TTL | TTL expiry (env doesn't change at runtime) |
default — YAML-declared default | Never cached | n/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
- The value is written to MongoDB.
- The local cache entry is invalidated on the writing worker.
- A
settings.updatedevent is published to the event bus. - Every worker subscribed with a unique GroupID receives the event (fan-out, not round-robin) and invalidates its local cache entry.
- On next read, the cache miss triggers a MongoDB re-query.
Cache TTL
| Variable | Default | Notes |
|---|---|---|
SETTINGS_CACHE_TTL_SECONDS | 30 | Maximum 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=0is fine — local invalidation is authoritative. - Multi-binary dev (
TRANSPORT=memory, separate cmd/api + cmd/admin-api): leave at default30. 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 to300for 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 flag | Public API (/settings/...) | Admin API (/admin/settings/...) |
|---|---|---|
| (none) | 404 | raw value |
public: true | resolved value (DB > ENV > default) | raw value |
sensitive: true | 404 | masked sentinel •••••••• |
public: true AND sensitive: true | rejected 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:
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.
| Method | Path | Description |
|---|---|---|
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.
| Method | Path | Description |
|---|---|---|
GET | /admin/settings/schemas | List settings schemas from all registered modules |
POST | /admin/settings/get | Bulk-get multiple settings by (module, key) pairs (masks sensitive) |
PUT | /admin/settings/{module}/{key} | Set or update a setting (mask passthrough) |
POST | /admin/settings/save | Alternative 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.
| Method | Path | Description |
|---|---|---|
GET | /internal/settings/schemas | Every module's settings schema, as {settings, groups, subgroups} |
POST | /internal/settings/refresh-cache | Drop 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
| Topic | Trigger | Payload |
|---|---|---|
settings.updated | Any save or delete | { moduleId, key, value, deleted } |
Environment variables
| Variable | Default | Description |
|---|---|---|
SETTINGS_COLLECTION | _settings | MongoDB collection name |
SETTINGS_CACHE_TTL_SECONDS | 30 | Resolution-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
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)
| Node | Kind | Description |
|---|---|---|
redelay/settings-get | transform | Read a setting by module + key; has NotFound output port |
redelay/settings-set | action | Write a setting value |
redelay/settings-delete | action | Remove a setting |
storage
Object storage module backed by MinIO (S3-compatible).
Implements modules.ObjectStorage — inject into other modules for file uploads and downloads.
Methods
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
| Variable | Default | Description |
|---|---|---|
MINIO_ENDPOINT | localhost:9000 | MinIO server endpoint |
MINIO_ACCESS_KEY | minioadmin | Access key ID |
MINIO_SECRET_KEY | minioadmin | Secret access key |
MINIO_BUCKET | redelay | Default bucket name |
MINIO_USE_SSL | false | Use 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
| Method | Path | Description |
|---|---|---|
GET | /notifications/stream | SSE stream for authenticated user |
GET | /notifications/ | List notifications (paginated) |
POST | /notifications/{id}/read | Mark notification as read |
Environment variables
| Variable | Default | Description |
|---|---|---|
NOTIFICATIONS_COLLECTION | notifications | MongoDB collection name |
Event schema (notification.created)
{
"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
| Method | Path | Description |
|---|---|---|
POST | /verification/request | Create token and emit verification email |
POST | /verification/verify | Consume 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():
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
| Variable | Default | Description |
|---|---|---|
VERIFICATION_COLLECTION | verification_tokens | MongoDB collection name |
VERIFICATION_TOKEN_TTL_MINUTES | 30 | Token expiry in minutes |
VERIFICATION_TOKEN_BYTES | 32 | Random bytes for token generation |
links
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
| Method | Path | Description |
|---|---|---|
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
| Variable | Default | Description |
|---|---|---|
LINKS_COLLECTION | short_links | MongoDB collection name |
LINKS_CODE_LENGTH | 6 | Generated 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
| Variable | Default | Description |
|---|---|---|
METRICS_PATH | /metrics | HTTP path for the metrics endpoint |
Exposed metrics
| Metric | Type | Labels |
|---|---|---|
http_request_duration_seconds | Histogram | method, 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:
| Backend | Source | Real | Not available |
|---|---|---|---|
| ClickHouse | system.query_log | total, avg/max ms, error rate, slow %, top query kinds | — |
| Redis | INFO commandstats / stats | total ops, avg ms, top commands | max, error rate, slow % |
| MongoDB | serverStatus | total ops, ops/min | avg/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)
| Method | Path | Description |
|---|---|---|
| GET | /user/activities | The caller's recent activity |
| GET | /user/rate-limits | The caller's current rate-limit status |
| GET | /user/security/profile | Suspicious count, restriction state, known devices |
| GET | /user/security/suspicious-activities | The caller's flagged activity |
| GET | /user/security/trusted-devices | List trusted devices |
| POST | /user/security/trusted-devices | Trust 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
| Variable | Default | Description |
|---|---|---|
SECURITY_ACTIVITY_COLLECTION | user_activities | Activity log collection |
SECURITY_DEVICE_COLLECTION | trusted_devices | Trusted-device collection |
SECURITY_SUSPICIOUS_COLLECTION | suspicious_activities | Flagged-activity collection |
SECURITY_RESTRICTION_COLLECTION | security_restrictions | Active restrictions |
SECURITY_ACTIVITY_RETENTION_DAYS | 30 | Activity-log TTL |
SECURITY_RESTRICT_THRESHOLD | 10 | Suspicious events before auto-restriction |
SECURITY_USER_RATE_LIMIT | 1000 | Per-user request budget |
SECURITY_IP_RATE_LIMIT | 200 | Per-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
| Method | Path | Description |
|---|---|---|
| GET | /admin/events/definitions | Known event types — live registry + runtime-registered |
| POST | /admin/events/definition | Register a runtime event definition (?definition=<json>) |
| POST | /admin/events/publish | Publish an event → bus + log + webhook dispatch |
| GET | /admin/events/logs | Browse the event log (entity_type/action/limit filters) |
| GET | /admin/events/subscriptions | List webhook subscriptions |
| POST | /admin/events/subscription | Register 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.
audit module.Environment variables
| Variable | Default | Description |
|---|---|---|
EVENTADMIN_DEFINITION_COLLECTION | event_definitions | Runtime definitions |
EVENTADMIN_SUBSCRIPTION_COLLECTION | event_subscriptions | Webhook subscriptions |
EVENTADMIN_LOG_COLLECTION | event_logs | Event log |
EVENTADMIN_LOG_TTL_HOURS | 168 | Event-log retention (7 days) |
EVENTADMIN_WEBHOOK_TIMEOUT_SECONDS | 5 | Per-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:
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.DB | Store | Sink |
|---|---|---|
| non-nil | mongostore (indexes created automatically) | mongots time-series sink (falls back to noop on init error) |
| nil | memstore | noop |
App code registers node handlers on the engine exposed by the module:
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.
| Method | Path | Description |
|---|---|---|
GET | /flows | List flows (cursor-paginated) |
POST | /flows | Create 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}/versions | Save a new immutable version — {document, note, labels} |
GET | /flows/{id}/versions | List 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:
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:
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:
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:
{
"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.