Central state machine and WebSocket handler for Project Aether — a distributed agent orchestration mesh.
The Core Orchestrator is the nerve centre of the Aether platform. It accepts real-time context frames from mesh clients over WebSocket, validates and persists them to Redis, and exposes Prometheus metrics for observability.
┌──────────────────────────────────────────────────────────────────────────┐
│ Aether Core Orchestrator │
│ │
│ ┌──────────┐ ┌────────────────────┐ ┌──────────────────────────┐ │
│ │ Client │───▶│ WebSocket Stream │───▶│ FrameProcessor │ │
│ │ (mesh) │ │ /v1/mesh/stream/ │ │ validate → persist │ │
│ └──────────┘ │ {session_id} │ │ state via Redis │ │
│ └────────────────────┘ └──────────┬───────────────┘ │
│ │ │
│ ┌──────────┐ ┌────────────────────┐ ┌──────────▼───────────────┐ │
│ │Prometheus│◀───│ TelemetryMiddleware│ │ Redis (asyncio) │ │
│ │ Scraper │ │ /metrics │ │ aether:frames:{sid} │ │
│ └──────────┘ └────────────────────┘ │ aether:state:{sid}:{f} │ │
│ │ aether:count:{sid} │ │
│ ┌──────────┐ ┌────────────────────┐ └──────────────────────────┘ │
│ │ K8s │───▶│ Health Check │ │
│ │ Probe │ │ /v1/health │ │
│ └──────────┘ └────────────────────┘ │
└──────────────────────────────────────────────────────────────────────────┘
| Component | Location | Responsibility |
|---|---|---|
Settings |
app/core/config.py |
Pydantic BaseSettings singleton; reads from env / .env |
configure_logging |
app/core/logging.py |
Structured JSON logging via structlog |
ContextFrame |
app/models/context.py |
Pydantic v2 domain model for inbound frames |
MeshConnectionManager |
app/services/connection_manager.py |
Thread-safe WebSocket registry |
FrameProcessor |
app/services/frame_processor.py |
Validation → Redis persistence pipeline |
TelemetryMiddleware |
app/middleware/telemetry.py |
Prometheus counter/histogram instrumentation |
create_app |
app/main.py |
FastAPI application factory with lifespan |
| Method | Path | Description |
|---|---|---|
GET |
/v1/health |
Service and dependency health check |
GET |
/metrics |
Prometheus metrics (text exposition format) |
GET |
/docs |
Swagger UI (DEBUG mode only) |
GET |
/redoc |
ReDoc UI (DEBUG mode only) |
| Path | Protocol | Description |
|---|---|---|
/v1/mesh/stream/{session_id} |
JSON text frames | Bidirectional mesh stream for a single session |
ws://localhost:8000/v1/mesh/stream/<session_id>
session_id must be unique across the mesh at any given moment. Duplicate session IDs are refused with close code 1008 Policy Violation.
{
"session_id": "sess_01HX3Q9YNFM8GV2DKTA7W0ZBE",
"source_app": "Browser_Extension_Snippet",
"timestamp": "2026-06-08T12:00:00Z",
"frame_id": "550e8400-e29b-41d4-a716-446655440000",
"payloads": [
{
"type": "dom_snapshot",
"content": "<html>...</html>",
"metadata": {
"url": "https://example.com",
"tab_id": "42"
}
}
]
}| Field | Type | Required | Description |
|---|---|---|---|
session_id |
string | Yes | Must match the URL path parameter |
source_app |
string | Yes | Originating application identifier |
timestamp |
ISO-8601 | No | Defaults to server UTC time if omitted |
frame_id |
UUID | No | Auto-generated if omitted |
payloads |
array | Yes | At least one payload required |
{
"status": "ACK",
"session_id": "sess_01HX3Q9YNFM8GV2DKTA7W0ZBE",
"frame_id": "550e8400-e29b-41d4-a716-446655440000",
"processed_frames": 7,
"payload_count": 1,
"state": "DONE",
"timestamp": "2026-06-08T12:00:00.123Z"
}{
"status": "REJECTED",
"code": "VALIDATION_ERROR",
"message": "The submitted frame is not valid JSON or does not conform to the ContextFrame schema.",
"state": "REJECTED",
"timestamp": "2026-06-08T12:00:00.123Z"
}| Error code | Cause |
|---|---|
VALIDATION_ERROR |
Malformed JSON or schema violation |
SESSION_ID_MISMATCH |
Frame session_id differs from URL path param |
REDIS_UNAVAILABLE |
Redis write failure during processing |
RECEIVED → VALIDATED → PROCESSING → DONE
↘ REJECTED
| Variable | Default | Description |
|---|---|---|
APP_NAME |
aether-core-orchestrator |
Service name in logs and health payloads |
VERSION |
1.0.0 |
Semantic version |
DEBUG |
false |
Enable Swagger UI and verbose logging |
REDIS_URL |
redis://localhost:6379 |
Redis connection URL |
REDIS_PASSWORD |
(empty) | Redis AUTH password |
REDIS_DB |
0 |
Redis logical database (0–15) |
MAX_CONNECTIONS |
1000 |
Max concurrent WebSocket sessions |
RATE_LIMIT_RPS |
100.0 |
Per-session frame rate limit (req/s) |
LOG_LEVEL |
INFO |
Logging verbosity |
CORS_ORIGINS |
["*"] |
Comma-separated list of allowed CORS origins |
FRAME_TTL_SECONDS |
3600 |
Redis key TTL for stored frames (seconds) |
- Python 3.11+
- Redis 7.x (local or Docker)
make(optional but recommended)
git clone <repo-url>
cd aether-core-orchestrator
python3.11 -m venv .venv
source .venv/bin/activatemake install-dev
# or manually:
pip install -r requirements-dev.txtcp .env.example .env # edit as neededMinimum .env for local development:
DEBUG=true
REDIS_URL=redis://localhost:6379
LOG_LEVEL=DEBUGdocker run -d --name redis -p 6379:6379 redis:7-alpinemake run
# or:
uvicorn app.main:app --host 0.0.0.0 --port 8000 --reloadThe service is now available at http://localhost:8000.
# Full suite
make test
# With coverage report
make test-cov
# Unit tests only
make test-unit
# Integration tests only
make test-integrationTests do not require a live Redis instance — all Redis I/O is intercepted by AsyncMock fixtures defined in tests/conftest.py.
make docker-build
# or:
docker build -t aether-core-orchestrator:latest .make docker-run
# or:
docker run --rm \
-p 8000:8000 \
-e REDIS_URL=redis://host.docker.internal:6379 \
aether-core-orchestrator:latestThe root-level docker-compose.yml in the aether/ workspace starts the complete platform including Redis, the API gateway, and all worker services.
# From the workspace root:
docker compose up -d| Metric | Type | Labels | Description |
|---|---|---|---|
aether_http_requests_total |
Counter | method, path, status_code |
Total HTTP requests |
aether_http_request_duration_seconds |
Histogram | method, path |
HTTP request latency |
aether_websocket_connections_active |
Gauge | — | Active WebSocket sessions |
aether_websocket_connections_total |
Counter | — | Cumulative WebSocket connections |
aether_frames_processed_total |
Counter | session_id |
Successfully processed frames |
aether_frames_rejected_total |
Counter | error_code |
Rejected frames by error type |
Metrics are available at GET /metrics in the Prometheus text exposition format.
# Lint (check only)
make lint
# Format + auto-fix
make formatThe project uses ruff for both linting and formatting (target: Python 3.11, line length: 99).
aether-core-orchestrator/
├── app/
│ ├── main.py # FastAPI app factory + lifespan
│ ├── api/v1/
│ │ ├── router.py # v1 APIRouter aggregator
│ │ └── endpoints/
│ │ ├── health.py # GET /v1/health
│ │ └── stream.py # WS /v1/mesh/stream/{session_id}
│ ├── core/
│ │ ├── config.py # Pydantic BaseSettings singleton
│ │ └── logging.py # structlog JSON/console setup
│ ├── models/
│ │ └── context.py # ContextFrame, Payload, AckResponse, …
│ ├── services/
│ │ ├── connection_manager.py # Thread-safe WebSocket registry
│ │ └── frame_processor.py # Validate → Redis persistence
│ └── middleware/
│ └── telemetry.py # Prometheus middleware + /metrics
├── tests/
│ ├── conftest.py # Shared fixtures (mock_redis, async_client, …)
│ ├── unit/test_models.py # Domain model unit tests
│ └── integration/test_endpoints.py # HTTP + WebSocket integration tests
├── Dockerfile # Multi-stage production image
├── requirements.txt # Production dependencies
├── requirements-dev.txt # Dev/test dependencies
├── pyproject.toml # Build system, pytest, coverage, ruff config
└── Makefile # Developer convenience targets