21 KiB
Future Improvements
Broadway + Oban: High-Volume Metric Ingestion
Status: Planning
Target: When reaching 100+ agents or experiencing API bottlenecks Effort: Medium (2-3 weeks) Priority: Medium (not urgent with current scale)
Problem Statement
Current architecture uses synchronous agent metric submission:
- Agent POSTs metrics → API validates → Direct DB insert → 200 OK
- At scale (100s of agents, 1000s of devices), this creates:
- Thundering herd: All agents polling on similar intervals
- API timeouts: DB writes block the HTTP response
- Poor throughput: One DB transaction per agent request
- No backpressure: Peaks overwhelm the system
Current Performance Profile
500 agents × 60s interval = 8.3 req/sec average
Peak (synchronized): 500 concurrent requests
Each request: 10-50ms (DB insert + processing)
Risk: Thundering herd, timeouts during peaks
Proposed Solution: Broadway Message Processing Pipeline
Architecture Overview
┌─────────────────────────────────────────────────────────────────┐
│ CURRENT (Oban) │
├─────────────────────────────────────────────────────────────────┤
│ DevicePollerWorker (Oban) ──┐ │
│ DeviceMonitorWorker (Oban) ─┼─→ Direct SNMP → DB writes │
│ Maintenance Workers (Oban) ─┘ │
└─────────────────────────────────────────────────────────────────┘
┌─────────────────────────────────────────────────────────────────┐
│ WITH BROADWAY (Proposed) │
├─────────────────────────────────────────────────────────────────┤
│ Agent POST /metrics ──→ Message Queue ──→ 202 Accepted │
│ ↓ │
│ Broadway Pipeline │
│ ├─ Consumer (50 workers) │
│ ├─ Batcher (1000 msgs) │
│ └─ Bulk DB Insert │
│ │
│ Oban continues: │
│ ├─ DevicePollerWorker (scheduled SNMP) │
│ ├─ DeviceMonitorWorker (health checks) │
│ └─ Maintenance workers (cleanup, etc.) │
└─────────────────────────────────────────────────────────────────┘
Role Separation
Broadway (async, high-volume, batchable):
- ✅ Agent metric submissions (
POST /api/v1/agent/metrics) - ✅ Agent heartbeats (optional - could batch these too)
- ✅ High-throughput data ingestion
- ✅ Batch database operations
Oban (scheduled, low-frequency, complex logic):
- ✅ DevicePollerWorker (scheduled SNMP polling)
- ✅ DeviceMonitorWorker (health checks)
- ✅ Maintenance workers (cleanup, evaluators)
- ✅ Discovery jobs (complex multi-step operations)
Performance Characteristics After Broadway
Agent POST: <5ms (just enqueue)
Broadway pipeline:
- 50 processor workers (parse/validate)
- 10 batcher workers (1000 rows/batch)
- Single INSERT for 1000 rows: ~50-100ms
Throughput: 10k-50k metrics/sec (limited by DB, not API)
Backpressure: Queue absorbs burst traffic
Message Queue Options
Option A: PostgreSQL-backed (Recommended Start)
Pros:
- ✅ No new infrastructure (use existing PostgreSQL)
- ✅ Transactional guarantees
- ✅ Simple deployment
- ✅ Easy development/testing
Cons:
- ❌ Not ideal for extremely high throughput (10k+ msg/sec)
- ❌ Adds load to primary database
Implementation: Custom Broadway producer reading from agent_metrics_queue table
When to use: Initial implementation, <100 agents, proof of concept
Option B: RabbitMQ (Recommended Production)
Pros:
- ✅ Battle-tested for high throughput
- ✅ Built-in clustering and persistence
- ✅ Management UI for monitoring
- ✅ Proven at scale (100k+ msg/sec)
Cons:
- ❌ Adds infrastructure dependency
- ❌ Another service to monitor
- ❌ Requires Kubernetes setup
Implementation: broadway_rabbitmq package with existing RabbitMQ cluster
When to use: 100+ agents, need guaranteed throughput, production workloads
Option C: AWS SQS
Pros:
- ✅ Fully managed, scales automatically
- ✅ No infrastructure to maintain
- ✅ Pay-per-use pricing
Cons:
- ❌ Requires AWS
- ❌ Potential latency (network calls)
- ❌ Visibility timeout complexity
Implementation: off_broadway_sqs package
When to use: Running on AWS, want fully managed solution
Implementation Details
1. Dependencies
# mix.exs
defp deps do
[
{:broadway, "~> 1.0"},
# Choose one:
{:broadway_rabbitmq, "~> 0.8"}, # Option B
{:off_broadway_sqs, "~> 0.7"}, # Option C
# Option A: custom producer (see below)
# ... existing deps
]
end
2. Broadway Pipeline
# lib/towerops/pipelines/agent_metrics_pipeline.ex
defmodule Towerops.Pipelines.AgentMetricsPipeline do
use Broadway
alias Broadway.Message
alias Towerops.Snmp
def start_link(_opts) do
Broadway.start_link(__MODULE__,
name: __MODULE__,
producer: [
module: {
BroadwayRabbitMQ.Producer,
queue: "agent_metrics",
connection: [
host: Application.get_env(:towerops, :rabbitmq_host),
username: Application.get_env(:towerops, :rabbitmq_username),
password: Application.get_env(:towerops, :rabbitmq_password)
],
on_failure: :reject_and_requeue
},
concurrency: 1
],
processors: [
default: [
concurrency: 50 # Parallel message processing
]
],
batchers: [
default: [
batch_size: 1000, # Max messages per batch
batch_timeout: 2_000, # Max wait time (ms)
concurrency: 10 # Parallel batch processors
]
]
)
end
@impl true
def handle_message(_, %Message{data: data} = message, _context) do
# Decode and validate the message
# No DB writes here - just prepare data
metric = Jason.decode!(data)
message
|> Message.update_data(fn _ -> metric end)
|> Message.put_batcher(:default)
end
@impl true
def handle_batch(:default, messages, _batch_info, _context) do
# Extract all metrics from the batch
metrics = Enum.map(messages, & &1.data)
# Single multi-row insert for entire batch
case Snmp.insert_metrics_batch(metrics) do
{:ok, _} ->
messages # All successful
{:error, reason} ->
# Mark all messages as failed (will be retried)
Enum.map(messages, &Message.failed(&1, reason))
end
end
end
3. Batch Insert Function
# lib/towerops/snmp.ex
defmodule Towerops.Snmp do
alias Towerops.Repo
alias Towerops.Snmp.SensorReading
@doc """
Insert a batch of sensor readings in a single transaction.
Expected input format from agents:
[
%{
"equipment_id" => "uuid",
"sensor_id" => "uuid",
"value" => 42.5,
"timestamp" => "2026-02-01T12:00:00Z"
},
...
]
"""
def insert_metrics_batch(metrics) when is_list(metrics) do
now = DateTime.utc_now()
# Transform to Ecto schema format
entries = Enum.map(metrics, fn metric ->
%{
equipment_id: metric["equipment_id"],
sensor_id: metric["sensor_id"],
value: Decimal.new(metric["value"]),
timestamp: parse_timestamp(metric["timestamp"]),
inserted_at: now,
updated_at: now
}
end)
# Single INSERT with multiple rows
Repo.insert_all(SensorReading, entries,
on_conflict: :nothing,
returning: false
)
end
defp parse_timestamp(ts) when is_binary(ts) do
case DateTime.from_iso8601(ts) do
{:ok, dt, _} -> dt
_ -> DateTime.utc_now()
end
end
end
4. Async API Endpoint
# lib/towerops_web/controllers/api/v1/agent_controller.ex
defmodule ToweropsWeb.API.V1.AgentController do
use ToweropsWeb, :controller
require Logger
# Current synchronous version (keep during transition)
def metrics_sync(conn, %{"metrics" => metrics}) do
# ... existing implementation
end
# New async version with Broadway
def metrics_async(conn, %{"metrics" => metrics}) do
agent = conn.assigns.current_agent
# Quick validation only
with :ok <- validate_agent_active(agent),
:ok <- validate_metrics_schema(metrics) do
# Enqueue to message queue (non-blocking)
message = Jason.encode!(%{
agent_id: agent.id,
organization_id: agent.organization_id,
metrics: metrics,
received_at: DateTime.utc_now()
})
# This is fast - just writes to queue
:ok = publish_to_queue("agent_metrics", message)
# Return immediately
conn
|> put_status(:accepted)
|> json(%{status: "accepted", queued: length(metrics)})
else
{:error, reason} ->
conn
|> put_status(:bad_request)
|> json(%{error: reason})
end
end
defp publish_to_queue(queue, message) do
# RabbitMQ example
AMQP.Basic.publish(
ToweropsWeb.RabbitMQ.channel(),
"",
queue,
message,
persistent: true
)
end
defp validate_agent_active(%{active: true}), do: :ok
defp validate_agent_active(_), do: {:error, "agent inactive"}
defp validate_metrics_schema(metrics) when is_list(metrics) do
# Quick schema validation only
if Enum.all?(metrics, &valid_metric?/1) do
:ok
else
{:error, "invalid metric schema"}
end
end
defp valid_metric?(%{"equipment_id" => _, "sensor_id" => _, "value" => _}), do: true
defp valid_metric?(_), do: false
end
5. Application Supervision Tree
# lib/towerops/application.ex
defmodule Towerops.Application do
def start(_type, _args) do
children = [
# ... existing children (Repo, PubSub, Endpoint, Oban, etc.)
# Add Broadway pipeline
Towerops.Pipelines.AgentMetricsPipeline,
]
opts = [strategy: :one_for_one, name: Towerops.Supervisor]
Supervisor.start_link(children, opts)
end
end
6. Router Updates
# lib/towerops_web/router.ex
scope "/api/v1/agent", ToweropsWeb.API.V1 do
pipe_through [:api, :api_agent_auth]
# Phase 1: Dual endpoints during migration
post "/metrics", AgentController, :metrics_sync # Current
post "/metrics/async", AgentController, :metrics_async # New
# Phase 3: After full migration
# post "/metrics", AgentController, :metrics_async
end
Migration Path
Phase 1: Dual Mode (Parallel Run)
Duration: 2-4 weeks Goal: Validate Broadway pipeline with real traffic
- Deploy both endpoints (
/metricsand/metrics/async) - Update 10% of agents to use
/metrics/async - Monitor:
- Queue depth
- Processing latency
- Error rates
- Database performance
Success Criteria:
- Zero data loss
- <100ms p95 latency for queue insertion
- <5s p95 latency for batch processing
- No queue backlog growth
Phase 2: Gradual Rollout
Duration: 2-4 weeks Goal: Full traffic migration
- Increase to 25% of agents
- Monitor for 1 week
- Increase to 50% of agents
- Monitor for 1 week
- Increase to 100% of agents
Rollback Plan:
- Keep sync endpoint active
- Simple config change in agent to switch back
- No data loss (queue persists messages)
Phase 3: Deprecate Synchronous Endpoint
Duration: 1 week Goal: Clean up legacy code
- Switch
/metricsroute to async handler - Monitor for 1 week
- Remove
metrics_sync/2function - Update documentation
Infrastructure Requirements
PostgreSQL-backed (Phase 1)
Database Schema:
CREATE TABLE agent_metrics_queue (
id BIGSERIAL PRIMARY KEY,
payload JSONB NOT NULL,
inserted_at TIMESTAMP NOT NULL DEFAULT NOW(),
processed_at TIMESTAMP,
retry_count INTEGER DEFAULT 0,
error TEXT
);
CREATE INDEX idx_agent_metrics_queue_unprocessed
ON agent_metrics_queue (inserted_at)
WHERE processed_at IS NULL;
No additional services required
RabbitMQ (Production)
Development (docker-compose.yml):
services:
rabbitmq:
image: rabbitmq:3.13-management-alpine
ports:
- "5672:5672" # AMQP
- "15672:15672" # Management UI
environment:
RABBITMQ_DEFAULT_USER: towerops
RABBITMQ_DEFAULT_PASS: ${RABBITMQ_PASSWORD}
volumes:
- rabbitmq_data:/var/lib/rabbitmq
volumes:
rabbitmq_data:
Production (Kubernetes):
# k8s/rabbitmq/helmrelease.yaml
apiVersion: helm.toolkit.fluxcd.io/v2beta1
kind: HelmRelease
metadata:
name: rabbitmq
namespace: towerops
spec:
chart:
spec:
chart: rabbitmq
version: 12.x.x
sourceRef:
kind: HelmRepository
name: bitnami
values:
auth:
username: towerops
existingPasswordSecret: rabbitmq-password
persistence:
enabled: true
size: 20Gi
clustering:
enabled: true
replicaCount: 3
metrics:
enabled: true
serviceMonitor:
enabled: true
Secrets (1Password → Kubernetes):
# Create secret from 1Password
kubectl create secret generic rabbitmq-password \
--from-literal=rabbitmq-password=$(op read "op://towerops/rabbitmq/password") \
-n towerops
Monitoring & Observability
Telemetry Metrics
# lib/towerops_web/telemetry.ex
defp periodic_measurements do
[
# Existing metrics...
# Broadway pipeline metrics
{Broadway, :metrics, [Towerops.Pipelines.AgentMetricsPipeline]},
# Queue depth (RabbitMQ)
{ToweropsWeb.Metrics, :queue_depth, ["agent_metrics"]},
# Processing rate
{ToweropsWeb.Metrics, :processing_rate, []},
]
end
LiveDashboard Additions
Add to /admin/dashboard:
- Queue Depth - Messages waiting to be processed
- Processing Rate - Messages/second throughput
- Batch Sizes - Average messages per batch
- Failed Messages - Retry queue depth
- Latency - p50/p95/p99 processing time
Alerts
Add to monitoring system (Prometheus/Grafana):
- Queue depth > 10,000 messages (backlog warning)
- Processing rate < 100 msg/sec (pipeline degraded)
- Failed message rate > 5% (data quality issue)
- Consumer crashes (Broadway supervisor restarts)
Testing Strategy
Unit Tests
# test/towerops/pipelines/agent_metrics_pipeline_test.exs
defmodule Towerops.Pipelines.AgentMetricsPipelineTest do
use Towerops.DataCase
alias Broadway.Message
alias Towerops.Pipelines.AgentMetricsPipeline
describe "handle_message/3" do
test "parses valid metric JSON" do
data = Jason.encode!(%{
"equipment_id" => "uuid",
"sensor_id" => "uuid",
"value" => 42.5
})
message = %Message{data: data}
result = AgentMetricsPipeline.handle_message(:default, message, %{})
assert result.batcher == :default
assert result.data["value"] == 42.5
end
end
describe "handle_batch/4" do
test "inserts batch of metrics" do
messages = [
%Message{data: %{"equipment_id" => "uuid1", ...}},
%Message{data: %{"equipment_id" => "uuid2", ...}}
]
result = AgentMetricsPipeline.handle_batch(:default, messages, %{}, %{})
assert length(result) == 2
assert Enum.all?(result, &match?(%Message{status: :ok}, &1))
end
end
end
Integration Tests
# test/towerops_web/controllers/api/v1/agent_controller_test.exs
describe "POST /api/v1/agent/metrics/async" do
test "enqueues metrics and returns 202", %{conn: conn, agent: agent} do
metrics = [%{
"equipment_id" => equipment.id,
"sensor_id" => sensor.id,
"value" => 42.5,
"timestamp" => DateTime.utc_now() |> DateTime.to_iso8601()
}]
conn = post(conn, ~p"/api/v1/agent/metrics/async", %{"metrics" => metrics})
assert %{"status" => "accepted", "queued" => 1} = json_response(conn, 202)
# Verify message in queue
assert_queue_length("agent_metrics", 1)
end
end
Load Tests
# test/load/agent_metrics_load_test.exs
# Use :fuse or similar to simulate 500 concurrent agents
defmodule Towerops.Load.AgentMetricsTest do
use ExUnit.Case
@tag timeout: :infinity
@tag :load_test
test "handles 500 concurrent agent submissions" do
agents = create_agents(500)
# Simulate all agents posting at once
tasks = Enum.map(agents, fn agent ->
Task.async(fn ->
submit_metrics(agent, generate_metrics(10))
end)
end)
results = Task.await_many(tasks, 30_000)
# All should succeed
assert Enum.all?(results, &match?({:ok, _}, &1))
# Verify queue processed everything
wait_for_queue_empty("agent_metrics", timeout: 60_000)
# Verify all metrics in database
assert Repo.aggregate(SensorReading, :count) == 5000
end
end
Performance Tuning
Database Optimization
Batch Insert Settings:
# config/runtime.exs
config :towerops, Towerops.Repo,
# Increase connection pool for Broadway workers
pool_size: 20, # Up from 10
# Optimize for batch inserts
prepare: :unnamed,
timeout: 30_000
TimescaleDB Tuning:
-- Increase batch size for continuous aggregates
ALTER MATERIALIZED VIEW sensor_readings_1h SET (
timescaledb.materialized_only = true,
timescaledb.compress_after = '7 days'
);
-- Add indexes for common queries
CREATE INDEX CONCURRENTLY idx_sensor_readings_equipment_timestamp
ON sensor_readings (equipment_id, timestamp DESC);
Broadway Tuning
# Processors: CPU-bound (parsing, validation)
# Start with: processors_concurrency = num_cores * 2
processors: [default: [concurrency: System.schedulers_online() * 2]]
# Batchers: I/O-bound (database writes)
# Start with: batch_size = 1000, tune based on DB performance
batchers: [
default: [
batch_size: 1000, # Increase if DB can handle larger batches
batch_timeout: 2_000, # Decrease for lower latency
concurrency: 10 # Increase if batchers become bottleneck
]
]
RabbitMQ Tuning
# Prefetch: How many unacked messages per consumer
# Higher = better throughput, lower = better distribution
connection: [
prefetch_count: 100 # Start with 100, tune based on message size
]
Cost Analysis
Infrastructure Costs (Estimated)
PostgreSQL-backed:
- $0/month (uses existing database)
- Adds ~10-20% CPU load to DB server
RabbitMQ (Self-hosted Kubernetes):
- 3 nodes × 2 vCPU × 4 GB RAM = ~$120/month
- 20 GB persistent storage = ~$10/month
- Total: ~$130/month
AWS SQS:
- 1M requests = $0.40
- 500 agents × 1 req/min × 60 min × 24 hrs × 30 days = 21.6M requests/month
- Total: ~$8.64/month (plus data transfer costs)
Development Costs
- Phase 1 (PostgreSQL-backed): 1 week development + 1 week testing
- Phase 2 (RabbitMQ migration): 1 week development + 1 week testing
- Phase 3 (Cleanup): 2-3 days
Total: 4-5 weeks engineering time
Decision Points
When to Implement?
Implement Broadway when you hit any of these thresholds:
- Agent Count: 100+ agents (current: <10)
- API Latency: p95 > 500ms for
/api/v1/agent/metrics - Timeout Rate: >1% of agent requests timing out
- Database Load: DB CPU > 70% during agent polling windows
- Business Need: SLA requires <100ms API response time
Which Message Queue?
| Criteria | PostgreSQL | RabbitMQ | AWS SQS |
|---|---|---|---|
| Agents | <100 | 100-1000 | Any |
| Throughput | <1k msg/sec | <100k msg/sec | Unlimited |
| Infrastructure | None | Kubernetes | AWS only |
| Ops Complexity | Low | Medium | Low |
| Cost | $0 | ~$130/mo | ~$9/mo |
Recommended Path:
- Start with PostgreSQL-backed (proof of concept)
- Migrate to RabbitMQ when hitting 100 agents
- Consider SQS only if already on AWS
References
- Broadway Documentation
- Building Custom Producers
- Broadway RabbitMQ
- OffBroadway SQS
- Batch Insert Performance
Related Documents
CLAUDE.md- Project-specific development guidelinesAGENTS.md- Elixir/Phoenix/LiveView patternsdocs/API.md- Agent API documentation