From beb906753fe0de6db2d2775bdfd593c3f7c581c4 Mon Sep 17 00:00:00 2001 From: Hampus Date: Fri, 14 Aug 2026 19:55:53 +0200 Subject: [PATCH] fix(api): eliminate idle Postgres I/O on self-hosted instances (#1593) --- fluxer_api/src/api/Config.ts | 3 + fluxer_api/src/api/admin/AdminRepository.ts | 32 +-------- fluxer_api/src/api/admin/IAdminRepository.ts | 9 --- fluxer_api/src/api/app/APILifecycle.ts | 65 +------------------ fluxer_api/src/api/config/APIConfig.ts | 3 + .../api/database/PostgresKvQueryExecutor.ts | 26 ++++---- .../src/api/middleware/ServiceMiddleware.ts | 23 +------ fluxer_api/src/api/risk/RiskCacheLoaders.ts | 23 ------- fluxer_api/src/api/risk/RiskCacheManager.ts | 64 ------------------ fluxer_api/src/api/risk/RiskToolboxFactory.ts | 6 +- fluxer_api/src/api/risk/RiskTypes.ts | 1 - .../__tests__/DeterministicRiskEngine.test.ts | 1 - .../risk/adapters/DisposableDomainChecker.ts | 13 ++-- fluxer_api/src/api/worker/CronScheduler.ts | 12 +++- fluxer_api/src/api/worker/WorkerMain.ts | 50 ++++++++------ .../tasks/SyncDisposableEmailDomains.ts | 12 +--- packages/config/src/ConfigLoader.ts | 1 + packages/config/src/MasterConfig.ts | 3 + .../src/config_loader/EnvironmentOverrides.ts | 1 + 19 files changed, 80 insertions(+), 268 deletions(-) delete mode 100644 fluxer_api/src/api/risk/RiskCacheLoaders.ts delete mode 100644 fluxer_api/src/api/risk/RiskCacheManager.ts diff --git a/fluxer_api/src/api/Config.ts b/fluxer_api/src/api/Config.ts index 6fc813abd..76bdf7408 100644 --- a/fluxer_api/src/api/Config.ts +++ b/fluxer_api/src/api/Config.ts @@ -276,6 +276,9 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig { ipinfoApiKey: master.integrations.risk_integration.ipinfo_api_key || undefined, accountPolicyDsl: master.integrations.risk_integration.account_policy_dsl, }, + blocklistFeeds: { + enabled: master.integrations.blocklist_feeds.enabled ?? !master.instance.self_hosted, + }, captcha: { enabled: master.integrations.captcha.enabled, provider: master.integrations.captcha.provider, diff --git a/fluxer_api/src/api/admin/AdminRepository.ts b/fluxer_api/src/api/admin/AdminRepository.ts index 418a6f520..8df71ce85 100644 --- a/fluxer_api/src/api/admin/AdminRepository.ts +++ b/fluxer_api/src/api/admin/AdminRepository.ts @@ -2,7 +2,7 @@ import {getSameIpDecisionKey} from '@fluxer/ip_utils/src/IpAddress'; import {createUserID} from '../BrandedTypes'; -import {deleteOneOrMany, fetchMany, fetchOne, fetchPage, upsertOne} from '../database/CassandraQueryExecution'; +import {deleteOneOrMany, fetchMany, fetchOne, upsertOne} from '../database/CassandraQueryExecution'; import type { AdminAuditLogRow, BannedAvatarHashRow, @@ -30,13 +30,7 @@ import { } from '../Tables'; import {parseIpBanEntry, tryParseSingleIp} from '../utils/IpRangeUtils'; import {canonicalizeStoredPhrase} from '../utils/PhraseBlocklistNormalization'; -import type { - AdminAuditLog, - BannedIpEntry, - BannedIpKind, - DisposableEmailDomainPage, - IAdminRepository, -} from './IAdminRepository'; +import type {AdminAuditLog, BannedIpEntry, BannedIpKind, IAdminRepository} from './IAdminRepository'; const FETCH_AUDIT_LOG_BY_ID_QUERY = AdminAuditLogs.select({ where: AdminAuditLogs.where.eq('log_id'), @@ -51,8 +45,6 @@ const IS_EMAIL_BANNED_QUERY = BannedEmails.select({ const IS_EMAIL_DOMAIN_SUSPICIOUS_QUERY = SuspiciousEmailDomains.select({ where: SuspiciousEmailDomains.where.eq('domain'), }); -const createLoadSuspiciousEmailDomainsQuery = (limit?: number) => - limit ? SuspiciousEmailDomains.select({limit}) : SuspiciousEmailDomains.select(); const IS_EMAIL_DOMAIN_DISPOSABLE_QUERY = DisposableEmailDomains.select({ where: DisposableEmailDomains.where.eq('domain'), }); @@ -273,13 +265,6 @@ export class AdminRepository implements IAdminRepository { await deleteOneOrMany(SuspiciousEmailDomains.deleteByPk({domain: domainLower})); } - async listSuspiciousEmailDomains(limit?: number): Promise> { - const rows = await fetchMany<{ - domain: string; - }>(createLoadSuspiciousEmailDomainsQuery(limit).bind({})); - return rows.map((row) => row.domain); - } - async isEmailDomainDisposable(domain: string): Promise { const domainLower = domain.toLowerCase(); if (isAccountPolicyContactDomainReputationExempt(domainLower)) return false; @@ -306,19 +291,6 @@ export class AdminRepository implements IAdminRepository { return rows.map((row) => row.domain); } - async listDisposableEmailDomainsPage(limit: number, pageState?: string | null): Promise { - const page = await fetchPage<{ - domain: string; - }>(createLoadDisposableEmailDomainsQuery().bind({}), undefined, { - pageSize: limit, - pageState, - }); - return { - domains: page.rows.map((row) => row.domain), - pageState: page.pageState, - }; - } - async isPhraseBanned(phrase: string): Promise { const phraseLower = canonicalizeStoredPhrase(phrase); const result = await fetchOne<{ diff --git a/fluxer_api/src/api/admin/IAdminRepository.ts b/fluxer_api/src/api/admin/IAdminRepository.ts index a2d512d73..45d55742a 100644 --- a/fluxer_api/src/api/admin/IAdminRepository.ts +++ b/fluxer_api/src/api/admin/IAdminRepository.ts @@ -32,11 +32,6 @@ export interface BannedIpEntry { createdAt: Date | null; } -export interface DisposableEmailDomainPage { - domains: Array; - pageState: string | null; -} - export abstract class IAdminRepository { abstract createAuditLog(log: AdminAuditLogRow): Promise; @@ -68,8 +63,6 @@ export abstract class IAdminRepository { abstract removeSuspiciousEmailDomain(domain: string): Promise; - abstract listSuspiciousEmailDomains(limit?: number): Promise>; - abstract isEmailDomainDisposable(domain: string): Promise; abstract addDisposableEmailDomain(domain: string): Promise; @@ -78,8 +71,6 @@ export abstract class IAdminRepository { abstract listDisposableEmailDomains(limit?: number): Promise>; - abstract listDisposableEmailDomainsPage(limit: number, pageState?: string | null): Promise; - abstract isPhraseBanned(phrase: string): Promise; abstract banPhrase(phrase: string): Promise; diff --git a/fluxer_api/src/api/app/APILifecycle.ts b/fluxer_api/src/api/app/APILifecycle.ts index be19b2eee..2cfd66fee 100644 --- a/fluxer_api/src/api/app/APILifecycle.ts +++ b/fluxer_api/src/api/app/APILifecycle.ts @@ -13,11 +13,7 @@ import type {ILogger} from '../ILogger'; import {JobLedgerRepository} from '../jobs/JobLedgerRepository'; import {startAbuseReplicationSubscriber, stopAbuseReplicationSubscriber} from '../middleware/AbusiveIpAutoBanner'; import {ipBanCache} from '../middleware/IpBanMiddleware'; -import { - getRiskCacheManagerInstance, - initializeServiceSingletons, - shutdownReportService, -} from '../middleware/ServiceMiddleware'; +import {initializeServiceSingletons, shutdownReportService} from '../middleware/ServiceMiddleware'; import { ensureVoiceResourcesInitialized, getKVClient, @@ -40,57 +36,6 @@ import {JetStreamWorkerQueue} from '../worker/JetStreamWorkerQueue'; import {WorkerService} from '../worker/WorkerService'; let jsConnectionManager: JetStreamConnectionManager | null = null; -let riskCacheRefreshInterval: NodeJS.Timeout | null = null; -let riskCacheRefreshInFlight = false; - -const RISK_CACHE_REFRESH_INTERVAL_MS = 5 * 60 * 1000; - -async function refreshRiskCache(logger: ILogger, source: 'startup' | 'interval'): Promise { - if (riskCacheRefreshInFlight) { - return; - } - riskCacheRefreshInFlight = true; - try { - const result = await getRiskCacheManagerInstance().refresh(); - if (result.subtaskErrors.length > 0) { - logger.warn({source, errors: result.subtaskErrors}, 'Risk cache refresh completed with errors'); - return; - } - logger.info( - { - source, - disposableDomainCount: result.disposableDomainCount, - }, - source === 'startup' ? 'Risk cache initialized on API startup' : 'Risk cache refresh complete on API', - ); - } catch (error) { - if (source === 'startup') { - logger.warn({error}, 'Risk cache initialisation failed on API startup'); - return; - } - logger.warn({error}, 'Periodic risk cache refresh failed on API'); - } finally { - riskCacheRefreshInFlight = false; - } -} - -function startRiskCacheRefreshLoop(logger: ILogger): void { - if (riskCacheRefreshInterval) { - return; - } - riskCacheRefreshInterval = setInterval(() => { - void refreshRiskCache(logger, 'interval'); - }, RISK_CACHE_REFRESH_INTERVAL_MS); -} - -function stopRiskCacheRefreshLoop(): void { - if (!riskCacheRefreshInterval) { - return; - } - clearInterval(riskCacheRefreshInterval); - riskCacheRefreshInterval = null; -} - export function createInitializer(config: APIConfig, logger: ILogger): () => Promise { return async (): Promise => { try { @@ -171,8 +116,6 @@ export function createInitializer(config: APIConfig, logger: ILogger): () => Pro logger.info('Profile substring blocklist cache initialized'); await initializeServiceSingletons(); logger.info('Service singletons initialized'); - await refreshRiskCache(logger, 'startup'); - startRiskCacheRefreshLoop(logger); if (!config.dev.testModeEnabled) { jsConnectionManager = new JetStreamConnectionManager({ url: config.nats.jetStreamUrl, @@ -285,12 +228,6 @@ export function createShutdown(logger: ILogger): () => Promise { } catch (error) { logger.error({error}, 'Error shutting down search service'); } - try { - stopRiskCacheRefreshLoop(); - logger.info('Risk cache refresh loop shut down'); - } catch (error) { - logger.error({error}, 'Error shutting down risk cache refresh loop'); - } try { ipBanCache.shutdown(); logger.info('IP ban cache shut down'); diff --git a/fluxer_api/src/api/config/APIConfig.ts b/fluxer_api/src/api/config/APIConfig.ts index 0c7f2afcc..82b0e2e7f 100644 --- a/fluxer_api/src/api/config/APIConfig.ts +++ b/fluxer_api/src/api/config/APIConfig.ts @@ -178,6 +178,9 @@ export interface APIConfig { ipinfoApiKey?: string; accountPolicyDsl?: unknown; }; + blocklistFeeds: { + enabled: boolean; + }; captcha: { enabled: boolean; provider: 'hcaptcha' | 'turnstile' | 'none'; diff --git a/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts b/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts index 4af2b32b4..db05aea5e 100644 --- a/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts +++ b/fluxer_api/src/api/database/PostgresKvQueryExecutor.ts @@ -561,14 +561,14 @@ WHERE NOT $6`, private async patch(meta: KvQueryMeta, params: CassandraParams, db: PostgresQueryable): Promise { const key = rowKeyFromParams(meta, params); - const existing = await this.getRow(meta, key, db); - const base = existing ?? paramsRow(params, (meta.pkColumns ?? meta.table.primaryKey) as ReadonlyArray); + const stored = await this.getStoredRow(meta, key, db); + const base = stored?.row ?? paramsRow(params, (meta.pkColumns ?? meta.table.primaryKey) as ReadonlyArray); const next = {...base}; for (const column of meta.patchKeys ?? []) { next[column] = column in params ? params[column] : null; } const ttl = ttlExpiresAt(meta, params); - const expiresAt = ttl === undefined ? await this.getExpiresAt(meta, key, db) : ttl; + const expiresAt = ttl === undefined ? (stored?.expiresAt ?? null) : ttl; await db.query( `INSERT INTO ${this.table} (table_name, partition_key, row_key, row_data, expires_at, updated_at) VALUES ($1, $2, $3, $4::jsonb, $5, now()) @@ -592,20 +592,20 @@ DO UPDATE SET partition_key = EXCLUDED.partition_key, row_data = EXCLUDED.row_da ]); } - private async getRow(meta: KvQueryMeta, key: string, db: PostgresQueryable): Promise { - const result = await db.query( - `SELECT row_key, row_data FROM ${this.table} WHERE table_name = $1 AND row_key = $2 AND (expires_at IS NULL OR expires_at > now()) LIMIT 1`, + private async getStoredRow( + meta: KvQueryMeta, + key: string, + db: PostgresQueryable, + ): Promise<{row: Row; expiresAt: Date | null} | null> { + const result = await db.query( + `SELECT row_key, row_data, expires_at FROM ${this.table} WHERE table_name = $1 AND row_key = $2 AND (expires_at IS NULL OR expires_at > now()) LIMIT 1`, [meta.table.name, key], ); const row = result.rows[0]; - return row ? decodeRow(row.row_data) : null; + return row ? {row: decodeRow(row.row_data), expiresAt: row.expires_at ?? null} : null; } - private async getExpiresAt(meta: KvQueryMeta, key: string, db: PostgresQueryable): Promise { - const result = await db.query<{expires_at: Date | null}>( - `SELECT expires_at FROM ${this.table} WHERE table_name = $1 AND row_key = $2 AND (expires_at IS NULL OR expires_at > now()) LIMIT 1`, - [meta.table.name, key], - ); - return result.rows[0]?.expires_at ?? null; + private async getRow(meta: KvQueryMeta, key: string, db: PostgresQueryable): Promise { + return (await this.getStoredRow(meta, key, db))?.row ?? null; } } diff --git a/fluxer_api/src/api/middleware/ServiceMiddleware.ts b/fluxer_api/src/api/middleware/ServiceMiddleware.ts index 425b8c94e..4e8e013a4 100644 --- a/fluxer_api/src/api/middleware/ServiceMiddleware.ts +++ b/fluxer_api/src/api/middleware/ServiceMiddleware.ts @@ -60,8 +60,6 @@ import {CassandraHistoricalOutcomeRepository} from '../risk/HistoricalOutcomeRep import {buildIpInfoCache, buildIpInfoRequestAuditLogger} from '../risk/IpInfoCacheFactory'; import {CassandraRegistrationEventsRepository} from '../risk/RegistrationEventsRepository'; import {CassandraRiskAssessmentRepository} from '../risk/RiskAssessmentRepository'; -import {buildRiskCacheLoaders} from '../risk/RiskCacheLoaders'; -import {RiskCacheManager} from '../risk/RiskCacheManager'; import {createRiskToolbox} from '../risk/RiskToolboxFactory'; import {CassandraSuspiciousIpRepository} from '../risk/SuspiciousIpRepository'; import {RpcService} from '../rpc/RpcService'; @@ -192,24 +190,6 @@ export function shutdownReportService(): void { } } -let _riskCacheManager: RiskCacheManager | null = null; - -function getRiskCacheManager(): RiskCacheManager { - if (!_riskCacheManager) { - _riskCacheManager = new RiskCacheManager({ - logger: Logger, - ...buildRiskCacheLoaders({ - adminRepository: getAdminRepository(), - }), - }); - } - return _riskCacheManager; -} - -export function getRiskCacheManagerInstance(): RiskCacheManager { - return getRiskCacheManager(); -} - let _inboundSmsChallengeService: InboundSmsChallengeService | null = null; function getInboundSmsChallengeService(): InboundSmsChallengeService { @@ -311,13 +291,12 @@ function getRegistrationRiskEvaluator(): IRegistrationRiskEvaluator { _registrationRiskEvaluator = noopRegistrationRiskEvaluator; return _registrationRiskEvaluator; } - const cacheManager = getRiskCacheManager(); const ipInfoService = getIpInfoService(); const ipInfoChecker = Config.risk.ipinfoApiKey ? createIpInfoChecker({ipInfoService}) : undefined; const cacheService = getCacheService(); const reverseDnsLookup = createReverseDnsLookup({cacheService}); const toolbox = createRiskToolbox({ - disposableDomainsRef: cacheManager.disposableDomainsRef, + adminRepository: getAdminRepository(), ipInfoChecker, reverseDnsLookup, ipInfoService, diff --git a/fluxer_api/src/api/risk/RiskCacheLoaders.ts b/fluxer_api/src/api/risk/RiskCacheLoaders.ts deleted file mode 100644 index af5f09cfc..000000000 --- a/fluxer_api/src/api/risk/RiskCacheLoaders.ts +++ /dev/null @@ -1,23 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -import type {IAdminRepository} from '../admin/IAdminRepository'; -import type {RiskCacheLoaders} from './RiskCacheManager'; - -interface BuildRiskCacheLoadersDeps { - adminRepository: Pick; -} - -export function buildRiskCacheLoaders(deps: BuildRiskCacheLoadersDeps): RiskCacheLoaders { - return { - loadDisposableDomains: async () => { - const [suspicious, disposable] = await Promise.all([ - deps.adminRepository.listSuspiciousEmailDomains(), - deps.adminRepository.listDisposableEmailDomains(), - ]); - const set = new Set(); - for (const d of suspicious) set.add(d.toLowerCase()); - for (const d of disposable) set.add(d.toLowerCase()); - return set; - }, - }; -} diff --git a/fluxer_api/src/api/risk/RiskCacheManager.ts b/fluxer_api/src/api/risk/RiskCacheManager.ts deleted file mode 100644 index ae83b69cc..000000000 --- a/fluxer_api/src/api/risk/RiskCacheManager.ts +++ /dev/null @@ -1,64 +0,0 @@ -// SPDX-License-Identifier: AGPL-3.0-or-later - -interface MutableRef { - current: T; -} - -export type ReadonlyRiskCacheRef = { - readonly current: T; -}; - -interface RiskCacheManagerLogger { - info(payload: object, msg: string): void; - warn(payload: object, msg: string): void; -} - -export interface RiskCacheLoaders { - loadDisposableDomains: () => Promise>; -} - -interface RiskCacheManagerOptions extends RiskCacheLoaders { - logger?: RiskCacheManagerLogger; -} - -interface RiskCacheRefreshResult { - disposableDomainCount: number; - subtaskErrors: ReadonlyArray<{ - step: string; - error: string; - }>; -} - -export class RiskCacheManager { - private readonly _disposable: MutableRef> = {current: new Set()}; - readonly disposableDomainsRef: ReadonlyRiskCacheRef> = this._disposable; - private readonly logger: RiskCacheManagerLogger | undefined; - private readonly loaders: RiskCacheLoaders; - - constructor(opts: RiskCacheManagerOptions) { - this.logger = opts.logger; - this.loaders = { - loadDisposableDomains: opts.loadDisposableDomains, - }; - } - - async refresh(): Promise { - const subtaskErrors: Array<{ - step: string; - error: string; - }> = []; - try { - const set = await this.loaders.loadDisposableDomains(); - this._disposable.current = set; - this.logger?.info({count: set.size}, 'RiskCacheManager: loaded disposable domains from DB'); - } catch (err) { - const message = err instanceof Error ? err.message : String(err); - this.logger?.warn({step: 'disposable_domains', err: message}, 'RiskCacheManager: subtask failed'); - subtaskErrors.push({step: 'disposable_domains', error: message}); - } - return { - disposableDomainCount: this._disposable.current.size, - subtaskErrors, - }; - } -} diff --git a/fluxer_api/src/api/risk/RiskToolboxFactory.ts b/fluxer_api/src/api/risk/RiskToolboxFactory.ts index 843c46968..23e9a0a70 100644 --- a/fluxer_api/src/api/risk/RiskToolboxFactory.ts +++ b/fluxer_api/src/api/risk/RiskToolboxFactory.ts @@ -2,6 +2,7 @@ import type {ICacheService} from '@pkgs/cache/src/ICacheService'; import type {IpInfoService} from '@pkgs/geoip/src/IpInfoService'; +import type {IAdminRepository} from '../admin/IAdminRepository'; import {createDisposableDomainChecker} from './adapters/DisposableDomainChecker'; import {createDnsMxChecker, type MxResolver, NodeDnsMxResolver} from './adapters/DnsMxChecker'; import {createDomainAgeChecker} from './adapters/DomainAgeChecker'; @@ -13,13 +14,12 @@ import {analyzeRegistrationTiming} from './adapters/RegistrationTimingAnalyzer'; import {analyzeUserAgent} from './adapters/UserAgentAnalyzer'; import {createVelocityAdapter, type IRegistrationEventsRepository} from './adapters/VelocityAdapter'; import type {IRiskHistoryRepository} from './HistoricalOutcomeRepository'; -import type {ReadonlyRiskCacheRef} from './RiskCacheManager'; import type {RiskToolbox} from './RiskToolbox'; import type {IpInfoAnonymousResult, ReverseDnsResult} from './RiskTypes'; import type {ISuspiciousIpRepository} from './SuspiciousIpRepository'; interface RiskToolboxFactoryOptions { - disposableDomainsRef: ReadonlyRiskCacheRef>; + adminRepository: Pick; ipInfoChecker?: (ip: string) => Promise; reverseDnsLookup?: (ip: string) => Promise; ipInfoService: IpInfoService; @@ -32,7 +32,7 @@ interface RiskToolboxFactoryOptions { } export function createRiskToolbox(opts: RiskToolboxFactoryOptions): RiskToolbox { - const checkDomainDisposable = createDisposableDomainChecker({disposableDomainsRef: opts.disposableDomainsRef}); + const checkDomainDisposable = createDisposableDomainChecker({adminRepository: opts.adminRepository}); const lookupGeoIpCity = createGeoIpCityAdapter({ipInfoService: opts.ipInfoService}); const lookupGeoIpAsn = createGeoIpAsnAdapter({ipInfoService: opts.ipInfoService}); const checkMx = createDnsMxChecker({ diff --git a/fluxer_api/src/api/risk/RiskTypes.ts b/fluxer_api/src/api/risk/RiskTypes.ts index 336ba47d6..201800543 100644 --- a/fluxer_api/src/api/risk/RiskTypes.ts +++ b/fluxer_api/src/api/risk/RiskTypes.ts @@ -97,7 +97,6 @@ export interface EmailSyntaxResult { export interface DisposableCheckResult { domain: string; isDisposable: boolean; - listSize: number; } export interface MxCheckResult { diff --git a/fluxer_api/src/api/risk/__tests__/DeterministicRiskEngine.test.ts b/fluxer_api/src/api/risk/__tests__/DeterministicRiskEngine.test.ts index d8439d733..a017edc2b 100644 --- a/fluxer_api/src/api/risk/__tests__/DeterministicRiskEngine.test.ts +++ b/fluxer_api/src/api/risk/__tests__/DeterministicRiskEngine.test.ts @@ -54,7 +54,6 @@ function createToolbox( checkDomainDisposable: async ({domain}) => ({ domain, isDisposable: false, - listSize: 0, }), checkMx: async ({domain}) => ({ domain, diff --git a/fluxer_api/src/api/risk/adapters/DisposableDomainChecker.ts b/fluxer_api/src/api/risk/adapters/DisposableDomainChecker.ts index 7776b0ee4..423039b6d 100644 --- a/fluxer_api/src/api/risk/adapters/DisposableDomainChecker.ts +++ b/fluxer_api/src/api/risk/adapters/DisposableDomainChecker.ts @@ -1,21 +1,22 @@ // SPDX-License-Identifier: AGPL-3.0-or-later +import type {IAdminRepository} from '../../admin/IAdminRepository'; import type {DisposableCheckResult} from '../RiskTypes'; interface DisposableDomainCheckerContext { - disposableDomainsRef: { - readonly current: ReadonlySet; - }; + adminRepository: Pick; } export function createDisposableDomainChecker(ctx: DisposableDomainCheckerContext) { return async function checkDomainDisposable(args: {domain: string}): Promise { const domain = args.domain.toLowerCase().trim(); - const set = ctx.disposableDomainsRef.current; + const [suspicious, disposable] = await Promise.all([ + ctx.adminRepository.isEmailDomainSuspicious(domain), + ctx.adminRepository.isEmailDomainDisposable(domain), + ]); return { domain, - isDisposable: set.has(domain), - listSize: set.size, + isDisposable: suspicious || disposable, }; }; } diff --git a/fluxer_api/src/api/worker/CronScheduler.ts b/fluxer_api/src/api/worker/CronScheduler.ts index 8a7038404..08ae96d9e 100644 --- a/fluxer_api/src/api/worker/CronScheduler.ts +++ b/fluxer_api/src/api/worker/CronScheduler.ts @@ -11,6 +11,7 @@ interface CronDefinition { taskType: WorkerTaskName; payload: WorkerJobPayload; cronExpression: string; + ledger: boolean; lastFired: number; } @@ -89,12 +90,19 @@ export class CronScheduler { this.kvClient = kvClient; } - upsert(id: string, taskType: WorkerTaskName, payload: WorkerJobPayload, cronExpression: string): void { + upsert( + id: string, + taskType: WorkerTaskName, + payload: WorkerJobPayload, + cronExpression: string, + options: {ledger: boolean}, + ): void { this.definitions.set(id, { id, taskType, payload, cronExpression, + ledger: options.ledger, lastFired: 0, }); } @@ -133,7 +141,7 @@ export class CronScheduler { if (!acquired) { continue; } - await this.workerService.addJob(def.taskType, def.payload, {jobKey}); + await this.workerService.addJob(def.taskType, def.payload, {jobKey, skipLedger: !def.ledger}); this.logger.debug({cronId: def.id, taskType: def.taskType}, 'Cron job fired'); } catch (error) { this.logger.error({err: error, cronId: def.id, taskType: def.taskType}, 'Failed to enqueue cron job'); diff --git a/fluxer_api/src/api/worker/WorkerMain.ts b/fluxer_api/src/api/worker/WorkerMain.ts index 3bc2c65d2..b42f1c8c8 100644 --- a/fluxer_api/src/api/worker/WorkerMain.ts +++ b/fluxer_api/src/api/worker/WorkerMain.ts @@ -41,21 +41,27 @@ const SEARCH_REQUIRED_TASKS = new Set([ ]); function registerCronJobs(cron: CronScheduler): void { - cron.upsert('processAssetDeletionQueue', 'processAssetDeletionQueue', {}, '0 */5 * * * *'); - cron.upsert('processBunnyPurgeQueue', 'processBunnyPurgeQueue', {}, '*/10 * * * * *'); - cron.upsert('processPendingBulkMessageDeletions', 'processPendingBulkMessageDeletions', {}, '0 */10 * * * *'); - cron.upsert('userProcessPendingDeletions', 'userProcessPendingDeletions', {}, '0 * * * * *'); - cron.upsert('processPremiumStateReconciliationQueue', 'processPremiumStateReconciliationQueue', {}, '0 * * * * *'); - cron.upsert('processExpiredPremiumSweep', 'processExpiredPremiumSweep', {}, '0 0 * * * *'); - cron.upsert('processInactivityDeletions', 'processInactivityDeletions', {}, '0 0 */6 * * *'); - cron.upsert('expireAttachments', 'expireAttachments', {}, '0 0 */12 * * *'); - cron.upsert('prunePostgresKvTtl', 'prunePostgresKvTtl', {}, '0 */5 * * * *'); - cron.upsert('syncDiscoveryIndex', 'syncDiscoveryIndex', {}, '0 */15 * * * *'); - cron.upsert('syncDisposableEmailDomains', 'syncDisposableEmailDomains', {}, '0 */30 * * * *'); - cron.upsert('syncUrlBlocklists', 'syncUrlBlocklists', {}, '0 0 */6 * * *'); - cron.upsert('syncFileShaBlocklists', 'syncFileShaBlocklists', {}, '0 0 */12 * * *'); - cron.upsert('flushUserActivityBuffer', 'flushUserActivityBuffer', {}, '*/10 * * * * *'); - Logger.info('Cron jobs registered successfully'); + cron.upsert('processAssetDeletionQueue', 'processAssetDeletionQueue', {}, '0 */5 * * * *', {ledger: false}); + cron.upsert('processBunnyPurgeQueue', 'processBunnyPurgeQueue', {}, '*/10 * * * * *', {ledger: false}); + cron.upsert('processPendingBulkMessageDeletions', 'processPendingBulkMessageDeletions', {}, '0 */10 * * * *', { + ledger: false, + }); + cron.upsert('userProcessPendingDeletions', 'userProcessPendingDeletions', {}, '0 * * * * *', {ledger: false}); + cron.upsert('processPremiumStateReconciliationQueue', 'processPremiumStateReconciliationQueue', {}, '0 * * * * *', { + ledger: false, + }); + cron.upsert('processExpiredPremiumSweep', 'processExpiredPremiumSweep', {}, '0 0 * * * *', {ledger: false}); + cron.upsert('processInactivityDeletions', 'processInactivityDeletions', {}, '0 0 */6 * * *', {ledger: false}); + cron.upsert('expireAttachments', 'expireAttachments', {}, '0 0 */12 * * *', {ledger: false}); + cron.upsert('prunePostgresKvTtl', 'prunePostgresKvTtl', {}, '0 */5 * * * *', {ledger: false}); + cron.upsert('syncDiscoveryIndex', 'syncDiscoveryIndex', {}, '0 */15 * * * *', {ledger: false}); + if (Config.blocklistFeeds.enabled) { + cron.upsert('syncDisposableEmailDomains', 'syncDisposableEmailDomains', {}, '0 0 */6 * * *', {ledger: true}); + cron.upsert('syncUrlBlocklists', 'syncUrlBlocklists', {}, '0 0 */6 * * *', {ledger: true}); + cron.upsert('syncFileShaBlocklists', 'syncFileShaBlocklists', {}, '0 0 */12 * * *', {ledger: true}); + } + cron.upsert('flushUserActivityBuffer', 'flushUserActivityBuffer', {}, '*/10 * * * * *', {ledger: false}); + Logger.info({blocklistFeeds: Config.blocklistFeeds.enabled}, 'Cron jobs registered successfully'); } function workerLanesRequireSearch(activeWorkerLanes: ReadonlyArray): boolean { @@ -191,10 +197,16 @@ export async function startWorkerMain(): Promise { setInjectedWorkerService(workerService); dependencies = await initializeWorkerDependencies(snowflakeService); setWorkerDependencies(dependencies); - const didClaimEmailSync = await dependencies.kvClient.setnx('sync:email_domains:initialized', '1'); - if (didClaimEmailSync) { - Logger.info('Triggering initial disposable email domain sync'); - await workerService.addJob('syncDisposableEmailDomains', {}); + if (Config.blocklistFeeds.enabled) { + const didClaimEmailSync = await dependencies.kvClient.setnx( + 'sync:email_domains:initialized', + '1', + ms('6 hours') / 1000, + ); + if (didClaimEmailSync) { + Logger.info('Triggering initial disposable email domain sync'); + await workerService.addJob('syncDisposableEmailDomains', {}); + } } cron = new CronScheduler(workerService, Logger, dependencies.kvClient); registerCronJobs(cron); diff --git a/fluxer_api/src/api/worker/tasks/SyncDisposableEmailDomains.ts b/fluxer_api/src/api/worker/tasks/SyncDisposableEmailDomains.ts index f6321108b..b843a594c 100644 --- a/fluxer_api/src/api/worker/tasks/SyncDisposableEmailDomains.ts +++ b/fluxer_api/src/api/worker/tasks/SyncDisposableEmailDomains.ts @@ -21,7 +21,6 @@ const SOURCES = [ 'https://raw.githubusercontent.com/vrittech/disposable-email/main/disposable_domains.txt', 'https://raw.githubusercontent.com/martenson/disposable-email-domains/master/disposable_email_blocklist.conf', ]; -const CURRENT_DOMAIN_PAGE_SIZE = 10000; const WRITE_PROGRESS_INTERVAL = 500; const DOMAIN_REGEX = /^[a-z0-9]([a-z0-9-]*[a-z0-9])?(\.[a-z0-9]([a-z0-9-]*[a-z0-9])?)+$/; @@ -120,16 +119,7 @@ function normaliseDomain(raw: string): string | null { async function loadCurrentDisposableEmailDomains(): Promise> { const {adminRepository} = getWorkerDependencies(); - const currentSet = new Set(); - let pageState: string | null = null; - do { - const page = await adminRepository.listDisposableEmailDomainsPage(CURRENT_DOMAIN_PAGE_SIZE, pageState); - for (const domain of page.domains) { - currentSet.add(domain); - } - pageState = page.pageState; - } while (pageState !== null); - return currentSet; + return new Set(await adminRepository.listDisposableEmailDomains()); } async function throwIfCancelled(helpers: WorkerTaskHelpers): Promise { diff --git a/packages/config/src/ConfigLoader.ts b/packages/config/src/ConfigLoader.ts index 761b4ee9e..6cc7838e7 100644 --- a/packages/config/src/ConfigLoader.ts +++ b/packages/config/src/ConfigLoader.ts @@ -231,6 +231,7 @@ function defaultConfig(): MasterConfig { api_key: '', pull_zone_id: 0, }, + blocklist_feeds: {}, risk_integration: { enabled: false, ipinfo_api_key: '', diff --git a/packages/config/src/MasterConfig.ts b/packages/config/src/MasterConfig.ts index f02b4c378..fa4665d35 100644 --- a/packages/config/src/MasterConfig.ts +++ b/packages/config/src/MasterConfig.ts @@ -271,6 +271,9 @@ export interface MasterConfig { api_key: string; pull_zone_id: number; }; + blocklist_feeds: { + enabled?: boolean; + }; risk_integration: { enabled: boolean; ipinfo_api_key: string; diff --git a/packages/config/src/config_loader/EnvironmentOverrides.ts b/packages/config/src/config_loader/EnvironmentOverrides.ts index 63e3b7611..eca8943ba 100644 --- a/packages/config/src/config_loader/EnvironmentOverrides.ts +++ b/packages/config/src/config_loader/EnvironmentOverrides.ts @@ -282,6 +282,7 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record = { FLUXER_KLIPY_API_KEY: {path: ['integrations', 'klipy', 'api_key']}, FLUXER_YOUTUBE_API_KEY: {path: ['integrations', 'youtube', 'api_key']}, FLUXER_BUNNY_PURGE_ENABLED: {path: ['integrations', 'bunny', 'purge_enabled'], parse: parseEnvValue}, + FLUXER_BLOCKLIST_FEEDS_ENABLED: {path: ['integrations', 'blocklist_feeds', 'enabled'], parse: parseEnvValue}, FLUXER_BUNNY_API_KEY: {path: ['integrations', 'bunny', 'api_key']}, FLUXER_BUNNY_PULL_ZONE_ID: {path: ['integrations', 'bunny', 'pull_zone_id'], parse: parseEnvValue}, FLUXER_RISK_INTEGRATION_ENABLED: {path: ['integrations', 'risk_integration', 'enabled'], parse: parseEnvValue},