From 74c206f515ede9793db8e7cd31a64843c2006d71 Mon Sep 17 00:00:00 2001 From: Juan Mrad Date: Thu, 16 Apr 2026 13:02:43 -0500 Subject: [PATCH] [Kysely] migrate rule-engine queries and related jobs to Kysely (phase 1) (#225) * [Vulnerabilities] Upgrade Kysely to latest * fix lint * code review * [Kysely] migrate rule-engine queries and related jobs to Kysely (phase 1) * fixes * fix lint by organizing errors to a file for simplifications * lint fix again * fix test * [Kysely] Remove knex migrate backtest pagination and takeLast to Kysely (#226) * [Kysely] Remove knex migrate backtest pagination and takeLast to Kysely * code revie fix * simplify enum uses --- server/graphql/datasources/LocationBankApi.ts | 5 +- server/graphql/datasources/RuleApi.ts | 128 +++++++------- server/iocContainer/index.ts | 39 ++--- server/models/OrgModel.ts | 68 -------- server/models/rules/RuleModel.ts | 92 ++-------- server/models/rules/ruleTypes.ts | 124 ++++++++++++++ server/package-lock.json | 121 +------------ server/package.json | 1 - server/plugins/warehouse/IWarehouseAdapter.ts | 2 +- server/rule_engine/ActionPublisher.test.ts | 2 + server/rule_engine/RuleEngine.ts | 14 +- server/rule_engine/ruleEngineQueries.ts | 160 +++++++----------- server/services/combinedDbTypes.ts | 10 +- server/services/coreAppTables.ts | 59 +++++++ .../moderationConfigService/dbTypes.ts | 20 +++ .../moderationConfigService/errors.ts | 49 ++++++ .../services/moderationConfigService/index.ts | 7 + .../moderationConfigService.test.ts | 1 + .../moderationConfigService.ts | 81 +++------ .../modules/ActionOperations.ts | 31 ++++ .../modules/PolicyOperations.ts | 46 ++++- .../modules/RuleReadOperations.ts | 147 ++++++++++++++++ .../moderationConfigService/types/actions.ts | 3 + .../moderationConfigService/types/policies.ts | 4 +- .../signalExecutionService.ts | 2 +- server/services/policyActionPenalties.ts | 73 ++++++++ .../detectRulePassRateAnomaliesJob.test.ts | 141 ++++++++++----- .../detectRulePassRateAnomaliesJob.ts | 90 ++++++---- .../signalsService/signals/SignalBase.ts | 2 +- .../services/userManagementService/dbTypes.ts | 3 +- .../userStatisticsService/computeUserScore.ts | 2 +- .../userStatisticsService.ts | 2 +- .../userStrikeService.test.ts | 3 + server/utils/kyselyTransactionWithRetry.ts | 36 ++++ server/utils/sql.test.ts | 73 +++++--- server/utils/sql.ts | 94 +++++----- server/workers_jobs/RunUserRulesJob.ts | 12 +- 37 files changed, 1056 insertions(+), 691 deletions(-) create mode 100644 server/models/rules/ruleTypes.ts create mode 100644 server/services/coreAppTables.ts create mode 100644 server/services/moderationConfigService/errors.ts create mode 100644 server/services/moderationConfigService/modules/RuleReadOperations.ts create mode 100644 server/services/policyActionPenalties.ts create mode 100644 server/utils/kyselyTransactionWithRetry.ts diff --git a/server/graphql/datasources/LocationBankApi.ts b/server/graphql/datasources/LocationBankApi.ts index eeed5db..65b915e 100644 --- a/server/graphql/datasources/LocationBankApi.ts +++ b/server/graphql/datasources/LocationBankApi.ts @@ -7,10 +7,7 @@ import { type LocationBank as TLocationBank } from '../../models/banks/LocationB import { isUniqueConstraintError } from '../../models/errors.js'; import { type LocationArea } from '../../models/types/locationArea.js'; import { type User } from '../../models/UserModel.js'; -// TODO: delete the import below when we move the location bank mutation logic -// into the moderation config service, which is where it should be. -// eslint-disable-next-line import/no-restricted-paths -import { makeLocationBankNameExistsError } from '../../services/moderationConfigService/moderationConfigService.js'; +import { makeLocationBankNameExistsError } from '../../services/moderationConfigService/index.js'; import { type PlacesApiService } from '../../services/placesApiService/index.js'; import { patchInPlace, safePick } from '../../utils/misc.js'; import { diff --git a/server/graphql/datasources/RuleApi.ts b/server/graphql/datasources/RuleApi.ts index 0cbe6b0..2449e60 100644 --- a/server/graphql/datasources/RuleApi.ts +++ b/server/graphql/datasources/RuleApi.ts @@ -3,7 +3,7 @@ import { type Exception } from '@opentelemetry/api'; import { makeEnumLike } from '@roostorg/types'; -import { sql, type Kysely } from 'kysely'; +import { type Kysely } from 'kysely'; import Sequelize from 'sequelize'; import { uid } from 'uid'; @@ -29,18 +29,19 @@ import { makeRuleHasRunningBacktestsError, makeRuleIsMissingContentTypeError, makeRuleNameExistsError, - // TODO: delete the import below when we move the rule mutation logic into the - // moderation config service, which is where it should be. - // eslint-disable-next-line import/no-restricted-paths -} from '../../services/moderationConfigService/moderationConfigService.js'; +} from '../../services/moderationConfigService/index.js'; import { isSignalId, signalIsExternal, type SignalId, } from '../../services/signalsService/index.js'; -import { type DataWarehousePublicSchema } from '../../storage/dataWarehouse/warehouseSchema.js'; +import { + type DataWarehousePublicSchema, + warehouseDateToDate, +} from '../../storage/dataWarehouse/warehouseSchema.js'; import { toCorrelationId } from '../../utils/correlationIds.js'; import { + type JsonOf, jsonParse, jsonStringify, tryJsonParse, @@ -84,7 +85,7 @@ export type RuleExecutionResult = { userId?: string; userTypeId?: string; content: string; - result: ConditionSetWithResultAsLogged; + result: ConditionSetWithResultAsLogged | null; environment: RuleStatus; passed: boolean; ruleId: string; @@ -270,7 +271,6 @@ class RuleAPI { private readonly warehouse: Kysely; constructor( - private readonly knex: Dependencies['Knex'], dialect: Dependencies['DataWarehouseDialect'], public readonly ruleInsights: Dependencies['RuleActionInsights'], private readonly actionStats: Dependencies['ActionStatisticsService'], @@ -714,45 +714,41 @@ class RuleAPI { // (no cursor, after cursor, before cursor) x (sort asc, desc). // But our pagination helpers let us handle reasonably simply, in steps. // First, we must define the result query if we weren't doing any pagination: - const allResultsQuery = this.knex('RULE_EXECUTIONS') - .select({ - // This select is aliasing each column to the corresponding object key, - // so we have to do fewer renames from the warehouse ALL_CAPS_SNAKE_CASE - // when we return the final result. - date: 'DS', - ts: 'TS', - contentId: 'ITEM_ID', - contentType: 'ITEM_TYPE_NAME', - userId: 'ITEM_CREATOR_ID', - content: 'ITEM_DATA', - result: 'RESULT', - }) - .where( - 'CORRELATION_ID', - toCorrelationId({ type: 'backtest', id: backtestId }), - ); + const correlationId = toCorrelationId({ + type: 'backtest', + id: backtestId, + }); - // Now, we can filter down the results to those that satisfy - // the cursor's before/after requirements, if there is a cursor. - // Note that how we do this filtering depends on how the results are sorted, - // because the sorting conceptually happens "before" pagination, and it - // effects what's "before" and what's "after" a given cursor. - // - // Specifically, if the results are sorted descending and we're looking for - // values _after_ the cursor, then we're looking for timestamp values that - // are less than the cursor. Similarly, if we're sorting ascending and - // looking for items before the cursor, then those items must have ts values - // less than the cursor. In the other cases, it's the opposite. - const filteredResultsQuery = !cursor - ? allResultsQuery - : allResultsQuery.andWhere( - 'TS', - (sortByTs === SortOrder.DESC && cursor.direction === 'after') || - (sortByTs === SortOrder.ASC && cursor.direction === 'before') - ? '<' - : '>', - new Date(cursor.value.ts), - ); + let filteredResultsQuery = this.warehouse + .selectFrom('RULE_EXECUTIONS') + .select([ + 'DS as date', + 'TS as ts', + 'ITEM_ID as contentId', + 'ITEM_TYPE_NAME as itemTypeName', + 'ITEM_TYPE_ID as itemTypeId', + 'ITEM_CREATOR_ID as userId', + 'ITEM_CREATOR_TYPE_ID as userTypeId', + 'ITEM_DATA as content', + 'RESULT as result', + 'ENVIRONMENT as environment', + 'PASSED as passed', + 'RULE_ID as ruleId', + 'RULE as ruleName', + 'TAGS as tags', + ]) + .where('CORRELATION_ID', '=', correlationId); + + if (cursor) { + filteredResultsQuery = filteredResultsQuery.where( + 'TS', + (sortByTs === SortOrder.DESC && cursor.direction === 'after') || + (sortByTs === SortOrder.ASC && cursor.direction === 'before') + ? '<' + : '>', + new Date(cursor.value.ts), + ); + } const desiredSort = { column: 'ts', @@ -766,16 +762,35 @@ class RuleAPI { // have to use our helper that implements "takeLast" in SQL. const finalQuery = takeFrom === 'start' - ? filteredResultsQuery.orderBy([desiredSort]).limit(count) - : takeLast(filteredResultsQuery, [desiredSort], count); - - const results = ( - await sql`${sql.raw(finalQuery.toString())}`.execute(this.warehouse) - ).rows; - - return results.map((it: any) => ({ - node: { ...it, result: it.result ? jsonParse(it.result) : null }, - cursor: { ts: new Date(it.ts).valueOf() }, + ? filteredResultsQuery + .orderBy('ts', desiredSort.order) + .limit(count) + : takeLast(this.warehouse, filteredResultsQuery, [desiredSort], count); + + const results = await finalQuery.execute(); + + return results.map>((it) => ({ + node: { + date: warehouseDateToDate(it.date).toISOString(), + ts: warehouseDateToDate(it.ts).toISOString(), + contentId: it.contentId, + itemTypeName: it.itemTypeName ?? '', + itemTypeId: it.itemTypeId, + userId: it.userId ?? undefined, + userTypeId: it.userTypeId ?? undefined, + content: (it.content ?? '') as string, + result: it.result + ? jsonParse( + it.result as JsonOf, + ) + : null, + environment: it.environment as RuleStatus, + passed: it.passed, + ruleId: it.ruleId, + ruleName: it.ruleName ?? '', + tags: [...it.tags], + }, + cursor: { ts: warehouseDateToDate(it.ts).valueOf() }, })); } @@ -910,7 +925,6 @@ class RuleAPI { export default inject( [ - 'Knex', 'DataWarehouseDialect', 'RuleActionInsights', 'ActionStatisticsService', diff --git a/server/iocContainer/index.ts b/server/iocContainer/index.ts index dacaf8a..0035885 100644 --- a/server/iocContainer/index.ts +++ b/server/iocContainer/index.ts @@ -5,8 +5,6 @@ import opentelemetry from '@opentelemetry/api'; import { makeDateString, type ItemIdentifier } from '@roostorg/types'; import { types as scyllaTypes } from 'cassandra-driver'; import IORedis, { type Cluster } from 'ioredis'; -import * as knexPkg from 'knex'; -import { type Knex } from 'knex'; import { Kysely, PostgresDialect, @@ -20,7 +18,6 @@ import { type JsonObject, type ReadonlyDeep } from 'type-fest'; import { v1 as uuidv1 } from 'uuid'; import makeDb from '../models/index.js'; -import { type PolicyActionPenalties } from '../models/OrgModel.js'; import type { IActionExecutionsAdapter } from '../plugins/warehouse/queries/IActionExecutionsAdapter.js'; import type { IActionStatisticsAdapter } from '../plugins/warehouse/queries/IActionStatisticsAdapter.js'; import type { IContentApiRequestsAdapter } from '../plugins/warehouse/queries/IContentApiRequestsAdapter.js'; @@ -39,6 +36,10 @@ import { makeItemSubmissionBulkWrite, type ItemSubmissionBulkWrite, } from '../queues/itemSubmissionQueue.js'; +import { + getPolicyActionPenaltiesForOrg, + type PolicyActionPenalties, +} from '../services/policyActionPenalties.js'; import makeActionPublisher, { type ActionPublisher, type ActionTargetItem, @@ -50,7 +51,6 @@ import { makeGetItemTypesForOrgEventuallyConsistent, makeGetLocationBankLocationsEventuallyConsistent, makeGetPoliciesForRulesEventuallyConsistent, - makeGetSequelizeItemTypeEventuallyConsistent, makeGetTextBankStringsEventuallyConsistent, makeRecordRuleActionLimitUsage, type GetActionsForRuleEventuallyConsistent, @@ -58,7 +58,6 @@ import { type GetItemTypesForOrgEventuallyConsistent, type GetLocationBankLocationsBankEventuallyConsistent, type GetPoliciesForRulesEventuallyConsistent, - type GetSequelizeItemTypeEventuallyConsistent, type GetTextBankStringsEventuallyConsistent, type RecordRuleActionLimitUsage, } from '../rule_engine/ruleEngineQueries.js'; @@ -334,7 +333,6 @@ export interface Dependencies { itemSubmissionQueueBulkWrite: ItemSubmissionBulkWrite; itemSubmissionRetryQueueBulkWrite: ItemSubmissionBulkWrite; - Knex: Knex; IORedis: IORedis.Redis | Cluster; // Loggers @@ -397,7 +395,6 @@ export interface Dependencies { (input: { orgId: string; bankId: string }) => Promise >; - getSequelizeItemTypeEventuallyConsistent: GetSequelizeItemTypeEventuallyConsistent; getItemTypesForOrgEventuallyConsistent: GetItemTypesForOrgEventuallyConsistent; getItemTypeEventuallyConsistent: GetItemTypeEventuallyConsistent; getEnabledRulesForItemTypeEventuallyConsistent: GetEnabledRulesForItemTypeEventuallyConsistent; @@ -467,9 +464,7 @@ export default async function getBottle() { // // - 'KyselyPg' is for issuing raw pg queries w/o sequelize (e.g., the queries // that some of the our "services" issue to pg, to the non-public schemas). - // These queries go to our primary db, which accepts writes. Using knex for - // query building is deprecated in favor of kysely, because the latter offers - // better typings. + // These queries go to our primary db, which accepts writes. // // - KyselyPgReadReplica gives us the same type safety, but sends queries to our // replicas, for when we only need reads and we're ok w/ eventual consistency. @@ -618,15 +613,6 @@ export default async function getBottle() { makeItemSubmissionBulkWrite(container.IORedis, ITEM_SUBMISSION_DLQ_NAME), ); - // Legacy service deprecated in favor of kysely. - bottle.value( - 'Knex', - knexPkg.default.knex({ - client: 'pg', - connection: getPgMasterConnectionInfo, - }), - ); - // Loggers register(bottle, 'RuleExecutionLogger', makeRuleExecutionLogger); @@ -1345,11 +1331,14 @@ export default async function getBottle() { bottle.factory( 'getPolicyActionPenaltiesEventuallyConsistent', (container) => { - const Org = container.OrgModel; + const moderationConfigService = container.ModerationConfigService; return cached({ async producer(orgId) { - return Org.getPolicyActionPenaltiesEventuallyConsistent(orgId); + return getPolicyActionPenaltiesForOrg( + moderationConfigService, + orgId, + ); }, directives: { freshUntilAge: 60 }, }); @@ -1393,12 +1382,6 @@ export default async function getBottle() { }); }); - register( - bottle, - 'getSequelizeItemTypeEventuallyConsistent', - makeGetSequelizeItemTypeEventuallyConsistent, - ); - register( bottle, 'getEnabledRulesForItemTypeEventuallyConsistent', @@ -1567,7 +1550,6 @@ export default async function getBottle() { 'itemSubmissionQueueBulkWrite', 'itemSubmissionRetryQueueBulkWrite', 'Sequelize', - 'Knex', 'IORedis', // Storage abstractions 'DataWarehouse', @@ -1576,7 +1558,6 @@ export default async function getBottle() { 'ReportingAnalyticsAdapter', 'KyselyPg', 'KyselyPgReadReplica', - 'getSequelizeItemTypeEventuallyConsistent', 'getEnabledRulesForItemTypeEventuallyConsistent', 'getPoliciesForRulesEventuallyConsistent', 'getActionsForRuleEventuallyConsistent', diff --git a/server/models/OrgModel.ts b/server/models/OrgModel.ts index 49b50b3..2d3deac 100644 --- a/server/models/OrgModel.ts +++ b/server/models/OrgModel.ts @@ -6,7 +6,6 @@ import sequelize, { type Sequelize, } from 'sequelize'; -import { UserPenaltySeverity } from '../services/moderationConfigService/index.js'; import { validateUrl } from '../utils/url.js'; import { type LocationBank } from './banks/LocationBankModel.js'; import { type DataTypes } from './index.js'; @@ -18,11 +17,6 @@ import { type User } from './UserModel.js'; const { Model } = sequelize; export type Org = InstanceType>; -export type PolicyActionPenalties = { - actionId: string; - policyId: string; - penalties: number[]; -}; /** * Data Model for Organizations @@ -62,25 +56,6 @@ export default function makeOrgModel( Org.hasMany(models.LocationBank, { as: 'LocationBanks' }); Org.hasMany(models.Policy, { as: 'policies' }); } - - static async getPolicyActionPenaltiesEventuallyConsistent(orgId: string) { - const [actions, policies] = await Promise.all([ - sequelize.models.Action.findAll({ where: { orgId } }), - sequelize.models.Policy.findAll({ where: { orgId } }), - ]); - - return (policies as Policy[]).flatMap((policy) => - (actions as SequelizeAction[]).map( - (action): PolicyActionPenalties => ({ - actionId: action.id, - policyId: policy.id, - penalties: [ - computeActionPolicyPenalty(action.penalty, policy.penalty), - ], - }), - ), - ); - } } /* Fields */ @@ -137,46 +112,3 @@ export default function makeOrgModel( return Org; } - -/** - * Computes the severity of the penalty we should apply for a given - * (action, policy) pair. The general idea is to make the penalties - * increase exponentially as severity levels increase, but the rate - * of increase can't be so high that a (severe, severe) penalty is - * 50x higher than a (high, high) penalty. - * - * The easiest way to achieve this exponential behavior is at the individual - * severity levels, rather than trying to multiply the action penalty - * by the severity penalty to compound their magnitudes. So the severity - * levels apply penalty magnitudes as follows: - * - * NONE = 0 - * LOW = 1 - * MEDIUM = 3 - * HIGH = 9 - * SEVERE = 27 - * - * To get the penalty value for an (action, policy) pair, we just add the - * penalty values of the action and penalty because the exponential nature - * of these penalties has already been taken into account. - */ -function computeActionPolicyPenalty( - actionPenalty: UserPenaltySeverity, - policyPenalty: UserPenaltySeverity, -) { - // Type annotation makes sure that every possible severity has a score. - const penaltySeverityMap: { [k in UserPenaltySeverity]: number } = { - [UserPenaltySeverity.NONE]: 0, - [UserPenaltySeverity.LOW]: 1, - [UserPenaltySeverity.MEDIUM]: 3, - [UserPenaltySeverity.HIGH]: 9, - [UserPenaltySeverity.SEVERE]: 27, - }; - - // If the action has no penalty (e.g., "Send to Moderation", "Restore - // Content"), we never apply any penalty, regardless of the policy penalty. - // Otherwise, the penalty accounts for both the action + policy penalties. - return actionPenalty === UserPenaltySeverity.NONE - ? 0 - : penaltySeverityMap[actionPenalty] + penaltySeverityMap[policyPenalty]; -} diff --git a/server/models/rules/RuleModel.ts b/server/models/rules/RuleModel.ts index 70ecccf..b5726d2 100644 --- a/server/models/rules/RuleModel.ts +++ b/server/models/rules/RuleModel.ts @@ -1,4 +1,3 @@ -import { type ScalarType, type TaggedScalar } from '@roostorg/types'; import _ from 'lodash'; import sequelize, { Sequelize, @@ -16,20 +15,25 @@ import { RuleStatus, RuleType, type ConditionSet, - type LeafCondition, } from '../../services/moderationConfigService/index.js'; -import { type SerializableError } from '../../utils/errors.js'; import { getUtcDateOnlyString } from '../../utils/time.js'; -import { - type NonEmptyArray, - type WithUndefined, -} from '../../utils/typescript-types.js'; import { type DataTypes } from '../index.js'; import { type User } from '../UserModel.js'; import { type SequelizeAction } from './ActionModel.js'; -import { type TaggedItemData } from './item-type-fields.js'; import { type RuleLatestVersion } from './RuleLatestVersionModel.js'; +export { + ConditionCompletionOutcome, + ConditionFailureOutcome, + type ConditionOutcome, + type ConditionCompletionMetadata, + type ConditionFailureMetadata, + type ConditionResult, + type ConditionWithResult, + type ConditionSetWithResult, + type LeafConditionWithResult, +} from './ruleTypes.js'; + const { Model, Op } = sequelize; const { without } = _; @@ -37,78 +41,6 @@ export type Rule = InstanceType>; export type RuleWithLatestVersion = Rule & Required>; -export enum ConditionCompletionOutcome { - PASSED = 'PASSED', - FAILED = 'FAILED', - INAPPLICABLE = 'INAPPLICABLE', -} - -// We might add more kinds of errors here later... -export enum ConditionFailureOutcome { - ERRORED = 'ERRORED', -} - -export type ConditionOutcome = - | ConditionCompletionOutcome - | ConditionFailureOutcome; - -// NB: For legacy reasons, score is stored as a string and matchingValue holds -// only a single string value. In the future, its easy to imagine condition -// completions having much richer metadata, which could vary by signal type. -export type ConditionCompletionMetadata = { - score?: string; - matchedValue?: string; -}; - -export type ConditionFailureMetadata = { - error?: SerializableError; -}; - -type ConditionResultCommonMetadata = { - signalInputValues?: (TaggedScalar | TaggedItemData)[]; -}; - -// A completion outcome can have optional metadata about the signal's score etc, -// while a failure outcome can have an error. In a completion outcome, the -// failure metadata must be undefined (except in IS_UNAVAILABLE comparisons), -// and vice-versa. -// prettier-ignore -export type ConditionResult = - | ({ outcome: ConditionCompletionOutcome } - & ConditionCompletionMetadata - // Conditions that completed successfully can optionally have - // some of the failure metadata (for now just the `error`) if the - // condition used a signal result IS_UNAVAILABLE comparison. - & Partial> - & ConditionResultCommonMetadata) - | ({ outcome: ConditionFailureOutcome; } - & ConditionFailureMetadata - // Failed condition must not have completion data - & WithUndefined - & ConditionResultCommonMetadata) - -export type ConditionWithResult = - | LeafConditionWithResult - | ConditionSetWithResult; - -export type ConditionSetWithResult = Omit & { - conditions: - | NonEmptyArray - | NonEmptyArray; - result?: ConditionResult; -}; - -// A useful type for logging. Gives us the result of the condition, -// with the context (i.e., the Condition definition) to interpret that result. -// The result is still optional because the conditions execution can be skipped -// if the result of a parent condition can be determined without running the -// child condition. Just leaving out the result property in that case is a very -// nice way to model "skipped", because it means that the original condition -// object can be used as a ConditionWithResult in that case. -export type LeafConditionWithResult = LeafCondition & { - result?: ConditionResult; -}; - /** * Data Model for Rules. Rules are comprised of * ContentType inputs, Conditions, and Actions. diff --git a/server/models/rules/ruleTypes.ts b/server/models/rules/ruleTypes.ts new file mode 100644 index 0000000..4c3f4ce --- /dev/null +++ b/server/models/rules/ruleTypes.ts @@ -0,0 +1,124 @@ +import { type ScalarType, type TaggedScalar } from '@roostorg/types'; + +import { + type RuleAlarmStatus, + RuleStatus, + type RuleType, + type ConditionSet, + type LeafCondition, + type Action, + type Policy, +} from '../../services/moderationConfigService/index.js'; +import { type SerializableError } from '../../utils/errors.js'; +import { + type NonEmptyArray, + type WithUndefined, +} from '../../utils/typescript-types.js'; +import { type User } from '../UserModel.js'; +import { type TaggedItemData } from './item-type-fields.js'; + +export enum ConditionCompletionOutcome { + PASSED = 'PASSED', + FAILED = 'FAILED', + INAPPLICABLE = 'INAPPLICABLE', +} + +export enum ConditionFailureOutcome { + ERRORED = 'ERRORED', +} + +export type ConditionOutcome = + | ConditionCompletionOutcome + | ConditionFailureOutcome; + +export type ConditionCompletionMetadata = { + score?: string; + matchedValue?: string; +}; + +export type ConditionFailureMetadata = { + error?: SerializableError; +}; + +type ConditionResultCommonMetadata = { + signalInputValues?: (TaggedScalar | TaggedItemData)[]; +}; + +// prettier-ignore +export type ConditionResult = + | ({ outcome: ConditionCompletionOutcome } + & ConditionCompletionMetadata + & Partial> + & ConditionResultCommonMetadata) + | ({ outcome: ConditionFailureOutcome; } + & ConditionFailureMetadata + & WithUndefined + & ConditionResultCommonMetadata) + +export type ConditionWithResult = + | LeafConditionWithResult + | ConditionSetWithResult; + +export type ConditionSetWithResult = Omit & { + conditions: + | NonEmptyArray + | NonEmptyArray; + result?: ConditionResult; +}; + +export type LeafConditionWithResult = LeafCondition & { + result?: ConditionResult; +}; + +export type RuleLatestVersionRow = { + ruleId: string; + version: string; +}; + +/** + * Rule row fields shared by the rule engine (no GraphQL resolver methods). + */ +export type PlainRuleWithLatestVersion = { + id: string; + name: string; + description: string | null; + statusIfUnexpired: Exclude; + status: RuleStatus; + tags: string[]; + maxDailyActions: number | null; + dailyActionsRun: number; + lastActionDate: string | null; + createdAt: Date; + updatedAt: Date; + orgId: string; + creatorId: string; + expirationTime: Date | null; + conditionSet: ConditionSet; + alarmStatus: RuleAlarmStatus; + alarmStatusSetAt: Date; + ruleType: RuleType; + parentId: string | null; + latestVersion: RuleLatestVersionRow; +}; + +export function computeRuleStatusFromRow( + expirationTime: Date | null, + statusIfUnexpired: Exclude, +): RuleStatus { + if (expirationTime && expirationTime.valueOf() < Date.now()) { + return RuleStatus.EXPIRED; + } + return statusIfUnexpired; +} + +export type RuleGraphqlMethods = { + getCreator(): Promise; + getActions(): Promise; + getPolicies(): Promise; +}; + +/** GraphQL parent for Rule / ContentRule / UserRule / RuleInsights. */ +export type Rule = PlainRuleWithLatestVersion & RuleGraphqlMethods; + +/** @deprecated Use {@link PlainRuleWithLatestVersion} directly. Remove after Kysely migration. */ +export type RuleWithLatestVersion = PlainRuleWithLatestVersion; diff --git a/server/package-lock.json b/server/package-lock.json index 0523b33..c1dcb79 100644 --- a/server/package-lock.json +++ b/server/package-lock.json @@ -69,7 +69,6 @@ "helmet": "^4.6.0", "ioredis": "^5.2.4", "jsonwebtoken": "^9.0.3", - "knex": "^2.3.0", "kysely": "^0.28.16", "latlon-geohash": "^2.0.0", "lodash": "^4.18.1", @@ -13141,11 +13140,6 @@ "resolved": "https://registry.npmjs.org/color-name/-/color-name-1.1.4.tgz", "integrity": "sha512-dOy+3AuW3a2wNbZHIuMZpTcgjGuLU/uBL/ubcZF9OXbDo8ff4O8yVp5Bf0efS8uEoYo5q4Fx7dY9OgQGXgAsQA==" }, - "node_modules/colorette": { - "version": "2.0.19", - "resolved": "https://registry.npmjs.org/colorette/-/colorette-2.0.19.tgz", - "integrity": "sha512-3tlv/dIP7FWvj3BsbHrGLJ6l/oKh1O3TcgBqMn+yyCagOxc23fyzDS6HypQbgxWbkpDnf52p1LuR4eWDQ/K9WQ==" - }, "node_modules/colors": { "version": "1.4.0", "resolved": "https://registry.npmjs.org/colors/-/colors-1.4.0.tgz", @@ -13165,14 +13159,6 @@ "node": ">= 0.8" } }, - "node_modules/commander": { - "version": "10.0.1", - "resolved": "https://registry.npmjs.org/commander/-/commander-10.0.1.tgz", - "integrity": "sha512-y4Mg2tXshplEbSGzx7amzPwKKOCGuoSRP/CjEdwwk0FOGlUbq6lKuoyDZTNZkmxHdJtp54hdfY/JUrdL7Xfdug==", - "engines": { - "node": ">=14" - } - }, "node_modules/comment-parser": { "version": "1.4.5", "resolved": "https://registry.npmjs.org/comment-parser/-/comment-parser-1.4.5.tgz", @@ -14558,14 +14544,6 @@ "integrity": "sha512-xbbCH5dCYU5T8LcEhhuh7HJ88HXuW3qsI3Y0zOZFKfZEHcpWiHU/Jxzk629Brsab/mMiHQti9wMP+845RPe3Vg==", "license": "MIT" }, - "node_modules/esm": { - "version": "3.2.25", - "resolved": "https://registry.npmjs.org/esm/-/esm-3.2.25.tgz", - "integrity": "sha512-U1suiZ2oDVWv4zPO56S0NcR5QriEahGtdN2OR6FiOG4WJvcjBVFB0qI4+eKoWFH483PKGuLuu6V8Z4T5g63UVA==", - "engines": { - "node": ">=6" - } - }, "node_modules/espree": { "version": "10.4.0", "resolved": "https://registry.npmjs.org/espree/-/espree-10.4.0.tgz", @@ -15426,6 +15404,7 @@ "version": "0.1.0", "resolved": "https://registry.npmjs.org/get-package-type/-/get-package-type-0.1.0.tgz", "integrity": "sha512-pjzuKtY64GYfWizNAJ0fr9VqttZkNiK2iS430LtIHzjBEr6bX8Am2zm4sW4Ro5wjWW5cAlRL1qAMTcXbjNAO2Q==", + "dev": true, "engines": { "node": ">=8.0.0" } @@ -15486,11 +15465,6 @@ "url": "https://github.com/privatenumber/get-tsconfig?sponsor=1" } }, - "node_modules/getopts": { - "version": "2.3.0", - "resolved": "https://registry.npmjs.org/getopts/-/getopts-2.3.0.tgz", - "integrity": "sha512-5eDf9fuSXwxBL6q5HX+dhDj+dslFGWzU5thZ9kNKUkcPtaPdatmUFKwHFrLb/uf/WpA4BHET+AX3Scl56cAjpA==" - }, "node_modules/glob": { "version": "7.2.3", "resolved": "https://registry.npmjs.org/glob/-/glob-7.2.3.tgz", @@ -16098,14 +16072,6 @@ "node": ">= 0.4" } }, - "node_modules/interpret": { - "version": "2.2.0", - "resolved": "https://registry.npmjs.org/interpret/-/interpret-2.2.0.tgz", - "integrity": "sha512-Ju0Bz/cEia55xDwUWEa8+olFpCiQoypjnQySseKtmjNrnps3P+xfpUmGr90T7yjlVJmOtybRvPXhKMbHr+fWnw==", - "engines": { - "node": ">= 0.10" - } - }, "node_modules/ioredis": { "version": "5.9.2", "resolved": "https://registry.npmjs.org/ioredis/-/ioredis-5.9.2.tgz", @@ -17608,64 +17574,6 @@ "node": ">=6" } }, - "node_modules/knex": { - "version": "2.5.1", - "resolved": "https://registry.npmjs.org/knex/-/knex-2.5.1.tgz", - "integrity": "sha512-z78DgGKUr4SE/6cm7ku+jHvFT0X97aERh/f0MUKAKgFnwCYBEW4TFBqtHWFYiJFid7fMrtpZ/gxJthvz5mEByA==", - "dependencies": { - "colorette": "2.0.19", - "commander": "^10.0.0", - "debug": "4.3.4", - "escalade": "^3.1.1", - "esm": "^3.2.25", - "get-package-type": "^0.1.0", - "getopts": "2.3.0", - "interpret": "^2.2.0", - "lodash": "^4.17.21", - "pg-connection-string": "2.6.1", - "rechoir": "^0.8.0", - "resolve-from": "^5.0.0", - "tarn": "^3.0.2", - "tildify": "2.0.0" - }, - "bin": { - "knex": "bin/cli.js" - }, - "engines": { - "node": ">=12" - }, - "peerDependenciesMeta": { - "better-sqlite3": { - "optional": true - }, - "mysql": { - "optional": true - }, - "mysql2": { - "optional": true - }, - "pg": { - "optional": true - }, - "pg-native": { - "optional": true - }, - "sqlite3": { - "optional": true - }, - "tedious": { - "optional": true - } - } - }, - "node_modules/knex/node_modules/resolve-from": { - "version": "5.0.0", - "resolved": "https://registry.npmjs.org/resolve-from/-/resolve-from-5.0.0.tgz", - "integrity": "sha512-qYg9KP24dD5qka9J47d0aVky0N+b4fTU89LN9iDnjB5waksiC49rvMB0PrUJQGoTmH50XPiqOvAjDfaijGxYZw==", - "engines": { - "node": ">=8" - } - }, "node_modules/kysely": { "version": "0.28.16", "resolved": "https://registry.npmjs.org/kysely/-/kysely-0.28.16.tgz", @@ -19312,17 +19220,6 @@ "node": ">= 6" } }, - "node_modules/rechoir": { - "version": "0.8.0", - "resolved": "https://registry.npmjs.org/rechoir/-/rechoir-0.8.0.tgz", - "integrity": "sha512-/vxpCXddiX8NGfGO/mTafwjq4aFa/71pvamip0++IQk3zG8cbCj0fifNPrjjF1XMXUne91jL9OoxmdykoEtifQ==", - "dependencies": { - "resolve": "^1.20.0" - }, - "engines": { - "node": ">= 10.13.0" - } - }, "node_modules/redis-errors": { "version": "1.2.0", "resolved": "https://registry.npmjs.org/redis-errors/-/redis-errors-1.2.0.tgz", @@ -20621,14 +20518,6 @@ "node": ">=6" } }, - "node_modules/tarn": { - "version": "3.0.2", - "resolved": "https://registry.npmjs.org/tarn/-/tarn-3.0.2.tgz", - "integrity": "sha512-51LAVKUSZSVfI05vjPESNc5vwqqZpbXCsU+/+wxlOrUjk2SnFTt97v9ZgQrD4YmxYW1Px6w2KjaDitCfkvgxMQ==", - "engines": { - "node": ">=8.0.0" - } - }, "node_modules/teeny-request": { "version": "10.1.2", "resolved": "https://registry.npmjs.org/teeny-request/-/teeny-request-10.1.2.tgz", @@ -20726,14 +20615,6 @@ "safe-buffer": "~5.1.0" } }, - "node_modules/tildify": { - "version": "2.0.0", - "resolved": "https://registry.npmjs.org/tildify/-/tildify-2.0.0.tgz", - "integrity": "sha512-Cc+OraorugtXNfs50hU9KS369rFXCfgGLpfCfvlc+Ud5u6VWmUQsOAa9HbTvheQdYnrdJqqv1e5oIqXppMYnSw==", - "engines": { - "node": ">=8" - } - }, "node_modules/tinyglobby": { "version": "0.2.15", "resolved": "https://registry.npmjs.org/tinyglobby/-/tinyglobby-0.2.15.tgz", diff --git a/server/package.json b/server/package.json index 0857495..94584a0 100644 --- a/server/package.json +++ b/server/package.json @@ -83,7 +83,6 @@ "helmet": "^4.6.0", "ioredis": "^5.2.4", "jsonwebtoken": "^9.0.3", - "knex": "^2.3.0", "kysely": "^0.28.16", "latlon-geohash": "^2.0.0", "lodash": "^4.18.1", diff --git a/server/plugins/warehouse/IWarehouseAdapter.ts b/server/plugins/warehouse/IWarehouseAdapter.ts index d17ce0f..58fe25c 100644 --- a/server/plugins/warehouse/IWarehouseAdapter.ts +++ b/server/plugins/warehouse/IWarehouseAdapter.ts @@ -5,7 +5,7 @@ import type { WarehouseQueryResult, WarehouseTransactionFn } from './types.js'; * * Adapters power operational workloads (transactions, point queries, etc.). * The interface intentionally mirrors a typical SQL client without enforcing - * a specific library (Kysely, Knex, pg, etc.). + * a specific library (Kysely, pg, etc.). */ export interface IWarehouseAdapter { /** Human friendly provider name for logging / diagnostics. */ diff --git a/server/rule_engine/ActionPublisher.test.ts b/server/rule_engine/ActionPublisher.test.ts index fbae061..36487e6 100644 --- a/server/rule_engine/ActionPublisher.test.ts +++ b/server/rule_engine/ActionPublisher.test.ts @@ -42,6 +42,7 @@ describe('ActionPublisher', () => { orgId: 'org-123', name: 'Action 1', applyUserStrikes: false, + penalty: 'NONE' as const, actionType: ActionType.CUSTOM_ACTION, callbackUrl: 'https://example.com/action1', callbackUrlHeaders: null, @@ -73,6 +74,7 @@ describe('ActionPublisher', () => { orgId: 'org-123', name: 'Action 2', applyUserStrikes: false, + penalty: 'NONE' as const, actionType: ActionType.CUSTOM_ACTION, callbackUrl: 'https://example.com/action2', callbackUrlHeaders: null, diff --git a/server/rule_engine/RuleEngine.ts b/server/rule_engine/RuleEngine.ts index b8f1cbe..c796242 100644 --- a/server/rule_engine/RuleEngine.ts +++ b/server/rule_engine/RuleEngine.ts @@ -7,11 +7,8 @@ import { } from '../condition_evaluator/conditionSet.js'; import { type Dependencies } from '../iocContainer/index.js'; import { inject } from '../iocContainer/utils.js'; -import { - ConditionCompletionOutcome, - type RuleWithLatestVersion, - type Rule as TRule, -} from '../models/rules/RuleModel.js'; +import { ConditionCompletionOutcome } from '../models/rules/RuleModel.js'; +import { type PlainRuleWithLatestVersion } from '../models/rules/ruleTypes.js'; import { evaluateAggregationRuntimeArgsForItem } from '../services/aggregationsService/index.js'; import { type ItemSubmission } from '../services/itemProcessingService/index.js'; import { @@ -182,13 +179,16 @@ class RuleEngine { * @param sync - whether the request should run synchronously */ async runRuleSet( - rules: ReadonlyDeep, + rules: ReadonlyDeep, evaluationContext: RuleEvaluationContext, environment: RuleEnvironment, executionsCorrelationId: RuleExecutionCorrelationId, sync?: boolean, ): Promise<{ - rulesToResults: Map, RuleExecutionResult>; + rulesToResults: Map< + ReadonlyDeep, + RuleExecutionResult + >; actions: readonly Action[]; }> { if (!rules.length) { diff --git a/server/rule_engine/ruleEngineQueries.ts b/server/rule_engine/ruleEngineQueries.ts index 2423254..c1a9d52 100644 --- a/server/rule_engine/ruleEngineQueries.ts +++ b/server/rule_engine/ruleEngineQueries.ts @@ -4,36 +4,25 @@ * the signal execution service) to enable those services to make the queries * they need to run a rule. Having these queries defined separately + injected * into the consumers gives us a cleaner place to add optimizations to the query - * logic (i.e., run the queries against replicas, add caching, etc) and makes + * logic (i.e. run the queries against replicas, add caching, etc) and makes * the consumers much more unit testable. */ -import { Op } from 'sequelize'; +import { sql, type Kysely } from 'kysely'; import { inject } from '../iocContainer/index.js'; import { type LocationArea } from '../models/types/locationArea.js'; -import { type Action } from '../services/moderationConfigService/index.js'; +import { type CombinedPg } from '../services/combinedDbTypes.js'; import { cached } from '../utils/caching.js'; import { jsonParse, jsonStringify } from '../utils/encoding.js'; +import { makeKyselyTransactionWithRetry } from '../utils/kyselyTransactionWithRetry.js'; import { getUtcDateOnlyString } from '../utils/time.js'; export const makeGetEnabledRulesForItemTypeEventuallyConsistent = inject( - ['getSequelizeItemTypeEventuallyConsistent'], - function (getSequelizeItemTypeEventuallyConsistent) { + ['ModerationConfigService'], + function (moderationConfigService) { return cached({ async producer(itemTypeId: string) { - // Getting the enabledRules is currently coupled to sequelize, so, - // annoyingly, we first have to convert the contentTypeId into a full - // contentType model object. However, we don't want to incur too much - // overhead for that, so we use a cached lookup. (Note: we can't just - // take a model object as the argument because our caching library - // requires the cache key to be a string; further, even if we could just - // accept the model object, we wouldn't want to, because we want to move - // away from this coupling to sequelize.) - const itemType = await getSequelizeItemTypeEventuallyConsistent({ - id: itemTypeId, - }); - - return itemType ? itemType.getEnabledRules() : null; + return moderationConfigService.getEnabledRulesForItemType(itemTypeId); }, directives: { freshUntilAge: 20 }, }); @@ -44,24 +33,6 @@ export type GetEnabledRulesForItemTypeEventuallyConsistent = ReturnType< typeof makeGetEnabledRulesForItemTypeEventuallyConsistent >; -export const makeGetSequelizeItemTypeEventuallyConsistent = inject( - ['ItemTypeModel'], - (ItemType) => { - return cached({ - async producer(key: { id: string } | { name: string; orgId: string }) { - return 'id' in key - ? ItemType.findByPk(key.id) - : ItemType.findOne({ where: { name: key.name, orgId: key.orgId } }); - }, - directives: { freshUntilAge: 10, maxStale: [0, 2, 2] }, - }); - }, -); - -export type GetSequelizeItemTypeEventuallyConsistent = ReturnType< - typeof makeGetSequelizeItemTypeEventuallyConsistent ->; - export const makeGetItemTypesForOrgEventuallyConsistent = inject( ['ModerationConfigService'], (moderationConfigService) => async (orgId: string) => @@ -74,18 +45,16 @@ export type GetItemTypesForOrgEventuallyConsistent = ReturnType< typeof makeGetItemTypesForOrgEventuallyConsistent >; -// TODO: this could probably be improved to increase cache hit rates, since -// rn the cache will only be used if all the ids have previously been fetched. export const makeGetPoliciesForRulesEventuallyConsistent = inject( - ['PolicyModel'], - function (Policy) { + ['ModerationConfigService'], + function (moderationConfigService) { return cached({ keyGeneration: { toString: (ids: readonly string[]) => jsonStringify([...ids].sort()), fromString: (it) => jsonParse(it), }, - async producer(key) { - return Policy.getPoliciesForRuleIds(key); + async producer(key: readonly string[]) { + return moderationConfigService.getPoliciesByRuleIds(key); }, directives: { freshUntilAge: 120 }, }); @@ -97,21 +66,11 @@ export type GetPoliciesForRulesEventuallyConsistent = ReturnType< >; export const makeGetActionsForRuleEventuallyConsistent = inject( - ['ActionModel'], - (Action) => { + ['ModerationConfigService'], + (moderationConfigService) => { return cached({ async producer(ruleId: string) { - // This generates a pretty slow/overly-complex query, but I think it's - // the best we can do with Sequelize. Eventually, we want to move off of - // Sequelize, so we don't fetch full model instances + we cast the - // result to be a plain data object, so that the rest of the code can't - // depend on getting an Action model instance back, as that won't always - // be the case. - return Action.findAll({ - where: { '$rules.id$': ruleId }, - include: [{ association: 'rules', attributes: ['id'] }], - raw: true, - }) as Promise; + return moderationConfigService.getActionsForRuleId(ruleId); }, directives: { freshUntilAge: 30 }, }); @@ -123,24 +82,25 @@ export type GetActionsForRuleEventuallyConsistent = ReturnType< >; export const makeGetLocationBankLocationsEventuallyConsistent = inject( - ['LocationBankLocationModel'], - (LocationBankLocation) => { + ['KyselyPgReadReplica'], + (db) => { return cached({ async producer(bankId: string) { - // NB: we use `raw: true` to get back plain JS objects, rather than - // sequelize model instances. We do that because, with model instances, - // every proprety access runs some extra getter code; see - // https://github.com/sequelize/sequelize/blob/e77dcf78b341b62c97dbb29f16ce7a23f46ddc53/src/model.js#L42 - // This ends up killing the performance of our hot-path code that checks - // whether a location is in each of these location areas. - // - // Meanwhile, we have to cast to LocationArea[] because findAll is still - // typed (incorrectly) to return the model instance, even when `raw: - // true` is provided. - return LocationBankLocation.findAll({ - where: { bankId }, - raw: true, - }) as Promise; + const rows = await db + .selectFrom('public.location_bank_locations') + .selectAll() + .where('bank_id', '=', bankId) + .execute(); + return rows.map( + (r) => + ({ + id: r.id, + name: r.name ?? undefined, + geometry: r.geometry as LocationArea['geometry'], + bounds: r.bounds as LocationArea['bounds'], + googlePlaceInfo: r.google_place_info as LocationArea['googlePlaceInfo'], + }) satisfies LocationArea, + ); }, directives(locations) { const numLocations = locations.length; @@ -197,40 +157,42 @@ export type GetImageBankEventuallyConsistent = ReturnType< >; export const makeRecordRuleActionLimitUsage = inject( - ['Sequelize', 'Tracer'], + ['KyselyPg', 'Tracer'], (db, tracer) => { - /** - * Record that each of the rules given by ruleIds has used up one of its - * daily action runs, against its maxDailyActions. - */ + const transactionWithRetry = makeKyselyTransactionWithRetry( + db as Kysely, + ); + async function recordRuleActionLimitUsage(ruleIds: readonly string[]) { if (ruleIds.length === 0) { return; } - const today = getUtcDateOnlyString(); - await db.transactionWithRetry(async () => { - // Using two queries like this isn't as efficient as, e.g., - // UPDATE `rules` - // SET `daily_actions_run` = - // IF(last_action_date != $1, 1, daily_actions_run + 1) - // SET `last_action_date` = $1 - // WHERE `id` IN (...); - // But it lets us keep the code in Sequelize, which is probably worth it. - await db.Rule.increment( - { dailyActionsRun: 1 }, - { where: { id: { [Op.in]: ruleIds }, lastActionDate: today } }, - ); - - await db.Rule.update( - { dailyActionsRun: 1, lastActionDate: today }, - { - where: { - id: { [Op.in]: ruleIds }, - lastActionDate: { [Op.ne]: today }, - }, - }, - ); + const today = String(getUtcDateOnlyString()); + await transactionWithRetry(async (trx) => { + await trx + .updateTable('public.rules') + .set({ + daily_actions_run: sql`daily_actions_run + 1`, + }) + .where('id', 'in', [...ruleIds]) + .where('last_action_date', '=', today) + .execute(); + + await trx + .updateTable('public.rules') + .set({ + daily_actions_run: 1, + last_action_date: today, + }) + .where('id', 'in', [...ruleIds]) + .where((eb) => + eb.or([ + eb('last_action_date', 'is', null), + eb('last_action_date', '!=', today), + ]), + ) + .execute(); }); } diff --git a/server/services/combinedDbTypes.ts b/server/services/combinedDbTypes.ts index 2a4eff2..80ff434 100644 --- a/server/services/combinedDbTypes.ts +++ b/server/services/combinedDbTypes.ts @@ -1,5 +1,11 @@ -import { type ModerationConfigServicePg } from './moderationConfigService/dbTypes.js'; import { type ApiKeyServicePg } from './apiKeyService/dbTypes.js'; +import { type CoreAppTablesPg } from './coreAppTables.js'; +import { type ModerationConfigServicePg } from './moderationConfigService/dbTypes.js'; import { type SigningKeyPairServicePg } from './signingKeyPairService/dbTypes.js'; +import { type UserManagementPg } from './userManagementService/dbTypes.js'; -export type CombinedPg = ModerationConfigServicePg & ApiKeyServicePg & SigningKeyPairServicePg; +export type CombinedPg = ModerationConfigServicePg & + ApiKeyServicePg & + SigningKeyPairServicePg & + UserManagementPg & + CoreAppTablesPg; diff --git a/server/services/coreAppTables.ts b/server/services/coreAppTables.ts new file mode 100644 index 0000000..80e0ab1 --- /dev/null +++ b/server/services/coreAppTables.ts @@ -0,0 +1,59 @@ +import { type Generated, type GeneratedAlways } from 'kysely'; + +/** Postgres enum for backtests.status (generated column — read-only in app). */ +export type BacktestStatusDb = 'RUNNING' | 'COMPLETE' | 'CANCELED'; + +export type CoreAppTablesPg = { + 'public.orgs': { + id: string; + email: string; + name: string; + website_url: string; + api_key_id: string | null; + created_at: Date; + updated_at: Date; + on_call_alert_email: string | null; + }; + 'public.location_banks': { + id: string; + name: string; + description: string | null; + org_id: string; + owner_id: string; + created_at: GeneratedAlways; + updated_at: Date; + full_places_api_responses: unknown[]; + }; + 'public.location_bank_locations': { + id: string; + bank_id: string; + geometry: unknown; + bounds: unknown | null; + name: string | null; + google_place_info: unknown | null; + created_at: GeneratedAlways; + updated_at: GeneratedAlways; + }; + 'public.backtests': { + id: string; + rule_id: string; + creator_id: string; + sample_desired_size: number; + sample_actual_size: Generated; + sample_start_at: Date; + sample_end_at: Date; + sampling_complete: Generated; + content_items_processed: Generated; + content_items_matched: Generated; + created_at: GeneratedAlways; + updated_at: Date; + cancelation_date: Date | null; + status: GeneratedAlways; + }; + 'public.users_and_favorite_rules': { + user_id: string; + rule_id: string; + created_at: GeneratedAlways; + updated_at: Date; + }; +}; diff --git a/server/services/moderationConfigService/dbTypes.ts b/server/services/moderationConfigService/dbTypes.ts index d5644e4..2c217cd 100644 --- a/server/services/moderationConfigService/dbTypes.ts +++ b/server/services/moderationConfigService/dbTypes.ts @@ -62,6 +62,20 @@ export type ModerationConfigServicePg = { created_at: GeneratedAlways; updated_at: GeneratedAlways; }; + 'public.rules_and_actions': { + action_id: string; + rule_id: string; + created_at: GeneratedAlways; + updated_at: GeneratedAlways; + sys_period: GeneratedAlways; + }; + 'public.rules_and_policies': { + policy_id: string; + rule_id: string; + created_at: Date; + updated_at: Date; + sys_period: GeneratedAlways; + }; 'public.actions_and_item_types': { action_id: string; item_type_id: string; @@ -100,9 +114,14 @@ export type ModerationConfigServicePg = { ACCEPT_APPEAL: { callback_url: null }; } >; + 'public.rules_latest_versions': { + rule_id: string; + version: string; + }; 'public.rules': { id: string; name: string; + description: string | null; status_if_unexpired: RuleStatus; tags: string[]; max_daily_actions: number | null; @@ -117,6 +136,7 @@ export type ModerationConfigServicePg = { alarm_status: Generated; alarm_status_set_at: Generated; rule_type: RuleType; + parent_id: string | null; }; 'public.policies': { id: string; diff --git a/server/services/moderationConfigService/errors.ts b/server/services/moderationConfigService/errors.ts new file mode 100644 index 0000000..c5eff23 --- /dev/null +++ b/server/services/moderationConfigService/errors.ts @@ -0,0 +1,49 @@ +import { + CoopError, + ErrorType, + type ErrorInstanceData, +} from '../../utils/errors.js'; + +export type RuleErrorType = + | 'RuleNameExistsError' + | 'RuleHasRunningBacktestsError' + | 'RuleIsMissingContentTypeError'; + +export const makeRuleNameExistsError = (data: ErrorInstanceData) => + new CoopError({ + status: 409, + type: [ErrorType.UniqueViolation], + title: 'A rule with that name already exists in this organization.', + name: 'RuleNameExistsError', + ...data, + }); + +export const makeRuleIsMissingContentTypeError = (data: ErrorInstanceData) => + new CoopError({ + status: 400, + type: [ErrorType.InvalidUserInput], + title: 'This rule must contain a content type on which to operate.', + name: 'RuleIsMissingContentTypeError', + ...data, + }); + +export const makeRuleHasRunningBacktestsError = (data: ErrorInstanceData) => + new CoopError({ + status: 409, + type: [ErrorType.AttemptingToMutateActiveRule], + title: + "This rule cannot be updated while it has running backtests, which are using the rule's current conditions.", + name: 'RuleHasRunningBacktestsError', + ...data, + }); + +export type LocationBankErrorType = 'LocationBankNameExistsError'; + +export const makeLocationBankNameExistsError = (data: ErrorInstanceData) => + new CoopError({ + status: 409, + type: [ErrorType.UniqueViolation], + title: 'A location bank with this name already exists', + name: 'LocationBankNameExistsError', + ...data, + }); diff --git a/server/services/moderationConfigService/index.ts b/server/services/moderationConfigService/index.ts index 6b409f1..ac9f0c7 100644 --- a/server/services/moderationConfigService/index.ts +++ b/server/services/moderationConfigService/index.ts @@ -44,3 +44,10 @@ export { ModerationConfigService, ModerationConfigErrorType, } from './moderationConfigService.js'; + +export { + makeRuleNameExistsError, + makeRuleIsMissingContentTypeError, + makeRuleHasRunningBacktestsError, + makeLocationBankNameExistsError, +} from './errors.js'; diff --git a/server/services/moderationConfigService/moderationConfigService.test.ts b/server/services/moderationConfigService/moderationConfigService.test.ts index 2df427c..07033d1 100644 --- a/server/services/moderationConfigService/moderationConfigService.test.ts +++ b/server/services/moderationConfigService/moderationConfigService.test.ts @@ -452,6 +452,7 @@ describe('ModerationConfigService', () => { "id": Any, "name": "Test Action", "orgId": Any, + "penalty": "NONE", } `, ); diff --git a/server/services/moderationConfigService/moderationConfigService.ts b/server/services/moderationConfigService/moderationConfigService.ts index d59149c..7b110c5 100644 --- a/server/services/moderationConfigService/moderationConfigService.ts +++ b/server/services/moderationConfigService/moderationConfigService.ts @@ -4,12 +4,7 @@ import { type JsonObject, type ReadonlyDeep } from 'type-fest'; import { type ConsumerDirectives } from '../../lib/cache/index.js'; import type { Invoker } from '../../models/types/permissioning.js'; -import { - CoopError, - ErrorType, - type ErrorInstanceData, -} from '../../utils/errors.js'; -import { __throw } from '../../utils/misc.js'; +import { type RuleErrorType, type LocationBankErrorType } from './errors.js'; import { type ModerationConfigServicePg } from './dbTypes.js'; import { type Action, type Policy } from './index.js'; import ActionOperations, { @@ -22,6 +17,7 @@ import MatchingBankOperations, { import PolicyOperations, { type PolicyErrorType, } from './modules/PolicyOperations.js'; +import RuleReadOperations from './modules/RuleReadOperations.js'; import UserStrikeOperations, { type UserStrikeThresholdErrorType, } from './modules/UserStrikeOperations.js'; @@ -35,6 +31,7 @@ import { type UserItemType, } from './types/itemTypes.js'; import type { PolicyType } from './types/policies.js'; +import { type PlainRuleWithLatestVersion } from '../../models/rules/ruleTypes.js'; export type ModerationConfigErrorType = | 'AttemptingToDeleteDefaultUserType' @@ -78,6 +75,7 @@ export class ModerationConfigService implements ReturnsModerationConfigTypes { private readonly itemTypeOps: ItemTypeOperations; private readonly userStrikeOps: UserStrikeOperations; private readonly matchingBankOps: MatchingBankOperations; + private readonly ruleReadOps: RuleReadOperations; constructor( pgQuery: Kysely, @@ -94,9 +92,9 @@ export class ModerationConfigService implements ReturnsModerationConfigTypes { onDeletePolicyId, ); this.itemTypeOps = new ItemTypeOperations(pgQuery, pgQueryReplica); - // TODO: Remove Rule API and replace with kysely this.userStrikeOps = new UserStrikeOperations(pgQuery, pgQueryReplica); this.matchingBankOps = new MatchingBankOperations(pgQuery, pgQueryReplica); + this.ruleReadOps = new RuleReadOperations(pgQuery, pgQueryReplica); } async getItemTypes(opts: { @@ -285,6 +283,28 @@ export class ModerationConfigService implements ReturnsModerationConfigTypes { return this.actionOps.getActions(opts); } + async getActionsForRuleId(ruleId: string) { + return this.actionOps.getActionsForRuleId({ + ruleId, + readFromReplica: true, + }); + } + + async getPoliciesByRuleIds(ruleIds: readonly string[]) { + return this.policyOps.getPoliciesByRuleIds({ + ruleIds, + readFromReplica: true, + }); + } + + async getEnabledRulesForItemType(itemTypeId: string) { + return this.ruleReadOps.getEnabledRulesForItemType(itemTypeId); + } + + async findEnabledUserRules(): Promise { + return this.ruleReadOps.findEnabledUserRules(); + } + async getPolicies(opts: { orgId: string; readFromReplica?: boolean }) { return this.policyOps.getPolicies(opts); } @@ -429,50 +449,3 @@ export class ModerationConfigService implements ReturnsModerationConfigTypes { } } -type RuleErrorType = - | 'RuleNameExistsError' - | 'RuleHasRunningBacktestsError' - | 'RuleIsMissingContentTypeError'; - -// TODO: throw this error as appropriate on failed rule creation/update. -export const makeRuleNameExistsError = (data: ErrorInstanceData) => - new CoopError({ - status: 409, - type: [ErrorType.UniqueViolation], - title: 'A rule with that name already exists in this organization.', - name: 'RuleNameExistsError', - ...data, - }); - -// TODO: throw this error as appropriate on failed rule creation/update. -export const makeRuleIsMissingContentTypeError = (data: ErrorInstanceData) => - new CoopError({ - status: 400, - type: [ErrorType.InvalidUserInput], - title: 'This rule must contain a content type on which to operate.', - name: 'RuleIsMissingContentTypeError', - ...data, - }); - -// TODO: throw this error as appropriate on failed rule creation/update. -export const makeRuleHasRunningBacktestsError = (data: ErrorInstanceData) => - new CoopError({ - status: 409, - type: [ErrorType.AttemptingToMutateActiveRule], - title: - "This rule cannot be updated while it has running backtests, which are using the rule's current conditions.", - name: 'RuleHasRunningBacktestsError', - ...data, - }); - -type LocationBankErrorType = 'LocationBankNameExistsError'; - -// TODO: throw this error as appropriate on failed bank creation/update. -export const makeLocationBankNameExistsError = (data: ErrorInstanceData) => - new CoopError({ - status: 409, - type: [ErrorType.UniqueViolation], - title: 'A location bank with this name already exists', - name: 'LocationBankNameExistsError', - ...data, - }); diff --git a/server/services/moderationConfigService/modules/ActionOperations.ts b/server/services/moderationConfigService/modules/ActionOperations.ts index c9dcdda..5cc42bc 100644 --- a/server/services/moderationConfigService/modules/ActionOperations.ts +++ b/server/services/moderationConfigService/modules/ActionOperations.ts @@ -29,6 +29,20 @@ const actionDbSelection = [ 'apply_user_strikes as applyUserStrikes', ] as const; +const actionJoinDbSelection = [ + 'a.id', + 'a.name', + 'a.description', + 'a.callback_url as callbackUrl', + 'a.callback_url_headers as callbackUrlHeaders', + 'a.callback_url_body as callbackUrlBody', + 'a.org_id as orgId', + 'a.penalty', + 'a.action_type as actionType', + 'a.applies_to_all_items_of_kind as appliesToAllItemsOfKind', + 'a.apply_user_strikes as applyUserStrikes', +] as const; + type ActionDbResult = FixKyselyRowCorrelation< ModerationConfigServicePg['public.actions'], typeof actionDbSelection @@ -105,12 +119,29 @@ export default class ActionOperations { return results.map((it) => this.#dbResultToAction(it)); } + async getActionsForRuleId(opts: { + ruleId: string; + readFromReplica?: boolean; + }) { + const { ruleId, readFromReplica } = opts; + const pgQuery = this.#getPgQuery(readFromReplica ?? true); + const results = (await pgQuery + .selectFrom('public.rules_and_actions as raa') + .innerJoin('public.actions as a', 'a.id', 'raa.action_id') + .select(actionJoinDbSelection) + .where('raa.rule_id', '=', ruleId) + .execute()) as ActionDbResult[]; + + return results.map((it) => this.#dbResultToAction(it)); + } + #dbResultToAction(it: ActionDbResult) { return { id: it.id, name: it.name, orgId: it.orgId, applyUserStrikes: it.applyUserStrikes, + penalty: it.penalty, ...(() => { switch (it.actionType) { case 'CUSTOM_ACTION': diff --git a/server/services/moderationConfigService/modules/PolicyOperations.ts b/server/services/moderationConfigService/modules/PolicyOperations.ts index 4828eff..bf8b643 100644 --- a/server/services/moderationConfigService/modules/PolicyOperations.ts +++ b/server/services/moderationConfigService/modules/PolicyOperations.ts @@ -36,7 +36,25 @@ const policyDbSelection = [ 'semantic_version as semanticVersion', 'user_strike_count as userStrikeCount', 'apply_user_strike_count_config_to_children as applyUserStrikeCountConfigToChildren', - 'penalty', // TODO: remove + 'penalty', +] as const; + +const policyJoinDbSelection = [ + 'rap.rule_id as ruleId', + 'p.id', + 'p.name', + 'p.org_id as orgId', + 'p.parent_id as parentId', + 'p.created_at as createdAt', + 'p.updated_at as updatedAt', + 'p.policy_text as policyText', + 'p.enforcement_guidelines as enforcementGuidelines', + 'p.sys_period as sysPeriod', + 'p.policy_type as policyType', + 'p.semantic_version as semanticVersion', + 'p.user_strike_count as userStrikeCount', + 'p.apply_user_strike_count_config_to_children as applyUserStrikeCountConfigToChildren', + 'p.penalty', ] as const; type PolicyDbResult = FixKyselyRowCorrelation< @@ -66,6 +84,32 @@ export default class PolicyOperations { return results.map((it) => this.#dbResultToPolicy(it)); } + async getPoliciesByRuleIds(opts: { + ruleIds: readonly string[]; + readFromReplica?: boolean; + }): Promise> { + const { ruleIds, readFromReplica } = opts; + if (ruleIds.length === 0) { + return {}; + } + const pgQuery = this.#getPgQuery(readFromReplica ?? true); + type Row = PolicyDbResult & { ruleId: string }; + const rows = (await pgQuery + .selectFrom('public.rules_and_policies as rap') + .innerJoin('public.policies as p', 'p.id', 'rap.policy_id') + .select(policyJoinDbSelection) + .where('rap.rule_id', 'in', [...ruleIds]) + .execute()) as Row[]; + + const out: Record = {}; + for (const row of rows) { + const { ruleId, ...policyFields } = row; + const policy = this.#dbResultToPolicy(policyFields as PolicyDbResult); + (out[ruleId] ??= []).push(policy); + } + return out; + } + async getPolicy(opts: { orgId: string; policyId: string; diff --git a/server/services/moderationConfigService/modules/RuleReadOperations.ts b/server/services/moderationConfigService/modules/RuleReadOperations.ts new file mode 100644 index 0000000..adc76f1 --- /dev/null +++ b/server/services/moderationConfigService/modules/RuleReadOperations.ts @@ -0,0 +1,147 @@ +import { type Kysely, sql } from 'kysely'; + +import { + type RuleAlarmStatus, + RuleStatus, + RuleType, + type ConditionSet, +} from '../index.js'; +import { type ModerationConfigServicePg } from '../dbTypes.js'; +import { getUtcDateOnlyString } from '../../../utils/time.js'; +import { + type PlainRuleWithLatestVersion, + computeRuleStatusFromRow, +} from '../../../models/rules/ruleTypes.js'; + +const ruleSelect = [ + 'r.id', + 'r.name', + 'r.description', + 'r.status_if_unexpired as statusIfUnexpired', + 'r.tags', + 'r.max_daily_actions as maxDailyActions', + 'r.daily_actions_run as dailyActionsRun', + 'r.last_action_date as lastActionDate', + 'r.created_at as createdAt', + 'r.updated_at as updatedAt', + 'r.org_id as orgId', + 'r.creator_id as creatorId', + 'r.expiration_time as expirationTime', + 'r.condition_set as conditionSet', + 'r.alarm_status as alarmStatus', + 'r.alarm_status_set_at as alarmStatusSetAt', + 'r.rule_type as ruleType', + 'r.parent_id as parentId', + 'rlv.version as latestVersionString', +] as const; + +type RuleRow = { + id: string; + name: string; + description: string | null; + statusIfUnexpired: Exclude; + tags: string[]; + maxDailyActions: number | null; + dailyActionsRun: number; + lastActionDate: string | null; + createdAt: Date; + updatedAt: Date; + orgId: string; + creatorId: string; + expirationTime: Date | null; + conditionSet: ConditionSet; + // Kysely returns the Postgres enum as a plain string; cast in rowToPlainRuleWithLatest. + alarmStatus: string; + alarmStatusSetAt: Date; + ruleType: RuleType; + parentId: string | null; + latestVersionString: string | null; +}; + +function enabledQuotaWhere(today: string) { + return sql`(r.max_daily_actions is null or r.last_action_date is distinct from ${today}::date or (r.max_daily_actions is not null and r.daily_actions_run < r.max_daily_actions))`; +} + +function rowToPlainRuleWithLatest(row: RuleRow): PlainRuleWithLatestVersion { + const status = computeRuleStatusFromRow(row.expirationTime, row.statusIfUnexpired); + const version = row.latestVersionString ?? ''; + return { + id: row.id, + name: row.name, + description: row.description, + statusIfUnexpired: row.statusIfUnexpired, + status, + tags: row.tags, + maxDailyActions: row.maxDailyActions, + dailyActionsRun: row.dailyActionsRun, + lastActionDate: row.lastActionDate, + createdAt: row.createdAt, + updatedAt: row.updatedAt, + orgId: row.orgId, + creatorId: row.creatorId, + expirationTime: row.expirationTime, + conditionSet: row.conditionSet, + alarmStatus: row.alarmStatus as RuleAlarmStatus, + alarmStatusSetAt: row.alarmStatusSetAt, + ruleType: row.ruleType, + parentId: row.parentId, + latestVersion: { ruleId: row.id, version }, + }; +} + +export default class RuleReadOperations { + constructor( + private readonly pgQuery: Kysely, + private readonly pgQueryReplica: Kysely, + ) {} + + async getEnabledRulesForItemType(itemTypeId: string) { + const today = String(getUtcDateOnlyString()); + const rows = (await this.pgQueryReplica + .selectFrom('public.rules as r') + .innerJoin('public.rules_and_item_types as rit', 'rit.rule_id', 'r.id') + .leftJoin('public.rules_latest_versions as rlv', 'rlv.rule_id', 'r.id') + .select(ruleSelect) + .where('rit.item_type_id', '=', itemTypeId) + .where((eb) => + eb.or([ + eb('r.expiration_time', 'is', null), + eb('r.expiration_time', '>', sql`now()`), + ]), + ) + .where('r.status_if_unexpired', 'in', [ + RuleStatus.LIVE, + RuleStatus.BACKGROUND, + ]) + .where(enabledQuotaWhere(today)) + .execute()) as RuleRow[]; + + return rows.map(rowToPlainRuleWithLatest); + } + + async findEnabledUserRules() { + const today = String(getUtcDateOnlyString()); + const rows = (await this.pgQueryReplica + .selectFrom('public.rules as r') + .leftJoin('public.rules_latest_versions as rlv', 'rlv.rule_id', 'r.id') + .select(ruleSelect) + .where('r.rule_type', '=', RuleType.USER) + .where( + sql`not exists (select 1 from public.rules_and_item_types rit where rit.rule_id = r.id)`, + ) + .where((eb) => + eb.or([ + eb('r.expiration_time', 'is', null), + eb('r.expiration_time', '>', sql`now()`), + ]), + ) + .where('r.status_if_unexpired', 'in', [ + RuleStatus.LIVE, + RuleStatus.BACKGROUND, + ]) + .where(enabledQuotaWhere(today)) + .execute()) as RuleRow[]; + + return rows.map(rowToPlainRuleWithLatest); + } +} diff --git a/server/services/moderationConfigService/types/actions.ts b/server/services/moderationConfigService/types/actions.ts index 34b26b8..a0bdfb0 100644 --- a/server/services/moderationConfigService/types/actions.ts +++ b/server/services/moderationConfigService/types/actions.ts @@ -6,6 +6,8 @@ import { type JsonObject, type ReadonlyDeep, type Simplify } from 'type-fest'; import { type TaggedUnionFromCases } from '../../../utils/typescript-types.js'; +import { type UserPenaltySeverity } from './shared.js'; + export const ActionType = makeEnumLike([ 'CUSTOM_ACTION', 'ENQUEUE_TO_MRT', @@ -28,6 +30,7 @@ type AnyAction = ReadonlyDeep< orgId: string; name: string; applyUserStrikes: boolean; + penalty: UserPenaltySeverity; } & TaggedUnionFromCases< { actionType: ActionType }, { diff --git a/server/services/moderationConfigService/types/policies.ts b/server/services/moderationConfigService/types/policies.ts index 00c9be7..2bbdf00 100644 --- a/server/services/moderationConfigService/types/policies.ts +++ b/server/services/moderationConfigService/types/policies.ts @@ -1,5 +1,7 @@ import { makeEnumLike } from '@roostorg/types'; +import { type UserPenaltySeverity } from './shared.js'; + export const PolicyType = makeEnumLike([ 'HATE', 'VIOLENCE', @@ -30,5 +32,5 @@ export type Policy = { semanticVersion: number; userStrikeCount: number; applyUserStrikeCountConfigToChildren: boolean; - penalty: string; // TODO: remove + penalty: UserPenaltySeverity; }; diff --git a/server/services/orgAwareSignalExecutionService/signalExecutionService.ts b/server/services/orgAwareSignalExecutionService/signalExecutionService.ts index 925d414..89ad951 100644 --- a/server/services/orgAwareSignalExecutionService/signalExecutionService.ts +++ b/server/services/orgAwareSignalExecutionService/signalExecutionService.ts @@ -4,7 +4,7 @@ import stringify from 'safe-stable-stringify'; import { type ReadonlyDeep } from 'type-fest'; import { inject } from '../../iocContainer/utils.js'; -import { type PolicyActionPenalties } from '../../models/OrgModel.js'; +import { type PolicyActionPenalties } from '../policyActionPenalties.js'; import { type MatchingValues } from '../../models/rules/matchingValues.js'; import { type LocationArea } from '../../models/types/locationArea.js'; import { jsonStringify } from '../../utils/encoding.js'; diff --git a/server/services/policyActionPenalties.ts b/server/services/policyActionPenalties.ts new file mode 100644 index 0000000..cdc3e36 --- /dev/null +++ b/server/services/policyActionPenalties.ts @@ -0,0 +1,73 @@ +import { + UserPenaltySeverity, + type ModerationConfigService, +} from './moderationConfigService/index.js'; + +export type PolicyActionPenalties = { + actionId: string; + policyId: string; + penalties: number[]; +}; + +/** + * Computes the severity of the penalty we should apply for a given + * (action, policy) pair. The general idea is to make the penalties + * increase exponentially as severity levels increase, but the rate + * of increase can't be so high that a (severe, severe) penalty is + * 50x higher than a (high, high) penalty. + * + * The easiest way to achieve this exponential behavior is at the individual + * severity levels, rather than trying to multiply the action penalty + * by the severity penalty to compound their magnitudes. So the severity + * levels apply penalty magnitudes as follows: + * + * NONE = 0 + * LOW = 1 + * MEDIUM = 3 + * HIGH = 9 + * SEVERE = 27 + * + * To get the penalty value for an (action, policy) pair, we just add the + * penalty values of the action and policy because the exponential nature + * of these penalties has already been taken into account. + * + * If the action has no penalty (e.g., "Send to Moderation", "Restore + * Content"), we never apply any penalty, regardless of the policy penalty. + * Otherwise, the penalty accounts for both the action + policy penalties. + */ +export function computeActionPolicyPenalty( + actionPenalty: UserPenaltySeverity, + policyPenalty: UserPenaltySeverity, +): number { + const penaltySeverityMap: { [k in UserPenaltySeverity]: number } = { + [UserPenaltySeverity.NONE]: 0, + [UserPenaltySeverity.LOW]: 1, + [UserPenaltySeverity.MEDIUM]: 3, + [UserPenaltySeverity.HIGH]: 9, + [UserPenaltySeverity.SEVERE]: 27, + }; + + return actionPenalty === UserPenaltySeverity.NONE + ? 0 + : penaltySeverityMap[actionPenalty] + penaltySeverityMap[policyPenalty]; +} + +export async function getPolicyActionPenaltiesForOrg( + moderationConfigService: ModerationConfigService, + orgId: string, +): Promise { + const [actions, policies] = await Promise.all([ + moderationConfigService.getActions({ orgId, readFromReplica: true }), + moderationConfigService.getPolicies({ orgId, readFromReplica: true }), + ]); + + return policies.flatMap((policy) => + actions.map((action) => ({ + actionId: action.id, + policyId: policy.id, + penalties: [ + computeActionPolicyPenalty(action.penalty, policy.penalty), + ], + })), + ); +} diff --git a/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.test.ts b/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.test.ts index d1a8673..f520910 100644 --- a/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.test.ts +++ b/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.test.ts @@ -1,25 +1,83 @@ -import _ from 'lodash'; - import getBottle, { type Dependencies, type PublicInterface, } from '../../iocContainer/index.js'; -import { type Rule as TRule } from '../../models/rules/RuleModel.js'; import { type NotificationsService } from '../../services/notificationsService/notificationsService.js'; import { type GetCurrentPeriodRuleAlarmStatuses } from '../../services/ruleAnomalyDetectionService/getCurrentPeriodRuleAlarmStatuses.js'; import createOrg from '../../test/fixtureHelpers/createOrg.js'; import createRule from '../../test/fixtureHelpers/createRule.js'; import createUser from '../../test/fixtureHelpers/createUser.js'; -import { mocked, type Mocked } from '../../test/mockHelpers/jestMocks.js'; +import { type Mocked } from '../../test/mockHelpers/jestMocks.js'; import { RuleAlarmStatus } from '../moderationConfigService/index.js'; import DetectRulePassRateAnomaliesJob from './detectRulePassRateAnomaliesJob.js'; +function makeMockKyselyForRules( + fakeRules: Array<{ + id: string; + orgId: string; + creatorId: string; + name: string; + alarmStatus: RuleAlarmStatus; + statusIfUnexpired: string; + }>, + orgRows: Array<{ id: string; on_call_alert_email: string | null }>, +) { + const updateExecute = jest.fn().mockResolvedValue(undefined); + const mockDb = { + selectFrom: jest.fn((table: string) => { + const chain: { + select: jest.Mock; + where: jest.Mock; + execute: jest.Mock; + } = { + select: jest.fn(), + where: jest.fn(), + execute: jest.fn(), + }; + chain.select.mockReturnValue(chain); + chain.where.mockReturnValue(chain); + chain.execute.mockImplementation(async () => { + if (table === 'public.rules') { + return fakeRules.map((r) => ({ + id: r.id, + org_id: r.orgId, + creator_id: r.creatorId, + name: r.name, + alarm_status: r.alarmStatus, + status_if_unexpired: r.statusIfUnexpired, + })); + } + if (table === 'public.orgs') { + return orgRows; + } + return []; + }); + return chain; + }), + updateTable: jest.fn(() => ({ + set: jest.fn().mockReturnValue({ + where: jest.fn().mockReturnValue({ + execute: updateExecute, + }), + }), + })), + __updateExecute: updateExecute, + }; + return mockDb; +} + describe('Detect Rule Anomalies', () => { describe('worker', () => { - let OrgModel: Dependencies['Sequelize']['Org'], - deleteMockData: () => Promise, - mockDummyRules: Mocked[], - mockRuleModel: Mocked, + let deleteMockData: () => Promise, + mockDummyRules: Array<{ + id: string; + orgId: string; + creatorId: string; + name: string; + alarmStatus: RuleAlarmStatus; + statusIfUnexpired: string; + }>, + mockKysely: ReturnType, mockGetCurrentPeriodRuleAlarmStatuses: GetCurrentPeriodRuleAlarmStatuses, mockNotificationsService: Mocked< PublicInterface, @@ -33,7 +91,6 @@ describe('Detect Rule Anomalies', () => { ModerationConfigService, ApiKeyService, } = (await getBottle()).container; - OrgModel = models.Org; // make some fake rules (w/ stable ids so we can match them in a snapshot) // in different initial alarm statuses, to test all 9 combinations [i.e., @@ -107,7 +164,15 @@ describe('Detect Rule Anomalies', () => { }), ]); - mockDummyRules = fakeRules.map((it) => mocked(it, ['save'])); + mockDummyRules = fakeRules.map((r) => ({ + id: r.id, + orgId: r.orgId, + creatorId: r.creatorId, + name: r.name, + alarmStatus: r.alarmStatus, + statusIfUnexpired: r.statusIfUnexpired, + })); + mockGetCurrentPeriodRuleAlarmStatuses = async () => { const newAlarmStatusByRule = // eslint-disable-next-line @typescript-eslint/consistent-type-assertions @@ -119,35 +184,27 @@ describe('Detect Rule Anomalies', () => { i % 3 === 0 ? RuleAlarmStatus.ALARM : i % 3 === 1 - ? RuleAlarmStatus.OK - : RuleAlarmStatus.INSUFFICIENT_DATA, + ? RuleAlarmStatus.OK + : RuleAlarmStatus.INSUFFICIENT_DATA, meta: { lastPeriodPassRate: 0.5, secondToLastPeriodPassRate: 0.4 }, }; }); return newAlarmStatusByRule; }; - mockNotificationsService = mocked( - { - async createNotifications(_it: any) {}, - async getNotificationsForUser() { - return []; - }, - }, - ['createNotifications'], - ); + mockNotificationsService = { + createNotifications: jest.fn(), + getNotificationsForUser: jest.fn(), + } as unknown as Mocked< + PublicInterface, + 'createNotifications' + >; - mockRuleModel = mocked(models.Rule, ['findAll']); - mockRuleModel.findAll.mockResolvedValue(mockDummyRules); + mockKysely = makeMockKyselyForRules(mockDummyRules, [ + { id: org.id, on_call_alert_email: null }, + { id: org2.id, on_call_alert_email: 'test@gmail.com' }, + ]); - // It might be nice if we could create the mock data at the start of a - // transaction, run the tests with that transaction open, and then just - // roll it back at the end to automatically delete and leave the db in a - // consistent/clean state. That's a little tricky, though, as it requires - // feeding the transacation object (or keeping a managed transaction - // callback open) all the way into calling the worker. So, instead, we - // settle for manually defining this compensating transacaction, which we - // call at the end. deleteMockData = async () => { await Promise.all(fakeRules.map(async (it) => it.destroy())); await Promise.all([ruleOwner.destroy(), ruleOwner2.destroy()]); @@ -163,17 +220,13 @@ describe('Detect Rule Anomalies', () => { test('should generate the proper notifications + update rules', async () => { const worker = DetectRulePassRateAnomaliesJob( - mockRuleModel, - OrgModel, + mockKysely as unknown as Dependencies['KyselyPg'], mockNotificationsService, mockGetCurrentPeriodRuleAlarmStatuses, jest.fn<() => Promise>(), ); await worker.run(); - // We should've sent 4 notifications: one for each of the two rules that - // was in alarm and transitioned to 'not alarm' (ok or insufficient data), - // and one for each of the rules that was in 'not alarm' and went to alarm. const mockCreateNotifications = mockNotificationsService.createNotifications; @@ -255,17 +308,11 @@ describe('Detect Rule Anomalies', () => { ] `); - // Except for the rules that stayed the same state (indexes 0, 4, 8), all - // the rules should've been saved with their new state. - const newRuleStatuses = await mockGetCurrentPeriodRuleAlarmStatuses(); - mockDummyRules.forEach((rule, i) => { - if (![0, 4, 8].includes(i)) { - expect(rule.save).toHaveBeenCalledTimes(1); - expect(rule.alarmStatus).toEqual(newRuleStatuses[rule.id].status); - } else { - expect(rule.save).toHaveBeenCalledTimes(0); - } - }); + await mockGetCurrentPeriodRuleAlarmStatuses(); + const expectedUpdates = mockDummyRules.filter( + (_rule, i) => ![0, 4, 8].includes(i), + ).length; + expect(mockKysely.updateTable).toHaveBeenCalledTimes(expectedUpdates); }); }); }); diff --git a/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.ts b/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.ts index 34cc2bd..9afbdda 100644 --- a/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.ts +++ b/server/services/ruleAnomalyDetectionService/detectRulePassRateAnomaliesJob.ts @@ -7,17 +7,20 @@ import { NotificationType } from '../notificationsService/notificationsService.j const { capitalize, keyBy } = lodash; +type OrgAlertRow = { + id: string; + on_call_alert_email: string | null; +}; + export default inject( [ - 'RuleModel', - 'OrgModel', + 'KyselyPg', 'NotificationsService', 'getCurrentPeriodRuleAlarmStatuses', 'closeSharedResourcesForShutdown', ], ( - Rule, - Org, + db, notificationsService, getCurrentPeriodRuleAlarmStatuses, sharedResourceShutdown, @@ -27,40 +30,54 @@ export default inject( const now = new Date(); const newAlarmStatusByRule = await getCurrentPeriodRuleAlarmStatuses(); - // TODO: at some point, we might have to chunk this, - // but we're very far from that right now. We also don't have to use a - // transaction, since there's no risk of concurrent updates to rule.alarmStatus. const ruleIds = Object.keys(newAlarmStatusByRule); - const rules = await Rule.findAll({ where: { id: ruleIds } }); + if (ruleIds.length === 0) { + return; + } + + const rules = await db + .selectFrom('public.rules') + .select([ + 'id', + 'org_id', + 'creator_id', + 'name', + 'alarm_status', + 'status_if_unexpired', + ]) + .where('id', 'in', ruleIds) + .execute(); + const alarmStatusChangedRules = rules.filter( - (rule) => rule.alarmStatus !== newAlarmStatusByRule[rule.id].status, + (rule) => rule.alarm_status !== newAlarmStatusByRule[rule.id].status, ); - const changedRuleOrgIds = alarmStatusChangedRules.map((it) => it.orgId); - const orgsForChangedRules = changedRuleOrgIds.length - ? keyBy( - await Org.findAll({ - where: { id: [...new Set(changedRuleOrgIds)] }, - }), - (it) => it.id, - ) - : {}; + const changedRuleOrgIds = alarmStatusChangedRules.map((it) => it.org_id); + let orgsForChangedRules: Record = {}; + if (changedRuleOrgIds.length > 0) { + const orgRows = (await db + .selectFrom('public.orgs') + .select(['id', 'on_call_alert_email']) + .where('id', 'in', [...new Set(changedRuleOrgIds)]) + .execute()) as OrgAlertRow[]; + orgsForChangedRules = keyBy(orgRows, (r) => r.id); + } - // Notify the creator of each rule, and the on call alert email (if any). const notifications = alarmStatusChangedRules .filter( - // Only alert if we're coming into an ALARM, or going out of one. - // If we transitioned (e.g.) from "OK" to "INSUFFICIENT_DATA", b/c a - // rule's conditions got updated, we don't care about that. Similarly, - // we don't care about INSUFFICIENT_DATA going to "OK". (it) => - it.alarmStatus === RuleAlarmStatus.ALARM || + it.alarm_status === RuleAlarmStatus.ALARM || newAlarmStatusByRule[it.id].status === RuleAlarmStatus.ALARM, ) - .map((rule) => { + .flatMap((rule) => { const ruleNowInAlarm = newAlarmStatusByRule[rule.id].status === RuleAlarmStatus.ALARM; + const orgRow = orgsForChangedRules[rule.org_id]; + if (!orgRow) { + return []; + } + return { type: ruleNowInAlarm ? NotificationType.RulePassRateIncreaseAnomalyStart @@ -76,19 +93,21 @@ export default inject( message: `${ ruleNowInAlarm ? `[Alarm Triggered - ${capitalize( - rule.statusIfUnexpired, + String(rule.status_if_unexpired), + )} Rule]` + : `[Alarm Cleared - ${capitalize( + String(rule.status_if_unexpired), )} Rule]` - : `[Alarm Cleared - ${capitalize(rule.statusIfUnexpired)} Rule]` } ${rule.name} has ${ ruleNowInAlarm ? 'started' : 'stopped' } passing at an anomalous rate.`, recipients: [ - { type: 'user_id', value: rule.creatorId }, - ...(orgsForChangedRules[rule.orgId].onCallAlertEmail + { type: 'user_id' as const, value: rule.creator_id }, + ...(orgRow.on_call_alert_email ? [ { type: 'email_address' as const, - value: orgsForChangedRules[rule.orgId].onCallAlertEmail!, + value: orgRow.on_call_alert_email, }, ] : []), @@ -97,9 +116,14 @@ export default inject( }); const ruleUpdateTasks = alarmStatusChangedRules.map(async (rule) => { - rule.alarmStatus = newAlarmStatusByRule[rule.id].status; - rule.alarmStatusSetAt = now; - return rule.save(); + await db + .updateTable('public.rules') + .set({ + alarm_status: newAlarmStatusByRule[rule.id].status, + alarm_status_set_at: now, + }) + .where('id', '=', rule.id) + .execute(); }); await Promise.all([ diff --git a/server/services/signalsService/signals/SignalBase.ts b/server/services/signalsService/signals/SignalBase.ts index 8be6176..7a51a41 100644 --- a/server/services/signalsService/signals/SignalBase.ts +++ b/server/services/signalsService/signals/SignalBase.ts @@ -6,7 +6,7 @@ import { } from '@roostorg/types'; import { type ReadonlyDeep, type Simplify } from 'type-fest'; -import { type PolicyActionPenalties } from '../../../models/OrgModel.js'; +import { type PolicyActionPenalties } from '../../policyActionPenalties.js'; import { type TaggedItemData } from '../../../models/rules/item-type-fields.js'; import { type CoopError } from '../../../utils/errors.js'; import { type Language } from '../../../utils/language.js'; diff --git a/server/services/userManagementService/dbTypes.ts b/server/services/userManagementService/dbTypes.ts index b6e4478..90f0abc 100644 --- a/server/services/userManagementService/dbTypes.ts +++ b/server/services/userManagementService/dbTypes.ts @@ -59,7 +59,7 @@ export type UserManagementPg = { 'public.users': { id: GeneratedAlways; email: string; - password: string; + password: string | null; first_name: string; last_name: string; role: UserRole; @@ -68,6 +68,7 @@ export type UserManagementPg = { created_at: GeneratedAlways; updated_at: GeneratedAlways; org_id: string; + login_methods: ('password' | 'saml')[]; }; 'public.invite_user_tokens': { id: GeneratedAlways; diff --git a/server/services/userStatisticsService/computeUserScore.ts b/server/services/userStatisticsService/computeUserScore.ts index 23992e5..08c9203 100644 --- a/server/services/userStatisticsService/computeUserScore.ts +++ b/server/services/userStatisticsService/computeUserScore.ts @@ -1,7 +1,7 @@ import _ from 'lodash'; import { type ReadonlyDeep } from 'type-fest'; -import { type PolicyActionPenalties } from '../../models/OrgModel.js'; +import { type PolicyActionPenalties } from '../policyActionPenalties.js'; import { jsonStringify, type JsonOf } from '../../utils/encoding.js'; import { unzip2 } from '../../utils/fp-helpers.js'; import { type UserActionStatistics } from './fetchUserActionStatistics.js'; diff --git a/server/services/userStatisticsService/userStatisticsService.ts b/server/services/userStatisticsService/userStatisticsService.ts index e180a33..1a81d45 100644 --- a/server/services/userStatisticsService/userStatisticsService.ts +++ b/server/services/userStatisticsService/userStatisticsService.ts @@ -7,7 +7,7 @@ import { sql, type Kysely } from 'kysely'; import { type ReadonlyDeep } from 'type-fest'; import { inject, type Dependencies } from '../../iocContainer/index.js'; -import { type PolicyActionPenalties } from '../../models/OrgModel.js'; +import { type PolicyActionPenalties } from '../policyActionPenalties.js'; import { initialUserScore, type UserScore } from './computeUserScore.js'; import { type UserStatisticsServicePg, diff --git a/server/services/userStrikeService/userStrikeService.test.ts b/server/services/userStrikeService/userStrikeService.test.ts index cc09751..01e3ac6 100644 --- a/server/services/userStrikeService/userStrikeService.test.ts +++ b/server/services/userStrikeService/userStrikeService.test.ts @@ -65,6 +65,7 @@ describe('Item Investigation Service', () => { name: 'testAction1', applyUserStrikes: true, orgId: 'fakeOrgId', + penalty: 'NONE' as const, callbackUrl: 'fakeCallbackUrl1', callbackUrlHeaders: null, callbackUrlBody: null, @@ -106,6 +107,7 @@ describe('Item Investigation Service', () => { name: 'testAction1', applyUserStrikes: false, orgId: 'fakeOrgId', + penalty: 'NONE' as const, callbackUrl: 'fakeCallbackUrl1', callbackUrlHeaders: null, callbackUrlBody: null, @@ -137,6 +139,7 @@ describe('Item Investigation Service', () => { name: 'testAction1', applyUserStrikes: false, orgId: 'fakeOrgId', + penalty: 'NONE' as const, callbackUrl: 'fakeCallbackUrl1', callbackUrlHeaders: null, callbackUrlBody: null, diff --git a/server/utils/kyselyTransactionWithRetry.ts b/server/utils/kyselyTransactionWithRetry.ts new file mode 100644 index 0000000..6da7fd2 --- /dev/null +++ b/server/utils/kyselyTransactionWithRetry.ts @@ -0,0 +1,36 @@ +import { type Kysely } from 'kysely'; + +import { safeGet } from './misc.js'; + +function isSerializationFailure(error: unknown): boolean { + return safeGet(error, ['code']) === '40001'; +} + +/** + * Like {@link server/models/sequelizeSetup.ts maketransactionWithRetry} but for Kysely. + */ +export function makeKyselyTransactionWithRetry(kysely: Kysely) { + return async function transactionWithRetry( + callback: (trx: Kysely) => Promise, + ): Promise { + let remainingTries = 3; + let lastError: unknown; + while (remainingTries > 0) { + remainingTries -= 1; + try { + return await kysely.transaction().execute(callback); + } catch (e: unknown) { + if (!isSerializationFailure(e)) { + throw e; + } + lastError = e; + } + } + + throw lastError; + }; +} + +export type KyselyTransactionWithRetry = ReturnType< + typeof makeKyselyTransactionWithRetry +>; diff --git a/server/utils/sql.test.ts b/server/utils/sql.test.ts index 05d6467..2703671 100644 --- a/server/utils/sql.test.ts +++ b/server/utils/sql.test.ts @@ -1,55 +1,76 @@ -import * as knexPkg from 'knex'; +import { Kysely, PostgresDialect, type PostgresQueryResult } from 'kysely'; import { takeLast } from './sql.js'; -const { knex: Knex } = knexPkg.default; +function makeCompileOnlyDb>>() { + return new Kysely({ + dialect: new PostgresDialect({ + pool: { + async connect() { + return { + query: jest.fn().mockResolvedValue({ + rows: [], + command: 'SELECT', + rowCount: 0, + } as PostgresQueryResult), + async release() {}, + }; + }, + async end() {}, + }, + }), + }); +} describe('Sql Helpers', () => { describe('takeLast', () => { test('should work for simple queries', () => { type User = { id: string; name: string; email: string }; + type TestDb = { users: User }; - const knex = Knex({ dialect: 'postgres' }); - const users = knex('users').select('id', 'name'); + const db = makeCompileOnlyDb(); + const users = db.selectFrom('users').select(['id', 'name']); - const result = takeLast(users, [{ column: 'id', order: 'desc' }], 2); + const result = takeLast(db, users, [{ column: 'id', order: 'desc' }], 2); - expect(result.toString()).toEqual( - 'select * from (select "id", "name" from "users" order by "id" asc limit 2) as "dc2d41a9-082e-48b0-a66f-345a22696b02" order by "id" desc', + expect(result.compile().sql).toEqual( + 'select * from (select "id", "name" from "users" order by "id" asc limit $1) as "dc2d41a9-082e-48b0-a66f-345a22696b02" order by "id" desc', ); + expect(result.compile().parameters).toEqual([2]); }); test('should work for arbitrarily complex queries', () => { - // We'll test this with the real query we use for backtesting. - // This is still only one case (notably, with no joins), but at least it - // uses aliases and a WHERE, so it'll give us a bit more confidence. - const knex = Knex({ dialect: 'postgres' }); - type Result = { - orgId: string; - ts: string; - content: string; - correlationId: string; + type RuleExecRow = { + ORG_ID: string; + TS: string; + CONTENT: string; + CORRELATION_ID: string; }; + type TestDb = { RULE_EXECUTIONS: RuleExecRow }; - const backtestResults = knex('RULE_EXECUTIONS') - .select({ - orgId: 'ORG_ID', - ts: 'TS', - content: 'CONTENT', - correlationId: 'CORRELATION_ID', - }) + const db = makeCompileOnlyDb(); + const backtestResults = db + .selectFrom('RULE_EXECUTIONS') + .select([ + 'ORG_ID as orgId', + 'TS as ts', + 'CONTENT as content', + 'CORRELATION_ID as correlationId', + ]) .where('CORRELATION_ID', '=', '47') - .andWhere('TS', '>', '2019-01-01'); + .where('TS', '>', '2019-01-01'); const result = takeLast( + db, backtestResults, [{ column: 'ts', order: 'asc' }], 50, ); - expect(result.toString()).toMatchInlineSnapshot( - `"select * from (select "ORG_ID" as "orgId", "TS" as "ts", "CONTENT" as "content", "CORRELATION_ID" as "correlationId" from "RULE_EXECUTIONS" where "CORRELATION_ID" = '47' and "TS" > '2019-01-01' order by "ts" desc limit 50) as "dc2d41a9-082e-48b0-a66f-345a22696b02" order by "ts" asc"`, + expect(result.compile().sql).toMatchInlineSnapshot( + `"select * from (select "ORG_ID" as "orgId", "TS" as "ts", "CONTENT" as "content", "CORRELATION_ID" as "correlationId" from "RULE_EXECUTIONS" where "CORRELATION_ID" = $1 and "TS" > $2 order by "ts" desc limit $3) as "dc2d41a9-082e-48b0-a66f-345a22696b02" order by "ts" asc"`, ); + expect(result.compile().parameters).toEqual(['47', '2019-01-01', 50]); }); }); }); diff --git a/server/utils/sql.ts b/server/utils/sql.ts index 2b4d393..5285183 100644 --- a/server/utils/sql.ts +++ b/server/utils/sql.ts @@ -1,7 +1,4 @@ -import * as knexPkg from 'knex'; -import type { Knex } from 'knex'; - -const { knex } = knexPkg.default; +import type { Kysely, SelectQueryBuilder } from 'kysely'; /** * When paginating backwards (i.e., the user's on page 5 and asks for page 4 @@ -17,10 +14,10 @@ const { knex } = knexPkg.default; * last 2" generically, you need to sort the items in reverse order, apply a * LIMIT, then reverse the result. * - * To do that generically, we need to use a query builder that'll let us build - * queries programmatically/work with them as data structures, so we use knex. - * This function, then implements the "take last n" operation given an unsorted - * knex select query and some sort criteria. + * This function implements that pattern on a Kysely select query. + * + * @param db Kysely instance used only to build the outer `select * from (…)`. + * It must use the same dialect as `unsortedSelectQuery`. * * @param unsortedSelectQuery The query that selects the set of items (without * them being sorted) from which we want to take the last n, after sorting. @@ -29,56 +26,47 @@ const { knex } = knexPkg.default; * returned by unsortedSelectQuery, before we can take the last n items. * NB: the column names provided here must refer to one of the columns * selected by `unsortedSelectQuery`, under that column's final alias. E.g., - * if unsortedSelectQuery is `SELECT a as "hello" from table`, then you can - * only provide "hello" as the sort criteria; not "a", and not some unselected - * column "b". + * if the select is `DS as date`, then sort criteria must use `date`, not + * `DS`. * * @param size How many items to take. * - * @param client The name of the knex client to use. This effects the - * SQL-dialect-specific settings that knex might apply to the generated query. - * NB: these dialect-specific settings potentially include data escaping rules - * that could be relevant for SQL injection. - * - * @param subqueryAlias Internally, this query generates a subquery, and SQL - * mandates that that subquery be given an alias. In theory, there's maybe - * some risk of that alias - * - * @returns A new knex query that selects the last n items, after sorting. + * @returns A Kysely query that selects the last n items, after sorting. */ -export function takeLast( - unsortedSelectQuery: Knex.QueryBuilder, - sortCriteria: { - column: (keyof T & string) | Knex.Raw; - order: 'desc' | 'asc'; - }[], +const SUBQUERY_ALIAS = 'dc2d41a9-082e-48b0-a66f-345a22696b02'; + +export function takeLast< + DB, + TB extends keyof DB, + O extends Record, +>( + db: Kysely, + unsortedSelectQuery: SelectQueryBuilder, + sortCriteria: readonly { column: keyof O & string; order: 'desc' | 'asc' }[], size: number, - client: string = 'pg', ) { - // SQL requires that the subquery we create have an alias. I don't _think_ - // there's risk of that name causing a naming conflict, but I haven't thought - // too hard about all the scoping implications, so, to be safe, we give this - // alias a very-unlikely-to-conflict name. - const subqueryAlias = 'dc2d41a9-082e-48b0-a66f-345a22696b02'; - - const inner = unsortedSelectQuery - .clone() - .orderBy( - sortCriteria.map((it) => ({ - // Cast here is because I think the knex typings are just wrong. - // They suggest that `column` has to be a string, but, actually, - // we can sort on arbitrary expressesions contained in a `knex.raw`. - column: it.column as keyof T & string, - order: it.order === 'desc' ? ('asc' as const) : ('desc' as const), - })), - ) - .limit(size); - - return knex({ client }) - .select('*') - .from(inner.as(subqueryAlias)) - .orderBy( - // Cast here is same as above. - sortCriteria as ((typeof sortCriteria)[number] & { column: string })[], + let inner = unsortedSelectQuery.clearOrderBy(); + for (const it of sortCriteria) { + inner = inner.orderBy( + it.column, + it.order === 'desc' ? 'asc' : 'desc', ); + } + inner = inner.limit(size); + + // Chaining `orderBy` in a loop widens `outer` to an incompatible union; the + // builder is still the same concrete Kysely select at runtime. + let outer = db.selectFrom(inner.as(SUBQUERY_ALIAS)).selectAll() as SelectQueryBuilder< + DB & { [K in typeof SUBQUERY_ALIAS]: O }, + typeof SUBQUERY_ALIAS, + O + >; + for (const it of sortCriteria) { + outer = outer.orderBy(it.column, it.order) as typeof outer; + } + return outer as SelectQueryBuilder< + DB & { [K in typeof SUBQUERY_ALIAS]: O }, + typeof SUBQUERY_ALIAS, + O + >; } diff --git a/server/workers_jobs/RunUserRulesJob.ts b/server/workers_jobs/RunUserRulesJob.ts index b30f637..2986b49 100644 --- a/server/workers_jobs/RunUserRulesJob.ts +++ b/server/workers_jobs/RunUserRulesJob.ts @@ -11,15 +11,21 @@ export default inject( 'RuleEngine', 'UserStatisticsService', 'closeSharedResourcesForShutdown', - 'RuleModel', + 'ModerationConfigService', 'getItemTypeEventuallyConsistent', ], - (RuleEngine, userStatisticsService, sharedResourceShutdown, Rule, getItemTypeEventuallyConsistent) => ({ + ( + RuleEngine, + userStatisticsService, + sharedResourceShutdown, + moderationConfigService, + getItemTypeEventuallyConsistent, + ) => ({ type: 'Job' as const, async run() { // TODO: we may have to do only some orgs per job run at some point. // For now, though, this is fine. - const userRules = await Rule.findEnabledUserRules(); + const userRules = await moderationConfigService.findEnabledUserRules(); if (!userRules.length) { return; -- 2.51.2