-
Notifications
You must be signed in to change notification settings - Fork 0
fix: recover redis clients and moderation workers after outages #155
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
lorenzocorallo
wants to merge
1
commit into
main
Choose a base branch
from
fix/incident-20260917-redis-recovery
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,28 @@ | ||
| import { setTimeout as delay } from "node:timers/promises" | ||
| import type { Worker } from "bullmq" | ||
|
|
||
| type RecoverableWorker = Pick<Worker, "run" | "isRunning" | "isPaused" | "closing"> | ||
|
|
||
| /** Restart an unexpectedly failed worker loop; respect deliberate shutdown and pause. */ | ||
| export async function runWorkerWithRecovery( | ||
| worker: RecoverableWorker, | ||
| signal: AbortSignal, | ||
| onError: (error: unknown) => void | ||
| ): Promise<void> { | ||
| while (!signal.aborted && !worker.closing && !worker.isPaused() && !worker.isRunning()) { | ||
| try { | ||
| await worker.run() | ||
| return | ||
| } catch (error) { | ||
| onError(error) | ||
| } | ||
|
|
||
| if (signal.aborted || worker.closing || worker.isPaused()) return | ||
| try { | ||
| await delay(1000, undefined, { signal }) | ||
| } catch (error) { | ||
| if (signal.aborted) return | ||
| throw error | ||
| } | ||
| } | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,106 @@ | ||
| import { createServer, type Server, type Socket } from "node:net" | ||
| import type { RedisClientOptions } from "redis" | ||
| import { afterEach, beforeEach, describe, expect, it, vi } from "vitest" | ||
|
|
||
| const fixtures = vi.hoisted(() => ({ | ||
| env: { REDIS_HOST: "127.0.0.1", REDIS_PORT: 0 }, | ||
| logger: { debug: vi.fn(), info: vi.fn(), error: vi.fn() }, | ||
| })) | ||
|
|
||
| vi.mock("@/env", () => ({ env: fixtures.env })) | ||
| vi.mock("@/logger", () => ({ logger: fixtures.logger })) | ||
| vi.mock("redis", async (importOriginal) => { | ||
| const actual = await importOriginal<typeof import("redis")>() | ||
| return { | ||
| ...actual, | ||
| createClient: (options: RedisClientOptions) => { | ||
| const strategy = options.socket?.reconnectStrategy | ||
| return actual.createClient({ | ||
| ...options, | ||
| // The test server implements only RPUSH. Retry decisions and the client | ||
| // are real; only the delay is shortened to keep outage tests fast. | ||
| disableClientInfo: true, | ||
| socket: { | ||
| ...options.socket, | ||
| reconnectStrategy: (retries, cause) => { | ||
| const delay = typeof strategy === "function" ? strategy(retries, cause) : strategy | ||
| return typeof delay === "number" ? 5 : (delay ?? 5) | ||
| }, | ||
| }, | ||
| }) | ||
| }, | ||
| } | ||
| }) | ||
|
|
||
| let server: Server | ||
| let sockets: Set<Socket> | ||
| let client: typeof import("@/redis")["redis"] | undefined | ||
|
|
||
| async function listen(port = 0) { | ||
| await new Promise<void>((resolve) => server.listen(port, "127.0.0.1", resolve)) | ||
| const address = server.address() | ||
| if (!address || typeof address === "string") throw new Error("Expected TCP address") | ||
| fixtures.env.REDIS_PORT = address.port | ||
| } | ||
|
|
||
| async function stopServer() { | ||
| for (const socket of sockets) socket.destroy() | ||
| if (server.listening) await new Promise<void>((resolve) => server.close(() => resolve())) | ||
| } | ||
|
|
||
| async function waitForRetry(attempt: number) { | ||
| await vi.waitFor(() => { | ||
| expect(fixtures.logger.debug).toHaveBeenCalledWith(`[REDIS] reconnect retry #${attempt}`) | ||
| }) | ||
| } | ||
|
|
||
| beforeEach(() => { | ||
| vi.resetModules() | ||
| vi.clearAllMocks() | ||
| sockets = new Set() | ||
| server = createServer((socket) => { | ||
| sockets.add(socket) | ||
| socket.on("close", () => sockets.delete(socket)) | ||
| let request = "" | ||
| socket.on("data", (chunk) => { | ||
| request += chunk.toString() | ||
| if (request.endsWith("$5\r\nentry\r\n")) { | ||
| socket.write(":1\r\n") | ||
| request = "" | ||
| } | ||
| }) | ||
| }) | ||
| }) | ||
|
|
||
| afterEach(async () => { | ||
| if (client?.isOpen) await client.disconnect() | ||
| client = undefined | ||
| await stopServer() | ||
| }) | ||
|
|
||
| describe("Redis outage recovery", () => { | ||
| it("runs RPUSH after Redis starts later than the initial retry budget", async () => { | ||
| await listen() | ||
| const port = fixtures.env.REDIS_PORT | ||
| await stopServer() | ||
| client = (await import("@/redis")).redis | ||
| await waitForRetry(3) | ||
|
|
||
| await listen(port) | ||
| await expect(client.rPush("moderation:test", "entry")).resolves.toBe(1) | ||
| expect(client.isReady).toBe(true) | ||
| }) | ||
|
|
||
| it("runs RPUSH after a connected server outlasts the reconnect budget", async () => { | ||
| await listen() | ||
| const port = fixtures.env.REDIS_PORT | ||
| client = (await import("@/redis")).redis | ||
| await vi.waitFor(() => expect(client?.isReady).toBe(true)) | ||
| await stopServer() | ||
| await waitForRetry(5) | ||
|
|
||
| await listen(port) | ||
| await expect(client.rPush("moderation:test", "entry")).resolves.toBe(1) | ||
| expect(client.isReady).toBe(true) | ||
| }) | ||
| }) |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,66 @@ | ||
| import { afterEach, beforeEach, describe, expect, it, vi } from "vitest" | ||
| import { runWorkerWithRecovery } from "@/utils/worker-recovery" | ||
|
|
||
| beforeEach(() => vi.useFakeTimers()) | ||
| afterEach(() => vi.useRealTimers()) | ||
|
|
||
| function fixture() { | ||
| return { | ||
| run: vi.fn<() => Promise<void>>(), | ||
| isRunning: vi.fn(() => false), | ||
| isPaused: vi.fn(() => false), | ||
| closing: undefined as Promise<void> | undefined, | ||
| } | ||
| } | ||
|
|
||
| describe("worker recovery", () => { | ||
| it("restarts a loop that stopped on a connection error and then processes work", async () => { | ||
| const worker = fixture() | ||
| const onError = vi.fn() | ||
| const processed = vi.fn() | ||
| const failure = new Error("Connection is closed") | ||
| worker.run.mockRejectedValueOnce(failure).mockImplementationOnce(async () => processed()) | ||
|
|
||
| const task = runWorkerWithRecovery(worker, new AbortController().signal, onError) | ||
| await vi.advanceTimersByTimeAsync(999) | ||
| expect(processed).not.toHaveBeenCalled() | ||
| await vi.advanceTimersByTimeAsync(1) | ||
| await task | ||
|
|
||
| expect(onError).toHaveBeenCalledWith(failure) | ||
| expect(processed).toHaveBeenCalledOnce() | ||
| expect(worker.run).toHaveBeenCalledTimes(2) | ||
| }) | ||
|
|
||
| it("cancels a pending restart when the module stops", async () => { | ||
| const worker = fixture() | ||
| const controller = new AbortController() | ||
| worker.run.mockRejectedValueOnce(new Error("Connection is closed")) | ||
| const task = runWorkerWithRecovery(worker, controller.signal, vi.fn()) | ||
| await vi.advanceTimersByTimeAsync(100) | ||
| controller.abort() | ||
| await task | ||
| await vi.advanceTimersByTimeAsync(2000) | ||
| expect(worker.run).toHaveBeenCalledOnce() | ||
| }) | ||
|
|
||
| it.each(["paused", "closing", "running"])("does not restart a %s worker", async (state) => { | ||
| const worker = fixture() | ||
| if (state === "paused") worker.isPaused.mockReturnValue(true) | ||
| if (state === "running") worker.isRunning.mockReturnValue(true) | ||
| if (state === "closing") worker.closing = Promise.resolve() | ||
| await runWorkerWithRecovery(worker, new AbortController().signal, vi.fn()) | ||
| expect(worker.run).not.toHaveBeenCalled() | ||
| }) | ||
|
|
||
| it("does not restart after a deliberate pause following an error", async () => { | ||
| const worker = fixture() | ||
| worker.run.mockImplementationOnce(async () => { | ||
| worker.isPaused.mockReturnValue(true) | ||
| throw new Error("Connection is closed") | ||
| }) | ||
| await runWorkerWithRecovery(worker, new AbortController().signal, vi.fn()) | ||
| await vi.advanceTimersByTimeAsync(2000) | ||
| expect(worker.run).toHaveBeenCalledOnce() | ||
| }) | ||
| }) |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🩺 Stability & Availability | 🟠 Major | ⚡ Quick win
🔎 Supported by static analysis
🏁 Script executed:
Repository: PoliNetworkOrg/telegram
Length of output: 5445
🏁 Script executed:
Repository: PoliNetworkOrg/telegram
Length of output: 1722
🌐 Web query:
node-redis redis 4.7.0 client quit behavior while reconnecting reconnectStrategy source💡 Result:
<search_synthesis>
In redis node-redis version 4.7.0, the behavior of client.quit while the client is in the process of reconnecting is designed to ensure a graceful closure, though it has historically been a subject of user confusion regarding timing and internal promise resolution [1][2]. When you call client.quit while the client is reconnecting: 1. Graceful Shutdown: The quit command is intended to wait for the client to complete its reconnection attempt [2]. Once the connection is successfully established, the client will process any pending commands in its offline queue, execute the QUIT command, and then finally close the socket connection [3][2]. 2. Offline Queue: If the offline queue is enabled (the default), commands sent during the reconnection period are held in memory [4][5]. Calling.quit informs the client that it should eventually terminate; it will work through the queue before closing [3][2]. If you have disabled the offline queue (using disableOfflineQueue: true), these commands will be rejected immediately instead of being queued [4]. 3. Historical Context and Known Issues: In versions prior to 4.x, there were identified issues where.quit did not correctly end the connection if called while the client was down, potentially leading to stalled states [1]. While 4.7.0 includes fixes to ensure the QUIT command properly terminates the connection, users have reported scenarios where.quit might appear to hang if the reconnection process itself is stuck or if the Promise returned by.quit does not resolve as expected [6][7]. 4. Reconnect Strategy: The reconnectStrategy is a configuration function (returning false, a number, or an Error) that dictates how the client behaves when the socket closes unexpectedly [8][9][10]. It does not directly affect the.quit method&
#39;s operation, other than determining whether the client enters the reconnecting state that.quit must then contend with [4][8]. Note on Future Versions: Be aware that in later versions of node-redis (specifically moving toward v5.x), the API has shifted. The.quit method has been replaced by.close for graceful shutdowns, and.disconnect has been replaced by.destroy for immediate, forceful shutdowns [11][7]. If you encounter persistent issues with.quit hanging in 4.7.0, some developers have used a Promise race (e.g., Promise.race([client.quit, someTimeout])) as a workaround to ensure the application continues execution [6].</search_synthesis>
<source_evidence>
Citations:
🏁 Script executed:
Repository: PoliNetworkOrg/telegram
Length of output: 1771
🌐 Web query:
node-redis v4.7.0 official source quit reconnectStrategy reconnecting💡 Result:
<search_synthesis>
In node-redis v4.7.0, the
reconnectStrategyis a configuration option that determines the client's behavior when the socket closes unexpectedly [1][2][3]. Key concepts for this version are: Reconnect Strategy When the socket closes unexpectedly—meaning it was not intentionally closed by calling a termination method—the client uses thereconnectStrategyto decide whether to attempt a reconnection [1][4][2]. * Supported values includefalse(do not reconnect), anumber(delay in milliseconds before reconnecting), or a function(retries, cause) => false | number | Error[1][2]. * By default, the client implements an exponential backoff strategy with a maximum of 2000 ms, plus a random jitter of 0–200 ms [1][5]. * While the client is reconnecting, it emits areconnectingevent [6][7]. Quit vs. Disconnect (and Deprecation) *.quit(): Historically used to gracefully close the connection by sending theQUITcommand to the server after flushing the command queue [6][8][9]. *.disconnect(): Historically used to forcibly close the connection immediately without waiting for the queue [6][8][9]. * Important Update: As of Redis 7.2, theQUITcommand is deprecated [10][9][7]. Consequently, in newer versions of node-redis (including v4.7.0 contexts), it is recommended to use.close()for graceful shutdown and.destroy()for immediate, forceful shutdown [10][8][7]. The older.quit()and.disconnect()methods exist as aliases but are deprecated [11]. When you call these closing methods (.quit()/.close()or.disconnect()/.destroy()), the client interprets the closure as intentional and does not trigger thereconnectStrategy[1][4][2].</search_synthesis>
<source_evidence>
Citations:
🌐 Web query:
node-redis 4.7.0 quit behavior during reconnect loop official source💡 Result:
<search_synthesis>
In node-redis, calling.quit while the client is in a reconnect loop is designed to stop the reconnection process and close the client [1][2][3]. When.quit is invoked, the client sets its internal state to closed, which prevents further automatic reconnection attempts [1][4]. Key points regarding this behavior include: 1. Expected Behavior: Historically, node-redis has been updated to ensure that.quit properly terminates the connection and halts reconnection logic, even if the connection was down at the time the method was called [2][3]. This fixes previous issues where the client would continue attempting to reconnect after.quit had been requested [3]. 2. Mechanism: When.quit is called, the client marks itself as closed [4]. If the client is currently reconnecting, it effectively cancels the pending reconnection cycle, ensuring the client reaches an inactive state [1][2]. 3. Recommended Modern Alternatives: As of recent versions, the QUIT command is considered deprecated in favor of methods that directly manage the socket connection [5]. You should prefer: -.close: Gracefully closes the connection, allowing pending commands to finish before closing [5]. -.destroy: Immediately closes the connection without waiting for pending commands [5]. For managing reconnection logic explicitly, you can configure the socket.reconnectStrategy option when creating the client [6][7]. Returning false from this function prevents automatic reconnections [6][8][7].
</search_synthesis>
<source_evidence>
Citations:
🏁 Script executed:
Repository: PoliNetworkOrg/telegram
Length of output: 2447
🌐 Web query:
site:github.com/redis/node-redis/blob 4.7.1 packages/client/lib/client/index.ts async quit socket.quit💡 Result:
<search_synthesis>
In the Node Redis client (node-redis), the
quitmethod is defined inpackages/client/lib/client/index.tsand leverages the underlyingRedisSocket'squitmethod [1][2][3][4]. The client-sideQUIT()method performs several cleanup operations, such as disposing of subscription credentials and clearing timers, before calling the socket'squitmethod [2][3][4]. The socket'squitmethod, defined inpackages/client/lib/client/socket.ts, is an asynchronous function that accepts a callback function returning a Promise [5]. It ensures the socket is not already closed, marks the socket as closed, awaits the resolution of the provided function (which executes the Redis QUIT command and performs associated client-side cleanup), and finally destroys the socket [5]. Historically,QUIT()has been marked as deprecated in favor ofclose()[2][3]. Theclose()method is designed to wait for pending commands to finish before closing the connection, whereasdestroy()terminates the connection immediately and rejects pending commands [2][3].</search_synthesis>
<source_evidence>
Citations:
Use
redis.disconnect()during shutdown. The lockfile resolvesredis@4.7.1with@redis/client@1.6.1. When the client is reconnecting,quit()queuesQUITand waits for a ready connection to process it. This reconnect strategy always retries, soredis.quit()can remain pending while Redis is unavailable.Promise.allSettled()then cannot reachprocess.exit(0), and deployment shutdown may wait until the process is force-killed.🤖 Prompt for AI Agents