diff --git a/fluxer_api/pkgs/cassandra/src/Client.ts b/fluxer_api/pkgs/cassandra/src/Client.ts index 1198c4511..1e148f2eb 100644 --- a/fluxer_api/pkgs/cassandra/src/Client.ts +++ b/fluxer_api/pkgs/cassandra/src/Client.ts @@ -6,6 +6,9 @@ import cassandra from 'cassandra-driver'; const distance = cassandra.types.distance; +const MAX_REQUESTS_PER_CONNECTION = 2048; +const READ_TIMEOUT_MS = 12000; + interface CassandraConfig { hosts: Array; port?: number | undefined; @@ -83,12 +86,15 @@ class CassandraClient implements ICassandraClient { port: this.config.port ?? 9042, }, pooling: { - maxRequestsPerConnection: 32768, + maxRequestsPerConnection: MAX_REQUESTS_PER_CONNECTION, coreConnectionsPerHost: { [distance.local]: 4, [distance.remote]: 2, }, }, + socketOptions: { + readTimeout: READ_TIMEOUT_MS, + }, encoding: { map: Map, set: Set, diff --git a/fluxer_api/src/api/database/CassandraQueryExecution.test.ts b/fluxer_api/src/api/database/CassandraQueryExecution.test.ts new file mode 100644 index 000000000..2761fb32f --- /dev/null +++ b/fluxer_api/src/api/database/CassandraQueryExecution.test.ts @@ -0,0 +1,36 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {ServiceUnavailableError} from '@fluxer/errors/src/domains/core/ServiceUnavailableError'; +import cassandra from 'cassandra-driver'; +import {describe, expect, it} from 'vitest'; +import {mapCassandraDriverError} from './CassandraQueryExecution'; + +function busyConnectionError(): cassandra.errors.BusyConnectionError { + return new cassandra.errors.BusyConnectionError('127.0.0.1:9042', 2048, 4); +} + +describe('mapCassandraDriverError', () => { + it('sheds a busy connection error as a 503', () => { + const mapped = mapCassandraDriverError(busyConnectionError()); + expect(mapped).toBeInstanceOf(ServiceUnavailableError); + expect((mapped as ServiceUnavailableError).status).toBe(503); + }); + + it('sheds a busy connection error nested in a no host available error as a 503', () => { + const mapped = mapCassandraDriverError( + new cassandra.errors.NoHostAvailableError({'127.0.0.1:9042': busyConnectionError()}), + ); + expect(mapped).toBeInstanceOf(ServiceUnavailableError); + expect((mapped as ServiceUnavailableError).status).toBe(503); + }); + + it('returns other no host available errors unchanged', () => { + const err = new cassandra.errors.NoHostAvailableError({'127.0.0.1:9042': new Error('connection refused')}); + expect(mapCassandraDriverError(err)).toBe(err); + }); + + it('returns unrelated errors unchanged', () => { + const err = new cassandra.errors.ResponseError(0x2200, 'invalid query'); + expect(mapCassandraDriverError(err)).toBe(err); + }); +}); diff --git a/fluxer_api/src/api/database/CassandraQueryExecution.ts b/fluxer_api/src/api/database/CassandraQueryExecution.ts index 0f28642e5..4569c9967 100644 --- a/fluxer_api/src/api/database/CassandraQueryExecution.ts +++ b/fluxer_api/src/api/database/CassandraQueryExecution.ts @@ -1,7 +1,8 @@ // SPDX-License-Identifier: AGPL-3.0-or-later +import {ServiceUnavailableError} from '@fluxer/errors/src/domains/core/ServiceUnavailableError'; import {getClient} from '@pkgs/cassandra/src/Client'; -import type cassandra from 'cassandra-driver'; +import cassandra from 'cassandra-driver'; import {Logger} from '../Logger'; import {getQueryType, logBatch, logQuery} from './CassandraDevLogger'; import {getIsDev} from './CassandraMetaRegistry'; @@ -16,6 +17,22 @@ import { const DEFAULT_MAX_PARTITION_KEYS_PER_QUERY = 100; +function isDriverOverloadError(err: unknown): boolean { + if (err instanceof cassandra.errors.BusyConnectionError) { + return true; + } + if (err instanceof cassandra.errors.NoHostAvailableError && err.innerErrors) { + return Object.values(err.innerErrors as Record).some( + (innerError) => innerError instanceof cassandra.errors.BusyConnectionError, + ); + } + return false; +} + +export function mapCassandraDriverError(err: unknown): unknown { + return isDriverOverloadError(err) ? new ServiceUnavailableError() : err; +} + export interface CassandraQueryExecutorForTesting { executeQuery, P extends CassandraParams = CassandraParams>( query: PreparedQuery

, @@ -121,7 +138,7 @@ export async function executeQuery, P extends Cassan } const errorMessage = err instanceof Error ? err.message : String(err); Logger.warn({error: errorMessage, query: cql, params: paramSummary}, 'Cassandra query failed'); - throw err; + throw mapCassandraDriverError(err); } } @@ -249,10 +266,14 @@ async function executeBatch(queries: Array, atomic = true): Promise< counter: false, }; const startTime = getIsDev() ? performance.now() : 0; - await getClient().batch( - queries.map(({query, params}) => ({query, params: normalizeInParams(query, params as CassandraParams)})), - options, - ); + try { + await getClient().batch( + queries.map(({query, params}) => ({query, params: normalizeInParams(query, params as CassandraParams)})), + options, + ); + } catch (err: unknown) { + throw mapCassandraDriverError(err); + } if (getIsDev()) { const durationMs = performance.now() - startTime; logBatch(queries, durationMs);