fix(api): eliminate idle Postgres I/O on self-hosted instances (#1593)

This commit is contained in:
Hampus
2026-08-14 19:55:53 +02:00
committed by GitHub
parent 8a9b12e6a1
commit beb906753f
19 changed files with 80 additions and 268 deletions
+3
View File
@@ -276,6 +276,9 @@ export function buildAPIConfigFromMaster(master: MasterConfig): APIConfig {
ipinfoApiKey: master.integrations.risk_integration.ipinfo_api_key || undefined, ipinfoApiKey: master.integrations.risk_integration.ipinfo_api_key || undefined,
accountPolicyDsl: master.integrations.risk_integration.account_policy_dsl, accountPolicyDsl: master.integrations.risk_integration.account_policy_dsl,
}, },
blocklistFeeds: {
enabled: master.integrations.blocklist_feeds.enabled ?? !master.instance.self_hosted,
},
captcha: { captcha: {
enabled: master.integrations.captcha.enabled, enabled: master.integrations.captcha.enabled,
provider: master.integrations.captcha.provider, provider: master.integrations.captcha.provider,
+2 -30
View File
@@ -2,7 +2,7 @@
import {getSameIpDecisionKey} from '@fluxer/ip_utils/src/IpAddress'; import {getSameIpDecisionKey} from '@fluxer/ip_utils/src/IpAddress';
import {createUserID} from '../BrandedTypes'; import {createUserID} from '../BrandedTypes';
import {deleteOneOrMany, fetchMany, fetchOne, fetchPage, upsertOne} from '../database/CassandraQueryExecution'; import {deleteOneOrMany, fetchMany, fetchOne, upsertOne} from '../database/CassandraQueryExecution';
import type { import type {
AdminAuditLogRow, AdminAuditLogRow,
BannedAvatarHashRow, BannedAvatarHashRow,
@@ -30,13 +30,7 @@ import {
} from '../Tables'; } from '../Tables';
import {parseIpBanEntry, tryParseSingleIp} from '../utils/IpRangeUtils'; import {parseIpBanEntry, tryParseSingleIp} from '../utils/IpRangeUtils';
import {canonicalizeStoredPhrase} from '../utils/PhraseBlocklistNormalization'; import {canonicalizeStoredPhrase} from '../utils/PhraseBlocklistNormalization';
import type { import type {AdminAuditLog, BannedIpEntry, BannedIpKind, IAdminRepository} from './IAdminRepository';
AdminAuditLog,
BannedIpEntry,
BannedIpKind,
DisposableEmailDomainPage,
IAdminRepository,
} from './IAdminRepository';
const FETCH_AUDIT_LOG_BY_ID_QUERY = AdminAuditLogs.select({ const FETCH_AUDIT_LOG_BY_ID_QUERY = AdminAuditLogs.select({
where: AdminAuditLogs.where.eq('log_id'), 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({ const IS_EMAIL_DOMAIN_SUSPICIOUS_QUERY = SuspiciousEmailDomains.select({
where: SuspiciousEmailDomains.where.eq('domain'), where: SuspiciousEmailDomains.where.eq('domain'),
}); });
const createLoadSuspiciousEmailDomainsQuery = (limit?: number) =>
limit ? SuspiciousEmailDomains.select({limit}) : SuspiciousEmailDomains.select();
const IS_EMAIL_DOMAIN_DISPOSABLE_QUERY = DisposableEmailDomains.select({ const IS_EMAIL_DOMAIN_DISPOSABLE_QUERY = DisposableEmailDomains.select({
where: DisposableEmailDomains.where.eq('domain'), where: DisposableEmailDomains.where.eq('domain'),
}); });
@@ -273,13 +265,6 @@ export class AdminRepository implements IAdminRepository {
await deleteOneOrMany(SuspiciousEmailDomains.deleteByPk({domain: domainLower})); await deleteOneOrMany(SuspiciousEmailDomains.deleteByPk({domain: domainLower}));
} }
async listSuspiciousEmailDomains(limit?: number): Promise<Array<string>> {
const rows = await fetchMany<{
domain: string;
}>(createLoadSuspiciousEmailDomainsQuery(limit).bind({}));
return rows.map((row) => row.domain);
}
async isEmailDomainDisposable(domain: string): Promise<boolean> { async isEmailDomainDisposable(domain: string): Promise<boolean> {
const domainLower = domain.toLowerCase(); const domainLower = domain.toLowerCase();
if (isAccountPolicyContactDomainReputationExempt(domainLower)) return false; if (isAccountPolicyContactDomainReputationExempt(domainLower)) return false;
@@ -306,19 +291,6 @@ export class AdminRepository implements IAdminRepository {
return rows.map((row) => row.domain); return rows.map((row) => row.domain);
} }
async listDisposableEmailDomainsPage(limit: number, pageState?: string | null): Promise<DisposableEmailDomainPage> {
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<boolean> { async isPhraseBanned(phrase: string): Promise<boolean> {
const phraseLower = canonicalizeStoredPhrase(phrase); const phraseLower = canonicalizeStoredPhrase(phrase);
const result = await fetchOne<{ const result = await fetchOne<{
@@ -32,11 +32,6 @@ export interface BannedIpEntry {
createdAt: Date | null; createdAt: Date | null;
} }
export interface DisposableEmailDomainPage {
domains: Array<string>;
pageState: string | null;
}
export abstract class IAdminRepository { export abstract class IAdminRepository {
abstract createAuditLog(log: AdminAuditLogRow): Promise<AdminAuditLog>; abstract createAuditLog(log: AdminAuditLogRow): Promise<AdminAuditLog>;
@@ -68,8 +63,6 @@ export abstract class IAdminRepository {
abstract removeSuspiciousEmailDomain(domain: string): Promise<void>; abstract removeSuspiciousEmailDomain(domain: string): Promise<void>;
abstract listSuspiciousEmailDomains(limit?: number): Promise<Array<string>>;
abstract isEmailDomainDisposable(domain: string): Promise<boolean>; abstract isEmailDomainDisposable(domain: string): Promise<boolean>;
abstract addDisposableEmailDomain(domain: string): Promise<void>; abstract addDisposableEmailDomain(domain: string): Promise<void>;
@@ -78,8 +71,6 @@ export abstract class IAdminRepository {
abstract listDisposableEmailDomains(limit?: number): Promise<Array<string>>; abstract listDisposableEmailDomains(limit?: number): Promise<Array<string>>;
abstract listDisposableEmailDomainsPage(limit: number, pageState?: string | null): Promise<DisposableEmailDomainPage>;
abstract isPhraseBanned(phrase: string): Promise<boolean>; abstract isPhraseBanned(phrase: string): Promise<boolean>;
abstract banPhrase(phrase: string): Promise<void>; abstract banPhrase(phrase: string): Promise<void>;
+1 -64
View File
@@ -13,11 +13,7 @@ import type {ILogger} from '../ILogger';
import {JobLedgerRepository} from '../jobs/JobLedgerRepository'; import {JobLedgerRepository} from '../jobs/JobLedgerRepository';
import {startAbuseReplicationSubscriber, stopAbuseReplicationSubscriber} from '../middleware/AbusiveIpAutoBanner'; import {startAbuseReplicationSubscriber, stopAbuseReplicationSubscriber} from '../middleware/AbusiveIpAutoBanner';
import {ipBanCache} from '../middleware/IpBanMiddleware'; import {ipBanCache} from '../middleware/IpBanMiddleware';
import { import {initializeServiceSingletons, shutdownReportService} from '../middleware/ServiceMiddleware';
getRiskCacheManagerInstance,
initializeServiceSingletons,
shutdownReportService,
} from '../middleware/ServiceMiddleware';
import { import {
ensureVoiceResourcesInitialized, ensureVoiceResourcesInitialized,
getKVClient, getKVClient,
@@ -40,57 +36,6 @@ import {JetStreamWorkerQueue} from '../worker/JetStreamWorkerQueue';
import {WorkerService} from '../worker/WorkerService'; import {WorkerService} from '../worker/WorkerService';
let jsConnectionManager: JetStreamConnectionManager | null = null; 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<void> {
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<void> { export function createInitializer(config: APIConfig, logger: ILogger): () => Promise<void> {
return async (): Promise<void> => { return async (): Promise<void> => {
try { try {
@@ -171,8 +116,6 @@ export function createInitializer(config: APIConfig, logger: ILogger): () => Pro
logger.info('Profile substring blocklist cache initialized'); logger.info('Profile substring blocklist cache initialized');
await initializeServiceSingletons(); await initializeServiceSingletons();
logger.info('Service singletons initialized'); logger.info('Service singletons initialized');
await refreshRiskCache(logger, 'startup');
startRiskCacheRefreshLoop(logger);
if (!config.dev.testModeEnabled) { if (!config.dev.testModeEnabled) {
jsConnectionManager = new JetStreamConnectionManager({ jsConnectionManager = new JetStreamConnectionManager({
url: config.nats.jetStreamUrl, url: config.nats.jetStreamUrl,
@@ -285,12 +228,6 @@ export function createShutdown(logger: ILogger): () => Promise<void> {
} catch (error) { } catch (error) {
logger.error({error}, 'Error shutting down search service'); 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 { try {
ipBanCache.shutdown(); ipBanCache.shutdown();
logger.info('IP ban cache shut down'); logger.info('IP ban cache shut down');
+3
View File
@@ -178,6 +178,9 @@ export interface APIConfig {
ipinfoApiKey?: string; ipinfoApiKey?: string;
accountPolicyDsl?: unknown; accountPolicyDsl?: unknown;
}; };
blocklistFeeds: {
enabled: boolean;
};
captcha: { captcha: {
enabled: boolean; enabled: boolean;
provider: 'hcaptcha' | 'turnstile' | 'none'; provider: 'hcaptcha' | 'turnstile' | 'none';
@@ -561,14 +561,14 @@ WHERE NOT $6`,
private async patch(meta: KvQueryMeta, params: CassandraParams, db: PostgresQueryable): Promise<void> { private async patch(meta: KvQueryMeta, params: CassandraParams, db: PostgresQueryable): Promise<void> {
const key = rowKeyFromParams(meta, params); const key = rowKeyFromParams(meta, params);
const existing = await this.getRow(meta, key, db); const stored = await this.getStoredRow(meta, key, db);
const base = existing ?? paramsRow(params, (meta.pkColumns ?? meta.table.primaryKey) as ReadonlyArray<string>); const base = stored?.row ?? paramsRow(params, (meta.pkColumns ?? meta.table.primaryKey) as ReadonlyArray<string>);
const next = {...base}; const next = {...base};
for (const column of meta.patchKeys ?? []) { for (const column of meta.patchKeys ?? []) {
next[column] = column in params ? params[column] : null; next[column] = column in params ? params[column] : null;
} }
const ttl = ttlExpiresAt(meta, params); 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( await db.query(
`INSERT INTO ${this.table} (table_name, partition_key, row_key, row_data, expires_at, updated_at) `INSERT INTO ${this.table} (table_name, partition_key, row_key, row_data, expires_at, updated_at)
VALUES ($1, $2, $3, $4::jsonb, $5, now()) 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<Row | null> { private async getStoredRow(
const result = await db.query<StoredRow>( meta: KvQueryMeta,
`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`, key: string,
db: PostgresQueryable,
): Promise<{row: Row; expiresAt: Date | null} | null> {
const result = await db.query<StoredRow & {expires_at: Date | null}>(
`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], [meta.table.name, key],
); );
const row = result.rows[0]; 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<Date | null> { private async getRow(meta: KvQueryMeta, key: string, db: PostgresQueryable): Promise<Row | null> {
const result = await db.query<{expires_at: Date | null}>( return (await this.getStoredRow(meta, key, db))?.row ?? 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;
} }
} }
@@ -60,8 +60,6 @@ import {CassandraHistoricalOutcomeRepository} from '../risk/HistoricalOutcomeRep
import {buildIpInfoCache, buildIpInfoRequestAuditLogger} from '../risk/IpInfoCacheFactory'; import {buildIpInfoCache, buildIpInfoRequestAuditLogger} from '../risk/IpInfoCacheFactory';
import {CassandraRegistrationEventsRepository} from '../risk/RegistrationEventsRepository'; import {CassandraRegistrationEventsRepository} from '../risk/RegistrationEventsRepository';
import {CassandraRiskAssessmentRepository} from '../risk/RiskAssessmentRepository'; import {CassandraRiskAssessmentRepository} from '../risk/RiskAssessmentRepository';
import {buildRiskCacheLoaders} from '../risk/RiskCacheLoaders';
import {RiskCacheManager} from '../risk/RiskCacheManager';
import {createRiskToolbox} from '../risk/RiskToolboxFactory'; import {createRiskToolbox} from '../risk/RiskToolboxFactory';
import {CassandraSuspiciousIpRepository} from '../risk/SuspiciousIpRepository'; import {CassandraSuspiciousIpRepository} from '../risk/SuspiciousIpRepository';
import {RpcService} from '../rpc/RpcService'; 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; let _inboundSmsChallengeService: InboundSmsChallengeService | null = null;
function getInboundSmsChallengeService(): InboundSmsChallengeService { function getInboundSmsChallengeService(): InboundSmsChallengeService {
@@ -311,13 +291,12 @@ function getRegistrationRiskEvaluator(): IRegistrationRiskEvaluator {
_registrationRiskEvaluator = noopRegistrationRiskEvaluator; _registrationRiskEvaluator = noopRegistrationRiskEvaluator;
return _registrationRiskEvaluator; return _registrationRiskEvaluator;
} }
const cacheManager = getRiskCacheManager();
const ipInfoService = getIpInfoService(); const ipInfoService = getIpInfoService();
const ipInfoChecker = Config.risk.ipinfoApiKey ? createIpInfoChecker({ipInfoService}) : undefined; const ipInfoChecker = Config.risk.ipinfoApiKey ? createIpInfoChecker({ipInfoService}) : undefined;
const cacheService = getCacheService(); const cacheService = getCacheService();
const reverseDnsLookup = createReverseDnsLookup({cacheService}); const reverseDnsLookup = createReverseDnsLookup({cacheService});
const toolbox = createRiskToolbox({ const toolbox = createRiskToolbox({
disposableDomainsRef: cacheManager.disposableDomainsRef, adminRepository: getAdminRepository(),
ipInfoChecker, ipInfoChecker,
reverseDnsLookup, reverseDnsLookup,
ipInfoService, ipInfoService,
@@ -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<IAdminRepository, 'listSuspiciousEmailDomains' | 'listDisposableEmailDomains'>;
}
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<string>();
for (const d of suspicious) set.add(d.toLowerCase());
for (const d of disposable) set.add(d.toLowerCase());
return set;
},
};
}
@@ -1,64 +0,0 @@
// SPDX-License-Identifier: AGPL-3.0-or-later
interface MutableRef<T> {
current: T;
}
export type ReadonlyRiskCacheRef<T> = {
readonly current: T;
};
interface RiskCacheManagerLogger {
info(payload: object, msg: string): void;
warn(payload: object, msg: string): void;
}
export interface RiskCacheLoaders {
loadDisposableDomains: () => Promise<ReadonlySet<string>>;
}
interface RiskCacheManagerOptions extends RiskCacheLoaders {
logger?: RiskCacheManagerLogger;
}
interface RiskCacheRefreshResult {
disposableDomainCount: number;
subtaskErrors: ReadonlyArray<{
step: string;
error: string;
}>;
}
export class RiskCacheManager {
private readonly _disposable: MutableRef<ReadonlySet<string>> = {current: new Set()};
readonly disposableDomainsRef: ReadonlyRiskCacheRef<ReadonlySet<string>> = 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<RiskCacheRefreshResult> {
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,
};
}
}
@@ -2,6 +2,7 @@
import type {ICacheService} from '@pkgs/cache/src/ICacheService'; import type {ICacheService} from '@pkgs/cache/src/ICacheService';
import type {IpInfoService} from '@pkgs/geoip/src/IpInfoService'; import type {IpInfoService} from '@pkgs/geoip/src/IpInfoService';
import type {IAdminRepository} from '../admin/IAdminRepository';
import {createDisposableDomainChecker} from './adapters/DisposableDomainChecker'; import {createDisposableDomainChecker} from './adapters/DisposableDomainChecker';
import {createDnsMxChecker, type MxResolver, NodeDnsMxResolver} from './adapters/DnsMxChecker'; import {createDnsMxChecker, type MxResolver, NodeDnsMxResolver} from './adapters/DnsMxChecker';
import {createDomainAgeChecker} from './adapters/DomainAgeChecker'; import {createDomainAgeChecker} from './adapters/DomainAgeChecker';
@@ -13,13 +14,12 @@ import {analyzeRegistrationTiming} from './adapters/RegistrationTimingAnalyzer';
import {analyzeUserAgent} from './adapters/UserAgentAnalyzer'; import {analyzeUserAgent} from './adapters/UserAgentAnalyzer';
import {createVelocityAdapter, type IRegistrationEventsRepository} from './adapters/VelocityAdapter'; import {createVelocityAdapter, type IRegistrationEventsRepository} from './adapters/VelocityAdapter';
import type {IRiskHistoryRepository} from './HistoricalOutcomeRepository'; import type {IRiskHistoryRepository} from './HistoricalOutcomeRepository';
import type {ReadonlyRiskCacheRef} from './RiskCacheManager';
import type {RiskToolbox} from './RiskToolbox'; import type {RiskToolbox} from './RiskToolbox';
import type {IpInfoAnonymousResult, ReverseDnsResult} from './RiskTypes'; import type {IpInfoAnonymousResult, ReverseDnsResult} from './RiskTypes';
import type {ISuspiciousIpRepository} from './SuspiciousIpRepository'; import type {ISuspiciousIpRepository} from './SuspiciousIpRepository';
interface RiskToolboxFactoryOptions { interface RiskToolboxFactoryOptions {
disposableDomainsRef: ReadonlyRiskCacheRef<ReadonlySet<string>>; adminRepository: Pick<IAdminRepository, 'isEmailDomainSuspicious' | 'isEmailDomainDisposable'>;
ipInfoChecker?: (ip: string) => Promise<IpInfoAnonymousResult>; ipInfoChecker?: (ip: string) => Promise<IpInfoAnonymousResult>;
reverseDnsLookup?: (ip: string) => Promise<ReverseDnsResult>; reverseDnsLookup?: (ip: string) => Promise<ReverseDnsResult>;
ipInfoService: IpInfoService; ipInfoService: IpInfoService;
@@ -32,7 +32,7 @@ interface RiskToolboxFactoryOptions {
} }
export function createRiskToolbox(opts: RiskToolboxFactoryOptions): RiskToolbox { 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 lookupGeoIpCity = createGeoIpCityAdapter({ipInfoService: opts.ipInfoService});
const lookupGeoIpAsn = createGeoIpAsnAdapter({ipInfoService: opts.ipInfoService}); const lookupGeoIpAsn = createGeoIpAsnAdapter({ipInfoService: opts.ipInfoService});
const checkMx = createDnsMxChecker({ const checkMx = createDnsMxChecker({
-1
View File
@@ -97,7 +97,6 @@ export interface EmailSyntaxResult {
export interface DisposableCheckResult { export interface DisposableCheckResult {
domain: string; domain: string;
isDisposable: boolean; isDisposable: boolean;
listSize: number;
} }
export interface MxCheckResult { export interface MxCheckResult {
@@ -54,7 +54,6 @@ function createToolbox(
checkDomainDisposable: async ({domain}) => ({ checkDomainDisposable: async ({domain}) => ({
domain, domain,
isDisposable: false, isDisposable: false,
listSize: 0,
}), }),
checkMx: async ({domain}) => ({ checkMx: async ({domain}) => ({
domain, domain,
@@ -1,21 +1,22 @@
// SPDX-License-Identifier: AGPL-3.0-or-later // SPDX-License-Identifier: AGPL-3.0-or-later
import type {IAdminRepository} from '../../admin/IAdminRepository';
import type {DisposableCheckResult} from '../RiskTypes'; import type {DisposableCheckResult} from '../RiskTypes';
interface DisposableDomainCheckerContext { interface DisposableDomainCheckerContext {
disposableDomainsRef: { adminRepository: Pick<IAdminRepository, 'isEmailDomainSuspicious' | 'isEmailDomainDisposable'>;
readonly current: ReadonlySet<string>;
};
} }
export function createDisposableDomainChecker(ctx: DisposableDomainCheckerContext) { export function createDisposableDomainChecker(ctx: DisposableDomainCheckerContext) {
return async function checkDomainDisposable(args: {domain: string}): Promise<DisposableCheckResult> { return async function checkDomainDisposable(args: {domain: string}): Promise<DisposableCheckResult> {
const domain = args.domain.toLowerCase().trim(); 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 { return {
domain, domain,
isDisposable: set.has(domain), isDisposable: suspicious || disposable,
listSize: set.size,
}; };
}; };
} }
+10 -2
View File
@@ -11,6 +11,7 @@ interface CronDefinition {
taskType: WorkerTaskName; taskType: WorkerTaskName;
payload: WorkerJobPayload; payload: WorkerJobPayload;
cronExpression: string; cronExpression: string;
ledger: boolean;
lastFired: number; lastFired: number;
} }
@@ -89,12 +90,19 @@ export class CronScheduler {
this.kvClient = kvClient; 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, { this.definitions.set(id, {
id, id,
taskType, taskType,
payload, payload,
cronExpression, cronExpression,
ledger: options.ledger,
lastFired: 0, lastFired: 0,
}); });
} }
@@ -133,7 +141,7 @@ export class CronScheduler {
if (!acquired) { if (!acquired) {
continue; 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'); this.logger.debug({cronId: def.id, taskType: def.taskType}, 'Cron job fired');
} catch (error) { } catch (error) {
this.logger.error({err: error, cronId: def.id, taskType: def.taskType}, 'Failed to enqueue cron job'); this.logger.error({err: error, cronId: def.id, taskType: def.taskType}, 'Failed to enqueue cron job');
+31 -19
View File
@@ -41,21 +41,27 @@ const SEARCH_REQUIRED_TASKS = new Set<string>([
]); ]);
function registerCronJobs(cron: CronScheduler): void { function registerCronJobs(cron: CronScheduler): void {
cron.upsert('processAssetDeletionQueue', 'processAssetDeletionQueue', {}, '0 */5 * * * *'); cron.upsert('processAssetDeletionQueue', 'processAssetDeletionQueue', {}, '0 */5 * * * *', {ledger: false});
cron.upsert('processBunnyPurgeQueue', 'processBunnyPurgeQueue', {}, '*/10 * * * * *'); cron.upsert('processBunnyPurgeQueue', 'processBunnyPurgeQueue', {}, '*/10 * * * * *', {ledger: false});
cron.upsert('processPendingBulkMessageDeletions', 'processPendingBulkMessageDeletions', {}, '0 */10 * * * *'); cron.upsert('processPendingBulkMessageDeletions', 'processPendingBulkMessageDeletions', {}, '0 */10 * * * *', {
cron.upsert('userProcessPendingDeletions', 'userProcessPendingDeletions', {}, '0 * * * * *'); ledger: false,
cron.upsert('processPremiumStateReconciliationQueue', 'processPremiumStateReconciliationQueue', {}, '0 * * * * *'); });
cron.upsert('processExpiredPremiumSweep', 'processExpiredPremiumSweep', {}, '0 0 * * * *'); cron.upsert('userProcessPendingDeletions', 'userProcessPendingDeletions', {}, '0 * * * * *', {ledger: false});
cron.upsert('processInactivityDeletions', 'processInactivityDeletions', {}, '0 0 */6 * * *'); cron.upsert('processPremiumStateReconciliationQueue', 'processPremiumStateReconciliationQueue', {}, '0 * * * * *', {
cron.upsert('expireAttachments', 'expireAttachments', {}, '0 0 */12 * * *'); ledger: false,
cron.upsert('prunePostgresKvTtl', 'prunePostgresKvTtl', {}, '0 */5 * * * *'); });
cron.upsert('syncDiscoveryIndex', 'syncDiscoveryIndex', {}, '0 */15 * * * *'); cron.upsert('processExpiredPremiumSweep', 'processExpiredPremiumSweep', {}, '0 0 * * * *', {ledger: false});
cron.upsert('syncDisposableEmailDomains', 'syncDisposableEmailDomains', {}, '0 */30 * * * *'); cron.upsert('processInactivityDeletions', 'processInactivityDeletions', {}, '0 0 */6 * * *', {ledger: false});
cron.upsert('syncUrlBlocklists', 'syncUrlBlocklists', {}, '0 0 */6 * * *'); cron.upsert('expireAttachments', 'expireAttachments', {}, '0 0 */12 * * *', {ledger: false});
cron.upsert('syncFileShaBlocklists', 'syncFileShaBlocklists', {}, '0 0 */12 * * *'); cron.upsert('prunePostgresKvTtl', 'prunePostgresKvTtl', {}, '0 */5 * * * *', {ledger: false});
cron.upsert('flushUserActivityBuffer', 'flushUserActivityBuffer', {}, '*/10 * * * * *'); cron.upsert('syncDiscoveryIndex', 'syncDiscoveryIndex', {}, '0 */15 * * * *', {ledger: false});
Logger.info('Cron jobs registered successfully'); 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<WorkerLaneDefinition>): boolean { function workerLanesRequireSearch(activeWorkerLanes: ReadonlyArray<WorkerLaneDefinition>): boolean {
@@ -191,10 +197,16 @@ export async function startWorkerMain(): Promise<void> {
setInjectedWorkerService(workerService); setInjectedWorkerService(workerService);
dependencies = await initializeWorkerDependencies(snowflakeService); dependencies = await initializeWorkerDependencies(snowflakeService);
setWorkerDependencies(dependencies); setWorkerDependencies(dependencies);
const didClaimEmailSync = await dependencies.kvClient.setnx('sync:email_domains:initialized', '1'); if (Config.blocklistFeeds.enabled) {
if (didClaimEmailSync) { const didClaimEmailSync = await dependencies.kvClient.setnx(
Logger.info('Triggering initial disposable email domain sync'); 'sync:email_domains:initialized',
await workerService.addJob('syncDisposableEmailDomains', {}); '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); cron = new CronScheduler(workerService, Logger, dependencies.kvClient);
registerCronJobs(cron); registerCronJobs(cron);
@@ -21,7 +21,6 @@ const SOURCES = [
'https://raw.githubusercontent.com/vrittech/disposable-email/main/disposable_domains.txt', 'https://raw.githubusercontent.com/vrittech/disposable-email/main/disposable_domains.txt',
'https://raw.githubusercontent.com/martenson/disposable-email-domains/master/disposable_email_blocklist.conf', '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 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])?)+$/; 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<Set<string>> { async function loadCurrentDisposableEmailDomains(): Promise<Set<string>> {
const {adminRepository} = getWorkerDependencies(); const {adminRepository} = getWorkerDependencies();
const currentSet = new Set<string>(); return new Set(await adminRepository.listDisposableEmailDomains());
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;
} }
async function throwIfCancelled(helpers: WorkerTaskHelpers): Promise<void> { async function throwIfCancelled(helpers: WorkerTaskHelpers): Promise<void> {
+1
View File
@@ -231,6 +231,7 @@ function defaultConfig(): MasterConfig {
api_key: '', api_key: '',
pull_zone_id: 0, pull_zone_id: 0,
}, },
blocklist_feeds: {},
risk_integration: { risk_integration: {
enabled: false, enabled: false,
ipinfo_api_key: '', ipinfo_api_key: '',
+3
View File
@@ -271,6 +271,9 @@ export interface MasterConfig {
api_key: string; api_key: string;
pull_zone_id: number; pull_zone_id: number;
}; };
blocklist_feeds: {
enabled?: boolean;
};
risk_integration: { risk_integration: {
enabled: boolean; enabled: boolean;
ipinfo_api_key: string; ipinfo_api_key: string;
@@ -282,6 +282,7 @@ const NAMED_FLUXER_ENV_OVERRIDES: Record<string, NamedEnvOverride> = {
FLUXER_KLIPY_API_KEY: {path: ['integrations', 'klipy', 'api_key']}, FLUXER_KLIPY_API_KEY: {path: ['integrations', 'klipy', 'api_key']},
FLUXER_YOUTUBE_API_KEY: {path: ['integrations', 'youtube', 'api_key']}, FLUXER_YOUTUBE_API_KEY: {path: ['integrations', 'youtube', 'api_key']},
FLUXER_BUNNY_PURGE_ENABLED: {path: ['integrations', 'bunny', 'purge_enabled'], parse: parseEnvValue}, 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_API_KEY: {path: ['integrations', 'bunny', 'api_key']},
FLUXER_BUNNY_PULL_ZONE_ID: {path: ['integrations', 'bunny', 'pull_zone_id'], parse: parseEnvValue}, FLUXER_BUNNY_PULL_ZONE_ID: {path: ['integrations', 'bunny', 'pull_zone_id'], parse: parseEnvValue},
FLUXER_RISK_INTEGRATION_ENABLED: {path: ['integrations', 'risk_integration', 'enabled'], parse: parseEnvValue}, FLUXER_RISK_INTEGRATION_ENABLED: {path: ['integrations', 'risk_integration', 'enabled'], parse: parseEnvValue},