diff --git a/AGENTS.md b/AGENTS.md index 1a7d2a43..f5a3938f 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -2,7 +2,7 @@ This file is for AI agents working inside applications scaffolded by Arkstack. It assumes the agent has access to the generated project only, not to the Arkstack framework monorepo. -Read `SKILLS.md` first, then use these workflows to combine the available skills safely. +Read `SKILL.md` first, then use these workflows to combine the available skills safely. ## Operating Rules diff --git a/packages/realtime/README.md b/packages/realtime/README.md index 087280b2..ff3bdaa4 100644 --- a/packages/realtime/README.md +++ b/packages/realtime/README.md @@ -44,6 +44,28 @@ unsubscribe(); await realtime.disconnect(); ``` +### Client events + +Pusher private and presence channels can send ephemeral client events (often +called whispers). Subscribe to the channel before sending an event: + +```ts +const stopTyping = await realtime.listenForWhisper( + 'private-room.7', + 'typing', + ({ userId }) => console.log(`${userId} is typing`), +); + +await realtime.whisper('private-room.7', 'typing', { userId: user.id }); +stopTyping(); +``` + +The `client-` prefix is added automatically. Pusher sends client events over its +private/presence channel. Firebase uses Realtime Database to publish the same +ephemeral events across connected clients. Use `listen(channel, event, handler)` +and `trigger(channel, event, payload)` when you need the lower-level APIs with +exact event names. + Each `notification` matches the payload broadcast by the server: ```ts @@ -115,10 +137,16 @@ const realtime = createRealtime({ projectId: import.meta.env.VITE_FIREBASE_PROJECT_ID, appId: import.meta.env.VITE_FIREBASE_APP_ID, messagingSenderId: import.meta.env.VITE_FIREBASE_SENDER_ID, + databaseURL: import.meta.env.VITE_FIREBASE_DATABASE_URL, }, }); ``` +Firebase client events are written below `arkstack/client-events` in Realtime +Database and removed immediately after publishing. Configure Firebase Security +Rules for the channels your authenticated clients may read and write. Override +the root with `firebase.clientEventsPath` when needed. + ## Custom transport Provide `transportFactory` to bridge any backend (a raw WebSocket, SSE, a test double, …): @@ -146,6 +174,9 @@ const realtime = createRealtime({ transportFactory: () => transport }); - `createRealtime(config)` — create a `RealtimeClient`. Config: `transport` (`'pusher'` | `'firebase'`), `event` (default `notification`), `channelPrefix` (default `user.`), `pusher`/`firebase` credentials, or a custom `transportFactory`. - `client.subscribe(channel, handler)` / `client.forUser(userId, handler)` — subscribe; returns an unsubscribe function. +- `client.listen(channel, event, handler)` — listen for an arbitrary event. +- `client.listenForWhisper(channel, event, handler)` / `client.whisper(channel, event, payload)` — receive and send Pusher client events. +- `client.trigger(channel, event, payload)` — emit an exact event name through a transport that supports it. - `client.channelFor(userId)` — the per-user channel name. - `client.disconnect()` — tear down the transport connection. - `@arkstack/realtime/react` — `useNotifications(client, channel, { limit? })` → `{ notifications, latest, clear }`. diff --git a/packages/realtime/src/RealtimeClient.ts b/packages/realtime/src/RealtimeClient.ts index ef2381f0..61a5052c 100644 --- a/packages/realtime/src/RealtimeClient.ts +++ b/packages/realtime/src/RealtimeClient.ts @@ -1,4 +1,4 @@ -import type { NotificationHandler, RealtimeConfig, RealtimeSubscription, RealtimeTransport } from './types' +import type { NotificationHandler, RealtimeConfig, RealtimeEventHandler, RealtimeSubscription, RealtimeTransport } from './types' /** * Consumes Arkstack realtime notifications. Resolves a transport (Pusher, @@ -15,18 +15,23 @@ export class RealtimeClient { this.channelPrefix = config.channelPrefix ?? 'user.' } - /** The channel name a given user's notifications are broadcast on. */ - channelFor (userId: string | number): string { + /** + * The channel name a given user's notifications are broadcast on. + * + * @param userId + * @returns + */ + channelFor(userId: string | number): string { return `${this.channelPrefix}${userId}` } - private transport (): Promise { + private transport(): Promise { this.transportPromise ??= this.resolveTransport() return this.transportPromise } - private async resolveTransport (): Promise { + private async resolveTransport(): Promise { if (this.config.transportFactory) { return await this.config.transportFactory() } @@ -56,25 +61,93 @@ export class RealtimeClient { * @param channel The channel name (e.g. `user.7`). * @param handler Called with each incoming notification. */ - async subscribe (channel: string, handler: NotificationHandler): Promise<() => void> { + async subscribe(channel: string, handler: NotificationHandler): Promise<() => void> { + return await this.listen(channel, this.event, handler as RealtimeEventHandler) + } + + /** + * Listen for an arbitrary event on a channel. + * + * @param channel + * @param event + * @param handler + * @returns + */ + async listen( + channel: string, + event: string, + handler: RealtimeEventHandler, + ): Promise<() => void> { const transport = await this.transport() - const subscription: RealtimeSubscription = await transport.subscribe(channel, this.event, handler) + const subscription: RealtimeSubscription = await transport.subscribe( + channel, + event, + handler as RealtimeEventHandler, + ) return () => subscription.unsubscribe() } + /** + * Listen for a client event (`client-{event}`). + * + * @param channel + * @param event + * @param handler + * @returns + */ + async listenForWhisper( + channel: string, + event: string, + handler: RealtimeEventHandler, + ): Promise<() => void> { + return await this.listen(channel, this.clientEventName(event), handler) + } + + /** + * Emit a client event (`client-{event}`) on a subscribed channel. + * + * @param channel + * @param event + * @param payload + */ + async whisper(channel: string, event: string, payload: Payload): Promise { + await this.trigger(channel, this.clientEventName(event), payload) + } + + /** + * Emit an event through a transport that supports client-originated events. + * + * @param channel + * @param event + * @param payload + */ + async trigger(channel: string, event: string, payload: Payload): Promise { + const transport = await this.transport() + + if (!transport.trigger) { + throw new Error('Realtime: the configured transport does not support client events') + } + + await transport.trigger(channel, event, payload) + } + + private clientEventName(event: string): string { + return event.startsWith('client-') ? event : `client-${event}` + } + /** * Subscribe to a user's channel (`{channelPrefix}{userId}`). * * @param userId The user id. * @param handler Called with each incoming notification. */ - async forUser (userId: string | number, handler: NotificationHandler): Promise<() => void> { + async forUser(userId: string | number, handler: NotificationHandler): Promise<() => void> { return await this.subscribe(this.channelFor(userId), handler) } /** Tear down the underlying transport connection. */ - async disconnect (): Promise { + async disconnect(): Promise { if (!this.transportPromise) { return } @@ -85,5 +158,7 @@ export class RealtimeClient { } } -/** Create a {@link RealtimeClient}. */ -export const createRealtime = (config: RealtimeConfig = {}): RealtimeClient => new RealtimeClient(config) +/** Create a new {@link RealtimeClient}. */ +export const createRealtime = ( + config: RealtimeConfig = {} +): RealtimeClient => new RealtimeClient(config) diff --git a/packages/realtime/src/transports/firebase.ts b/packages/realtime/src/transports/firebase.ts index 3886c753..e90b41ca 100644 --- a/packages/realtime/src/transports/firebase.ts +++ b/packages/realtime/src/transports/firebase.ts @@ -1,4 +1,4 @@ -import type { FirebaseClientConfig, NotificationHandler, RealtimeTransport } from '../types' +import type { FirebaseClientConfig, RealtimeEventHandler, RealtimeTransport } from '../types' /** The slice of `firebase/messaging` this transport uses. */ interface FirebaseMessagePayload { @@ -7,6 +7,33 @@ interface FirebaseMessagePayload { type OnMessage = (messaging: unknown, next: (payload: FirebaseMessagePayload) => void) => () => void +interface DatabaseReference { + readonly key?: string | null +} + +interface DatabaseSnapshot { + val(): unknown +} + +interface FirebaseClientEvent { + sender: string + payload: unknown +} + +interface FirebaseDatabaseModule { + getDatabase(app: unknown, url?: string): unknown + ref(database: unknown, path: string): DatabaseReference + onChildAdded(reference: DatabaseReference, handler: (snapshot: DatabaseSnapshot) => void): () => void + push(reference: DatabaseReference, payload: unknown): Promise + remove(reference: DatabaseReference): Promise +} + +const pathSegment = (value: string): string => encodeURIComponent(value).replace(/\./g, '%2E') + +const isFirebaseClientEvent = (value: unknown): value is FirebaseClientEvent => { + return typeof value === 'object' && value !== null && 'sender' in value && 'payload' in value +} + /** * Realtime transport backed by [Firebase Cloud Messaging](https://firebase.google.com/docs/cloud-messaging/js/receive) * foreground messages. `firebase` is an optional peer dependency imported lazily. @@ -17,10 +44,12 @@ type OnMessage = (messaging: unknown, next: (payload: FirebaseMessagePayload) => export const createFirebaseTransport = async (config: FirebaseClientConfig): Promise => { const appSpecifier = 'firebase/app' const messagingSpecifier = 'firebase/messaging' + const databaseSpecifier = 'firebase/database' - const [appMod, messagingMod] = await Promise.all([ + const [appMod, messagingMod, databaseMod] = await Promise.all([ import(appSpecifier), import(messagingSpecifier), + import(databaseSpecifier), ]).catch(() => { throw new Error( 'The "firebase" package is required for the Firebase transport. Install it with `npm i firebase`.', @@ -32,13 +61,25 @@ export const createFirebaseTransport = async (config: FirebaseClientConfig): Pro projectId: config.projectId, appId: config.appId, messagingSenderId: config.messagingSenderId, + databaseURL: config.databaseURL, }) const messaging = messagingMod.getMessaging(app) const onMessage = messagingMod.onMessage as OnMessage + const messageUnsubscribers = new Set<() => void>() + const databaseUnsubscribers = new Set<() => void>() + const databaseApi = databaseMod as unknown as FirebaseDatabaseModule + const database = databaseApi.getDatabase(app, config.databaseURL) + const clientEventsPath = config.clientEventsPath ?? 'arkstack/client-events' + const clientId = globalThis.crypto?.randomUUID?.() ?? `${Date.now()}-${Math.random()}` + + const clientEventReference = (channel: string, event: string): DatabaseReference => databaseApi.ref( + database, + `${clientEventsPath}/${pathSegment(channel)}/${pathSegment(event)}`, + ) const transport: RealtimeTransport = { - subscribe(channel: string, event: string, handler: NotificationHandler) { + subscribe(channel: string, event: string, handler: RealtimeEventHandler) { const off = onMessage(messaging, (payload) => { if (payload.data?.event !== event || !payload.data.payload) { return @@ -50,10 +91,47 @@ export const createFirebaseTransport = async (config: FirebaseClientConfig): Pro /** Ignore malformed payloads. */ } }) + messageUnsubscribers.add(off) + const offClientEvent = event.startsWith('client-') + ? databaseApi.onChildAdded(clientEventReference(channel, event), (snapshot) => { + const clientEvent = snapshot.val() + + if (isFirebaseClientEvent(clientEvent) && clientEvent.sender !== clientId) { + handler(clientEvent.payload) + } + }) + : undefined + + if (offClientEvent) { + databaseUnsubscribers.add(offClientEvent) + } - return { channel, unsubscribe: off } + return { + channel, + unsubscribe() { + off() + messageUnsubscribers.delete(off) + if (offClientEvent) { + offClientEvent() + databaseUnsubscribers.delete(offClientEvent) + } + }, + } + }, + async trigger(channel: string, event: string, payload: unknown) { + const eventReference = await databaseApi.push(clientEventReference(channel, event), { + sender: clientId, + payload, + }) + + await databaseApi.remove(eventReference) + }, + disconnect() { + messageUnsubscribers.forEach((off) => off()) + messageUnsubscribers.clear() + databaseUnsubscribers.forEach((off) => off()) + databaseUnsubscribers.clear() }, - disconnect() { }, } return transport diff --git a/packages/realtime/src/transports/pusher.ts b/packages/realtime/src/transports/pusher.ts index 29b45869..9975dd85 100644 --- a/packages/realtime/src/transports/pusher.ts +++ b/packages/realtime/src/transports/pusher.ts @@ -1,9 +1,10 @@ -import type { NotificationHandler, PusherClientConfig, RealtimeTransport } from '../types' +import type { PusherClientConfig, RealtimeEventHandler, RealtimeTransport } from '../types' /** The slice of the `pusher-js` client this transport uses. */ interface PusherChannel { bind(event: string, handler: (data: unknown) => void): void unbind(event: string, handler: (data: unknown) => void): void + trigger(event: string, data: unknown): boolean } interface PusherClient { @@ -39,23 +40,49 @@ export const createPusherTransport = async ( authEndpoint: config.authEndpoint ?? `${config.apiBase}/realtime/auth`, auth: config.auth, }) + const channels = new Map() const transport: RealtimeTransport = { - subscribe(channel: string, event: string, handler: NotificationHandler) { - const subscription = client.subscribe(channel) + subscribe(channel: string, event: string, handler: RealtimeEventHandler) { + const active = channels.get(channel) + const subscription = active?.channel ?? client.subscribe(channel) const listener = (data: unknown) => handler(data as never) + channels.set(channel, { + channel: subscription, + subscriptions: (active?.subscriptions ?? 0) + 1, + }) + subscription.bind(event, listener) return { channel, unsubscribe() { subscription.unbind(event, listener) - client.unsubscribe(channel) + const current = channels.get(channel) + + if (!current || current.subscriptions <= 1) { + channels.delete(channel) + client.unsubscribe(channel) + } else { + current.subscriptions-- + } }, } }, + trigger(channel: string, event: string, payload: unknown) { + const subscription = channels.get(channel)?.channel + + if (!subscription) { + throw new Error(`Realtime: subscribe to channel "${channel}" before triggering client events`) + } + + if (!subscription.trigger(event, payload)) { + throw new Error(`Realtime: failed to trigger client event "${event}" on channel "${channel}"`) + } + }, disconnect() { + channels.clear() client.disconnect() }, } diff --git a/packages/realtime/src/types.ts b/packages/realtime/src/types.ts index 193d5d12..4f9f557b 100644 --- a/packages/realtime/src/types.ts +++ b/packages/realtime/src/types.ts @@ -19,6 +19,9 @@ export type RealtimeTransportName = 'pusher' | 'firebase' export type NotificationHandler = (notification: RealtimeNotification) => void +/** Handles an arbitrary realtime event payload. */ +export type RealtimeEventHandler = (payload: Payload) => void + /** A live subscription to one channel; call `unsubscribe()` to stop listening. */ export interface RealtimeSubscription { channel: string @@ -31,11 +34,33 @@ export interface RealtimeSubscription { * supplied via {@link RealtimeConfig.transportFactory} for custom backends/tests. */ export interface RealtimeTransport { + /** + * Subscribe to a realtime channel + * + * @param channel + * @param event + * @param handler + */ subscribe( channel: string, event: string, - handler: NotificationHandler, + handler: RealtimeEventHandler, ): RealtimeSubscription | Promise + /** + * Emit an event from the connected client, when supported by the transport. + * + * @param channel + * @param event + * @param payload + */ + trigger?( + channel: string, + event: string, + payload: unknown, + ): void | Promise + /** + * Disconnect from a connected realtime channel + */ disconnect(): void | Promise } @@ -64,6 +89,10 @@ export interface FirebaseClientConfig { projectId: string appId: string messagingSenderId: string + /** Realtime Database URL. Uses the Firebase project's default database when omitted. */ + databaseURL?: string + /** Root path used for ephemeral client events (default `arkstack/client-events`). */ + clientEventsPath?: string /** Web push VAPID key used when requesting a messaging token. */ vapidKey?: string } diff --git a/packages/realtime/tests/realtime.test.ts b/packages/realtime/tests/realtime.test.ts index 03cf882a..3ddb5667 100644 --- a/packages/realtime/tests/realtime.test.ts +++ b/packages/realtime/tests/realtime.test.ts @@ -1,13 +1,50 @@ import { describe, expect, it, vi } from 'vitest' import type { NotificationHandler, RealtimeNotification, RealtimeTransport } from '../src' -import { createRealtime } from '../src' +import { createFirebaseTransport, createRealtime } from '../src' + +const firebaseDatabase = vi.hoisted(() => { + const listeners = new Map void>>() + const removed: string[] = [] + + return { listeners, removed } +}) + +vi.mock('firebase/app', () => ({ + initializeApp: vi.fn(() => ({})), +})) + +vi.mock('firebase/messaging', () => ({ + getMessaging: vi.fn(() => ({})), + onMessage: vi.fn(() => vi.fn()), +})) + +vi.mock('firebase/database', () => ({ + getDatabase: vi.fn(() => ({})), + ref: vi.fn((_database: unknown, path: string) => ({ path })), + onChildAdded: vi.fn((reference: { path: string }, handler: (snapshot: { val(): unknown }) => void) => { + const handlers = firebaseDatabase.listeners.get(reference.path) ?? new Set() + handlers.add(handler) + firebaseDatabase.listeners.set(reference.path, handlers) + + return () => handlers.delete(handler) + }), + push: vi.fn(async (reference: { path: string }, payload: unknown) => { + firebaseDatabase.listeners.get(reference.path)?.forEach((handler) => handler({ val: () => payload })) + + return { path: `${reference.path}/generated-key` } + }), + remove: vi.fn(async (reference: { path: string }) => { + firebaseDatabase.removed.push(reference.path) + }), +})) /** A fake transport that lets a test push notifications into subscribers. */ const fakeTransport = () => { const handlers = new Map>() const unsubscribed: string[] = [] let disconnected = false + const triggered: Array<{ channel: string, event: string, payload: unknown }> = [] const transport: RealtimeTransport = { subscribe (channel, _event, handler) { @@ -23,6 +60,9 @@ const fakeTransport = () => { }, } }, + trigger (channel, event, payload) { + triggered.push({ channel, event, payload }) + }, disconnect () { disconnected = true }, @@ -32,7 +72,7 @@ const fakeTransport = () => { handlers.get(channel)?.forEach((h) => h(notification)) } - return { transport, emit, unsubscribed, isDisconnected: () => disconnected } + return { transport, emit, triggered, unsubscribed, isDisconnected: () => disconnected } } const sample = (over: Partial = {}): RealtimeNotification => ({ @@ -87,6 +127,31 @@ describe('RealtimeClient', () => { expect(unsubscribed).toContain('user.7') }) + it('listens for and emits client events with the client prefix', async () => { + const { transport, triggered } = fakeTransport() + const subscribe = vi.spyOn(transport, 'subscribe') + const client = createRealtime({ transportFactory: () => transport }) + const typing: Array<{ userId: number }> = [] + + await client.listenForWhisper<{ userId: number }>('private-room.7', 'typing', (payload) => typing.push(payload)) + await client.whisper('private-room.7', 'typing', { userId: 7 }) + + expect(subscribe).toHaveBeenCalledWith('private-room.7', 'client-typing', expect.any(Function)) + expect(triggered).toEqual([{ + channel: 'private-room.7', + event: 'client-typing', + payload: { userId: 7 }, + }]) + }) + + it('reports transports that do not support client events', async () => { + const { transport } = fakeTransport() + delete transport.trigger + const client = createRealtime({ transportFactory: () => transport }) + + await expect(client.whisper('private-room.7', 'typing', {})).rejects.toThrow(/does not support client events/i) + }) + it('resolves the transport once and disconnects it', async () => { const factory = vi.fn(() => fakeTransport().transport) const client = createRealtime({ transportFactory: factory }) @@ -107,3 +172,45 @@ describe('RealtimeClient', () => { await expect(client.subscribe('user.7', () => undefined)).rejects.toThrow(/pusher/i) }) }) + +describe('Firebase client events', () => { + it('publishes client events through Realtime Database with channel isolation', async () => { + firebaseDatabase.listeners.clear() + firebaseDatabase.removed.length = 0 + const receiver = await createFirebaseTransport({ + apiKey: 'key', + projectId: 'project', + appId: 'app', + messagingSenderId: 'sender', + }) + const sender = await createFirebaseTransport({ + apiKey: 'key', + projectId: 'project', + appId: 'app', + messagingSenderId: 'sender', + }) + const received: Array<{ userId: number }> = [] + const otherChannel: Array<{ userId: number }> = [] + const subscription = await receiver.subscribe( + 'private-room.7', + 'client-typing', + (payload) => received.push(payload as { userId: number }), + ) + await receiver.subscribe( + 'private-room.8', + 'client-typing', + (payload) => otherChannel.push(payload as { userId: number }), + ) + + await sender.trigger?.('private-room.7', 'client-typing', { userId: 7 }) + expect(received).toEqual([{ userId: 7 }]) + expect(otherChannel).toEqual([]) + expect(firebaseDatabase.removed).toEqual([ + 'arkstack/client-events/private-room%2E7/client-typing/generated-key', + ]) + + subscription.unsubscribe() + await sender.trigger?.('private-room.7', 'client-typing', { userId: 8 }) + expect(received).toEqual([{ userId: 7 }]) + }) +})