fix(cassandra): shorten the read timeout below the rpc deadline (#2198)

This commit is contained in:
Hampus
2026-08-31 00:06:29 +02:00
committed by GitHub
parent b5496097d2
commit f1400ae58e
4 changed files with 18 additions and 4 deletions
+8 -2
View File
@@ -7,7 +7,10 @@ import cassandra from 'cassandra-driver';
const distance = cassandra.types.distance;
const MAX_REQUESTS_PER_CONNECTION = 2048;
const READ_TIMEOUT_MS = 12000;
const CONNECT_TIMEOUT_MS = 5000;
const DEFAULT_READ_TIMEOUT_MS = 5000;
export const BACKGROUND_READ_TIMEOUT_MS = 12000;
interface CassandraConfig {
hosts: Array<string>;
@@ -16,6 +19,7 @@ interface CassandraConfig {
localDc: string;
username?: string | undefined;
password?: string | undefined;
readTimeoutMs?: number | undefined;
}
interface CassandraClientOptions {
@@ -66,6 +70,7 @@ class CassandraClient implements ICassandraClient {
localDc: config.localDc,
username: config.username,
password: config.password,
readTimeoutMs: config.readTimeoutMs,
};
this.logger = options.logger ?? NoopLogger;
this.client = null;
@@ -93,7 +98,8 @@ class CassandraClient implements ICassandraClient {
},
},
socketOptions: {
readTimeout: READ_TIMEOUT_MS,
connectTimeout: CONNECT_TIMEOUT_MS,
readTimeout: this.config.readTimeoutMs ?? DEFAULT_READ_TIMEOUT_MS,
},
encoding: {
map: Map,
@@ -164,6 +164,7 @@ export async function fetchPage<T = Record<string, unknown>, P extends Cassandra
options: {
pageSize: number;
pageState?: string | null;
readTimeout?: number | undefined;
},
): Promise<PagedQueryResult<T>> {
const {cql, params: boundRaw} = normalizeExecuteArgs(queryOrPrepared, params);
@@ -193,6 +194,7 @@ export async function fetchPage<T = Record<string, unknown>, P extends Cassandra
prepare: true,
fetchSize: options.pageSize,
pageState: options.pageState ?? undefined,
readTimeout: options.readTimeout,
});
return {
rows: (result.rows as Array<T>) ?? [],
@@ -1,6 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {UserFlags} from '@fluxer/constants/src/UserConstants';
import {BACKGROUND_READ_TIMEOUT_MS} from '@pkgs/cassandra/src/Client';
import {createUserID, type UserID} from '../../../../BrandedTypes';
import {fetchMany, fetchOne, fetchPage, upsertOne} from '../../../../database/CassandraQueryExecution';
import {Db, type DbOp, nextVersion} from '../../../../database/CassandraTypes';
@@ -94,7 +95,11 @@ export class UserDataRepository {
users: Array<User>;
pageState: string | null;
}> {
const result = await fetchPage<UserRow>(FETCH_ALL_USERS_SCAN_CQL, {}, {pageSize: limit, pageState});
const result = await fetchPage<UserRow>(
FETCH_ALL_USERS_SCAN_CQL,
{},
{pageSize: limit, pageState, readTimeout: BACKGROUND_READ_TIMEOUT_MS},
);
return {
users: result.rows.map((user) => new User(user)),
pageState: result.pageState,
+2 -1
View File
@@ -1,7 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
import {setupGracefulShutdown} from '@fluxer/hono/src/Server';
import {initCassandra, shutdownCassandra} from '@pkgs/cassandra/src/Client';
import {BACKGROUND_READ_TIMEOUT_MS, initCassandra, shutdownCassandra} from '@pkgs/cassandra/src/Client';
import {JetStreamConnectionManager} from '@pkgs/nats/src/JetStreamConnectionManager';
import {getDefaultPostgresClient, initPostgres, shutdownPostgres} from '@pkgs/postgres/src/Client';
import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask';
@@ -168,6 +168,7 @@ export async function startWorkerMain(): Promise<void> {
localDc: Config.cassandra.localDc,
username: Config.cassandra.username || undefined,
password: Config.cassandra.password || undefined,
readTimeoutMs: BACKGROUND_READ_TIMEOUT_MS,
});
cassandraInitialized = true;
Logger.info('Cassandra client initialised for worker backend');