towerops/docs/FUTURE_IMPROVEMENTS.md

21 KiB
Raw Blame History

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

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


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

  1. Deploy both endpoints (/metrics and /metrics/async)
  2. Update 10% of agents to use /metrics/async
  3. 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

  1. Increase to 25% of agents
  2. Monitor for 1 week
  3. Increase to 50% of agents
  4. Monitor for 1 week
  5. 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

  1. Switch /metrics route to async handler
  2. Monitor for 1 week
  3. Remove metrics_sync/2 function
  4. 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:

  1. Agent Count: 100+ agents (current: <10)
  2. API Latency: p95 > 500ms for /api/v1/agent/metrics
  3. Timeout Rate: >1% of agent requests timing out
  4. Database Load: DB CPU > 70% during agent polling windows
  5. 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:

  1. Start with PostgreSQL-backed (proof of concept)
  2. Migrate to RabbitMQ when hitting 100 agents
  3. Consider SQS only if already on AWS

References


  • CLAUDE.md - Project-specific development guidelines
  • AGENTS.md - Elixir/Phoenix/LiveView patterns
  • docs/API.md - Agent API documentation