Nexa-Flow: Plug-in Protocol Industrial IoT Data Acquisition Service

Nexa-Flow is a high-throughput industrial IoT data acquisition and monitoring service built on .NET 8. Its core architectural choice is a plug-in protocol model: five protocols (Modbus TCP, Modbus RTU, MQTT, TCP Listener, Virtual) each run with their own IDeviceProtocolHandler + IProtocolRuntime + IHostedService triple, but every read converges into a single protocol-agnostic DeviceReadPipeline that handles MongoDB time-series writes, in-memory cache updates, alarm/calculation dirty-set notifications and SignalR fan-out. Adding a new protocol therefore requires no changes to downstream code. The service ships with a hash-tracked SQL artifact migrator that automatically applies versioned SQL files at startup (replacing manual psql -f steps), a hybrid PostgreSQL + MongoDB time-series data model, NCalc-based runtime formula evaluation for virtual tags, AES-256-GCM credential encryption (MQTT broker passwords, flow export keys), JWT with hashed refresh-token rotation, AES-protected secrets, liveness/readiness split health probes for HA orchestration, and a fail-fast production startup. Two major refactors hardened the polling pipeline: (1) DI cleanup that split ConnectionManager into IHostedService + IProtocolRuntime pairs, and (2) a performance pass replacing GUID-prefix substring scans with ConcurrentDictionary lookups, full-table scans with dirty-set patterns and active-view UNION ALL views, a SemaphoreSlim(64) read budget for thread-pool starvation safety, and an unbounded Task.Run pattern with a channel-based bounded fire-and-forget worker pool. CI enforces this via anti-pattern and warning-regression lints with baselines.

Solo Senior .NET Backend Engineer & System Architect

  1. Plug-in Protocol Architecture – Designed an IProtocolRegistry + IDeviceProtocolHandler + IProtocolRuntime + IHostedService model so that Modbus TCP, Modbus RTU, MQTT, TCP Listener and Virtual devices share a single DeviceReadPipeline(List<BsonDocument>) contract; cache, alarm, calculation and SignalR flows light up automatically for any new protocol.
  2. Performance Refactor (Phase 2) – Replaced GUID-prefix substring scans with ConcurrentDictionary O(1) lookups, converted full-table scans into a dirty-set pattern + UNION ALL active views (vw_active_tags / vw_active_enabled_tags / vw_active_devices), introduced a SemaphoreSlim(64) global read budget and a channel-based bounded fire-and-forget worker pool to eliminate unbounded Task.Run and thread-pool starvation.
  3. DI Cleanup (Phase 1) – Migrated static-field dependencies into DI and split ConnectionManagers into clean IHostedService + IProtocolRuntime pairs so runtime CRUD (device added/removed/updated) reflects in polling lists instantly without service restarts.
  4. Protocol Implementations – Built Modbus TCP (NModbus + 64-concurrency budget, FC1-4), Modbus RTU (serial port queue + read-back verification), Virtual Device (formula-driven synthetic tags), TCP Listener (FixedLengthFromRules / PerRecvChunk framing for push-only frames) and MQTT (MQTTnet pub/sub, JSONPath payload extraction, write-template substitution, hot-reload subscriptions).
  5. Hash-Tracked SQL Artifact Migrator – Wrote a startup migrator that loads Migrations/Sql/*.sql as embedded resources, tracks each file by SHA-256 in _sql_artifact_migrations and re-applies automatically on hash change — eliminating manual psql -f steps in deployment.
  6. Security & Secrets – Implemented AES-256-GCM encryption for MQTT broker passwords and flow export keys, JWT access + DB-hashed refresh tokens, forwarded headers for reverse-proxy IP correctness, and production fail-fast on missing identity seed or export key.
  7. Operational Hardening – Split health endpoints into /health (liveness, process only) and /ready (readiness, Postgres SELECT 1 + Mongo admin ping + NexaFlow init), added a 15-retry DB readiness barrier on startup, and made migrations HA-safe via Database__RunStartupTasks flag.
  8. CI & Code Quality – Wrote tools/check-antipatterns.ps1 (blocks new GetList().Where, unbounded Task.Run, full-scan patterns via grandfathered baselines), tools/check-warnings.ps1 (warning-count regression guard), and xUnit tests for AesEncryptor (round-trip + tamper detection), RefreshTokenHasher and RequestIdentityNormalizer.

Features

  1. Plug-in protocol architecture — 5 protocols (Modbus TCP, Modbus RTU, MQTT, TCP Listener, Virtual) on one shared DeviceReadPipeline contract
  2. Hash-tracked automatic SQL artifact migration (zero manual psql steps)
  3. Dirty-set pattern for alarm/calculation engines (no full-table scans)
  4. Channel-based bounded fire-and-forget worker pool (no unbounded Task.Run)
  5. SemaphoreSlim(64) global read budget protecting against thread-pool starvation
  6. PostgreSQL + MongoDB time-series hybrid storage with active-view UNION ALL views
  7. NCalc runtime formula engine for virtual tags and calculations
  8. AES-256-GCM encryption for broker passwords and export keys
  9. JWT access + hashed refresh-token rotation
  10. Liveness / readiness split health probes for HA orchestration
  11. SignalR real-time fan-out for SCADA-style dashboards, alarms and monitoring
  12. Custom hourly/daily reporting with automated PDF export
  13. Alarm rules, service/maintenance tracking, virtual tag computation
  14. Docker Compose multi-service deployment (Postgres, Mongo, Seq, app) with non-root container
  15. CI anti-pattern lint + compiler warning regression guard (baseline-based)
  16. Dynamic role-based access control for users, devices and tags

Technologies

  1. C#
  2. .NET 8 Web API
  3. Entity Framework Core
  4. NPoco
  5. PostgreSQL
  6. MongoDB (Time-Series)
  7. SignalR
  8. NModbus (Modbus TCP/RTU)
  9. MQTTnet (MQTT)
  10. TCP Listener (push frames)
  11. NCalc (formula engine)
  12. System.Threading.Channels
  13. Docker / Docker Compose
  14. Windows Service
  15. Serilog + Seq
  16. JWT + Refresh Token (hashed)
  17. AES-256-GCM
  18. xUnit
  19. GitHub Actions
Associated with