Files
WRNexusJS/packages/pubsub/test/redis-driver.test.ts
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

74 lines
2.8 KiB
TypeScript

import { afterEach, expect, test } from "bun:test";
import { redisDriver } from "../src/redis.ts";
const bun = Bun as typeof Bun & { connect: (options: unknown) => Promise<unknown> };
const originalConnect = bun.connect;
afterEach(() => {
bun.connect = originalConnect;
});
test("rediss URLs enable TLS and decode credentials", () => {
const options: Array<Record<string, unknown>> = [];
bun.connect = ((value: Record<string, unknown>) => {
options.push(value);
return new Promise(() => {});
}) as typeof bun.connect;
const driver = redisDriver("rediss://user:p%40ss@redis.example.com:6380/2");
expect(options).toHaveLength(2);
expect(options[0]).toMatchObject({ hostname: "redis.example.com", port: 6380, tls: true });
driver.close();
});
test("rejects invalid Redis URLs before connecting", () => {
expect(() => redisDriver("http://localhost:6379")).toThrow("redis:// or rediss://");
expect(() => redisDriver("redis://localhost/not-a-db")).toThrow("database");
});
test("bounds writes queued while Redis is unavailable", () => {
bun.connect = (() => new Promise(() => {})) as typeof bun.connect;
const driver = redisDriver("redis://localhost:6379");
for (let index = 0; index < 1000; index++) driver.publish("topic", index);
expect(() => driver.publish("topic", "overflow")).toThrow("queue is full");
driver.close();
});
test("validates reconnect and backpressure options", () => {
bun.connect = (() => new Promise(() => {})) as typeof bun.connect;
expect(() => redisDriver(undefined, { maxPending: 0 })).toThrow("maxPending");
expect(() => redisDriver(undefined, { reconnectDelayMs: 20, reconnectMaxDelayMs: 10 })).toThrow(
"reconnectMaxDelayMs",
);
});
test("reconnects and replays subscriptions after a socket closes", async () => {
const connections: Array<Record<string, any>> = [];
bun.connect = ((options: Record<string, any>) => {
connections.push(options);
return Promise.resolve({});
}) as typeof bun.connect;
const driver = redisDriver("redis://localhost:6379", {
reconnectDelayMs: 1,
reconnectMaxDelayMs: 1,
});
expect(connections).toHaveLength(2);
const firstWrites: string[] = [];
connections[0].socket.open({
write: (bytes: Uint8Array) => firstWrites.push(new TextDecoder().decode(bytes)),
end() {},
});
driver.subscribe("order:*", () => {});
expect(firstWrites.some((write) => write.includes("PSUBSCRIBE"))).toBe(true);
connections[0].socket.close();
await new Promise((resolve) => setTimeout(resolve, 10));
expect(connections.length).toBeGreaterThanOrEqual(3);
const replayed: string[] = [];
connections[2].socket.open({
write: (bytes: Uint8Array) => replayed.push(new TextDecoder().decode(bytes)),
end() {},
});
expect(replayed.some((write) => write.includes("PSUBSCRIBE"))).toBe(true);
driver.close();
});