@finueva/pub-sub
v0.8.0
Published
Finueva Pub Sub Web SDK
Readme
@finueva/pub-sub
Use @finueva/pub-sub to receive live events, manage browser push, issue tickets, and publish events.
Pub Sub is experimental. It attempts delivery immediately. It does not retain, replay, order, retry, deduplicate, or acknowledge events.
Install the SDK
pnpm add @finueva/pub-subNode.js tools and server applications require Node.js 22 or later.
Opaque channel signaling
The next minor SDK adds client.createSignalingClient(...) to the canonical createPubSubClient result. It uses the same configured service origin and WebSocket factory for Cloudflare or an explicitly composed Node relay. The trusted-server client's issueChannelTicket({ userId, channel, publish, subscribe, expiresAt }) supplies a single-use grant after Application authorization; expiresAt is an absolute millisecond deadline at most 60 seconds ahead.
const signaling = await client.createSignalingClient({
userId: authorizedOwner,
channel: authorizedChannel,
getTicket: (signal) => requestApplicationChannelTicket(signal),
onPeers: (self, peers) => peerProvider.updatePeers(self, peers),
onEvent: ({ peerId, payload }) => peerProvider.receive(peerId, payload),
});
signaling.publish(opaquePayload, optionalRecipientPeerId);
await signaling.closed;userId is the room's owner namespace, not the connecting participant's personal identity. Participants authorized for that same owner/channel receive distinct signed peerId connection identities. Other owners with the same channel name are isolated. The current owner-room capacity is four connections. Data and peer presence end at capability expiry; the Application must reauthorize and reconnect. Pub Sub does not interpret the payload or implement WebRTC or resource policy.
connectChannel is also available from @finueva/pub-sub/channels. For a server-declared alternate Node route, pass its advertised absolute websocketUrl; it must share the configured service origin and contain no query, fragment or credentials. The server verifies the ticket's signed exact path. Both runtimes use channel/ticket subprotocols, never a ticket query parameter. See the repository's channel architecture for the Node provider contract and route composition. Consume the published minor release after service rollout; source availability does not mean these methods exist in an older installed SDK.
API summary
Browser client:
| Method | Description |
| --------------- | -------------------------------------------------- |
| connect() | Open a live event connection. |
| enablePush() | Enable push for the current browser installation. |
| disablePush() | Disable push for the current browser installation. |
| onPush() | Listen for push receipt and notification taps. |
Server client:
| Method | Description |
| --------------- | ------------------------------------------------- |
| issueTicket() | Issue a ticket for one subscriber operation. |
| publish() | Publish an event to one application-defined user. |
Create a browser client
Provide the authenticated user's ID and a getTicket function. The function must request a ticket from your application server.
import { createPubSubClient } from "@finueva/pub-sub";
const pubSub = await createPubSubClient({
userId: currentUser.id,
async getTicket({ action, signal }) {
const response = await fetch("/api/pub-sub/ticket", {
method: "POST",
credentials: "include",
headers: { "content-type": "application/json" },
body: JSON.stringify({ action }),
signal,
});
if (!response.ok) throw new Error("Pub Sub ticket unavailable");
return ((await response.json()) as { ticket: string }).ticket;
},
});Your server must get the user ID from authenticated session data. Do not accept a user ID from the request body. Never send a Pub Sub API token to browser code.
Production is the default environment. Without explicit public configuration, client creation makes one GET to /v1/configuration on the embedded production service origin, without cookies or Vault headers. It accepts a complete public bundle within three seconds and 16 KiB; redirects, failures, and invalid or partial responses preserve the complete embedded bundle without retry. Passing serviceOrigin, variables, or push.firebase skips this request and merges caller values over the selected embedded bundle. Pass environment: "staging" to select the staging service and fallback bundle. The service returns its cached configuration snapshot, not a fresh OpenBao read. Browser access requires an origin allowed by the service.
Receive live events
Call connect() and process events with onEvent.
const connection = await pubSub.connect({
onEvent(event) {
if (event.type === "order.updated") {
refreshOrder(event.data);
}
},
});
// Close the connection on sign-out or account change.
connection.close();
await connection.closed;event.data is unknown. Validate it before use. The SDK does not reconnect. Call connect() again to create a new connection. Events sent while disconnected are not replayed.
Registered hosted applications must set deliveryMode: "acknowledged" in connect(). This selects finueva.pubsub.v2 without falling back to v1. The SDK validates each bounded frame and admits it to the callback before returning matching transport credit. A synchronous callback failure closes the connection without credit. Credit confirms transport admission, not asynchronous persistence or message read status.
The hosted service closes slow connections when outstanding frame, byte, attachment or deadline bounds are reached. Recover current authorized state through the application's API after reconnect or a cursor gap. No live event body is retained for transport replay. Independent legacy application namespaces retain their separate existing protocol.
Enable browser push
Install the listener in the service worker that controls your application origin.
import { installPubSubServiceWorker } from "@finueva/pub-sub/service-worker";
installPubSubServiceWorker({
notificationIcon: "/notification.png",
openPath: "/inbox",
});Call enablePush() after a user action, such as a button click.
const stopPushEvents = pubSub.onPush((event) => {
if (event.kind === "tapped") {
openEvent(event.eventId, event.eventType);
}
});
const enabled = await pubSub.enablePush();
// Disable push on opt-out, sign-out, or account change.
await pubSub.disablePush();
await enabled.unlisten();
await stopPushEvents();Push events contain an event ID, an event type, and optional display text. Fetch current application data from an authorized API after the application opens.
Create a server client
Import the server entry point only from server code.
import { createPubSubServerClient } from "@finueva/pub-sub/server";
const pubSub = createPubSubServerClient({
apiToken: process.env.PUB_SUB_API_TOKEN!,
});Issue a ticket
const ticket = await pubSub.issueTicket({
userId: authenticatedUser.id,
action: "connect",
});Publish an event
const result = await pubSub.publish({
userId: recipient.id,
event: {
id: crypto.randomUUID(),
type: "order.updated",
occurredAt: new Date().toISOString(),
data: { orderId: order.id },
push: {
title: "Order updated",
body: "Open the application for details.",
},
},
});Omit push to send the event only to live connections. Keep private information out of push display text because the operating system can show it on a lock screen.
Handle errors
SDK operations throw PubSubError. Use its stable code property to select retry or user-interface behavior.
import { PubSubError } from "@finueva/pub-sub";
try {
await pubSub.connect({
onEvent(event) {
handleApplicationEvent(event);
},
});
} catch (error) {
if (error instanceof PubSubError && error.code === "OPERATION_TIMED_OUT") {
showRetryAction();
}
}Publishing is not idempotent. A retry after an unknown result can deliver the same event more than once.
Learn more
License
MIT
Private presence watches (unpublished candidate)
The server entry point provides registerPresenceWatch, renewPresenceWatch and cancelPresenceWatch. Configure the registered application alias as well as its apiToken. Generate a new watch UUID v4 and generation UUID for each new ordered selection, then retain the returned handle for renew/cancel. The Application must freshly authorize the audience and every selected subject before each operation; Pub Sub registration credentials do not establish organization membership. Returned timestamps are immutable server metadata, not an authorization proof. These methods are server-only; complete Application presence and Node provider parity remain unfinished.
setPresenceVisibility submits an ordered batch of at most fifty current Application-owned policy revisions. It requires the registered application alias, refuses changed or partial responses, and never accepts a caller clock or lease. A failed multi-shard operation may leave bounded grants on successful shards; retry from freshly read policy and retain the Application's disclosure checks. Missing or expired controls are hidden. This candidate does not establish complete Application presence or Node parity.
wakePresencePolicy submits the Application's current canonical policy identity: positive sequence, availabilityRevision, dndRevision, a UUID v4 changeId and reason. The Application registration must carry the operator-granted presence_policy_wake capability. A receipt means a live durable obligation was accepted, never delivery. PRESENCE_POLICY_CONFLICT means the command contradicts the retained identity (stale sequence, changed fields without a new sequence, reused change ID or regressed revision); do not retry it, re-read canonical state instead. PRESENCE_POLICY_RETIRED means the identical command's work already completed or expired; recover with a fresh sequence and change ID. PRESENCE_POLICY_RATE_LIMITED means a later retry of the same command may succeed. Unavailable or ambiguous outcomes keep the same command for a bounded retry. The client performs no automatic retries.
