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.
When to use it
Section titled “When to use it”- 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.
Architecture
Section titled “Architecture” ┌─────────────────────────────────────────────────────────────┐ │ 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.
Quick start
Section titled “Quick start”# 1. Run with an explicit config filego run ./server/cmd/workflow/main.go --config server/configs/workflow.yaml
# 2. Development mode – uses hardcoded defaults, no config file requiredgo run ./server/cmd/workflow/main.go --dev
# 3. Lite mode – SQLite only, no Redis requiredgo run ./server/cmd/workflow/main.go --config server/configs/workflow.yaml --liteMigrations 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.
Configuration
Section titled “Configuration”YAML reference
Section titled “YAML reference”| Key | Type | Default | Description |
|---|---|---|---|
mode | standard | lite | standard | lite uses SQLite and skips Redis. |
aether.address | string | localhost:50051 | Gateway gRPC address. |
aether.implementation | string | aether-workflow | Client identifier registered with the gateway. |
aether.workspace | string | _system | Workspace the client operates under. |
aether.tls.cert_file | path | — | Client TLS certificate. |
aether.tls.key_file | path | — | Client TLS private key. |
aether.tls.ca_file | path | — | CA certificate for gateway verification. |
aether.credentials.api_key | string | — | API key for gateway authentication. |
postgres.host | string | localhost | PostgreSQL host. |
postgres.port | int | 5432 | PostgreSQL port. |
postgres.database | string | aether | Database name. |
postgres.user | string | aether | Database user. |
postgres.password | string | — | Database password. |
postgres.ssl_mode | string | disable | PostgreSQL SSL mode. |
postgres.max_connections | int | 10 | Max open DB connections. |
postgres.max_idle_connections | int | 5 | Max idle DB connections. |
sqlite.path | path | workflow.db | SQLite file path (lite mode only). |
redis.cluster | list | ["localhost:6379"] | Redis addresses (single or cluster). |
redis.password | string | — | Redis password. |
workflow.rule_cache_ttl | duration | 2m | How long compiled rules stay cached. |
workflow.rule_cache_size | int | 2048 | Max compiled rule entries in cache. |
workflow.scheduler_poll_interval | duration | 1s | How often the scheduler polls for due triggers. |
workflow.dag_monitor_interval | duration | 5s | How often the DAG monitor checks running executions. |
workflow.step_default_timeout | duration | 5m | Per-step timeout when none is declared. |
workflow.dag_default_timeout | duration | 1h | Whole-DAG timeout when none is declared. |
workflow.max_concurrent_executions | int | 100 | Maximum simultaneous DAG executions. |
admin.enabled | bool | true | Enables the admin REST API. |
admin.port | int | 31881 | Admin REST API listen port. |
admin.api_key | string | — | API key protecting the admin surface (optional). |
logging.level | string | info | Log level: debug, info, warn, error. |
logging.format | string | auto | json or console (auto-detected from TTY). |
Environment variable overrides
Section titled “Environment variable overrides”The following variables are applied after the YAML file is parsed. See Environment Reference for the full list.
| Variable | Overrides |
|---|---|
WORKFLOW_MODE | mode |
WORKFLOW_ADMIN_ENABLED | admin.enabled |
WORKFLOW_ADMIN_PORT | admin.port |
WORKFLOW_ADMIN_API_KEY | admin.api_key |
AETHER_ADDRESS | aether.address |
AETHER_WORKSPACE | aether.workspace |
AETHER_API_KEY | aether.credentials.api_key |
POSTGRES_HOST | postgres.host |
POSTGRES_PORT | postgres.port |
POSTGRES_USER | postgres.user |
POSTGRES_PASSWORD | postgres.password |
POSTGRES_DATABASE | postgres.database |
REDIS_ADDR | redis.cluster[0] (single address) |
REDIS_PASSWORD | redis.password |
AETHER_LOG_LEVEL | logging.level |
SQLITE_PATH | sqlite.path |
Admin REST API
Section titled “Admin REST API”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.
Endpoint families
Section titled “Endpoint families”| Prefix | Methods | Purpose |
|---|---|---|
/api/v1/rules | GET, POST | List or create event routing rules. |
/api/v1/rules/{id} | GET, PUT, DELETE | Read, update, or delete a rule by integer ID. |
/api/v1/workflows | GET, POST | List or create DAG workflow definitions. |
/api/v1/workflows/{id} | GET, DELETE | Read or deactivate a workflow definition. |
/api/v1/schedules | GET, POST | List or create scheduled triggers. |
/api/v1/schedules/{id} | DELETE | Delete a schedule. |
/api/v1/executions | GET | List executions (filter with ?status=running). |
/api/v1/executions/{id} | GET | Get execution detail including step states. |
/api/v1/executions/{id}/cancel | POST | Cancel a running execution. |
/api/v1/statemachines | GET, POST | List or create state machine definitions. |
/api/v1/statemachines/{id} | GET, DELETE | Read or deactivate a state machine. |
/api/v1/statemachines/{id}/instances | GET, POST | List or create instances of a state machine. |
/api/v1/statemachines/{mid}/instances/{iid} | GET | Get a specific instance. |
/api/v1/statemachines/{mid}/instances/{iid}/event | POST | Send an event to an instance. |
/health | GET | Liveness check; returns {"status":"ok"}. |
Example curl commands
Section titled “Example curl commands”# List all active routing rulescurl 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 executioncurl -s -X POST http://localhost:31881/api/v1/executions/<exec-id>/cancel \ -H "X-API-Key: $WORKFLOW_ADMIN_API_KEY"
# Check livenesscurl http://localhost:31881/healthMigrations
Section titled “Migrations”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:
| Migration | Contents |
|---|---|
001_workflow_schema.sql | Core tables: workflow_rules, workflow_definitions, workflow_executions, workflow_step_states, workflow_schedules. |
002_statemachine_schema.sql | State machine tables: workflow_state_machines, workflow_state_machine_instances. |
003_schedule_enhancements.sql | Adds max_concurrent and active_task_id columns to workflow_schedules. |
SQLite runs an equivalent set of migrations from a parallel embedded directory.
Operational notes
Section titled “Operational notes”Leader election
Section titled “Leader election”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.
Restart safety
Section titled “Restart safety”- 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.
Single active model
Section titled “Single active model”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.
Troubleshooting
Section titled “Troubleshooting”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.