// SPDX-License-Identifier: AGPL-3.0-or-later import {getClient} from '@pkgs/cassandra/src/Client'; import type cassandra from 'cassandra-driver'; import {Logger} from '../Logger'; import {getQueryType, logBatch, logQuery} from './CassandraDevLogger'; import {getIsDev} from './CassandraMetaRegistry'; import type {CassandraParams, KvQueryMeta, PreparedQuery, QueryTemplate} from './CassandraTypes'; import { assertNoUndefinedParams, chunkArray, isUnsafePreparedStatement, normalizeExecuteArgs, normalizeInParams, } from './CassandraTypes'; const DEFAULT_MAX_PARTITION_KEYS_PER_QUERY = 100; export interface CassandraQueryExecutorForTesting { executeQuery, P extends CassandraParams = CassandraParams>( query: PreparedQuery

, ): Promise>; executePagedQuery?, P extends CassandraParams = CassandraParams>( query: PreparedQuery

, options: { pageSize: number; pageState?: string | null; }, ): Promise>; executeBatch(queries: Array<{query: string; params: object; meta?: KvQueryMeta}>, atomic?: boolean): Promise; reset?(): void; shutdown?(): Promise; } let injectedExecutorForTesting: CassandraQueryExecutorForTesting | null = null; let configuredExecutor: CassandraQueryExecutorForTesting | null = null; function activeExecutor(): CassandraQueryExecutorForTesting | null { return injectedExecutorForTesting ?? configuredExecutor; } export function setDatabaseQueryExecutor(executor: CassandraQueryExecutorForTesting | null): void { configuredExecutor = executor; } export function hasDatabaseQueryExecutor(): boolean { return activeExecutor() !== null; } export function setCassandraQueryExecutorForTesting(executor: CassandraQueryExecutorForTesting | null): void { injectedExecutorForTesting = executor; } export function hasCassandraQueryExecutorForTesting(): boolean { return injectedExecutorForTesting !== null; } export function resetCassandraQueryExecutorForTesting(): void { injectedExecutorForTesting?.reset?.(); } export async function shutdownCassandraQueryExecutorForTesting(): Promise { await injectedExecutorForTesting?.shutdown?.(); injectedExecutorForTesting = null; } export interface PagedQueryResult { rows: Array; pageState: string | null; } async function collectSelectRows(queryType: string, result: cassandra.types.ResultSet): Promise> { if (queryType !== 'SELECT' || !result.pageState) { return (result.rows ?? []) as Array; } const rows: Array = []; for await (const row of result) { rows.push(row as T); } return rows; } export async function executeQuery, P extends CassandraParams = CassandraParams>( queryOrPrepared: string | PreparedQuery

, params?: P, ): Promise> { const {cql, params: boundRaw} = normalizeExecuteArgs(queryOrPrepared, params); const bound = normalizeInParams(cql, boundRaw); if (isUnsafePreparedStatement(cql)) { throw new Error('Cannot prepare a statement that looks like `SELECT *`'); } assertNoUndefinedParams(bound as Record); const executor = activeExecutor(); if (executor) { return executor.executeQuery({ cql, params: bound as P, kvMeta: typeof queryOrPrepared === 'string' ? undefined : queryOrPrepared.kvMeta, }); } const startTime = getIsDev() ? performance.now() : Date.now(); const queryType = getQueryType(cql); try { const result = await getClient().execute(cql, bound, {prepare: true}); const rows = await collectSelectRows(queryType, result); if (getIsDev()) { const durationMs = performance.now() - startTime; logQuery(queryType, cql, bound as Record, durationMs, rows.length); } return rows; } catch (err: unknown) { const paramSummary: Record = {}; for (const [k, v] of Object.entries(bound as Record)) { if (typeof v === 'string') paramSummary[k] = {type: 'string', len: v.length}; else if (typeof v === 'bigint') paramSummary[k] = {type: 'bigint'}; else if (typeof v === 'number') paramSummary[k] = {type: 'number'}; else if (typeof v === 'boolean') paramSummary[k] = {type: 'boolean'}; else if (v instanceof Buffer) paramSummary[k] = {type: 'buffer', len: v.length}; else if (v instanceof Set) paramSummary[k] = {type: 'set', size: (v as Set).size}; else if (v instanceof Map) paramSummary[k] = {type: 'map', size: (v as Map).size}; else if (v instanceof Date) paramSummary[k] = {type: 'date'}; else if (Array.isArray(v)) paramSummary[k] = {type: 'array', len: v.length}; else if (v === null) paramSummary[k] = {type: 'null'}; else paramSummary[k] = {type: typeof v}; } const errorMessage = err instanceof Error ? err.message : String(err); Logger.warn({error: errorMessage, query: cql, params: paramSummary}, 'Cassandra query failed'); throw err; } } export async function fetchOne, P extends CassandraParams = CassandraParams>( queryOrPrepared: PreparedQuery

| string, params?: P, ): Promise { const [row] = await executeQuery(queryOrPrepared, params); return row ?? null; } export async function fetchMany, P extends CassandraParams = CassandraParams>( queryOrPrepared: PreparedQuery

| string, params?: P, ): Promise> { return executeQuery(queryOrPrepared, params); } export async function fetchPage, P extends CassandraParams = CassandraParams>( queryOrPrepared: PreparedQuery

| string, params: P | undefined, options: { pageSize: number; pageState?: string | null; }, ): Promise> { const {cql, params: boundRaw} = normalizeExecuteArgs(queryOrPrepared, params); const bound = normalizeInParams(cql, boundRaw); if (isUnsafePreparedStatement(cql)) { throw new Error('Cannot prepare a statement that looks like `SELECT *`'); } assertNoUndefinedParams(bound as Record); const executor = activeExecutor(); if (executor) { const preparedQuery = { cql, params: bound as P, kvMeta: typeof queryOrPrepared === 'string' ? undefined : queryOrPrepared.kvMeta, }; if (executor.executePagedQuery) { return executor.executePagedQuery(preparedQuery, options); } const rows = await executor.executeQuery(preparedQuery); return { rows: rows.slice(0, options.pageSize), pageState: null, }; } const result = await getClient().execute(cql, bound, { prepare: true, fetchSize: options.pageSize, pageState: options.pageState ?? undefined, }); return { rows: (result.rows as Array) ?? [], pageState: result.pageState ?? null, }; } export async function fetchManyInChunks< T = Record, V = unknown, P extends CassandraParams = CassandraParams, >( query: QueryTemplate

| PreparedQuery

| string, values: Array, paramsFactory: (chunk: Array) => P, chunkSize = DEFAULT_MAX_PARTITION_KEYS_PER_QUERY, ): Promise> { if (values.length === 0) return []; const chunks = chunkArray(values, chunkSize); const results = await Promise.all( chunks.map(async (chunk) => { const params = paramsFactory(chunk); if (typeof query === 'string') { return executeQuery(query, params); } if ((query as PreparedQuery

).params !== undefined) { return executeQuery(query as PreparedQuery

); } return executeQuery((query as QueryTemplate

).bind(params)); }), ); return results.flat(); } export async function upsertOne

( queryOrPrepared: PreparedQuery

| string, params?: P, ): Promise { await executeQuery(queryOrPrepared, params); } export async function deleteOneOrMany

( queryOrPrepared: PreparedQuery

| string, params?: P, ): Promise { await executeQuery(queryOrPrepared, params); } interface BatchQuery { query: string; params: object; meta?: KvQueryMeta; } async function executeBatch(queries: Array, atomic = true): Promise { if (queries.length === 0) return; for (const {query} of queries) { if (isUnsafePreparedStatement(query)) { throw new Error('Cannot prepare a statement that looks like `SELECT *`'); } } for (const {params} of queries) { assertNoUndefinedParams(params as Record); } const executor = activeExecutor(); if (executor) { await executor.executeBatch(queries, atomic); return; } const options = { prepare: true, logged: atomic, counter: false, }; const startTime = getIsDev() ? performance.now() : 0; await getClient().batch( queries.map(({query, params}) => ({query, params: normalizeInParams(query, params as CassandraParams)})), options, ); if (getIsDev()) { const durationMs = performance.now() - startTime; logBatch(queries, durationMs); } } export class BatchBuilder { private queries: Array = []; add(query: string, params: object, meta?: KvQueryMeta): this { this.queries.push({query, params, meta}); return this; } addPrepared(q: PreparedQuery): this { this.queries.push({query: q.cql, params: q.params, meta: q.kvMeta}); return this; } addIf(condition: boolean, query: string, params: object, meta?: KvQueryMeta): this { if (condition) this.queries.push({query, params, meta}); return this; } addPreparedIf(condition: boolean, q: PreparedQuery): this { if (condition) this.queries.push({query: q.cql, params: q.params, meta: q.kvMeta}); return this; } async execute(atomic = true): Promise { if (this.queries.length === 0) return; await executeBatch(this.queries, atomic); } async executeChunked(chunkSize: number, atomic = false): Promise { if (this.queries.length === 0) return; for (let i = 0; i < this.queries.length; i += chunkSize) { await executeBatch(this.queries.slice(i, i + chunkSize), atomic); } } getQueries(): Array { return this.queries; } }