nxt-device-messaging
Purpose
nxt-device-messaging is a standalone Fastify service that owns reliable, prioritized, retrying command delivery to addressable field devices. A billing or operations app POSTs a job with a correlationId it chose; this service queues it, sends it through the plugin that already speaks that device's network, retries with exponential backoff, and reports outcomes on a signed webhook.
A "message" here is a command or read to physical hardware (read credit, deliver an STS token, set a power limit) — not SMS, chat, or product notifications.
Scope
- In scope:
- Command API: enqueue, inspect in-flight jobs, cancel, synchronous STS mint, optional plugin provisioning.
- Vendor ingress (
POST /ingress/:pluginId) so network servers (ChirpStack, …) never callback into the adopter app. - Plugin catalog: CALIN HTTP V1/V2, CALIN over ChirpStack (LoRaWAN), nxt-sts mint, stub PUSH/PULL for local tests.
- Redis/Valkey-backed job, retry, and in-flight state. No relational database.
- Outbound HMAC-signed webhook events and Prometheus metrics.
- Out of scope:
- Running ChirpStack, a CALIN cloud, or meter firmware — those stay with the operator / vendor.
- Billing, wallet, or meter-domain business logic (owned by the adopter; in this suite, nxt-backend when metering is enabled).
- Running more than one process against a given Valkey. A second replica races engine timers and can split LoRaWAN ACK vs uplink.
- Living inside
nxt-backend. Callers talk HTTP + webhook.
What this service does in production
- Accepts
POST /message/enqueueand takes responsibility until it can emit a terminal webhook (DELIVERY_SUCCESSFULorDELIVERY_FAILED). - Dispatches through a PUSH plugin (send, then wait for the network server to call ingress — typical of LoRaWAN) or a PULL plugin (create a vendor task and poll status).
- Mints STS tokens synchronously via
POST /token/generate(usually thenxt-stsplugin). Mint is not queued and has no delivery webhook. Putting a token on a meter is still an enqueue (DELIVER_PREEXISTING_TOKEN, or a plugin that mints then delivers). - Turns ChirpStack (and similar) callbacks at
/ingress/calin-chirpstackinto the same webhook events the adopter already handles. - Stores queue state in Valkey/Redis so a process restart does not drop jobs.
/healthzis liveness only — it does not check Redis.
Current package version: 0.1.2. Pin GHCR tags (ghcr.io/nxtgrid/nxt-device-messaging:vX.Y.Z); images are linux/amd64 and linux/arm64.
Primary workflows
- Enqueue and deliver: adopter POSTs enqueue with
pluginId+correlationId→ job lands in Redis (QUEUED) → engine tick (1 s) sends via plugin → webhookSENT_TO_NSon first send → laterDELIVERY_SUCCESSFUL/DELIVERY_FAILED. After success the Redis record is removed, soGET /message/:correlationId404s; the webhook is the outcome channel. - LoRaWAN ingress: ChirpStack HTTP integration POSTs to
/ingress/calin-chirpstack(optionalX-API-KEYviaCALIN_CHIRPSTACK_INGRESS_API_KEY) → plugin correlates up/ack in-process → same webhook envelope as an enqueue job. Unsolicited device events (JOIN_NETWORK,READ_REPORT) have nocorrelationId. - STS mint then deliver:
POST /token/generatewithpluginId: "nxt-sts"(sidecarNXT_STS_URL, oftenhttp://nxt-sts:8080) returns{ "token": "<20 digits>" }→ adopter enqueuesDELIVER_PREEXISTING_TOKENwithrequestData.token, or uses a plugin that mints and delivers in one command. - Local stub path: enable
stub-push/stub-pull, run Valkey +pnpm dev— command is accepted and leavesQUEUED; no real meter, no terminal success.
Setup and run
-
Repository: github.com/nxtgrid/nxt-device-messaging
-
Prerequisites: Node.js 24.x, pnpm 11 (Corepack), Docker for Valkey.
-
Local (stubs, host process + Valkey only):
corepack enable && pnpm installcp .env.example .envdocker compose up -d valkeypnpm devListens on
PORT(default 3100)..env.examplepointsDEVICE_MESSAGING_CONFIG_PATHatconfig.example.json(both stubs). PlaceholdereventWebhook.urlproduces failed outbound POSTs in logs; they do not stop delivery. -
Compose (service + Valkey): copy
.env, uncomment the config volume indocker-compose.yml, thendocker compose up --build. Compose setsREDIS_HOST=valkey. -
Production: one replica, durable Valkey, pinned GHCR tag, non-empty
DEVICE_MESSAGING_API_KEYor private network / reverse-proxy auth, JSON config artifact with the plugins you actually run. Empty bundledplugins[](config.default.json) means enqueue fails until you enable at least one.
Deployment
Platform-specific deployment guides live under this repository section:
- DigitalOcean App Platform — deploy from GHCR or build from the GitHub repository; one replica, Valkey with TLS, config via
DEVICE_MESSAGING_CONFIG_JSON.
Source runbook: docs/deployment/ in the repository (same-app Valkey / STS / webhook bindables).
APIs and interfaces
Command and token routes use Authorization: Bearer <DEVICE_MESSAGING_API_KEY> when that env is set. Unset/empty is local-only. Ingress and ops probes are not Bearer-authenticated.
| Method | Path | Auth | Role |
|---|---|---|---|
POST | /message/enqueue | Bearer | Queue a command. 201 + message body. correlationId is the adopter's. Higher priority is more urgent. |
GET | /message/:correlationId | Bearer | In-flight inspect. 404 after successful delivery (record deleted) or unknown id. |
POST | /message/cancel | Bearer | Cancel one id. Result: CANCELLED / NOT_CANCELLABLE (already in-flight) / NOT_FOUND. No webhook. |
POST | /messages/cancel | Bearer | Same, many ids. |
POST | /token/generate | Bearer | Sync mint. Discriminated body: pluginId, issueDateString, device, type (TOP_UP_KWH, CLEAR_CREDIT, CLEAR_TAMPER, SET_POWER_LIMIT) plus type-specific payload. 200 { "token": "…" }. Live shape: /swagger. |
POST | /plugin/provisioning | Bearer | Sync vendor provision/deprovision. calin-api-v1 has no facet (400). |
POST | /ingress/:pluginId | none (optional plugin X-API-KEY) | Vendor → service. PUSH plugins only. Opaque JSON body. |
GET | /healthz | none | Liveness { "ok": true }. Not Redis readiness. |
GET | /metrics | none | Prometheus text. |
GET | /swagger, /v3/api-docs | none | Live contract. |
Enqueue commandType is a closed set (src/lib/device-message/command-types.ts). Plugins accept a subset; unsupported type → 400. READ_VOLTAGE / READ_CURRENT also need phase (A/B/C). CALIN HTTP plugins require device.relayNode.id. JOIN_NETWORK and READ_REPORT are ingress-only.
Delivery pipeline (Redis): QUEUED → SENT_TO_NS → DELIVERED_TO_NS → SENT_TO_DEVICE → DELIVERY_SUCCESSFUL, or TO_RETRY → QUEUED / DELIVERY_FAILED. Webhook fires on first send (SENT_TO_NS), terminals, and unsolicited ingress. Delivery success means the radio/HTTP path completed; message.response.status is meter execution (EXECUTION_SUCCESS vs EXECUTION_FAILURE).
Outbound webhook: JSON to eventWebhook.url. Headers X-Device-Messaging-Event-Id and, if DEVICE_MESSAGING_WEBHOOK_SECRET is set, X-Device-Messaging-Signature: sha256=<hex> (HMAC-SHA256 of the raw body). Treat eventId as the idempotency key (retries reuse it). 2xx = stored; retries on network/408/429/5xx; other 4xx not retried; after six attempts the event sits in a Redis DLQ for seven days.
TypeScript/Zod wire types: npm package @nxtgrid/device-messaging-contract (not an HTTP client; ingress is not exported).
Integrations and dependencies
- Upstream / sidecars: Valkey or Redis (required). nxt-sts when
{ "id": "nxt-sts" }is inplugins[](NXT_STS_URL). ChirpStack when usingcalin-chirpstack. CALIN HTTP APIs forcalin-api-v1/calin-api-v2. - Downstream consumers: any HTTP adopter that enqueues and handles the webhook.
nxt-backendis the intended consumer when metering is enabled (HTTP + HMAC webhook, types from@nxtgrid/device-messaging-contract); that capability is not imported yet onnxt-backendmain. - Stack: Node 24, Fastify, Zod (plain modules, no DI container). License MPL-2.0.
- Config split: JSON artifact (
$schemaVersion: "1") for plugins, webhook URL, retries — not passwords. Env for secrets, Redis,PORT, and which artifact to load. Precedence:DEVICE_MESSAGING_CONFIG_JSON→_URL→_PATH→ bundledconfig.default.json. Plugin listed inplugins[]with missing env → boot fails. Request for a plugin you did not enable → that request fails, process stays up.
Operations notes
- Runtime: one HTTP process + one engine interval (
ENGINE_TICK_INTERVAL_MS= 1000). Timeouts/backoffs are observed on the next tick (whole seconds are the useful unit).engine.enabled: falseis ingest/inspect only. - One replica. A second process against the same Redis competes on timers and can split ChirpStack ACK vs uplink across in-memory correlators.
- Image
HEALTHCHECKprobes/healthzonPORT. PaaS ignores it — set the platform HTTP probe. - Logging: pretty stdout by default;
"logging": { "stdout": "json" }for aggregators.
Change impact map:
- if enqueue/webhook shapes change, update
@nxtgrid/device-messaging-contractconsumers (nxt-backendmetering, any adopter). - if
nxt-ststoken types or URL change, mint viaPOST /token/generate/ sidecarNXT_STS_URLbreaks. - if ChirpStack integration URL or ingress API key changes, LoRaWAN callbacks 401/miss and jobs stall in-flight.
- if you add a second replica against the same Redis, deliveries and correlators silently split.
Failure modes to check first:
- enqueue 400
Unknown or disabled pluginId— id not inplugins[], ornxt-stsused for enqueue (mint-only). - boot crash naming an env key — plugin enabled without its secrets.
GET /message/:id404 after a job you think succeeded — expected; look at the webhook / DLQ, not polling.- webhook unsigned or rejected — secret unset vs caller verifying HMAC on re-serialized JSON instead of raw bytes; placeholder URL in local
config.example.json. /healthz200 but nothing delivers — Redis down orengine.enabled: false; liveness does not cover that.- ChirpStack ACK without data (or the reverse) on more than one replica — correlator is in-process only.
Source of truth
- Repository: github.com/nxtgrid/nxt-device-messaging
- Operator + plugin chooser:
README.md - Integrator (auth, webhook HMAC, event set):
docs/guides/integrating.md - HTTP composition:
src/app.ts - Routes:
src/http/message-routes.ts,token-routes.ts,ingress-routes.ts,provisioning-routes.ts,auth.ts - Command vocabulary:
src/lib/device-message/command-types.ts - Engine tick:
src/engine/timers.ts - Webhook envelope:
src/engine/webhook/event-schema.ts,sign.ts - Plugin catalog:
src/plugins/catalog.ts - Wire package:
packages/contract/ - Config:
src/config/,config.default.json,config.example.json,.env.example - Why one replica:
docs/architecture/(single-replica note) - Deployment:
docs/deployment/(hub: DigitalOcean App Platform)