Sign In

@msgmesh/sdk

Package Overview
Dependencies
Maintainers
1
Versions
10
Alerts
File Explorer

Advanced tools

Socket logo

Install Socket

Detect and block malicious and high-risk dependencies

Install

@msgmesh/sdk

MsgMesh TypeScript SDK — 多租戶事件總線的收發 / 即時(SSE·WS)/ 治理 client(Node 與瀏覽器通用)。

Source
npmnpm
Version
0.1.7
Version published
Weekly downloads
78
-49.35%
Maintainers
1
Weekly downloads
 
Created
Source

@msgmesh/sdk

English | 繁體中文

The TypeScript SDK for MsgMesh — a publish / consume / realtime (SSE / WebSocket) / governance client for the multi-tenant event bus, universal across Node and the browser (built on fetch). It is the single source of truth for the entire frontend ecosystem's HTTP contract: both the MCP server and the panel build on it.

Also available for Python: msgmesh — the same API surface in snake_case.

Quick start

Register an account in the panel and issue an API key (shown in plaintext only once), then:

import { MsgMesh } from "@msgmesh/sdk";

const mq = new MsgMesh({
  apiKey: process.env.MSGMESH_KEY,          // long-lived key, server-side only
  controlPlaneUrl: "https://cp.example.com",
  gatewayUrl: "https://gw.example.com",
  realtimeUrl: "https://rt.example.com",
});

await mq.createTopic("orders");
await mq.publish("orders", { hello: 1 });
const msgs = await mq.poll("orders", { group: "g1" });

Browsers / untrusted clients: use getToken, never embed an API key

Anything in the browser leaks — never put a long-lived API key in the frontend. Use a token-broker instead: your backend (which holds the key) calls POST /v1/tokens to exchange it for a short-lived dp token, and the frontend only ever holds the token. The SDK caches it, refetches before expiry, and rotates it on SSE reconnect.

const mq = new MsgMesh({
  // no apiKey; provide a function that fetches a short-lived token from your backend
  getToken: async () => (await fetch("/api/mm-token")).then((r) => r.json()), // { token, expires_in }
  gatewayUrl: "https://gw.example.com",
  realtimeUrl: "https://rt.example.com",
});
mq.stream("room.42", (data) => console.log(data));   // SSE; auto-rotates the token on expiry and reconnects
mq.streamWs("room.42", (data) => console.log(data)); // WebSocket; same interface, SDK-managed reconnect

Two options for realtime receive, with the same interface, each returning a stop function:

  • stream (SSE): backed by the browser-native EventSource (which auto-reconnects), over /…/sse. Browser only.
  • streamWs (WebSocket): the global WebSocket (browser-native; built into Node ≥ 22), over /…/ws. WebSocket has no native reconnect, so the SDK takes it over: reconnect with a 1s backoff after each drop, reset the failure counter on a successful connect, and stop once consecutive failures (never having connected) reach the limit — avoiding infinite reconnect to a revoked credential or a persistently unavailable endpoint; in getToken mode it also rotates the token before reconnecting. Revocation mid-connection is CLOSE 1008 authorization revoked (stops immediately); during the handshake (HTTP 401 → CloseEvent 1006) it is bounded by the failure limit. Good when SSE is blocked by a middlebox, or when you already have WebSocket infrastructure. Node < 22 has no global WebSocket and will throw — use subscribe (long-polling) instead.

Resume on reconnect (at-least-once, no gaps). Both stream (SSE) and streamWs (WebSocket) resume across reconnects: each message carries a <partition>-<offset> cursor, the SDK tracks the last one seen, and on reconnect it asks the server to replay from there — so messages dropped during a disconnect are backfilled, not lost. Delivery is at-least-once: the SDK dedupes per-partition by cursor, so a rare overlap is suppressed rather than delivered twice. If the server can't cover the gap (older than the replay window), it emits a resync signal — pass onResync to be told to re-fetch a snapshot. All of this is transparent: onMessage still receives the raw value string, no API change. (Resume requires the platform's realtime resume tier; against an older server the stream degrades gracefully to live-tail.)

Rooms

A single topic can be split into multiple rooms (room = Kafka record key), decoupling "number of rooms" from "number of topics". Two layers:

① Routing — publish with publish(topic, body, { key: roomId }) to target a room, and subscribe with the optional room (the fourth argument opts, same for stream / streamWs) to receive only that room:

mq.stream("chat", (data) => console.log(data), undefined, { room: "room-42" });   // only room-42
mq.streamWs("chat", (data) => console.log(data), undefined, { room: "room-42" });
await mq.publish("chat", { text: "hi" }, { key: "room-42" });                     // publish to room-42

Omitting room = receive all messages on the topic (backward compatible). Routing only filters — it does not enforce isolation; a malicious client can switch to someone else's room and eavesdrop on other rooms in the same topic. For real isolation, see ②.

② Isolation (platform-enforced) — add the optional rooms to a credential's capabilities and the platform enforces that the credential can only send/receive the named rooms (403 on overreach). rooms omitted/empty = all rooms (backward compatible); non-empty = only these. The typical approach: the backend holds an all-rooms key and downscopes it via POST /v1/tokens to mint a short-lived "this room only" token for the frontend (a downscope may only narrow, must be a subset of the key's capabilities, 403 on overreach):

// Backend token-broker: downscope an all-rooms key to a short-lived "chat / room-42 only"
// token, returned to the frontend as getToken
const r = await fetch(`${controlPlaneUrl}/v1/tokens`, {
  method: "POST",
  headers: { Authorization: `Bearer ${process.env.MSGMESH_KEY}`, "Content-Type": "application/json" },
  body: JSON.stringify({
    ttl_seconds: 600,
    capabilities: [{ ops: ["subscribe", "publish"], topics: ["chat"], rooms: ["room-42"] }],
  }),
});
const { token, expires_in } = await r.json();   // return to the frontend; it connects SSE/WS via getToken

You can also mint a persistent room-scoped key with createKey("key", { capabilities: [{ ops, topics, rooms }] }). Platform enforcement points: subscribe (SSE/WS) must carry a ?room within the allowed set (omitting it = wanting all rooms, also 403); publish ?key must be within the allowed set.

⚠️ A room-scoped credential can only use realtime (SSE/WS) + publish to its rooms; it cannot poll / consume / DLQ. Those are a whole-topic firehose (the consumer-group offset would consume other rooms; one group per room = read amplification) and can't be cleanly per-room filtered, so a room-restricted credential always gets 403 (use realtime SSE/WS ?room=). Use an unrestricted credential when you need poll/consume.

Room isolation security notes (must read)

  • Isolation strength = the scope of the token you issue. Isolation only exists when the backend downscopes an all-rooms key into a room-scoped token for the frontend. Never put an unrestricted credential (a full key, or a token without rooms) into the frontend / untrusted clients — that lets anyone change room and see all rooms, so isolation is meaningless.
  • The platform does not verify "who the sender is." Room isolation governs "which rooms you can send/receive," not "who you are in the room." Within a room, anyone holding that room's token can impersonate any sender in the payload. To prevent in-room impersonation: mint a token per user on the backend and stamp / verify the sender there, don't let untrusted clients self-report identity.
  • Note: presence (online count) is currently per-topic, not per-room (only leaks an aggregate number); short-lived tokens are bearer tokens — leaking one = usable for that room until TTL expires (so keep the TTL short and don't log it).

Production configuration (must read)

  • Always set the service URLs explicitly: controlPlaneUrl (governance API), gatewayUrl (send/receive), realtimeUrl (SSE/WS/presence). When unset, the SDK falls back to local-dev defaults (http://localhost:8080/8081/8082), which are for local use only; when a production call can't connect, the error message appends a "set controlPlaneUrl/gatewayUrl/realtimeUrl" hint.
  • Credential handling: an API key is returned in plaintext only once, at creation — never commit it to a repo or write it to logs. Never put any API key in the browser; use getToken instead.
  • Scopes: publishing needs producer, consuming needs consumer, administration needs admin (which covers everything); a data-plane key that both publishes and consumes, or needs fine-grained access, can use the neutral key scope (which requires capabilities). getToken mode is data-plane only (send/receive); calling governance endpoints returns 401/403.

Error handling

Non-2xx responses throw a typed error by status code (all inherit MsgMeshError and carry status/code/path):

import { ValidationError, AuthError, NotFoundError, RateLimitError } from "@msgmesh/sdk";

try {
  await mq.createTopic("Bad Name!");
} catch (e) {
  if (e instanceof ValidationError) console.error("invalid argument:", e.message);
  else if (e instanceof RateLimitError) console.error("rate limited, retry later");
  else throw e;
}
StatusTypecode
400 / 422ValidationErrorvalidation
401 / 403AuthErrorauth
404NotFoundErrornot_found
429RateLimitErrorrate_limit
otherMsgMeshErrorserver

API overview

  • Topics: createTopic / listTopics / deleteTopic
  • Send/receive: publish / poll / subscribe (polling) / stream (SSE, browser) / streamWs (WebSocket, browser + Node ≥ 22) / getPresence
  • Keys: listKeys (returns capabilities / name) / createKey (accepts scope + capabilities) / deleteKey
  • Webhooks: listWebhooks / createWebhook / deleteWebhook / reactivateWebhook
  • Schemas: registerSchema / listSchemas / getLatestSchema / deleteSchema
  • Functions: registerFunction / getFunction / deleteFunction (JavaScript / WASM)
  • Plan: getPlan / setPlan; usage: getUsage
  • Settings: getSettings / setStrictTopics (data-plane topic gate toggle)
  • Billing (crypto PAYG prepaid): getBilling / getDepositAddresses / getDeposits / getLedger / getUsageDebits / getDepositStatus
  • Misc: getSnippet / getDocs / getAudit, DLQ dlqPeek / dlqReplay

Registration and admin (finance / tenant governance) go through panel sessions, not this SDK.

Development

npm test -w @msgmesh/sdk && npm run build -w @msgmesh/sdk

@msgmesh/sdk · 繁體中文

English | 繁體中文

MsgMesh 的 TypeScript SDK(Node / 瀏覽器通用,基於 fetch)。是整個前端生態的唯一 HTTP 真相源——MCP server 與面板都複用它。

也有 Python 版:msgmesh(相同 API 面、snake_case)。

快速開始

先在面板註冊帳號、簽發一把 API key(明文僅顯示一次),再:

import { MsgMesh } from "@msgmesh/sdk";

const mq = new MsgMesh({
  apiKey: process.env.MSGMESH_KEY,          // 伺服器端用長期 key
  controlPlaneUrl: "https://cp.example.com",
  gatewayUrl: "https://gw.example.com",
  realtimeUrl: "https://rt.example.com",
});

await mq.createTopic("orders");
await mq.publish("orders", { hello: 1 });
const msgs = await mq.poll("orders", { group: "g1" });

瀏覽器/不可信端:用 getToken,不要放 API key

瀏覽器裡的東西都會外洩,不要把長期 API key 放進前端。改採 token-broker:後端(持 key)呼叫 POST /v1/tokens 換一張短期 dp token,前端只拿 token;SDK 會自動快取、將過期前重取,SSE 重連時亦換新。

const mq = new MsgMesh({
  // 不放 apiKey;給一個「去我後端拿短期 token」的函式
  getToken: async () => (await fetch("/api/mm-token")).then((r) => r.json()), // { token, expires_in }
  gatewayUrl: "https://gw.example.com",
  realtimeUrl: "https://rt.example.com",
});
mq.stream("room.42", (data) => console.log(data));   // SSE;token 過期自動換新後重連
mq.streamWs("room.42", (data) => console.log(data)); // WebSocket;同介面,SDK 自管重連

即時接收兩種選擇,介面一致、皆回傳停止函式:

  • stream(SSE):靠瀏覽器原生 EventSource(有原生自動重連),走 /…/sse僅瀏覽器
  • streamWs(WebSocket):全域 WebSocket(瀏覽器原生;Node ≥ 22 內建),走 /…/wsWebSocket 無原生重連,故由 SDK 接管:每次斷線退避 1s 重連,成功連上即重置失敗計數;連續失敗達上限(未曾連上)即停止(避免對已撤銷的憑證或持續不可用的端點無限重連),getToken 模式另在重連前換新 token。撤權若在連線中發生為 CLOSE 1008 authorization revoked(立即停);若在握手期(HTTP 401→CloseEvent 1006)則由上限收口。適合 SSE 被中間層擋掉、或已有 WS 基礎設施的場景。Node < 22 無全域 WebSocket 會拋錯,改用 subscribe(長輪詢)。

多房間(rooms)

一個 topic 內可再切多個房間(room = Kafka record key),脫鉤「房間數」與「topic 數」。分兩層:

① 路由——發佈時用 publish(topic, body, { key: roomId }) 指定房間,訂閱時傳選用 room(第四參數 opts,stream/streamWs 皆同)只收該房間:

mq.stream("chat", (data) => console.log(data), undefined, { room: "room-42" });   // 只收 room-42
mq.streamWs("chat", (data) => console.log(data), undefined, { room: "room-42" });
await mq.publish("chat", { text: "hi" }, { key: "room-42" });                     // 發到 room-42

省略 room=收該 topic 全部訊息(向後相容)。路由本身只做過濾、無強制隔離——惡意 client 可改成別人的 room 偷聽同 topic 其他房間。要真隔離看 ②。

② 隔離(平台強制)——把憑證的 capabilities 加上選用 rooms,平台即強制該憑證只能收發指定房間(逾越 403)。rooms 省略/空 = 所有房間(向後相容);非空 = 僅限這些。典型作法是後端持一把全房間金鑰,向 POST /v1/tokens 降權簽出「只准某房間」的短期 token 給前端(降權只准更窄、須為金鑰能力子集,逾越 403):

// 後端 token-broker:用全房間 key 降權鑄「只准 chat / room-42」的短期 token,回給前端當 getToken
const r = await fetch(`${controlPlaneUrl}/v1/tokens`, {
  method: "POST",
  headers: { Authorization: `Bearer ${process.env.MSGMESH_KEY}`, "Content-Type": "application/json" },
  body: JSON.stringify({
    ttl_seconds: 600,
    capabilities: [{ ops: ["subscribe", "publish"], topics: ["chat"], rooms: ["room-42"] }],
  }),
});
const { token, expires_in } = await r.json();   // 回給前端;前端以 getToken 用它連 SSE/WS

也可用 createKey("key", { capabilities: [{ ops, topics, rooms }] }) 簽一把常駐 room-scoped 鍵。平台強制點:訂閱(SSE/WS)必須帶允許集內的 ?room(不帶=想收全部房間,一樣 403);發佈的 ?key 必須 ∈ 允許集。

⚠️ room-scoped 憑證只能走即時(SSE/WS)+ 對其房間 publish;不能 poll / consume / DLQ。後者是整個 topic 的 firehose(consumer-group offset 會吃掉別房間、每房一 group = 讀取放大),無法乾淨 per-room 過濾,受限房間憑證一律 403(use realtime SSE/WS ?room=)。需要 poll/consume 時請改用不限房間的憑證。

房間隔離的安全須知(必讀)

  • 隔離強度 = 你發的 token 範圍。 只有在「後端用全房間金鑰降權鑄 room-scoped token 給前端」時才有隔離。別把不限房間的憑證(全權 key、或沒有 rooms 的 token)放進前端——那樣 client 改個 room 就能看到所有房間,隔離形同虛設。
  • 平台不驗「發訊者是誰」。 房間隔離管的是「能收發哪些房間」,不是「你是房裡的誰」。同一房內,任何持該房 token 的人都能在 payload 裡冒充任何 sender。要防房內冒名:後端為每個使用者各自鑄 token、並由後端戳上 / 驗證 sender,別讓前端自報身分。
  • 附帶:presence(在線數)目前是 per-topic 非 per-room(只洩漏聚合數字);短期 token 為 bearer,洩漏 = 該房 ≤TTL 可用(故 TTL 短、勿記進 log)。

生產環境設定(必讀)

  • 務必明確指定服務 URL:controlPlaneUrl(治理 API)、gatewayUrl(收發)、realtimeUrl(SSE/WS/presence)。 未指定時 SDK 退回本機開發預設值(http://localhost:8080/8081/8082),只適合本機; 生產環境連不上時錯誤訊息會附「請指定 controlPlaneUrl/gatewayUrl/realtimeUrl」提示。
  • 憑證保管:API key 只在建立時回傳一次明文,不要寫進 repo 或日誌。瀏覽器端勿放任何 API key,改用 getToken
  • scope:發訊需 producer、收訊需 consumer、管理需 admin(通吃);又推又收/細粒度的資料面鍵可用中性 key scope(強制附 capabilities)。getToken 模式僅供資料面(收發),呼叫治理端點會 401/403。

錯誤處理

非 2xx 回應會依狀態碼拋型別化錯誤(都繼承 MsgMeshError,帶 status/code/path):

import { ValidationError, AuthError, NotFoundError, RateLimitError } from "@msgmesh/sdk";

try {
  await mq.createTopic("Bad Name!");
} catch (e) {
  if (e instanceof ValidationError) console.error("參數不合法:", e.message);
  else if (e instanceof RateLimitError) console.error("被限流,稍後重試");
  else throw e;
}
狀態碼型別code
400 / 422ValidationErrorvalidation
401 / 403AuthErrorauth
404NotFoundErrornot_found
429RateLimitErrorrate_limit
其他MsgMeshErrorserver

API 一覽

  • Topics:createTopic / listTopics / deleteTopic
  • 收發:publish / poll / subscribe(輪詢)/ stream(SSE,瀏覽器)/ streamWs(WebSocket,瀏覽器 + Node ≥ 22)/ getPresence
  • Keys:listKeys(回傳含 capabilities/name)/ createKey(可帶 scope+capabilities)/ deleteKey
  • Webhooks:listWebhooks / createWebhook / deleteWebhook / reactivateWebhook
  • Schemas:registerSchema / listSchemas / getLatestSchema / deleteSchema
  • Functions:registerFunction / getFunction / deleteFunction(JavaScript / WASM)
  • 方案:getPlan / setPlan;用量:getUsage
  • 設定:getSettings / setStrictTopics(資料面 topic 閘門開關)
  • 帳務(加密貨幣 PAYG 預付):getBilling / getDepositAddresses / getDeposits / getLedger / getUsageDebits / getDepositStatus
  • 其他:getSnippet / getDocs / getAudit、DLQ dlqPeek / dlqReplay

註冊、超管(finance/租戶治理)走面板 session,不在本 SDK。

開發

npm test -w @msgmesh/sdk && npm run build -w @msgmesh/sdk

Keywords

msgmesh

FAQs

Package last updated on 24 Jul 2026

Related posts