From 6f466f21564d1fe5f0740208a31f1dce7253b66e Mon Sep 17 00:00:00 2001 From: Dmitry Kropachev Date: Mon, 29 Jun 2026 22:53:00 -0400 Subject: [PATCH] Drain discovery responses before non-2xx errors --- src/discovery.ts | 9 +++ test/discovery.test.ts | 121 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 130 insertions(+) diff --git a/src/discovery.ts b/src/discovery.ts index 75a6ae9..6ec112c 100644 --- a/src/discovery.ts +++ b/src/discovery.ts @@ -250,6 +250,7 @@ export class AlternatorDiscovery { }); if (response.response.statusCode < 200 || response.response.statusCode >= 300) { + await drainResponseBody(response.response.body); throw new Error(`/localnodes returned HTTP ${response.response.statusCode}`); } @@ -271,6 +272,14 @@ export class AlternatorDiscovery { } } +async function drainResponseBody(body: unknown): Promise { + try { + await bodyToString(body); + } catch (_error) { + return; + } +} + function queryToRequestQuery(query: LocalNodesQuery): Record { const result: Record = {}; if (query.dc) { diff --git a/test/discovery.test.ts b/test/discovery.test.ts index f0c3c5b..4affc5c 100644 --- a/test/discovery.test.ts +++ b/test/discovery.test.ts @@ -3,6 +3,8 @@ import { AlternatorDynamoDBClient, routing } from "../src/index.js"; import { AlternatorDynamoDBClient as EdgeAlternatorDynamoDBClient } from "../src/edge.js"; import { RecordingHandler } from "./helpers.js"; import { ListTablesCommand } from "@aws-sdk/client-dynamodb"; +import { createServer, type Server } from "node:http"; +import type { AddressInfo } from "node:net"; const missingDatacenterQuery = { dc: "__alternator_client_missing_dc__" }; const missingRackQuery = { rack: "__alternator_client_missing_rack__" }; @@ -274,4 +276,123 @@ describe("Alternator discovery", () => { expect(handler.requests[1]?.hostname).toBe("edge-node"); expect(handler.requests[1]?.headers.connection).toBeUndefined(); }); + + it("keeps the discovery socket reusable after non-2xx responses", async () => { + let requests = 0; + let connections = 0; + const server = createServer((request, response) => { + expect(request.url).toBe("/localnodes"); + requests += 1; + response.setHeader("content-type", "application/json"); + if (requests === 1) { + response.statusCode = 500; + response.end(JSON.stringify({ error: "temporary failure" })); + return; + } + response.end(JSON.stringify(["node-a.internal"])); + }); + server.on("connection", () => { + connections += 1; + }); + const address = await listen(server); + const client = new AlternatorDynamoDBClient({ + seeds: [address.address], + port: address.port, + discovery: { + background: false, + timeoutMs: 500, + }, + connection: { + keepAlive: true, + maxSockets: 1, + }, + }); + + try { + await client.alternator.refreshNodes(); + await expect(client.alternator.refreshNodes()).resolves.toEqual([ + { + host: "node-a.internal", + scheme: "http", + port: address.port, + url: `http://node-a.internal:${address.port}`, + }, + ]); + expect(connections).toBe(1); + } finally { + client.destroy(); + await close(server); + } + }); + + it("keeps the DynamoDB socket reusable after repeated non-2xx responses", async () => { + let requests = 0; + let connections = 0; + const server = createServer((request, response) => { + expect(request.method).toBe("POST"); + expect(request.url).toBe("/"); + request.resume(); + request.on("end", () => { + requests += 1; + response.setHeader("content-type", "application/x-amz-json-1.0"); + if (requests < 3) { + response.statusCode = 400; + response.end(JSON.stringify({ __type: "ValidationException", message: "bad" })); + return; + } + response.end(JSON.stringify({ TableNames: [] })); + }); + }); + server.on("connection", () => { + connections += 1; + }); + const address = await listen(server); + const client = new AlternatorDynamoDBClient({ + seeds: [address.address], + port: address.port, + discovery: { + background: false, + }, + connection: { + keepAlive: true, + maxSockets: 1, + }, + maxAttempts: 1, + }); + + try { + await expect(client.send(new ListTablesCommand({}))).rejects.toThrow(/bad/); + await expect(client.send(new ListTablesCommand({}))).rejects.toThrow(/bad/); + await expect(client.send(new ListTablesCommand({}))).resolves.toMatchObject({ + TableNames: [], + }); + expect(requests).toBe(3); + expect(connections).toBe(1); + } finally { + client.destroy(); + await close(server); + } + }); }); + +function listen(server: Server): Promise { + return new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(0, "127.0.0.1", () => { + server.off("error", reject); + resolve(server.address() as AddressInfo); + }); + }); +} + +function close(server: Server): Promise { + return new Promise((resolve, reject) => { + server.close((error) => { + if (error) { + reject(error); + return; + } + resolve(); + }); + }); +}