Skip to main content

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/enqueue and takes responsibility until it can emit a terminal webhook (DELIVERY_SUCCESSFUL or DELIVERY_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 the nxt-sts plugin). 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-chirpstack into the same webhook events the adopter already handles.
  • Stores queue state in Valkey/Redis so a process restart does not drop jobs. /healthz is 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 → webhook SENT_TO_NS on first send → later DELIVERY_SUCCESSFUL / DELIVERY_FAILED. After success the Redis record is removed, so GET /message/:correlationId 404s; the webhook is the outcome channel.
  • LoRaWAN ingress: ChirpStack HTTP integration POSTs to /ingress/calin-chirpstack (optional X-API-KEY via CALIN_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 no correlationId.
  • STS mint then deliver: POST /token/generate with pluginId: "nxt-sts" (sidecar NXT_STS_URL, often http://nxt-sts:8080) returns { "token": "<20 digits>" } → adopter enqueues DELIVER_PREEXISTING_TOKEN with requestData.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 leaves QUEUED; 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 install
    cp .env.example .env
    docker compose up -d valkey
    pnpm dev

    Listens on PORT (default 3100). .env.example points DEVICE_MESSAGING_CONFIG_PATH at config.example.json (both stubs). Placeholder eventWebhook.url produces failed outbound POSTs in logs; they do not stop delivery.

  • Compose (service + Valkey): copy .env, uncomment the config volume in docker-compose.yml, then docker compose up --build. Compose sets REDIS_HOST=valkey.

  • Production: one replica, durable Valkey, pinned GHCR tag, non-empty DEVICE_MESSAGING_API_KEY or private network / reverse-proxy auth, JSON config artifact with the plugins you actually run. Empty bundled plugins[] (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.

MethodPathAuthRole
POST/message/enqueueBearerQueue a command. 201 + message body. correlationId is the adopter's. Higher priority is more urgent.
GET/message/:correlationIdBearerIn-flight inspect. 404 after successful delivery (record deleted) or unknown id.
POST/message/cancelBearerCancel one id. Result: CANCELLED / NOT_CANCELLABLE (already in-flight) / NOT_FOUND. No webhook.
POST/messages/cancelBearerSame, many ids.
POST/token/generateBearerSync 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/provisioningBearerSync vendor provision/deprovision. calin-api-v1 has no facet (400).
POST/ingress/:pluginIdnone (optional plugin X-API-KEY)Vendor → service. PUSH plugins only. Opaque JSON body.
GET/healthznoneLiveness { "ok": true }. Not Redis readiness.
GET/metricsnonePrometheus text.
GET/swagger, /v3/api-docsnoneLive 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 in plugins[] (NXT_STS_URL). ChirpStack when using calin-chirpstack. CALIN HTTP APIs for calin-api-v1 / calin-api-v2.
  • Downstream consumers: any HTTP adopter that enqueues and handles the webhook. nxt-backend is the intended consumer when metering is enabled (HTTP + HMAC webhook, types from @nxtgrid/device-messaging-contract); that capability is not imported yet on nxt-backend main.
  • 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 → bundled config.default.json. Plugin listed in plugins[] 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: false is 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 HEALTHCHECK probes /healthz on PORT. 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-contract consumers (nxt-backend metering, any adopter).
  • if nxt-sts token types or URL change, mint via POST /token/generate / sidecar NXT_STS_URL breaks.
  • 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 in plugins[], or nxt-sts used for enqueue (mint-only).
  • boot crash naming an env key — plugin enabled without its secrets.
  • GET /message/:id 404 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.
  • /healthz 200 but nothing delivers — Redis down or engine.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)