npm package discovery and stats viewer.

Discover Tips

  • General search

    [free text search, go nuts!]

  • Package details

    pkg:[package-name]

  • User packages

    @[username]

Sponsor

Optimize Toolset

I’ve always been into building performant and accessible sites, but lately I’ve been taking it extremely seriously. So much so that I’ve been building a tool to help me optimize and monitor the sites that I build to make sure that I’m making an attempt to offer the best experience to those who visit them. If you’re into performant, accessible and SEO friendly sites, you might like it too! You can check it out at Optimize Toolset.

About

Hi, 👋, I’m Ryan Hefner  and I built this site for me, and you! The goal of this site was to provide an easy way for me to check the stats on my npm packages, both for prioritizing issues and updates, and to give me a little kick in the pants to keep up on stuff.

As I was building it, I realized that I was actually using the tool to build the tool, and figured I might as well put this out there and hopefully others will find it to be a fast and useful way to search and browse npm packages as I have.

If you’re interested in other things I’m working on, follow me on Twitter or check out the open source projects I’ve been publishing on GitHub.

I am also working on a Twitter bot for this site to tweet the most popular, newest, random packages from npm. Please follow that account now and it will start sending out packages soon–ish.

Open Software & Tools

This site wouldn’t be possible without the immense generosity and tireless efforts from the people who make contributions to the world and share their work via open source initiatives. Thank you 🙏

© 2026 – Pkg Stats / Ryan Hefner

sse-to-mqtt-node

v0.1.1

Published

Bridge Server-Sent Events (SSE) streams to MQTT topics, as a library or a config-driven CLI

Readme

sse-to-mqtt-node

Bridge Server-Sent Events (SSE) streams to MQTT topics.

Open any number of SSE connections (GET or POST, each with its own body, headers and topic rule) and republish every event to an MQTT broker. Use it as a CLI driven by a JSON config file, or as a typed library in your own Node.js app.

  • Spec-compliant SSE parsing: multi-line data:, event:, id:, retry:, and Last-Event-ID on reconnect
  • Topic templates such as vessels/{name}/{imoNumber} filled from message fields
  • Automatic reconnects with exponential backoff, jitter, retry limits and connect timeouts
  • Optional OAuth2 client credentials with token caching and refresh on 401
  • Typed lifecycle events, a pluggable logger (silent by default), and zero HTTP dependencies (built-in fetch)

Requires Node.js 22 or later.

CLI

npm install -g sse-to-mqtt-node
# or run without installing: npx sse-to-mqtt-node --config connections.json

Create connections.json:

[
  { "name": "ticker", "method": "GET", "url": "/ticker" },
  {
    "name": "orders",
    "method": "POST",
    "url": "/orders/stream",
    "body": { "region": "eu" },
    "topic": "{name}/{orderId}",
    "qos": 1
  }
]

Set the environment (a .env file in the working directory is loaded automatically) and start it:

export STREAMING_ENDPOINT=https://api.example.com
export MQTT_BROKER_URL=mqtt://localhost:1883
export MQTT_TOPIC=example
sse-to-mqtt-node --config connections.json

Events from /ticker are published to example/ticker; events from /orders/stream go to example/orders/<orderId>.

| Variable | Required | Description | | --- | --- | --- | | STREAMING_ENDPOINT | yes | Base URL of the SSE server. Connection urls are resolved against it. | | MQTT_BROKER_URL | yes | e.g. mqtt://localhost:1883, mqtts://broker:8883 | | MQTT_TOPIC | yes | Base topic prefixed to every published topic | | MQTT_USERNAME, MQTT_PASSWORD | no | Broker credentials | | CONNECTIONS_CONFIG | no | Config path, if --config is not given | | AUTHENTICATION_URL | no | Enables OAuth2 client credentials; then CLIENT_ID and CLIENT_SECRET are required | | CLIENT_ID, CLIENT_SECRET, CLIENT_SCOPE | no | OAuth2 client credentials (CLIENT_SCOPE is optional) | | LOG_LEVEL | no | debug, info (default), warn or error |

The CLI stops cleanly on SIGINT / SIGTERM.

Connection config

Each entry in the JSON array:

| Field | Type | Description | | --- | --- | --- | | name | string | Unique name, available as {name} in topics | | method | "GET" | "POST" | HTTP method | | body | object | JSON request body; only allowed for POST | | url | string | Absolute URL, or a path resolved against the endpoint. Defaults to the endpoint | | headers | object | Extra request headers | | topic | string | string[] | Topic template relative to MQTT_TOPIC. Defaults to "{name}" | | qos | 0 | 1 | 2 | MQTT QoS (default 0) | | retain | boolean | MQTT retain flag (default false) |

The file is validated at startup; mistakes stop the CLI with a message naming the connection and field.

Topic templates

  • {name} is the connection name.
  • {field} or {nested.field} reads a field from the event data parsed as JSON.
  • If a placeholder has no value (missing, null, empty or an object), the event is skipped. Use this to filter.
  • Values have +, #, / and NUL replaced with _, so message content can't add topic levels or wildcards. Literal wildcards in the template itself are rejected.
  • The payload is the raw event data. Only templates that reference fields parse it as JSON, so non-JSON streams work with {name}-only topics.

Library

npm install sse-to-mqtt-node
import { SseToMqttBridge, createConsoleLogger } from 'sse-to-mqtt-node';

const bridge = new SseToMqttBridge({
  endpoint: 'https://api.example.com',
  mqtt: { brokerUrl: 'mqtt://localhost:1883', baseTopic: 'example' },
  logger: createConsoleLogger('info'),
  connections: [
    { name: 'ticker', method: 'GET', url: '/ticker' },
    {
      name: 'orders',
      method: 'POST',
      url: '/orders/stream',
      body: { region: 'eu' },
      topic: '{name}/{orderId}',
      // Reshape or drop events before publishing (return undefined to skip)
      transform: (data) => {
        const order = JSON.parse(data);
        return order.status === 'test' ? undefined : { id: order.orderId, total: order.total };
      }
    }
  ]
});

bridge.on('published', (connection, topic) => console.log(`${connection} -> ${topic}`));
bridge.on('error', (error, connection) => console.error(connection, error.message));

await bridge.start();
// later
await bridge.stop();

Connection topic can also be a function (data, connectionName) => string | string[] | undefined for full control.

Options

| Option | Description | | --- | --- | | endpoint | Base URL; each connection's url is resolved against it | | connections | Connections as described above; transform and function topics are library-only | | mqtt | { brokerUrl, baseTopic, username?, password?, qos?, retain?, clientOptions? }, or { client, baseTopic } to reuse an existing mqtt.js client (left open on stop()) | | tokenProvider | Anything with getBearerToken(): Promise<string> (and optionally invalidate()), e.g. BearerTokenProvider | | headers | Headers sent on every connection | | retry | { initialDelayMs = 2000, maxDelayMs = 30000, maxRetries = Infinity, jitter = 0.2 } | | connectTimeoutMs | Fail an attempt if no response headers arrive in time (default 30000) | | logger | { debug, info, warn, error }; compatible with console, pino and winston. Silent by default |

Events

| Event | Arguments | | --- | --- | | connected | connection | | disconnected | connection, error? | | reconnecting | connection, delayMs, attempt | | gaveUp | connection (after maxRetries) | | message | connection, event ({ data, event, id? }) | | published | connection, topic, payload | | error | error, connection: only emitted when you listen for it, so an unhandled error never crashes your process |

Authentication

import { BearerTokenProvider, BodyType } from 'sse-to-mqtt-node';

const tokenProvider = new BearerTokenProvider({
  url: 'https://auth.example.com/oauth/token',
  bodyType: BodyType.FormUrlEncoded, // or BodyType.Json
  body: { grant_type: 'client_credentials', client_id: '...', client_secret: '...' }
});

Tokens are cached until expires_in (minus expiryMarginMs, default 30s), concurrent requests share one token fetch, and a 401 from the SSE server triggers a fresh token. Use tokenFields if the token isn't in access_token or token.

Building blocks

SseDataProvider (a single reconnecting SSE connection), SseParser, MqttPublisher and loadConnectionsConfig are exported too, for custom pipelines.

Docker

The repository's Dockerfile runs the CLI with CONNECTIONS_CONFIG=config/connections.json. Mount your own config and pass the environment:

docker build -t sse-to-mqtt-node .
docker run --env-file .env -v "$PWD/connections.json:/app/config/connections.json:ro" sse-to-mqtt-node

Development

npm install
npm run typecheck && npm run lint && npm test
MQTT_TEST_URL=mqtt://localhost:1883 npm test   # also run broker integration tests
npm run build

A broker for the integration tests: docker run -p 1883:1883 eclipse-mosquitto:2 mosquitto -c /mosquitto-no-auth.conf

License

MIT