125 lines
3.5 KiB
TypeScript
125 lines
3.5 KiB
TypeScript
import type { Handler, PubSub, PubSubDriver } from "./index.ts";
|
|
|
|
export interface MessageEnvelope<T = unknown> {
|
|
id: string;
|
|
topic: string;
|
|
data: T;
|
|
timestamp: number;
|
|
attempts: number;
|
|
}
|
|
|
|
export interface ResilientPubSubOptions {
|
|
retries?: number;
|
|
retryDelayMs?: number;
|
|
onError?: (error: unknown, envelope: MessageEnvelope) => void;
|
|
}
|
|
|
|
export function createResilientPubSub(
|
|
driver: PubSubDriver,
|
|
options: ResilientPubSubOptions = {},
|
|
): PubSub {
|
|
const retries = options.retries ?? 2;
|
|
const retryDelayMs = options.retryDelayMs ?? 100;
|
|
if (!Number.isInteger(retries) || retries < 0) {
|
|
throw new RangeError("pubsub retries must be a non-negative integer");
|
|
}
|
|
if (!Number.isFinite(retryDelayMs) || retryDelayMs < 0) {
|
|
throw new RangeError("pubsub retryDelayMs must be a non-negative number");
|
|
}
|
|
|
|
return {
|
|
async publish(topic, message) {
|
|
if (!topic.trim()) throw new TypeError("pubsub topic cannot be empty");
|
|
const envelope: MessageEnvelope = {
|
|
id: crypto.randomUUID(),
|
|
topic,
|
|
data: message,
|
|
timestamp: Date.now(),
|
|
attempts: 0,
|
|
};
|
|
let lastError: unknown;
|
|
for (let attempt = 0; attempt <= retries; attempt++) {
|
|
envelope.attempts = attempt + 1;
|
|
try {
|
|
await driver.publish(topic, envelope);
|
|
return;
|
|
} catch (error) {
|
|
lastError = error;
|
|
if (attempt < retries && retryDelayMs > 0) {
|
|
await new Promise((resolve) => setTimeout(resolve, retryDelayMs * 2 ** attempt));
|
|
}
|
|
}
|
|
}
|
|
options.onError?.(lastError, envelope);
|
|
throw lastError;
|
|
},
|
|
|
|
subscribe<T>(pattern: string, handler: Handler<T>) {
|
|
if (!pattern.trim()) throw new TypeError("pubsub pattern cannot be empty");
|
|
return driver.subscribe(pattern, async (value, topic) => {
|
|
const envelope = value as MessageEnvelope<T>;
|
|
if (envelope && typeof envelope === "object" && "data" in envelope) {
|
|
await handler(envelope.data, topic);
|
|
} else {
|
|
await handler(value as T, topic);
|
|
}
|
|
});
|
|
},
|
|
};
|
|
}
|
|
|
|
export interface PresenceMember {
|
|
id: string;
|
|
metadata?: Record<string, unknown>;
|
|
joinedAt: number;
|
|
expiresAt: number;
|
|
}
|
|
|
|
export class PresenceChannel {
|
|
readonly #members = new Map<string, PresenceMember>();
|
|
|
|
constructor(
|
|
private readonly ttlMs = 30_000,
|
|
private readonly now: () => number = Date.now,
|
|
) {
|
|
if (!Number.isFinite(ttlMs) || ttlMs <= 0) {
|
|
throw new RangeError("presence ttlMs must be a positive number");
|
|
}
|
|
}
|
|
|
|
touch(id: string, metadata?: Record<string, unknown>): PresenceMember {
|
|
if (!id.trim()) throw new TypeError("presence member id cannot be empty");
|
|
const current = this.#members.get(id);
|
|
const timestamp = this.now();
|
|
const member = {
|
|
id,
|
|
metadata: metadata ?? current?.metadata,
|
|
joinedAt: current?.joinedAt ?? timestamp,
|
|
expiresAt: timestamp + this.ttlMs,
|
|
};
|
|
this.#members.set(id, member);
|
|
return structuredClone(member);
|
|
}
|
|
|
|
leave(id: string): boolean {
|
|
return this.#members.delete(id);
|
|
}
|
|
|
|
list(): PresenceMember[] {
|
|
this.prune();
|
|
return [...this.#members.values()].map((member) => structuredClone(member));
|
|
}
|
|
|
|
prune(): number {
|
|
const timestamp = this.now();
|
|
let removed = 0;
|
|
for (const [id, member] of this.#members) {
|
|
if (member.expiresAt <= timestamp) {
|
|
this.#members.delete(id);
|
|
removed++;
|
|
}
|
|
}
|
|
return removed;
|
|
}
|
|
}
|