import { afterEach, expect, test } from "bun:test"; import { redisDriver } from "../src/redis.ts"; const bun = Bun as typeof Bun & { connect: (options: unknown) => Promise }; const originalConnect = bun.connect; afterEach(() => { bun.connect = originalConnect; }); test("rediss URLs enable TLS and decode credentials", () => { const options: Array> = []; bun.connect = ((value: Record) => { 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> = []; bun.connect = ((options: Record) => { 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(); });