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
- 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.
- 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.
- 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.
- 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).
- 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.
- 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.
- 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.
- 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
- Plug-in protocol architecture — 5 protocols (Modbus TCP, Modbus RTU, MQTT, TCP Listener, Virtual) on one shared DeviceReadPipeline contract
- Hash-tracked automatic SQL artifact migration (zero manual psql steps)
- Dirty-set pattern for alarm/calculation engines (no full-table scans)
- Channel-based bounded fire-and-forget worker pool (no unbounded Task.Run)
- SemaphoreSlim(64) global read budget protecting against thread-pool starvation
- PostgreSQL + MongoDB time-series hybrid storage with active-view UNION ALL views
- NCalc runtime formula engine for virtual tags and calculations
- AES-256-GCM encryption for broker passwords and export keys
- JWT access + hashed refresh-token rotation
- Liveness / readiness split health probes for HA orchestration
- SignalR real-time fan-out for SCADA-style dashboards, alarms and monitoring
- Custom hourly/daily reporting with automated PDF export
- Alarm rules, service/maintenance tracking, virtual tag computation
- Docker Compose multi-service deployment (Postgres, Mongo, Seq, app) with non-root container
- CI anti-pattern lint + compiler warning regression guard (baseline-based)
- Dynamic role-based access control for users, devices and tags
Technologies
- C#
- .NET 8 Web API
- Entity Framework Core
- NPoco
- PostgreSQL
- MongoDB (Time-Series)
- SignalR
- NModbus (Modbus TCP/RTU)
- MQTTnet (MQTT)
- TCP Listener (push frames)
- NCalc (formula engine)
- System.Threading.Channels
- Docker / Docker Compose
- Windows Service
- Serilog + Seq
- JWT + Refresh Token (hashed)
- AES-256-GCM
- xUnit
- GitHub Actions