This is the diagram for the prototype analytic engine data map.
Note: This is AI generated and at this time, not fully reviewed for production
flowchart LR
A([Analytics Engine App<br/>Rust + Axum + SeaORM])
W1[Calculation Worker 1]
W2[Calculation Worker 2]
R1[(Redis Primary Cache)]
R2[(Redis Replica Cache)]
C1[(Citus Coordinator)]
C2[(Citus Coordinator Standby)]
P1[(Postgres / Citus Worker 1)]
P2[(Postgres / Citus Worker 2)]
P3[(Postgres / Citus Worker 3)]
A -->|API requests / DB queries| C1
A -->|Cache reads / writes| R1
W1 -->|Analytical jobs| C1
W2 -->|Analytical jobs| C1
W1 -->|Progress / cache| R1
W2 -->|Progress / cache| R1
R1 -->|Replication| R2
C1 -->|Streaming replication| C2
C1 -->|Shard queries| P1
C1 -->|Shard queries| P2
C1 -->|Shard queries| P3
classDef app fill:#f7dc6f,stroke:#7d6608,stroke-width:2px,color:#000;
classDef worker fill:#7dcea0,stroke:#145a32,stroke-width:2px,color:#000;
classDef redis fill:#f5b041,stroke:#a04000,stroke-width:2px,color:#000;
classDef db fill:#85c1e9,stroke:#1b4f72,stroke-width:2px,color:#000;
class A app;
class W1,W2 worker;
class R1,R2 redis;
class C1,C2,P1,P2,P3 db;
flowchart LR
S[(Ingest Data Warehouse)]
A[Analytics Engine App<br/>Axum API<br/>Calculation Worker]
C[(Analytics Database<br/>Citus Coordinator)]
W1[(Citus Worker 1<br/>Warehouse Shards)]
W2[(Citus Worker 2<br/>Warehouse Shards)]
S -->|Periodic versioned export| C
A -->|SeaORM / SQL| C
C -->|Distribute imported rows| W1
C -->|Distribute imported rows| W2
C -->|Distributed queries| W1
C -->|Distributed queries| W2
W1 -->|Partial results| C
W2 -->|Partial results| C
C -->|Combined result| A
classDef app fill:#f7dc6f,stroke:#7d6608,stroke-width:2px,color:#000;
classDef db fill:#85c1e9,stroke:#1b4f72,stroke-width:2px,color:#000;
class A app;
class S,C,W1,W2 db;
Cloudflair, CDN, etc.
Multiple servers will ultimately fill this role
HAProxy and the regional router can ultimately occupy different servers.
flowchart LR
EDGE["Managed Edge"]
subgraph ROUTER["One Ubuntu 24.04 Routing Server"]
HAP["HAProxy<br/>0.0.0.0:443"]
APP["mt-regional-router<br/>127.0.0.1:8080"]
HAP --> APP
end
AUTH["Authentication Service"]
GAME["Regional Game Services"]
TELEMETRY["Regional Telemetry Service"]
EDGE --> HAP
APP --> AUTH
APP --> GAME
APP --> TELEMETRY
flowchart TB
%% =========================================================
%% CLIENT / EXTERNAL ACCESS
%% =========================================================
subgraph CLIENTS["Clients and Administrative Access"]
WEB["Web / Android Clients"]
ADMIN["Administrator VPN / Bastion Host"]
ETL["Authorized Data Import Processes"]
end
%% =========================================================
%% EDGE / APPLICATION ENTRY
%% =========================================================
subgraph EDGE["Application Entry Tier"]
DNS["Internal / Public DNS"]
APP_LB["Application Load Balancer<br/>HAProxy or Nginx<br/>HTTPS :443"]
end
WEB -->|"HTTPS :443"| DNS
DNS --> APP_LB
%% =========================================================
%% APPLICATION TIER
%% =========================================================
subgraph APP_TIER["MT Analytics Engine Application Tier"]
API1["Analytics Engine API 1<br/>Rust + Axum + SeaORM"]
API2["Analytics Engine API 2<br/>Rust + Axum + SeaORM"]
WORKER1["Calculation Worker 1<br/>Rust + SeaORM"]
WORKER2["Calculation Worker 2<br/>Rust + SeaORM"]
WORKERN["Calculation Worker N<br/>Rust + SeaORM"]
SCHEDULER["Optional Scheduler / Dispatcher<br/>Starts periodic or queued work"]
end
APP_LB -->|"HTTP :8080"| API1
APP_LB -->|"HTTP :8080"| API2
SCHEDULER -->|"Create calculation jobs"| API1
SCHEDULER -->|"Create calculation jobs"| API2
%% =========================================================
%% REDIS CACHE TIER
%% =========================================================
subgraph CACHE_TIER["Redis Cache and Coordination Tier"]
REDIS_ENDPOINT["Redis Endpoint<br/>redis.internal :6379"]
REDIS_PRIMARY["Redis Primary"]
REDIS_REPLICA["Redis Replica"]
REDIS_SENTINEL["Redis Sentinel Quorum<br/>Optional HA discovery"]
end
REDIS_ENDPOINT --> REDIS_PRIMARY
REDIS_PRIMARY -.->|"Asynchronous replication"| REDIS_REPLICA
REDIS_SENTINEL -.->|"Health and role monitoring"| REDIS_PRIMARY
REDIS_SENTINEL -.->|"Health and role monitoring"| REDIS_REPLICA
API1 -->|"Cache reads/writes<br/>TCP :6379"| REDIS_ENDPOINT
API2 -->|"Cache reads/writes<br/>TCP :6379"| REDIS_ENDPOINT
WORKER1 -->|"Progress, locks, invalidation<br/>TCP :6379"| REDIS_ENDPOINT
WORKER2 -->|"Progress, locks, invalidation<br/>TCP :6379"| REDIS_ENDPOINT
WORKERN -->|"Progress, locks, invalidation<br/>TCP :6379"| REDIS_ENDPOINT
%% =========================================================
%% POSTGRESQL CONNECTION ENTRY
%% =========================================================
subgraph DB_ENTRY["Database Connection and Routing Tier"]
DB_DNS["warehouse-db.internal"]
DB_LB1["Database Router 1<br/>HAProxy / Keepalived"]
DB_LB2["Database Router 2<br/>HAProxy / Keepalived"]
DB_VIP["Database Virtual IP<br/>PostgreSQL :6432"]
PGBOUNCER1["PgBouncer 1<br/>Connection Pool"]
PGBOUNCER2["PgBouncer 2<br/>Connection Pool"]
end
DB_DNS --> DB_VIP
DB_VIP --> DB_LB1
DB_VIP --> DB_LB2
DB_LB1 -->|"Healthy active coordinator"| PGBOUNCER1
DB_LB2 -->|"Healthy active coordinator"| PGBOUNCER2
API1 -->|"SeaORM / PostgreSQL<br/>TCP :6432"| DB_DNS
API2 -->|"SeaORM / PostgreSQL<br/>TCP :6432"| DB_DNS
WORKER1 -->|"SeaORM / analytical SQL<br/>TCP :6432"| DB_DNS
WORKER2 -->|"SeaORM / analytical SQL<br/>TCP :6432"| DB_DNS
WORKERN -->|"SeaORM / analytical SQL<br/>TCP :6432"| DB_DNS
ETL -->|"Controlled warehouse load<br/>TCP :6432"| DB_DNS
%% =========================================================
%% CITUS COORDINATOR TIER
%% =========================================================
subgraph CITUS_COORDINATORS["Citus Coordinator Tier"]
COORD_PRIMARY["Citus Coordinator Primary<br/>PostgreSQL + Citus<br/>Metadata and query planning"]
COORD_STANDBY["Citus Coordinator Standby<br/>PostgreSQL streaming replica"]
FAILOVER["PostgreSQL HA Manager<br/>Patroni / repmgr / equivalent"]
end
PGBOUNCER1 -->|"PostgreSQL :5432"| COORD_PRIMARY
PGBOUNCER2 -->|"PostgreSQL :5432"| COORD_PRIMARY
COORD_PRIMARY -.->|"WAL streaming replication"| COORD_STANDBY
FAILOVER -.->|"Health, fencing and promotion"| COORD_PRIMARY
FAILOVER -.->|"Health, fencing and promotion"| COORD_STANDBY
%% =========================================================
%% CITUS WORKER SHARD TIER
%% =========================================================
subgraph CITUS_WORKERS["Citus Worker Shard Tier"]
subgraph WORKER_GROUP_1["Citus Worker Group 1"]
PGW1["Worker 1 Primary<br/>Shard Set A"]
PGW1R["Worker 1 Standby<br/>Replica of Shard Set A"]
end
subgraph WORKER_GROUP_2["Citus Worker Group 2"]
PGW2["Worker 2 Primary<br/>Shard Set B"]
PGW2R["Worker 2 Standby<br/>Replica of Shard Set B"]
end
subgraph WORKER_GROUP_3["Citus Worker Group 3"]
PGW3["Worker 3 Primary<br/>Shard Set C"]
PGW3R["Worker 3 Standby<br/>Replica of Shard Set C"]
end
subgraph WORKER_GROUP_N["Additional Worker Groups"]
PGWN["Worker N Primary<br/>Additional shards"]
PGWNR["Worker N Standby"]
end
end
COORD_PRIMARY -->|"Distributed SQL<br/>PostgreSQL :5432"| PGW1
COORD_PRIMARY -->|"Distributed SQL<br/>PostgreSQL :5432"| PGW2
COORD_PRIMARY -->|"Distributed SQL<br/>PostgreSQL :5432"| PGW3
COORD_PRIMARY -->|"Distributed SQL<br/>PostgreSQL :5432"| PGWN
PGW1 -.->|"WAL replication"| PGW1R
PGW2 -.->|"WAL replication"| PGW2R
PGW3 -.->|"WAL replication"| PGW3R
PGWN -.->|"WAL replication"| PGWNR
FAILOVER -.->|"Worker health and promotion"| PGW1
FAILOVER -.->|"Worker health and promotion"| PGW2
FAILOVER -.->|"Worker health and promotion"| PGW3
FAILOVER -.->|"Worker health and promotion"| PGWN
%% =========================================================
%% DATA DISTRIBUTION
%% =========================================================
subgraph DATA_MODEL["Logical Data Placement"]
LOCAL_TABLES["Coordinator-local tables<br/>analysis_job<br/>analysis_run<br/>users and configuration"]
REF_TABLES["Citus reference tables<br/>industry codes<br/>geography types<br/>period definitions"]
FACT_TABLES["Distributed warehouse facts<br/>Census / BLS / market facts<br/>Sharded by selected distribution key"]
RESULT_TABLES["Distributed TAM worksets and results<br/>Sharded by analysis_run_id"]
end
COORD_PRIMARY --- LOCAL_TABLES
COORD_PRIMARY --- REF_TABLES
PGW1 --- FACT_TABLES
PGW2 --- FACT_TABLES
PGW3 --- FACT_TABLES
PGWN --- FACT_TABLES
PGW1 --- RESULT_TABLES
PGW2 --- RESULT_TABLES
PGW3 --- RESULT_TABLES
PGWN --- RESULT_TABLES
%% =========================================================
%% BACKUP / OBSERVABILITY
%% =========================================================
subgraph OPERATIONS["Operations, Monitoring and Recovery"]
PROM["Prometheus"]
GRAFANA["Grafana"]
LOGS["Central Logging<br/>Loki / OpenSearch / equivalent"]
BACKUP["pgBackRest Repository<br/>Backup storage / object storage"]
ALERTS["Alert Manager"]
end
PROM -.->|"Metrics"| API1
PROM -.->|"Metrics"| API2
PROM -.->|"Metrics"| WORKER1
PROM -.->|"Metrics"| WORKER2
PROM -.->|"Metrics"| REDIS_PRIMARY
PROM -.->|"Metrics"| COORD_PRIMARY
PROM -.->|"Metrics"| PGW1
PROM -.->|"Metrics"| PGW2
PROM -.->|"Metrics"| PGW3
GRAFANA --> PROM
ALERTS --> PROM
API1 -.->|"Application logs"| LOGS
API2 -.->|"Application logs"| LOGS
WORKER1 -.->|"Worker logs"| LOGS
WORKER2 -.->|"Worker logs"| LOGS
COORD_PRIMARY -.->|"PostgreSQL logs"| LOGS
COORD_PRIMARY -.->|"Base backups and WAL archive"| BACKUP
COORD_STANDBY -.->|"Optional backup source"| BACKUP
PGW1 -.->|"Base backups and WAL archive"| BACKUP
PGW2 -.->|"Base backups and WAL archive"| BACKUP
PGW3 -.->|"Base backups and WAL archive"| BACKUP
PGWN -.->|"Base backups and WAL archive"| BACKUP
%% =========================================================
%% ADMINISTRATION
%% =========================================================
ADMIN -->|"SSH :22 through restricted network"| API1
ADMIN -->|"SSH :22 through restricted network"| WORKER1
ADMIN -->|"SSH :22 through restricted network"| DB_LB1
ADMIN -->|"SSH :22 through restricted network"| COORD_PRIMARY
ADMIN -->|"SSH :22 through restricted network"| PGW1
### Sequence
sequenceDiagram
autonumber
participant Client
participant API as Axum Analytics API
participant Redis
participant Coordinator as Citus Coordinator
participant Queue as PostgreSQL Job Queue
participant Worker as Rust Calculation Worker
participant Shards as Citus Worker Shards
Client->>API: POST /v1/analysis-runs
API->>Redis: Check duplicate/idempotency key
alt Existing completed calculation
Redis-->>API: Cached result
API-->>Client: 200 OK with result
else New calculation
API->>Coordinator: Begin transaction
Coordinator->>Queue: Insert analysis_run
Coordinator->>Queue: Insert pending analysis_job
Coordinator-->>API: Commit
API->>Redis: Store pending status
API-->>Client: 202 Accepted + analysis_run_id
end
loop Claim available jobs
Worker->>Coordinator: SELECT ... FOR UPDATE SKIP LOCKED
Coordinator->>Queue: Lease pending job
Queue-->>Worker: Job definition
end
Worker->>Coordinator: Execute staged TAM query
Coordinator->>Shards: Dispatch shard-local query fragments
Shards-->>Coordinator: Partial aggregates
Coordinator-->>Worker: Combined result
Worker->>Coordinator: Store durable result rows
Worker->>Coordinator: Mark job completed
Worker->>Redis: Cache summary and completed status
Client->>API: GET /v1/analysis-runs/{id}
API->>Redis: Read cached status/result
alt Cache hit
Redis-->>API: Completed summary
else Cache miss
API->>Coordinator: Read durable run and result
Coordinator-->>API: Run status and result
API->>Redis: Populate cache
end
API-->>Client: Result or current progress