Files
WRNexusJS/packages/realtime/src/streams.ts
T
Clintchiz 586a6db8ff
Quality / quality (ubuntu-latest) (push) Failing after 21s
Quality / quality (windows-latest) (push) Canceled after 0s
release: WRNexusJS 0.8.0
2026-08-02 23:18:51 +05:30

99 lines
3.4 KiB
TypeScript

export interface DatabaseChange<T = unknown> {
table: string;
operation: "insert" | "update" | "delete";
key?: string | number;
record?: T;
occurredAt: number;
}
export interface DatabaseChangeSource {
subscribe(handler: (change: DatabaseChange) => void | Promise<void>): () => void;
}
export function databaseChangeFeed(
source: DatabaseChangeSource,
publish: (topic: string, change: DatabaseChange) => void | Promise<void>,
options: { prefix?: string; allowTables?: string[] } = {},
): () => void {
const prefix = options.prefix ?? "db";
const allowed = options.allowTables ? new Set(options.allowTables) : null;
return source.subscribe(async (change) => {
if (!/^[A-Za-z_][A-Za-z0-9_]*$/.test(change.table)) return;
if (allowed && !allowed.has(change.table)) return;
await publish(`${prefix}:${change.table}`, structuredClone(change));
});
}
export interface FileStreamFrame {
streamId: string;
index: number;
total: number;
bytes: Uint8Array;
}
export function frameFileStream(
streamId: string,
bytes: Uint8Array,
options: { chunkBytes?: number; maxBytes?: number } = {},
): FileStreamFrame[] {
if (!/^[A-Za-z0-9_-]{1,128}$/.test(streamId)) throw new Error("Invalid stream id");
const chunkBytes = options.chunkBytes ?? 64 * 1024;
const maxBytes = options.maxBytes ?? 25 * 1024 * 1024;
if (!Number.isInteger(chunkBytes) || chunkBytes < 1024 || chunkBytes > 1024 * 1024)
throw new RangeError("chunkBytes must be between 1KiB and 1MiB");
if (bytes.byteLength > maxBytes) throw new Error("WRN-REALTIME-FILE-LIMIT");
const total = Math.max(1, Math.ceil(bytes.byteLength / chunkBytes));
return Array.from({ length: total }, (_, index) => ({
streamId,
index,
total,
bytes: bytes.slice(index * chunkBytes, (index + 1) * chunkBytes),
}));
}
export function createFileStreamReceiver(options: { maxBytes?: number; maxStreams?: number } = {}) {
const maxBytes = options.maxBytes ?? 25 * 1024 * 1024;
const maxStreams = options.maxStreams ?? 32;
const streams = new Map<string, Map<number, Uint8Array>>();
return {
accept(frame: FileStreamFrame): Uint8Array | null {
if (
!Number.isInteger(frame.total) ||
frame.total < 1 ||
frame.total > 25_600 ||
frame.index < 0 ||
frame.index >= frame.total
)
throw new Error("Invalid file stream frame");
let parts = streams.get(frame.streamId);
if (!parts) {
if (streams.size >= maxStreams) throw new Error("WRN-REALTIME-STREAM-CAPACITY");
streams.set(frame.streamId, (parts = new Map()));
}
parts.set(frame.index, frame.bytes.slice());
const size = [...parts.values()].reduce((sum, part) => sum + part.byteLength, 0);
if (size > maxBytes) {
streams.delete(frame.streamId);
throw new Error("WRN-REALTIME-FILE-LIMIT");
}
if (parts.size !== frame.total) return null;
const output = new Uint8Array(size);
let offset = 0;
for (let index = 0; index < frame.total; index++) {
const part = parts.get(index);
if (!part) return null;
output.set(part, offset);
offset += part.byteLength;
}
streams.delete(frame.streamId);
return output;
},
snapshot: () => ({
activeStreams: streams.size,
bufferedBytes: [...streams.values()]
.flatMap((parts) => [...parts.values()])
.reduce((sum, part) => sum + part.byteLength, 0),
}),
};
}