Skip to content

Workflow Server

The Workflow Server is a separate binary (cmd/workflow) that connects to the Aether gateway as a WorkflowEngineClient. Once connected, it receives every message the gateway delivers on the event.* topic family and routes them through a configurable rule engine. It never shares process space with the gateway, so it can be scaled, restarted, or replaced independently.

Inside the server, six cooperating components handle the full lifecycle of event-driven automation: the Router matches incoming events against stored rules using expr-lang expressions and dispatches transformed messages; the Scheduler polls for cron or one-shot triggers using robfig/cron; the DAGEngine executes multi-step directed-acyclic-graph workflows where each step can depend on the output of earlier steps; the StateMachineEngine manages long-running stateful objects with declared state transitions, entry/exit actions, and timeouts; the Executor sends resulting Aether tasks back to the gateway; and the LeaderElector ensures that only one instance drives scheduling and monitoring at a time.

  • Reacting to gateway events with conditional dispatch (replace custom event-bridge scripts).
  • Running recurring jobs on a cron schedule without external schedulers.
  • Coordinating multi-step tasks where each step depends on the result of the previous one.
  • Tracking stateful workflows that persist across many events and time windows.
┌─────────────────────────────────────────────────────────────┐
│ Aether Gateway (gRPC :50051) │
│ │
│ event.* topics ──────────────────────────────────────────► │
│ ▲ WorkflowEngineClient │
└────────────────────────────────┼────────────────────────────┘
│ gRPC bidirectional stream
┌────────────────────────────────┼────────────────────────────┐
│ Workflow Server │ │
│ │ │
│ handleMessage() │ │
│ ├── Router (rule match + transform) │
│ └── DAGEngine (event-triggered DAGs) │
│ │
│ Scheduler (cron / interval) │
│ StateMachineEngine (state + transitions) │
│ LeaderElector (Redis or single-node) │
│ │
│ AdminServer (REST :31881) │
│ │
│ PostgreSQL (workflow_schema_migrations) │
└─────────────────────────────────────────────────────────────┘

The server uses its own migration tracking table (workflow_schema_migrations) so it never conflicts with the gateway’s schema_migrations table. Both can share the same Postgres database.

Terminal window
# 1. Run with an explicit config file
go run ./server/cmd/workflow/main.go --config server/configs/workflow.yaml
# 2. Development mode – uses hardcoded defaults, no config file required
go run ./server/cmd/workflow/main.go --dev
# 3. Lite mode – SQLite only, no Redis required
go run ./server/cmd/workflow/main.go --config server/configs/workflow.yaml --lite

Migrations run automatically on startup. The binary exits non-zero if the config file is missing (unless --dev is passed) or if required fields fail validation.

KeyTypeDefaultDescription
modestandard | litestandardlite uses SQLite and skips Redis.
aether.addressstringlocalhost:50051Gateway gRPC address.
aether.implementationstringaether-workflowClient identifier registered with the gateway.
aether.workspacestring_systemWorkspace the client operates under.
aether.tls.cert_filepath—Client TLS certificate.
aether.tls.key_filepath—Client TLS private key.
aether.tls.ca_filepath—CA certificate for gateway verification.
aether.credentials.api_keystring—API key for gateway authentication.
postgres.hoststringlocalhostPostgreSQL host.
postgres.portint5432PostgreSQL port.
postgres.databasestringaetherDatabase name.
postgres.userstringaetherDatabase user.
postgres.passwordstring—Database password.
postgres.ssl_modestringdisablePostgreSQL SSL mode.
postgres.max_connectionsint10Max open DB connections.
postgres.max_idle_connectionsint5Max idle DB connections.
sqlite.pathpathworkflow.dbSQLite file path (lite mode only).
redis.clusterlist["localhost:6379"]Redis addresses (single or cluster).
redis.passwordstring—Redis password.
workflow.rule_cache_ttlduration2mHow long compiled rules stay cached.
workflow.rule_cache_sizeint2048Max compiled rule entries in cache.
workflow.scheduler_poll_intervalduration1sHow often the scheduler polls for due triggers.
workflow.dag_monitor_intervalduration5sHow often the DAG monitor checks running executions.
workflow.step_default_timeoutduration5mPer-step timeout when none is declared.
workflow.dag_default_timeoutduration1hWhole-DAG timeout when none is declared.
workflow.max_concurrent_executionsint100Maximum simultaneous DAG executions.
admin.enabledbooltrueEnables the admin REST API.
admin.portint31881Admin REST API listen port.
admin.api_keystring—API key protecting the admin surface (optional).
logging.levelstringinfoLog level: debug, info, warn, error.
logging.formatstringautojson or console (auto-detected from TTY).

The following variables are applied after the YAML file is parsed. See Environment Reference for the full list.

VariableOverrides
WORKFLOW_MODEmode
WORKFLOW_ADMIN_ENABLEDadmin.enabled
WORKFLOW_ADMIN_PORTadmin.port
WORKFLOW_ADMIN_API_KEYadmin.api_key
AETHER_ADDRESSaether.address
AETHER_WORKSPACEaether.workspace
AETHER_API_KEYaether.credentials.api_key
POSTGRES_HOSTpostgres.host
POSTGRES_PORTpostgres.port
POSTGRES_USERpostgres.user
POSTGRES_PASSWORDpostgres.password
POSTGRES_DATABASEpostgres.database
REDIS_ADDRredis.cluster[0] (single address)
REDIS_PASSWORDredis.password
AETHER_LOG_LEVELlogging.level
SQLITE_PATHsqlite.path

The admin server listens on port 31881 (default) and exposes a JSON REST API under /api/v1. When admin.api_key is set, every /api/v1 request must include either an Authorization: Bearer <key> or X-API-Key: <key> header. The /health endpoint is always unauthenticated.

PrefixMethodsPurpose
/api/v1/rulesGET, POSTList or create event routing rules.
/api/v1/rules/{id}GET, PUT, DELETERead, update, or delete a rule by integer ID.
/api/v1/workflowsGET, POSTList or create DAG workflow definitions.
/api/v1/workflows/{id}GET, DELETERead or deactivate a workflow definition.
/api/v1/schedulesGET, POSTList or create scheduled triggers.
/api/v1/schedules/{id}DELETEDelete a schedule.
/api/v1/executionsGETList executions (filter with ?status=running).
/api/v1/executions/{id}GETGet execution detail including step states.
/api/v1/executions/{id}/cancelPOSTCancel a running execution.
/api/v1/statemachinesGET, POSTList or create state machine definitions.
/api/v1/statemachines/{id}GET, DELETERead or deactivate a state machine.
/api/v1/statemachines/{id}/instancesGET, POSTList or create instances of a state machine.
/api/v1/statemachines/{mid}/instances/{iid}GETGet a specific instance.
/api/v1/statemachines/{mid}/instances/{iid}/eventPOSTSend an event to an instance.
/healthGETLiveness check; returns {"status":"ok"}.
Terminal window
# List all active routing rules
curl http://localhost:31881/api/v1/rules
# Create a routing rule (with API key authentication)
curl -s -X POST http://localhost:31881/api/v1/rules \
-H "X-API-Key: $WORKFLOW_ADMIN_API_KEY" \
-H "Content-Type: application/json" \
-d '{
"rule_name": "notify-on-failure",
"source_event": "task.failed",
"source_agent": "*",
"destination_template": "notify:\n channel: ops\n message: \"{{ .event_name }} from {{ .source_agent }}\"",
"workspace": "*"
}'
# Cancel a running execution
curl -s -X POST http://localhost:31881/api/v1/executions/<exec-id>/cancel \
-H "X-API-Key: $WORKFLOW_ADMIN_API_KEY"
# Check liveness
curl http://localhost:31881/health

The workflow server manages its own schema migrations independently from the gateway. It tracks applied migrations in the workflow_schema_migrations table using the same runner pattern as the gateway but with a separate table name to avoid conflicts.

Migrations run automatically at startup in the correct order:

MigrationContents
001_workflow_schema.sqlCore tables: workflow_rules, workflow_definitions, workflow_executions, workflow_step_states, workflow_schedules.
002_statemachine_schema.sqlState machine tables: workflow_state_machines, workflow_state_machine_instances.
003_schedule_enhancements.sqlAdds max_concurrent and active_task_id columns to workflow_schedules.

SQLite runs an equivalent set of migrations from a parallel embedded directory.

In standard mode the server uses a Redis-backed leader lock (workflow:leader). Only the leader instance runs the scheduler, DAG monitor, and state machine timeout monitor. All instances handle incoming events and admin API requests.

In lite mode (--lite), the server always assumes leadership (single-node only). Redis is not required.

  • The server reconnects to the gateway automatically on disconnect with exponential backoff (1 s → 30 s).
  • If the process is restarted while executions are running, the DAG monitor resumes monitoring them on the next poll cycle.
  • Leader election re-attempts every 5 seconds until the lock is acquired.

The scheduler and monitors only run on the leader. If the leader crashes without releasing its lock, the next instance acquires leadership after the Redis lock TTL expires. No manual intervention is required.

Server exits immediately with “config file not found” Pass --dev to use built-in development defaults, or verify the --config path is correct.

“configuration validation failed: aether.address is required” The aether.address YAML key (or AETHER_ADDRESS env var) must be set. Confirm the gateway is reachable before starting the workflow server.

Server connects but no events are processed Verify that agents are publishing to event.* topics and that the aether.workspace value matches the workspace used by those agents. The server only processes EVENT-typed messages.

Admin API returns 401 An API key is configured. Supply X-API-Key: <key> or Authorization: Bearer <key> in the request headers.

Schedules are not firing Check that the server has acquired leadership (workflow:leader key in Redis). In a multi-instance deployment only the leader fires schedules. If Redis is unavailable, all instances will fail to acquire leadership and no schedules will run.

SQLite mode with concurrent writes Lite mode sets max_open_conns=1 to serialize writes. It is not suitable for high-concurrency production deployments; use standard (PostgreSQL) mode instead.