Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions src/extensions/distributed/redis-runtime-provider.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import { register, unregister } from "../contracts.ts";
import { ensureRedisRuntimeProvider } from "./defaults.ts";
import {
captureRedisRuntimeProvider,
type NodeRedisClient,
type RedisClient,
type RedisRuntimeProvider,
RedisRuntimeProviderName,
Expand Down Expand Up @@ -48,6 +49,18 @@ function createProvider(): RedisRuntimeProvider {
}

describe("RedisRuntimeProvider", () => {
it("keeps error-only Redis event listeners source compatible", () => {
const on: NodeRedisClient["on"] = (
event: "error",
listener: (error: unknown) => void,
): unknown => {
void listener;
return event;
};

assertEquals(on("error", () => {}), "error");
});

it("captures provider methods without losing their receiver", async () => {
let receivedThis: unknown;
const provider = {
Expand Down
87 changes: 86 additions & 1 deletion src/proxy/routing-invalidation-redis.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,10 @@ async function settleWithin<T>(promise: Promise<T>, label: string): Promise<T> {

function createFakeRedisServer() {
const subscriptions = new Map<RoutingInvalidationRedisClient, Map<string, RedisListener>>();
const clientEventListeners = new Map<
RoutingInvalidationRedisClient,
Map<string, Set<(value?: unknown) => void>>
>();
const clients: RoutingInvalidationRedisClient[] = [];
let onPublish:
| ((channel: string, message: string) => void | Promise<void>)
Expand Down Expand Up @@ -71,9 +75,20 @@ function createFakeRedisServer() {
},
close: () => {
subscriptions.delete(client);
clientEventListeners.delete(client);
return Promise.resolve();
},
destroy: () => subscriptions.delete(client),
destroy: () => {
subscriptions.delete(client);
clientEventListeners.delete(client);
},
on: ((event: string, listener: (value?: unknown) => void) => {
const listeners = clientEventListeners.get(client) ?? new Map();
const eventListeners = listeners.get(event) ?? new Set();
eventListeners.add(listener);
listeners.set(event, eventListeners);
clientEventListeners.set(client, listeners);
}) as NonNullable<RoutingInvalidationRedisClient["on"]>,
};
clients.push(client);
return client;
Expand All @@ -82,6 +97,13 @@ function createFakeRedisServer() {
return {
clients,
createClient,
emitClientEvent(clientIndex: number, event: string, value?: unknown) {
const client = clients[clientIndex];
if (!client) throw new Error(`Missing Redis client ${clientIndex}`);
for (const listener of clientEventListeners.get(client)?.get(event) ?? []) {
listener(value);
}
},
publishRaw,
setOnPublish(listener: typeof onPublish) {
onPublish = listener;
Expand Down Expand Up @@ -143,6 +165,69 @@ function deferred(): { promise: Promise<void>; resolve: () => void } {
}

describe("proxy routing invalidation Redis bus", () => {
it("warns for a managed socket recycle and reports successful resubscription", async () => {
const redis = createFakeRedisServer();
const info: Array<{ message: string; extra?: Record<string, unknown> }> = [];
const warnings: Array<{ message: string; extra?: Record<string, unknown> }> = [];
const errors: Array<{ message: string; error?: Error }> = [];
const bus = await startProxyRoutingInvalidationBus({
redisUrl: "redis://example.test:6379",
replicaId: "replica-a",
createClient: redis.createClient,
integritySecret: createIntegritySecret(),
logger: {
info: (message, extra) => info.push({ message, extra }),
warn: (message, extra) => warnings.push({ message, extra }),
error: (message, error) => errors.push({ message, error }),
},
onInvalidate: () => {},
});

class SocketClosedUnexpectedlyError extends Error {
constructor() {
super("Socket closed unexpectedly");
}
}

redis.emitClientEvent(1, "ready");
redis.emitClientEvent(1, "error", new SocketClosedUnexpectedlyError());
redis.emitClientEvent(1, "ready");
redis.emitClientEvent(0, "ready");
redis.emitClientEvent(0, "error", new SocketClosedUnexpectedlyError());
redis.emitClientEvent(0, "ready");
redis.emitClientEvent(0, "error", new Error("Redis authentication failed"));

assertEquals(warnings, [
{
message: "Proxy routing invalidation Redis socket closed; reconnecting",
extra: { clientRole: "subscriber" },
},
{
message: "Proxy routing invalidation Redis socket closed; reconnecting",
extra: { clientRole: "publisher" },
},
]);
assertEquals(
info.filter((entry) => entry.message.includes("reconnected")),
[
{
message: "Proxy routing invalidation Redis subscriber reconnected and resubscribed",
extra: { clientRole: "subscriber" },
},
{
message: "Proxy routing invalidation Redis publisher reconnected",
extra: { clientRole: "publisher" },
},
],
);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
assertEquals(errors.map(({ message, error }) => ({ message, error: error?.message })), [{
message: "Proxy routing invalidation Redis error",
error: "Redis authentication failed",
}]);

await bus?.close();
});

it("fans out to every replica and waits for a distinct acknowledgement from each", async () => {
const redis = createFakeRedisServer();
const integritySecret = createIntegritySecret();
Expand Down
61 changes: 54 additions & 7 deletions src/proxy/routing-invalidation-redis.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ const DEFAULT_MAX_ENVELOPE_FUTURE_MS = 5_000;
const INTEGRITY_SECRET_ENV_VAR = "VERYFRONT_PROXY_ROUTING_INVALIDATION_SECRET";
const EVENT_SIGNATURE_DOMAIN = "vf-proxy-routing-invalidation:event:v1";
const ACK_SIGNATURE_DOMAIN = "vf-proxy-routing-invalidation:ack:v1";
const reflectApply = Reflect.apply;

type RedisListener = (message: string, channel: string) => void;
type SignatureDomain = "event" | "ack";
Expand Down Expand Up @@ -79,6 +80,18 @@ function encodedByteLength(value: string): number {
return new TextEncoder().encode(value).byteLength;
}

function observeRedisReady(
client: RoutingInvalidationRedisClient,
listener: () => void,
): void {
if (!client.on) return;
try {
reflectApply(client.on, client, ["ready", listener]);
} catch {
// Injected clients only guarantee support for error listeners.
}
}

function parseSignedEnvelope(message: string): SignedRoutingInvalidationEnvelope | null {
if (encodedByteLength(message) > MAX_SIGNED_ENVELOPE_BYTES) return null;
let parsed: unknown;
Expand Down Expand Up @@ -287,6 +300,12 @@ async function closeClient(
}
}

function isExpectedSocketRecycle(error: unknown): error is Error {
return error instanceof Error &&
error.constructor.name === "SocketClosedUnexpectedlyError" &&
error.message === "Socket closed unexpectedly";
}

export async function startProxyRoutingInvalidationBus(
options: StartProxyRoutingInvalidationBusOptions,
): Promise<ProxyRoutingInvalidationBus | null> {
Expand Down Expand Up @@ -391,14 +410,42 @@ export async function startProxyRoutingInvalidationBus(
return processing;
};

const logRedisError = (error: unknown) => {
options.logger?.error(
"Proxy routing invalidation Redis error",
error instanceof Error ? error : new Error(String(error)),
);
const observeRedisClient = (
client: RoutingInvalidationRedisClient,
clientRole: "publisher" | "subscriber",
): void => {
let hasBeenReady = false;
let recoveryPending = false;

observeRedisReady(client, () => {
if (hasBeenReady && recoveryPending) {
options.logger?.info(
clientRole === "subscriber"
? "Proxy routing invalidation Redis subscriber reconnected and resubscribed"
: "Proxy routing invalidation Redis publisher reconnected",
{ clientRole },
);
}
hasBeenReady = true;
recoveryPending = false;
});
client.on?.("error", (error) => {
recoveryPending ||= hasBeenReady;
if (hasBeenReady && isExpectedSocketRecycle(error)) {
options.logger?.warn(
"Proxy routing invalidation Redis socket closed; reconnecting",
{ clientRole },
);
return;
}
options.logger?.error(
"Proxy routing invalidation Redis error",
error instanceof Error ? error : new Error(String(error)),
);
});
};
publishClient.on?.("error", logRedisError);
subscribeClient.on?.("error", logRedisError);
observeRedisClient(publishClient, "publisher");
observeRedisClient(subscribeClient, "subscriber");

try {
await Promise.all([publishClient.connect(), subscribeClient.connect()]);
Expand Down