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
6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -221,8 +221,10 @@ Header optimization is disabled by default. When enabled, headers are
whitelisted, not removed by a strip list. The default
whitelist is `Host`, `X-Amz-Target`, `Content-Length`, `Accept-Encoding`, and
`Content-Encoding`; when credentials are configured, `Authorization` and
`X-Amz-Date` are also kept. The Alternator `User-Agent` is applied after this
filter, so it is kept unless `userAgent: false` is configured.
`X-Amz-Date` are also kept. Alternator does not use AWS session tokens, so
`sessionToken` is not sent even when provided in credentials. The Alternator
`User-Agent` is applied after this filter, so it is kept unless
`userAgent: false` is configured.

By default, the client replaces the AWS SDK `User-Agent` with the ScyllaDB
Alternator client identity:
Expand Down
22 changes: 20 additions & 2 deletions src/client-base.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import {
type ServiceOutputTypes,
} from "@aws-sdk/client-dynamodb";
import type { HttpHandler, HttpHandlerUserInput } from "@smithy/protocol-http";
import type { HttpHandlerOptions } from "@smithy/types";
import type { AwsCredentialIdentity, HttpHandlerOptions } from "@smithy/types";
import { DEFAULT_REGION, firstEndpointUrl, NO_AUTH_CREDENTIALS, normalizeConfig } from "./config.js";
import { AlternatorDiscovery } from "./discovery.js";
import { KeyRouteAffinityPlanner } from "./affinity.js";
Expand Down Expand Up @@ -137,7 +137,7 @@ function buildDynamoConfig(
...awsConfig,
endpoint: firstEndpointUrl(alternatorConfig),
region: region ?? DEFAULT_REGION,
credentials: credentials ?? NO_AUTH_CREDENTIALS,
credentials: dropSessionToken(credentials) ?? NO_AUTH_CREDENTIALS,
requestHandler,
};
if (alternatorConfig.headerOptimization.enabled) {
Expand All @@ -146,4 +146,22 @@ function buildDynamoConfig(
return dynamoConfig;
}

function dropSessionToken(
credentials: DynamoDBClientConfig["credentials"] | undefined,
): DynamoDBClientConfig["credentials"] | undefined {
if (credentials === undefined) {
return undefined;
}
if (typeof credentials === "function") {
return async (identityProperties?: Record<string, unknown>) =>
removeSessionToken(await credentials(identityProperties));
}
return removeSessionToken(credentials);
}

function removeSessionToken(credentials: AwsCredentialIdentity): AwsCredentialIdentity {
const { sessionToken: _sessionToken, ...withoutSessionToken } = credentials;
return withoutSessionToken;
}

export type AlternatorRequestHandler = HttpHandler<HttpHandlerOptions>;
46 changes: 41 additions & 5 deletions src/discovery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -250,7 +250,7 @@ export class AlternatorDiscovery {
});

if (response.response.statusCode < 200 || response.response.statusCode >= 300) {
await drainResponseBody(response.response.body);
await drainResponseBody(response.response.body, this.config.discovery.timeoutMs);
throw new Error(`/localnodes returned HTTP ${response.response.statusCode}`);
}

Expand All @@ -272,12 +272,48 @@ export class AlternatorDiscovery {
}
}

async function drainResponseBody(body: unknown): Promise<void> {
try {
await bodyToString(body);
} catch (_error) {
async function drainResponseBody(body: unknown, timeoutMs: number): Promise<void> {
const drain = bodyToString(body).then(
() => undefined,
() => undefined,
);
let timeout: ReturnType<typeof setTimeout> | undefined;
const timeoutPromise = new Promise<void>((resolve) => {
timeout = setTimeout(() => {
destroyResponseBody(body);
resolve();
}, Math.max(1, timeoutMs));
timeout.unref?.();
});

await Promise.race([drain, timeoutPromise]);
if (timeout) {
clearTimeout(timeout);
}
}

function destroyResponseBody(body: unknown): void {
if (typeof body !== "object" || body === null) {
return;
}

if ("destroy" in body && typeof (body as { destroy?: unknown }).destroy === "function") {
try {
(body as { destroy(): void }).destroy();
} catch (_error) {
return;
}
return;
}

if ("cancel" in body && typeof (body as { cancel?: unknown }).cancel === "function") {
try {
const cancellation = (body as { cancel(): unknown }).cancel();
void Promise.resolve(cancellation).catch(() => undefined);
} catch (_error) {
return;
}
}
}

function queryToRequestQuery(query: LocalNodesQuery): Record<string, string> {
Expand Down
44 changes: 44 additions & 0 deletions test/discovery.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -325,6 +325,50 @@ describe("Alternator discovery", () => {
}
});

it("bounds draining non-terminating non-2xx discovery bodies", async () => {
let requests = 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.write(JSON.stringify({ error: "temporary failure" }));
return;
}
response.end(JSON.stringify(["node-a.internal"]));
});
const address = await listen(server);
const client = new AlternatorDynamoDBClient({
seeds: [address.address, address.address],
port: address.port,
discovery: {
background: false,
timeoutMs: 20,
},
connection: {
keepAlive: true,
maxSockets: 1,
},
});

try {
await expect(client.alternator.refreshNodes()).resolves.toEqual([
{
host: "node-a.internal",
scheme: "http",
port: address.port,
url: `http://node-a.internal:${address.port}`,
},
]);
expect(requests).toBe(2);
} finally {
client.destroy();
server.closeAllConnections?.();
await close(server);
}
});

it("keeps the DynamoDB socket reusable after repeated non-2xx responses", async () => {
let requests = 0;
let connections = 0;
Expand Down
29 changes: 28 additions & 1 deletion test/middleware.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -156,7 +156,7 @@ describe("Alternator middleware", () => {
]);
});

it("uses the default optimized header whitelist with session credentials", async () => {
it("drops session tokens before signing Alternator requests", async () => {
const handler = new RecordingHandler(() => ({ TableNames: [] }));
const client = new AlternatorDynamoDBClient({
seeds: ["seed"],
Expand All @@ -177,6 +177,33 @@ describe("Alternator middleware", () => {
expect(headers["x-amz-date"]).toBeDefined();
expect(headers["x-amz-target"]).toBe("DynamoDB_20120810.ListTables");
expect(headers["x-amz-security-token"]).toBeUndefined();
expect(signedHeaderNames(headers.authorization)).toEqual([
"content-length",
"host",
"x-amz-date",
"x-amz-target",
]);
});

it("drops session tokens from credential providers before signing Alternator requests", async () => {
const handler = new RecordingHandler(() => ({ TableNames: [] }));
const client = new AlternatorDynamoDBClient({
seeds: ["seed"],
requestHandler: handler,
discovery: { background: false },
credentials: () => Promise.resolve({
accessKeyId: "key",
secretAccessKey: "secret",
sessionToken: "session-token",
}),
});

await client.send(new ListTablesCommand({}));

const headers = commandRequests(handler)[0]?.headers ?? {};
expect(headers.authorization).toContain("AWS4-HMAC-SHA256");
expect(headers["x-amz-security-token"]).toBeUndefined();
expect(signedHeaderNames(headers.authorization)).not.toContain("x-amz-security-token");
});

it("compresses JSON request bodies when enabled", async () => {
Expand Down