diff --git a/docker-compose.yaml b/docker-compose.yaml index a7c226f..6ca192a 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -72,6 +72,10 @@ services: scylla: image: scylladb/scylla:5.2 + # Cap Scylla's footprint so it doesn't grab nearly all available RAM + # on small hosts (it allocates ~all "available" memory at startup by + # default, which starves the rest of the stack on a VPS env). + command: --smp 1 --memory 4G --overprovisioned 1 --developer-mode 1 volumes: - scylla_data:/var/lib/scylla ports: diff --git a/server/.env.example b/server/.env.example index 261225f..49a5f5a 100644 --- a/server/.env.example +++ b/server/.env.example @@ -44,6 +44,16 @@ CLICKHOUSE_PASSWORD=clickhouse CLICKHOUSE_DATABASE=analytics CLICKHOUSE_PROTOCOL=http +# ClickHouse analytics insert retry policy. By default retries up to twice on +# transient network errors (ECONNRESET, socket hang up, etc.) with jittered +# exponential backoff starting at 100ms, capped at 1000ms (cap only reached +# at higher retry counts). After exhaustion, transient failures are logged +# and dropped; non-transient errors (schema, auth, payload) propagate. +# Set CLICKHOUSE_INSERT_MAX_RETRIES=0 to disable retries entirely. +# CLICKHOUSE_INSERT_MAX_RETRIES=2 +# CLICKHOUSE_INSERT_RETRY_INITIAL_MS=100 +# CLICKHOUSE_INSERT_RETRY_MAX_MS=1000 + HMA_SERVICE_URL=http://localhost:9876 # Scylla Cluster Details diff --git a/server/bin/run-worker-or-job.ts b/server/bin/run-worker-or-job.ts index c72db86..d0b6dae 100644 --- a/server/bin/run-worker-or-job.ts +++ b/server/bin/run-worker-or-job.ts @@ -69,5 +69,14 @@ process.on('uncaughtException', (err, _) => { process.exit(1); }); +// Log but don't exit; a stray rejection shouldn't kill the worker. +process.on('unhandledRejection', (reason) => { + // eslint-disable-next-line no-restricted-syntax + logErrorJson({ + message: 'UnhandledRejection', + error: reason instanceof Error ? reason : new Error(String(reason)), + }); +}); + process.once('SIGTERM', onFinish); process.once('SIGINT', onFinish); diff --git a/server/bin/www.ts b/server/bin/www.ts index 61889f3..1582748 100755 --- a/server/bin/www.ts +++ b/server/bin/www.ts @@ -132,6 +132,15 @@ process.on('uncaughtException', (err, _) => { process.exit(1); }); +// Log but don't exit; a stray rejection shouldn't kill the server. +process.on('unhandledRejection', (reason) => { + // eslint-disable-next-line no-restricted-syntax + logErrorJson({ + message: 'UnhandledRejection', + error: reason instanceof Error ? reason : new Error(String(reason)), + }); +}); + process.once('SIGTERM', shutdownOnce); process.once('SIGINT', shutdownOnce); diff --git a/server/iocContainer/utils.ts b/server/iocContainer/utils.ts index f41e045..053f0c8 100644 --- a/server/iocContainer/utils.ts +++ b/server/iocContainer/utils.ts @@ -691,3 +691,24 @@ export function safeGetEnvInt(varName: string, defaultValue: number): number { } return parsed; } + +/** + * Like `safeGetEnvInt` but allows `0`. Use when zero is a meaningful value + * (e.g. disabling retries, no timeout). + */ +export function safeGetEnvNonNegativeInt( + varName: string, + defaultValue: number, +): number { + const raw = process.env[varName]; + if (raw === undefined) return defaultValue; + const parsed = parseInt(raw, 10); + if (!Number.isInteger(parsed) || parsed < 0) { + // eslint-disable-next-line no-console + console.error( + `Invalid env var ${varName}: expected a non-negative integer, got ${jsonStringify(raw)}. Using default value ${defaultValue}.`, + ); + return defaultValue; + } + return parsed; +} diff --git a/server/plugins/analytics/adapters/ClickhouseAnalyticsAdapter.ts b/server/plugins/analytics/adapters/ClickhouseAnalyticsAdapter.ts index 66fb16d..754f386 100644 --- a/server/plugins/analytics/adapters/ClickhouseAnalyticsAdapter.ts +++ b/server/plugins/analytics/adapters/ClickhouseAnalyticsAdapter.ts @@ -2,6 +2,7 @@ import { createClient, type ClickHouseClient } from '@clickhouse/client'; import type SafeTracer from '../../../utils/SafeTracer.js'; import { jsonStringify, tryJsonParse } from '../../../utils/encoding.js'; +import { logErrorJson } from '../../../utils/logging.js'; import type { IAnalyticsAdapter } from '../IAnalyticsAdapter.js'; import { type AnalyticsEventInput, @@ -9,6 +10,10 @@ import { type AnalyticsWriteOptions, } from '../types.js'; import { formatClickhouseQuery } from '../../warehouse/utils/clickhouseSql.js'; +import { + isTransientNetworkError, + withClickhouseInsertRetries, +} from './clickhouseRetry.js'; export interface ClickhouseAnalyticsConnection { host: string; @@ -94,6 +99,10 @@ export class ClickhouseAnalyticsAdapter implements IAnalyticsAdapter { private readonly tracer?: SafeTracer; private readonly client: ClickHouseClient; private readonly defaultBatchSize: number; + private readonly insertWithRetries: (params: { + table: string; + values: AnalyticsEventInput[]; + }) => Promise; constructor(options: ClickhouseAnalyticsAdapterOptions) { this.tracer = options.tracer; @@ -112,6 +121,12 @@ export class ClickhouseAnalyticsAdapter implements IAnalyticsAdapter { ...(password ? { password } : {}), database: options.connection.database, }); + + this.insertWithRetries = withClickhouseInsertRetries( + async ({ table, values }: { table: string; values: AnalyticsEventInput[] }) => { + await this.client.insert({ table, values, format: 'JSONEachRow' }); + }, + ); } async writeEvents( @@ -130,11 +145,20 @@ export class ClickhouseAnalyticsAdapter implements IAnalyticsAdapter { this.normalizeRecord(row, table), ); - await this.client.insert({ - table, - values: normalizedBatch, - format: 'JSONEachRow', - }); + try { + await this.insertWithRetries({ table, values: normalizedBatch }); + } catch (err) { + // Only swallow transient network errors; let schema/auth/payload + // errors propagate so they're not silently dropped. + if (!isTransientNetworkError(err)) { + throw err; + } + // eslint-disable-next-line no-restricted-syntax + logErrorJson({ + message: `clickhouse.analytics.insert_failed table=${table} batchSize=${normalizedBatch.length}`, + error: err, + }); + } } } diff --git a/server/plugins/analytics/adapters/clickhouseRetry.ts b/server/plugins/analytics/adapters/clickhouseRetry.ts new file mode 100644 index 0000000..0eba7a8 --- /dev/null +++ b/server/plugins/analytics/adapters/clickhouseRetry.ts @@ -0,0 +1,46 @@ +import { + safeGetEnvInt, + safeGetEnvNonNegativeInt, +} from '../../../iocContainer/utils.js'; +import { withRetries } from '../../../utils/misc.js'; + +// Network errors we'll retry on. ClickHouse over HTTP can RST in-flight +// connections (remote restart, idle-socket reaper between us and CH, etc.); +// these are transient and worth one or two retries before giving up. +const RETRYABLE_ERROR_CODES = new Set([ + 'ECONNRESET', + 'ECONNREFUSED', + 'ETIMEDOUT', + 'EPIPE', + 'EAI_AGAIN', +]); + +export function isTransientNetworkError(err: unknown): boolean { + if (err == null || typeof err !== 'object') return false; + const code = (err as { code?: unknown }).code; + if (typeof code === 'string' && RETRYABLE_ERROR_CODES.has(code)) { + return true; + } + const message = (err as { message?: unknown }).message; + return typeof message === 'string' && message.includes('socket hang up'); +} + +export function withClickhouseInsertRetries( + fn: (this: void, ...args: Args) => Promise, +): (...args: Args) => Promise { + return withRetries( + { + maxRetries: safeGetEnvNonNegativeInt('CLICKHOUSE_INSERT_MAX_RETRIES', 2), + initialTimeMsBetweenRetries: safeGetEnvInt( + 'CLICKHOUSE_INSERT_RETRY_INITIAL_MS', + 100, + ), + maxTimeMsBetweenRetries: safeGetEnvInt( + 'CLICKHOUSE_INSERT_RETRY_MAX_MS', + 1000, + ), + isRetryableError: isTransientNetworkError, + }, + fn, + ); +} diff --git a/server/rule_engine/RuleEngine.ts b/server/rule_engine/RuleEngine.ts index 355d67a..3e01657 100644 --- a/server/rule_engine/RuleEngine.ts +++ b/server/rule_engine/RuleEngine.ts @@ -22,6 +22,7 @@ import { type CorrelationIdType, } from '../utils/correlationIds.js'; import { equalLengthZip } from '../utils/fp-helpers.js'; +import { logErrorJson } from '../utils/logging.js'; import { safePick } from '../utils/misc.js'; import type SafeTracer from '../utils/SafeTracer.js'; import { @@ -249,24 +250,34 @@ class RuleEngine { rules.map((it) => it.id), ); - const logRuleExecutionsPromise = this.ruleExecutionLogger.logRuleExecutions( - [...rulesToResults.entries()].map(([rule, result]) => ({ - orgId: org.id, - rule: { - id: rule.id, - name: rule.name, - version: rule.latestVersion.version, - tags: rule.tags, - }, - ruleInput, - environment, - result: result.conditionResults, - correlationId: executionsCorrelationId, - passed: result.passed, - policies: policiesByRule[rule.id] ?? [], - })), - sync, - ); + // Catch at construction; awaited below via Promise.all. + const logRuleExecutionsPromise = this.ruleExecutionLogger + .logRuleExecutions( + [...rulesToResults.entries()].map(([rule, result]) => ({ + orgId: org.id, + rule: { + id: rule.id, + name: rule.name, + version: rule.latestVersion.version, + tags: rule.tags, + }, + ruleInput, + environment, + result: result.conditionResults, + correlationId: executionsCorrelationId, + passed: result.passed, + policies: policiesByRule[rule.id] ?? [], + })), + sync, + ) + .catch((err) => { + this.tracer.logActiveSpanFailedIfAny(err); + // eslint-disable-next-line no-restricted-syntax + logErrorJson({ + message: `logRuleExecutions failed orgId=${org.id} environment=${environment} correlationId=${executionsCorrelationId}`, + error: err, + }); + }); if (!shouldRunActions) { await logRuleExecutionsPromise; diff --git a/server/services/reportingService/reportingRuleEngine.ts b/server/services/reportingService/reportingRuleEngine.ts index 4571e62..98952cc 100644 --- a/server/services/reportingService/reportingRuleEngine.ts +++ b/server/services/reportingService/reportingRuleEngine.ts @@ -5,6 +5,7 @@ import { type Dependencies } from '../../iocContainer/index.js'; import { RuleEnvironment } from '../../rule_engine/RuleEngine.js'; import { type RuleEvaluationContext } from '../../rule_engine/RuleEvaluator.js'; import { equalLengthZip } from '../../utils/fp-helpers.js'; +import { logErrorJson } from '../../utils/logging.js'; import { safePick } from '../../utils/misc.js'; import { type ItemSubmission } from '../itemProcessingService/makeItemSubmission.js'; import { type Action } from '../moderationConfigService/index.js'; @@ -126,8 +127,9 @@ export default class ReportingRuleEngine { // since, logically, each rule triggered the action. const { org, input: ruleInput } = evaluationContext; - const logRuleExecutionsPromise = - this.reportingRuleExecutionLogger.logReportingRuleExecutions( + // Catch at construction; awaited below via Promise.all. + const logRuleExecutionsPromise = this.reportingRuleExecutionLogger + .logReportingRuleExecutions( [...rulesToResults.entries()].map(([rule, result]) => ({ orgId: org.id, reportingRule: { @@ -145,7 +147,15 @@ export default class ReportingRuleEngine { [], policyIds: rule.policyIds, })), - ); + ) + .catch((err) => { + this.tracer.logActiveSpanFailedIfAny(err); + // eslint-disable-next-line no-restricted-syntax + logErrorJson({ + message: `logReportingRuleExecutions failed orgId=${org.id} environment=${environment} correlationId=${executionsCorrelationId}`, + error: err, + }); + }); if (!shouldRunActions) { await logRuleExecutionsPromise;