Skip to main content
Kenneth Nnorom
Project SeriesPart 4 of 4

Lightweight Network SLA Monitoring Platform (LNMP)

View Series Hub
Series Sequence & Navigation:
Technical LabNetwork Automation•

Building LNMP v3.0: Enterprise Architecture, Real-Time SSE Streams, and Dual-Driver Storage

How production scaling and improved Python architectural practices led to LNMP v3.0: deterministic 32s/28s sweeper concurrency, typed SQLAlchemy 2.0 repositories, Server-Sent Events, dual PostgreSQL/Redis storage, and multi-protocol synthetics.

Environment Specifications
Tools & Utilities
Python 3FastAPISQLAlchemy 2.0TimescaleDBRedisVue 3SSE
Operating Systems
Linux

In the previous parts of this series, we developed the Lightweight Network Monitoring Platform (LNMP) and refined its topology layering and time-series compression engines. While Version 2.0 stabilized baseline operations, running the platform under live multi-operator conditions exposed new architectural friction points: probe sweeper collisions at minute boundaries, 5-second polling database overhead on NOC screens, and tight coupling with raw SQL strings. At the same time, as I continued sharpening my Python programming skills and studying modern software architecture patterns, I realized there were cleaner, more decoupled ways to structure our background services and data models that were not obvious to me during the initial build of Version 2.0. LNMP Version 3.0 refactors the system into a featherlight monolith with hardware-adaptive toggles, featuring deterministic 32s/28s sweeper concurrency, typed SQLAlchemy 2.0 repositories, Server-Sent Events (SSE), dual PostgreSQL/Redis storage drivers, and multi-protocol synthetics.

LNMP v3.0 Architecture Overview

API & Streaming Layer

FastAPI Service API & Server-Sent Events (SSE)

Single Persistent Stream (/api/v1/events/stream) • 15s Keep-Alive Heartbeat • SQL-Level Pagination

Concurrency Engine

32s Probe / 28s Headroom Sweeper

5 Pings @ 8.0s • 0-2000ms Startup Jitter • Dynamic Endpoint Registry

Diagnostic & Synthetics

Multi-Protocol Synthetic Probes

TCP Handshakes • HTTP/HTTPS Status (SSRF Shielded) • SSL Expiry Tracking

Dual-Driver Storage
Postgres-Native vs. Redis Acceleration
SQLAlchemy 2.0 Repos
Async Declarative Mapped Models
Fleet SLA Reports
Chunked Cursor Streaming Export

Architectural Evolution: Version 2.0 vs. Version 3.0

The transition to Version 3.0 represents both a response to real-world deployment bottlenecks and an intentional refactoring toward modern software design practices:

Architectural DimensionLNMP Version 2.0 (Baseline)LNMP Version 3.0 (Current Production)
Database Access LayerRaw SQL strings embedded directly inside route handlers.Pure SQLAlchemy 2.0 Typed Repositories with declarative mapped models and database-level LIMIT/OFFSET pagination.
Storage ArchitectureFixed PostgreSQL storage backend.Dual-Driver Storage Architecture: Postgres-Native (standard edge mode) with automated Redis memory acceleration for high-concurrency NOC tiers.
Probe Sweeper Timing10 pings @ 6.0s (tight 60s window with frequent cycle overlaps during packet loss).Deterministic 32s Probe / 28s Headroom Sweeper (5 pings @ 8.0s + 0–2000ms randomized startup jitter).
Target ManagementDaemon restart required to apply endpoint additions or edits.Dynamic In-Memory Registry (EndpointRegistry) with non-blocking async task spawning and cancellation.
Dashboard TelemetryPeriodic client-side polling (HTTP GET every 5 seconds).Server-Sent Events (SSE) persistent stream (/api/v1/events/stream) with 15s keep-alive heartbeat.
Application ProbesICMP ping reachability only.Multi-Protocol Synthetics: Non-blocking TCP port checks, HTTP/HTTPS response validation, and SSL certificate expiry tracking.
Route DiagnosticsBasic system traceroute execution.High-Fidelity ICMP Diagnostics (-q 2 -w 3 -I) with automatic Layer-2 subnet bypass.
Fleet ReportingSingle-endpoint uptime calculation only.Dedicated Reports Console (/reports) with fleet-wide SLA summary cards and chunked streaming CSV export.
Design SystemScoped component CSS with color drift and dark hover bugs.Single-Source Global CSS (App.vue), high-contrast tactical NOC palette, and tabular-nums alignment.

Deterministic Sweeper Concurrency and the 32s / 28s Timing Budget

The 60-Second Sweeper Collision Problem

In LNMP v2.0, our ping sweeper executed 10 pings per target at 6.0-second intervals. Under perfect network conditions, this finished in around 54 seconds. However, when an endpoint experienced packet drops or high transit latency, each failed probe waited for the full socket timeout (up to 4.5 seconds):

Worst-Case v2.0 Execution=10×5.85s=58.5 seconds\text{Worst-Case v2.0 Execution} = 10 \times 5.85\text{s} = 58.5\text{ seconds}

When the next top-of-the-minute boundary arrived (:00.000), the previous cycle was still writing historical rows and closing sockets. The incoming cycle launched simultaneously, causing thread contention, connection pool exhaustion, and CPU spikes.

The 32s Probe / 28s Headroom Timing Formulation

In Version 3.0, we re-architected the execution budget in monitoring/engine.py and monitoring/ping.py around a strict mathematical ceiling:

Probe Count=5∣Probe Interval=8.0s∣Socket Timeout=1.0s\text{Probe Count} = 5 \quad | \quad \text{Probe Interval} = 8.0\text{s} \quad | \quad \text{Socket Timeout} = 1.0\text{s} Active Sweep Duration=(5×Probe Window)+Startup Jitter≈32.0 seconds\text{Active Sweep Duration} = (5 \times \text{Probe Window}) + \text{Startup Jitter} \approx 32.0\text{ seconds} Maintenance Headroom=60s−32s=28 seconds of idle buffer\text{Maintenance Headroom} = 60\text{s} - 32\text{s} = \mathbf{28\text{ seconds of idle buffer}}

60-Second Sweeper Execution Budget

Active Probe Window00s – 32s

5-Ping Concurrency Sweeper

  • • 0–2000ms randomized startup jitter (J=random(0,2000)msJ = \text{random}(0, 2000)\text{ms})
  • • Raw ICMP socket execution with Linux CAP_NET_RAW
  • • Instant sub-cycle differential RCA trigger on first packet drop
Maintenance Headroom32s – 60s

28s Idle Buffer & Rollups

  • • TimescaleDB continuous aggregate refreshes
  • • CSV report generation and background telemetry streaming
  • • Zero CPU spikes or overlapping socket collisions

Eliminating the Socket Thundering Herd

When monitoring hundreds of nodes, starting all ping tasks at :00.000 creates a momentary burst of raw socket requests that can overwhelm unprivileged network buffers.

To eliminate this, each endpoint task introduces a randomized delay before its first ping:

Ji∼U(0,2000) msJ_i \sim \mathcal{U}(0, 2000)\text{ ms}

This smooths socket allocation across the first two seconds of the minute window without altering the 60-second telemetry aggregation boundary.


Dynamic In-Memory Endpoint Registry

Zero-Downtime Hot-Reloading

In Version 2.0, adding, modifying, or removing a monitored target required a systemd service restart (systemctl restart netmon-engine). In production, restarting the daemon dropped in-flight ping cycles, caused false gaps in continuous aggregates, and temporarily interrupted reachability monitoring.

In Version 3.0, we introduced EndpointRegistry (monitoring/registry.py), a concurrent in-memory repository protected by an asynchronous lock:

import asyncio
from typing import Any, Dict, Optional
from uuid import UUID

class EndpointRegistry:
    """Thread-safe in-memory registry for dynamic target management."""
    
    def __init__(self):
        self._endpoints: Dict[UUID, Any] = {}
        self._tasks: Dict[UUID, asyncio.Task] = {}
        self._lock = asyncio.Lock()

    async def add_or_update(self, endpoint_data: Any) -> None:
        async with self._lock:
            ep_id = self._get_val(endpoint_data, "id")
            
            # Cancel existing worker task if updating
            if ep_id in self._tasks:
                self._tasks[ep_id].cancel()
                try:
                    await self._tasks[ep_id]
                except asyncio.CancelledError:
                    pass

            self._endpoints[ep_id] = endpoint_data
            
            # Spawn new autonomous worker coroutine
            if self._get_val(endpoint_data, "monitoring_enabled", True):
                self._tasks[ep_id] = asyncio.create_task(
                    self._run_endpoint_worker(ep_id)
                )

    async def remove(self, ep_id: UUID) -> None:
        async with self._lock:
            if ep_id in self._tasks:
                self._tasks[ep_id].cancel()
                self._tasks.pop(ep_id, None)
            self._endpoints.pop(ep_id, None)

    @staticmethod
    def _get_val(obj: Any, key: str, default: Any = None) -> Any:
        """Safe model-agnostic attribute extractor."""
        if isinstance(obj, dict):
            val = obj.get(key, default)
        else:
            val = getattr(obj, key, default)
        return val if val is not None else default

When an operator adds or pauses an endpoint via the web dashboard, FastAPI calls the registry in memory. The corresponding worker coroutine is spawned or canceled in sub-second time without restarting the daemon.


SQLAlchemy 2.0 Typed Repositories and SQL-Level Pagination

Eliminating Raw SQL Coupling

Earlier versions of LNMP relied on raw SQL query strings concatenated inside API router functions. While fast to write during prototyping, this approach made schema refactoring error-prone and caused memory bloat during pagination. To paginate a table, the backend fetched thousands of rows from PostgreSQL and applied Python list slicing [offset:offset+limit].

Version 3.0 adopts the modern SQLAlchemy 2.0 declarative mapping standard, introducing strongly typed models and repository patterns:

from uuid import UUID, uuid4
from sqlalchemy import String, Boolean, Float, Integer
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
from sqlalchemy.dialects.postgresql import UUID as PG_UUID

class Base(DeclarativeBase):
    pass

class Endpoint(Base):
    __tablename__ = "endpoints"

    id: Mapped[UUID] = mapped_column(
        PG_UUID(as_uuid=True), primary_key=True, default=uuid4
    )
    hostname: Mapped[str] = mapped_column(String(255), nullable=False, index=True)
    ip_address: Mapped[str] = mapped_column(String(64), nullable=False, index=True)
    device_type: Mapped[str] = mapped_column(String(32), default="SERVER")
    monitoring_enabled: Mapped[bool] = mapped_column(Boolean, default=True)
    packet_loss_threshold: Mapped[float] = mapped_column(Float, default=20.0)
    latency_threshold_ms: Mapped[float] = mapped_column(Float, default=100.0)
    is_deleted: Mapped[bool] = mapped_column(Boolean, default=False)

True SQL-Level Pagination

All list queries now execute database-level LIMIT and OFFSET clauses, returning standardized pagination envelopes:

from sqlalchemy import select, func
from sqlalchemy.ext.asyncio import AsyncSession

class EndpointRepository:
    def __init__(self, session: AsyncSession):
        self.session = session

    async def get_paginated(
        self, page: int = 1, page_size: int = 50
    ) -> dict:
        offset = (page - 1) * page_size
        
        # Exact database-level count
        count_stmt = select(func.count()).select_from(Endpoint).where(
            Endpoint.is_deleted == False
        )
        total_count = (await self.session.execute(count_stmt)).scalar() or 0
        
        # Paginated fetch
        query = (
            select(Endpoint)
            .where(Endpoint.is_deleted == False)
            .order_by(Endpoint.hostname.asc())
            .offset(offset)
            .limit(page_size)
        )
        result = await self.session.execute(query)
        items = result.scalars().all()
        
        total_pages = max(1, (total_count + page_size - 1) // page_size)
        
        return {
            "items": items,
            "total_count": total_count,
            "page": page,
            "page_size": page_size,
            "total_pages": total_pages
        }

Dual-Driver Storage Architecture

Different monitoring deployments have distinct hardware constraints. An edge appliance monitoring a remote branch substation might run on 1 vCPU with 1 GB RAM, where installing extra services is impractical. Conversely, a central enterprise Network Operations Center (NOC) with 50+ operators requires sub-millisecond session validation and high-throughput event distribution.

LNMP v3.0 introduces a Dual-Driver Storage Architecture (backend/app/services/driver_manager.py):

Hardware-Adaptive Dual-Driver Storage Model

Standard Mode (Edge Tier)

PostgreSQL-Native Driver

Zero extra daemon dependencies. Relational session state and native LISTEN/NOTIFY IPC event streams. Ideal for low-resource standalone edge gateways.

Accelerated Mode (Enterprise Tier)

Redis-Accelerated Driver

Sub-millisecond token session validation, in-memory telemetry buffers, and Redis Pub/Sub broadcasting for high-concurrency multi-operator video walls.

The active storage mode is managed via /etc/netmon/config.toml:

[redis]
enabled = true
host = "127.0.0.1"
port = 6379
performance_mode = true

If Redis acceleration is enabled but the Redis service becomes unreachable, the driver manager automatically logs a warning and falls back to PostgreSQL-Native operations without dropping incoming telemetry.


Real-Time Telemetry via Server-Sent Events (SSE)

Replacing 5-Second HTTP Polling

In v2.0, web dashboard clients polled GET /api/v1/endpoints/ every 5 seconds to update latency charts and status badges. In environments with 20 active operator dashboards and wall-mounted NOC displays, this generated over 240 heavy database queries per minute, creating unnecessary read load on TimescaleDB.

Version 3.0 replaces polling with a single persistent Server-Sent Events (SSE) stream at GET /api/v1/events/stream:

import asyncio
from fastapi import APIRouter, Request
from fastapi.responses import StreamingResponse

router = APIRouter(prefix="/api/v1/events", tags=["Events"])

@router.get("/stream")
async def stream_telemetry_events(request: Request):
    """Publish real-time telemetry updates and keep-alive heartbeats."""
    
    async def event_generator():
        queue = asyncio.Queue()
        # Register subscriber in event broker
        await request.app.state.broker.subscribe(queue)
        
        try:
            while True:
                # Check for client disconnect
                if await request.is_disconnected():
                    break

                try:
                    # Wait for telemetry event with 15s timeout
                    data = await asyncio.wait_for(queue.get(), timeout=15.0)
                    yield f"event: {data['type']}\ndata: {data['payload']}\n\n"
                except asyncio.TimeoutError:
                    # Send periodic keep-alive comment to prevent proxy timeouts
                    yield ": heartbeat\n\n"
        finally:
            await request.app.state.broker.unsubscribe(queue)

    return StreamingResponse(
        event_generator(),
        media_type="text/event-stream",
        headers={
            "Cache-Control": "no-cache",
            "Connection": "keep-alive",
            "X-Accel-Buffering": "no"
        }
    )

The browser establishes a single long-lived HTTP connection. State transitions, latency spikes, and root-cause notifications stream to clients instantly, reducing database read queries from polling to zero.


Multi-Protocol Synthetic Service Probes

A monitored server’s operating system kernel will continue answering ICMP Echo requests as long as the network interface is up, even if the underlying web application, reverse proxy, or database process has crashed.

LNMP v3.0 expands beyond pure ICMP reachability by adding three non-blocking synthetic service probes (monitoring/synthetic.py):

import asyncio
import ssl
import socket
from datetime import datetime, timezone
import httpx

async def check_tcp_port(host: str, port: int, timeout: float = 3.0) -> dict:
    """Validate TCP 3-way handshake and measure connect latency."""
    start = asyncio.get_event_loop().time()
    try:
        _, writer = await asyncio.wait_for(
            asyncio.open_connection(host, port), timeout=timeout
        )
        latency_ms = (asyncio.get_event_loop().time() - start) * 1000
        writer.close()
        await writer.wait_closed()
        return {"status": "OPEN", "latency_ms": round(latency_ms, 2)}
    except Exception as e:
        return {"status": "CLOSED", "error": str(e)}

async def check_http_status(url: str, timeout: float = 5.0) -> dict:
    """Validate HTTP response code with SSRF protection."""
    # Defend against AWS/GCP cloud metadata SSRF attempts
    if "169.254.169.254" in url or "127.0.0.1" in url:
        return {"status": "BLOCKED", "error": "SSRF Target Denied"}
        
    async with httpx.AsyncClient(verify=False, timeout=timeout) as client:
        try:
            resp = await client.get(url)
            return {"status_code": resp.status_code, "reachable": resp.is_success}
        except Exception as e:
            return {"status_code": 0, "error": str(e)}

def check_ssl_expiry(host: str, port: int = 443) -> dict:
    """Extract x509 certificate expiry metadata over TLS 1.2+."""
    ctx = ssl.create_default_context()
    with socket.create_connection((host, port), timeout=5.0) as sock:
        with ctx.wrap_socket(sock, server_hostname=host) as ssock:
            cert = ssock.getpeercert()
            expiry_str = cert['notAfter']
            expiry_date = datetime.strptime(expiry_str, "%b %d %H:%M:%S %Y %Z").replace(tzinfo=timezone.utc)
            days_left = (expiry_date - datetime.now(timezone.utc)).days
            return {"days_remaining": days_left, "expiry_date": expiry_date.isoformat()}
  • TCP Port Probe: Verifies service port availability (e.g. 5432 for PostgreSQL, 22 for SSH) and records connect latency.
  • HTTP/HTTPS Status Probe: Validates web endpoint health with explicit SSRF defense blocking cloud metadata IP ranges.
  • SSL Certificate Expiry: Tracks TLS certificate validity windows, warning operators before expiration outages occur.

Dedicated Fleet SLA Reports and Streaming CSV Export

Operators require verifiable historical SLA records for compliance audits and executive reviews. Fetching hundreds of thousands of historical telemetry rows into RAM before generating a file can exhaust server memory.

In backend/app/routers/reports.py, LNMP v3.0 streams historical records directly from a PostgreSQL server-side cursor to the client stream in batches of 1,000 rows:

from fastapi import APIRouter, Depends
from fastapi.responses import StreamingResponse
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession

router = APIRouter(prefix="/api/v1/reports", tags=["Reports"])

def sanitize_csv_cell(val: str) -> str:
    """Neutralize spreadsheet formula injection triggers."""
    if val and str(val).strip().startswith(("=", "+", "-", "@", "\t", "\r")):
        return f"'{val}"
    return str(val)

@router.get("/export/csv")
async def stream_sla_csv_export(db: AsyncSession = Depends(get_db)):
    """Stream raw telemetry records directly from database cursor."""
    
    async def csv_row_generator():
        # Yield CSV Header
        yield "Timestamp,Endpoint,IP,RTT_Avg_ms,Packet_Loss_Pct,State\n"
        
        query = select(HourlyTelemetry).order_by(HourlyTelemetry.bucket.desc()).execution_options(yield_per=1000)
        result = await db.stream(query)
        
        async for row in result.scalars():
            line = [
                row.bucket.isoformat(),
                sanitize_csv_cell(row.hostname),
                row.ip_address,
                f"{row.avg_rtt_ms:.2f}",
                f"{row.packet_loss_pct:.1f}",
                row.state
            ]
            yield ",".join(line) + "\n"

    return StreamingResponse(
        csv_row_generator(),
        media_type="text/csv",
        headers={"Content-Disposition": "attachment; filename=sla_telemetry_export.csv"}
    )

As specified in the OWASP CSV Injection defense guidelines, leading formula initiation characters (=, +, -, @) are sanitized on the fly with a single quote prefix, protecting spreadsheet users against formula execution.


Tactical NOC Design System and Single-Source Global CSS

In Version 2.0, CSS rules were duplicated across individual Vue components (DashboardView.vue, ReportsView.vue, EndpointCard.vue), causing subtle color mismatches and dark-on-dark text bugs during theme toggling.

In Version 3.0, all UI tokens were refactored into a single global stylesheet (frontend/src/App.vue):

  • Single Source of Truth: Centralized .status-pill, .status-dot, .dense-table, .table-card, and .btn-primary tokens.
  • Numeric Column Alignment: Applied font-variant-numeric: tabular-nums across all latency values, IP addresses, uptime percentages, and timestamps to eliminate column jitter during live telemetry refreshes.
  • Unified Theme Surface Variables: Mapped all elevated surfaces to explicit CSS properties (--bg-surface, --bg-surface-hover, --text-primary, --border-color), guaranteeing consistent contrast across dark and light modes.

Post-Implementation Incident Log and Live Upgrade Triage

During live production upgrade testing on branch 3.0.0, our test matrix surfaced 7 specific operational edge cases that were triaged and resolved:

Defect / Incident IdentifiedTechnical Root CauseResolution Applied
PostgreSQL Backup Permission Deniedsudo -u postgres psql -f /var/backups/netmon/backup.sql failed because the postgres user cannot access 750 root-owned directories.Replaced file flag with standard input streaming: cat <backup> | sudo -u postgres psql, which bypasses filesystem permission restrictions.
Endpoint Detail View 500 ErrorSQLAlchemy subquery latest_event_sub was joined with a constant True and omitted the endpoint_id grouping column.Explicitly projected EndpointEvent.endpoint_id and updated join to .outerjoin(latest_event_sub, Endpoint.id == latest_event_sub.c.endpoint_id).
Engine Registry Loop Freezeregistry.py assumed inputs were dictionaries and called .get(), crashing when receiving SQLAlchemy ORM model objects.Implemented safe model-agnostic helper _get_val(obj, key) supporting both dictionary keys and ORM object attributes.
Blank Latency Values on Cardslist_with_stats() omitted avg_rtt_ms in the subquery projection, causing latency to return null.Added EndpointEvent.avg_rtt_ms to the subquery select list and serialized it as a float in the response schema.
Microsoft Edge Autofill BlockedLogin form used empty action="", nested PrimeVue card slots, and omitted HTML5 autocomplete attributes.Restructured into semantic <form action="/api/v1/auth/login" method="post" autocomplete="on"> with explicit field IDs.
Light Theme Table Hover Turn Pitch BlackReportsView.vue referenced undefined fallback variable --bg-surface-elevated, #18181b, causing dark text on dark background in light mode.Bound all table hover states to global semantic token var(--bg-surface-hover).
Status Badges Missing Color ClassesEndpointCard.vue computed status classes but lacked local CSS definitions; Edit buttons were exposed to unauthenticated viewers.Centralized .status-pill classes in App.vue and guarded mutation actions with v-if="isAdmin".

Automated In-Place Upgrade Pipeline

To ensure zero-downtime upgrades on enterprise Linux hosts, Version 3.0 includes an updated deploy/upgrade.sh script:

# Execute automated in-place upgrade
sudo ./deploy/upgrade.sh

Upgrade Pipeline Workflow

  1. Pre-Upgrade Backup: Dumps an archive of the active database into /var/backups/netmon/backup_pre_v3.0_TIMESTAMP.sql. Aborts cleanly if the dump fails.
  2. Configuration Migration: Inspects /etc/netmon/config.toml and merges new v3.0 defaults (such as [redis] settings) without modifying existing database passwords or JWT secret keys.
  3. Daemon Pause: Pauses netmon-engine and netmon-api to prevent write collisions during database migrations.
  4. Code & Asset Build: Pulls the latest Git commits, updates Python virtualenv dependencies, and compiles the production Vue 3 frontend bundle (npm run build).
  5. Alembic Forward Migrations: Runs alembic upgrade head to construct new synthetic tables and indexes while keeping existing TimescaleDB compressed chunks intact.
  6. Daemon Restart & Health Verification: Restarts systemd units (systemctl restart netmon-*), reloads Nginx, and validates the /api/v1/version endpoint.

Conclusion & Open Source Repository

Building LNMP Version 3.0 has been a valuable lesson in balancing network performance engineering with disciplined software architecture. Moving from raw SQL strings and tight polling loops to typed repositories, deterministic concurrency budgets, and real-time event streams demonstrated that a monitoring system can deliver enterprise-grade reliability and low latency while remaining lightweight and simple to operate.

The complete open-source codebase for LNMP v3.0 (Production), including the asynchronous polling engine, FastAPI backend, Alembic migrations, and Vue 3 frontend, is available on GitHub:

Feedback & Discussion

Have questions, corrections, or perspectives to share? Connect directly to discuss systems and security.

Table of Contents (20 sections)
↑ ↓ navigate↵ select
Publications indexed