diff --git a/.github/ISSUE_TEMPLATE/config.yml b/.github/ISSUE_TEMPLATE/config.yml index c8fda80..f80aee2 100644 --- a/.github/ISSUE_TEMPLATE/config.yml +++ b/.github/ISSUE_TEMPLATE/config.yml @@ -1,9 +1,9 @@ blank_issues_enabled: false # maintainers can still file blank issues for convenience contact_links: - name: Discuss Coop - about: "See previously-asked questions and existing discussions about Coop, or start your own" + about: 'See previously-asked questions and existing discussions about Coop, or start your own' url: https://github.com/roostorg/coop/discussions - name: Chat with contributors - about: "Join the #coop channel on the ROOST Discord server" + about: 'Join the #coop channel on the ROOST Discord server' url: https://discord.gg/2r7uUqvsut diff --git a/docs/integrations/ncmec.md b/docs/integrations/ncmec.md index 3d27049..14e4e3b 100644 --- a/docs/integrations/ncmec.md +++ b/docs/integrations/ncmec.md @@ -233,27 +233,27 @@ Coop signs every request with your org's signing key. Verify the signature befor #### Response fields -| Field | Type | Description | -| ------------------------- | --------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `users` | Array | Must include an entry for every user in the request. | -| `users.id` | String | Must match the `id` from the request. | -| `users.typeId` | String | Must match the `typeId` from the request. | -| `users.screenName` | String | The user's screen name or username on your platform. | -| `users.email` | Array | Known email addresses for the user. `type` may be `Business`, `Home`, or `Work`. | +| Field | Type | Description | +| ------------------------- | --------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `users` | Array | Must include an entry for every user in the request. | +| `users.id` | String | Must match the `id` from the request. | +| `users.typeId` | String | Must match the `typeId` from the request. | +| `users.screenName` | String | The user's screen name or username on your platform. | +| `users.email` | Array | Known email addresses for the user. `type` may be `Business`, `Home`, or `Work`. | | `users.ipCaptureEvent` | Array | IP events associated with the user (e.g. logins, registrations). `eventName` may be `Login`, `Registration`, `Purchase`, `Upload`, `Other`, or `Unknown`. Coop also appends the user item's `ipAddress` field role (when mapped) as an `Unknown` event at the report's incident time. | -| `users.data` | Object | Raw item data for the user. | -| `media` | Array | Must include an entry for every media item in the request if present. | -| `media.id` | String | Must match the `id` from the request. | -| `media.typeId` | String | Must match the `typeId` from the request. | -| `media.missing` | Boolean | Set to `true` if the media is no longer available. If **all** media items are `missing: true`, no CyberTip is filed. | -| `media.publiclyAvailable` | Boolean | Whether the media was publicly accessible on your platform at the time of reporting. | -| `media.fileName` | String | Original filename of the media. | -| `media.additionalInfo` | Array\ | Additional context about the media to include in the NCMEC file details. | -| `media.ipCaptureEvent` | Array | IP events associated with this media item. Coop also appends the media item's `ipAddress` field role (when mapped) as an `Upload` event at the media's `createdAt`. | -| `media.fileDetails` | Object | Hash information for the file: `{ hash, hashType }`. | -| `additionalFiles` | Array | Extra files to upload to NCMEC as supplemental evidence (e.g. screenshots). | -| `messages` | Array | Message-level IP address data for conversation thread context. | -| `additionalInfo` | String | Top-level freeform additional information to include in the report. | +| `users.data` | Object | Raw item data for the user. | +| `media` | Array | Must include an entry for every media item in the request if present. | +| `media.id` | String | Must match the `id` from the request. | +| `media.typeId` | String | Must match the `typeId` from the request. | +| `media.missing` | Boolean | Set to `true` if the media is no longer available. If **all** media items are `missing: true`, no CyberTip is filed. | +| `media.publiclyAvailable` | Boolean | Whether the media was publicly accessible on your platform at the time of reporting. | +| `media.fileName` | String | Original filename of the media. | +| `media.additionalInfo` | Array\ | Additional context about the media to include in the NCMEC file details. | +| `media.ipCaptureEvent` | Array | IP events associated with this media item. Coop also appends the media item's `ipAddress` field role (when mapped) as an `Upload` event at the media's `createdAt`. | +| `media.fileDetails` | Object | Hash information for the file: `{ hash, hashType }`. | +| `additionalFiles` | Array | Extra files to upload to NCMEC as supplemental evidence (e.g. screenshots). | +| `messages` | Array | Message-level IP address data for conversation thread context. | +| `additionalInfo` | String | Top-level freeform additional information to include in the report. | > [!IMPORTANT] > If the response does not include an entry for every user and media item in the request, Coop will throw an error and not submit the CyberTip. Your endpoint must return a response entry for each requested user and media item. diff --git a/server/services/ncmecService/buildSubmitReportParamsFromDecision.ts b/server/services/ncmecService/buildSubmitReportParamsFromDecision.ts index 7c0e376..4e11086 100644 --- a/server/services/ncmecService/buildSubmitReportParamsFromDecision.ts +++ b/server/services/ncmecService/buildSubmitReportParamsFromDecision.ts @@ -150,8 +150,7 @@ export async function buildSubmitReportParamsFromDecision( // role mapped to a non-date column). Better to send the retry // timestamp than to permanently block the report. const createdAt = - roleCreatedAt !== undefined && - !Number.isNaN(Date.parse(roleCreatedAt)) + roleCreatedAt !== undefined && !Number.isNaN(Date.parse(roleCreatedAt)) ? roleCreatedAt : makeDateString(new Date().toISOString()); if (createdAt === undefined) { diff --git a/server/test/fixtureHelpers/createRule.ts b/server/test/fixtureHelpers/createRule.ts index e2fc9fb..038b5db 100644 --- a/server/test/fixtureHelpers/createRule.ts +++ b/server/test/fixtureHelpers/createRule.ts @@ -48,6 +48,8 @@ export default async function createRule( ruleType?: RuleType; status?: RuleStatus; conditionSet?: ConditionSet; + actionIds?: readonly string[]; + contentTypeIds?: readonly string[]; } = {}, ) { const ruleId = extra.id ?? uid(); @@ -72,9 +74,9 @@ export default async function createRule( orgId, ruleType, parentId: null, - actionIds: [], + actionIds: extra.actionIds ?? [], policyIds: [], - contentTypeIds: [], + contentTypeIds: extra.contentTypeIds ?? [], }).catch(logErrorAndThrow); // `kyselyCreateRule` always seeds `INSUFFICIENT_DATA`; patch when callers diff --git a/server/test/integ/report-flow.integ.test.ts b/server/test/integ/report-flow.integ.test.ts index 6cc8bdd..5184410 100644 --- a/server/test/integ/report-flow.integ.test.ts +++ b/server/test/integ/report-flow.integ.test.ts @@ -96,7 +96,6 @@ describe('Report flow (integration)', () => { try { await fn(); } catch (err) { - console.warn('[report-flow.integ] cleanup step failed', err); } }; diff --git a/server/test/integ/rule-change.integ.test.ts b/server/test/integ/rule-change.integ.test.ts new file mode 100644 index 0000000..e43df0b --- /dev/null +++ b/server/test/integ/rule-change.integ.test.ts @@ -0,0 +1,338 @@ +/** + * Integration test for #340: rule changes take effect. + * + * The rule engine reads enabled rules per item type via an eventually + * consistent cache (`getEnabledRulesForItemTypeEventuallyConsistent`, + * `freshUntilAge: 20s`). This test exercises that contract end-to-end against + * the real stack: + * + * 1. Newly created rule fires on the next matching item submission. No prior + * cache entry exists for the fresh item type, so the first submission's + * cache miss fetches the rule from Postgres and the rule engine matches. + * + * 2. Updated rule supersedes the original after the cache window expires. + * A submission made within the cache TTL still sees the old condition; + * one made past it sees the new condition. The test waits past + * `freshUntilAge` between the update and the post-update submission. + * + * Each test provisions its own item type + action so the rules cache (keyed by + * item type) starts empty for that test. Sharing an item type across tests + * leaks cached rule lists from earlier tests, masking whether a freshly created + * rule is actually being picked up. + * + * Run with: npm run test:integ + * Requires: `npm run up && npm run db:update` + */ +import { ScalarTypes } from '@roostorg/coop-types'; +import { type Kysely } from 'kysely'; +import { uid } from 'uid'; + +import { kyselyUpdateRule } from '../../graphql/datasources/ruleKyselyPersistence.js'; +import { type CombinedPg } from '../../services/combinedDbTypes.js'; +import { + ConditionConjunction, + RuleStatus, + RuleType, + type ConditionSet, +} from '../../services/moderationConfigService/index.js'; +import { SignalType } from '../../services/signalsService/index.js'; +import { jsonStringify } from '../../utils/encoding.js'; +import createActions from '../fixtureHelpers/createActions.js'; +import createContentItemTypes from '../fixtureHelpers/createContentItemTypes.js'; +import createOrg from '../fixtureHelpers/createOrg.js'; +import createRule from '../fixtureHelpers/createRule.js'; +import { + makeIntegrationServer, + type IntegrationServer, +} from './setupIntegrationServer.js'; +import { + assertNoActionExecution, + waitForActionExecution, + waitForItemInScylla, +} from './wait.js'; + +// Matches `freshUntilAge: 20` on `getEnabledRulesForItemTypeEventuallyConsistent` +// in `server/rule_engine/ruleEngineQueries.ts`. Anything past the TTL forces +// the next read to refetch, so a wait slightly longer than that proves a rule +// update reaches the engine without depending on an explicit cache-invalidate +// hook (which the rule update path does not currently provide). +const RULES_CACHE_TTL_SECONDS = 20; +const RULES_CACHE_REFRESH_BUFFER_MS = (RULES_CACHE_TTL_SECONDS + 5) * 1000; + +function makeTextContainsConditionSet( + contentTypeId: string, + keyword: string, +): ConditionSet { + return { + conditions: [ + { + input: { type: 'CONTENT_FIELD', name: 'text', contentTypeId }, + signal: { + id: jsonStringify({ type: SignalType.TEXT_MATCHING_CONTAINS_TEXT }), + type: SignalType.TEXT_MATCHING_CONTAINS_TEXT, + }, + matchingValues: { strings: [keyword] }, + }, + ], + conjunction: ConditionConjunction.AND, + }; +} + +describe('Rule changes take effect (integration)', () => { + const orgId = uid(); + let harness: IntegrationServer | undefined; + let apiKey: string; + let orgCleanup: (() => Promise) | undefined; + + beforeAll(async () => { + harness = await makeIntegrationServer(); + + const orgFixture = await createOrg( + { + KyselyPg: harness.deps.KyselyPg, + ModerationConfigService: harness.deps.ModerationConfigService, + ApiKeyService: harness.deps.ApiKeyService, + }, + orgId, + ); + apiKey = orgFixture.apiKey; + orgCleanup = orgFixture.cleanup; + }, 60_000); + + afterAll(async () => { + try { + await orgCleanup?.(); + } finally { + await harness?.shutdown(); + } + }, 30_000); + + /** + * Spins up a fresh item type + action pair so the per-item-type rules cache + * starts empty. Caller passes any rule cleanups to `cleanup`; we tear + * everything down in FK-safe order (rules → action → item type). + */ + async function provisionScenario(deps: IntegrationServer['deps']) { + const itemTypeFixture = await createContentItemTypes({ + moderationConfigService: deps.ModerationConfigService, + orgId, + extra: { + fields: [ + { + name: 'text', + type: ScalarTypes.STRING, + required: true, + container: null, + }, + ], + }, + }); + const itemTypeId = itemTypeFixture.itemTypes[0].id; + + const actionsFixture = await createActions({ + actionAPI: deps.ActionAPIDataSource, + itemTypeIds: [itemTypeId], + orgId, + numActions: 1, + }); + const actionId = actionsFixture.actions[0].id; + + return { + itemTypeId, + actionId, + async cleanup(ruleDestroyers: ReadonlyArray<() => Promise>) { + for (const destroy of ruleDestroyers) { + await destroy().catch(() => undefined); + } + await actionsFixture.cleanup(); + await itemTypeFixture.cleanup(); + }, + }; + } + + test('newly created rule fires on the next matching item submission', async () => { + if (!harness) throw new Error('harness was not initialized'); + const scenario = await provisionScenario(harness.deps); + let rule: Awaited> | undefined; + + try { + rule = await createRule(harness.deps.KyselyPg, orgId, { + name: `text-contains-trigger-${uid()}`, + status: RuleStatus.LIVE, + ruleType: RuleType.CONTENT, + conditionSet: makeTextContainsConditionSet( + scenario.itemTypeId, + 'forbidden', + ), + actionIds: [scenario.actionId], + contentTypeIds: [scenario.itemTypeId], + }); + + const itemId = uid(); + await harness.request + .post('/api/v1/items/async') + .set('x-api-key', apiKey) + .send({ + items: [ + { + id: itemId, + typeId: scenario.itemTypeId, + data: { text: 'this contains the forbidden keyword' }, + }, + ], + }) + .expect(202); + + // Gate on Scylla so we know the worker has picked up the submission before + // we start polling for the (potentially later) action execution row. + // `waitForItemInScylla` looks up by `item_identifier` only (per its doc + // comment), so guard against a cross-org id collision after the row comes + // back. + const scyllaRow = await waitForItemInScylla(harness.deps, { + orgId, + itemIdentifier: { id: itemId, typeId: scenario.itemTypeId }, + }); + expect(scyllaRow.org_id).toBe(orgId); + + const actionRow = await waitForActionExecution(harness.deps, { + orgId, + actionId: scenario.actionId, + itemIdentifier: { id: itemId, typeId: scenario.itemTypeId }, + }); + expect(actionRow.action_id).toBe(scenario.actionId); + } finally { + await scenario.cleanup(rule ? [rule.destroy] : []); + } + }, 60_000); + + test('updated rule applies to items submitted after the cache refreshes', async () => { + if (!harness) throw new Error('harness was not initialized'); + const scenario = await provisionScenario(harness.deps); + let rule: Awaited> | undefined; + + try { + rule = await createRule(harness.deps.KyselyPg, orgId, { + name: `text-contains-update-${uid()}`, + status: RuleStatus.LIVE, + ruleType: RuleType.CONTENT, + conditionSet: makeTextContainsConditionSet( + scenario.itemTypeId, + 'firstkeyword', + ), + actionIds: [scenario.actionId], + contentTypeIds: [scenario.itemTypeId], + }); + + // Pre-update: confirm the rule fires on the original condition so the + // post-update "no action" assertion has a meaningful baseline. + const preUpdateItemId = uid(); + await harness.request + .post('/api/v1/items/async') + .set('x-api-key', apiKey) + .send({ + items: [ + { + id: preUpdateItemId, + typeId: scenario.itemTypeId, + data: { text: 'matches firstkeyword pre-update' }, + }, + ], + }) + .expect(202); + + await waitForActionExecution(harness.deps, { + orgId, + actionId: scenario.actionId, + itemIdentifier: { id: preUpdateItemId, typeId: scenario.itemTypeId }, + }); + + // Update the rule's condition. The cached rule list for this item type + // is now stale; reads within `freshUntilAge` will still see the original + // condition. There is no explicit invalidation on the rule-update path. + const kysely = harness.deps.KyselyPg as Kysely; + await kysely + .transaction() + .setIsolationLevel('repeatable read') + .execute(async (trx) => { + await kyselyUpdateRule(trx, { + id: rule!.id, + orgId, + name: undefined, + description: undefined, + conditionSet: makeTextContainsConditionSet( + scenario.itemTypeId, + 'secondkeyword', + ), + tags: undefined, + ruleType: RuleType.CONTENT, + maxDailyActions: undefined, + expirationTime: undefined, + parentId: undefined, + actionIds: undefined, + policyIds: undefined, + contentTypeIds: undefined, + }); + }); + + // Wait past the cache TTL so the next read refetches and sees the update. + await new Promise((r) => setTimeout(r, RULES_CACHE_REFRESH_BUFFER_MS)); + + // Submission matching the *new* condition: action should fire. + const newMatchItemId = uid(); + await harness.request + .post('/api/v1/items/async') + .set('x-api-key', apiKey) + .send({ + items: [ + { + id: newMatchItemId, + typeId: scenario.itemTypeId, + data: { text: 'matches secondkeyword post-update' }, + }, + ], + }) + .expect(202); + + await waitForActionExecution(harness.deps, { + orgId, + actionId: scenario.actionId, + itemIdentifier: { id: newMatchItemId, typeId: scenario.itemTypeId }, + }); + + // Submission matching the *old* condition: should NOT fire under the + // updated rule. + const staleMatchItemId = uid(); + await harness.request + .post('/api/v1/items/async') + .set('x-api-key', apiKey) + .send({ + items: [ + { + id: staleMatchItemId, + typeId: scenario.itemTypeId, + data: { text: 'matches firstkeyword post-update' }, + }, + ], + }) + .expect(202); + + // Wait for Scylla so we know the worker started processing this + // submission. Same cross-org guard as in the test above. + // `assertNoActionExecution` then waits for `analytics.CONTENT_API_REQUESTS` + // (logged after `runEnabledRules`) before checking absence — that + // post-rules signal, not Scylla, is what makes the negative trustworthy. + const staleScyllaRow = await waitForItemInScylla(harness.deps, { + orgId, + itemIdentifier: { id: staleMatchItemId, typeId: scenario.itemTypeId }, + }); + expect(staleScyllaRow.org_id).toBe(orgId); + await assertNoActionExecution(harness.deps, { + orgId, + actionId: scenario.actionId, + itemIdentifier: { id: staleMatchItemId, typeId: scenario.itemTypeId }, + }); + } finally { + await scenario.cleanup(rule ? [rule.destroy] : []); + } + }, 120_000); +}); diff --git a/server/test/integ/wait.ts b/server/test/integ/wait.ts index eb3df7e..bafd905 100644 --- a/server/test/integ/wait.ts +++ b/server/test/integ/wait.ts @@ -106,6 +106,98 @@ export async function waitForItemInClickHouse( ); } +/** + * Polls `analytics.ACTION_EXECUTIONS` for a row where the given action fired on + * the given item submission. Returned by RuleEngine when a rule's conditions + * match and the rule has actions attached; the row is written via + * ActionPublisher after the worker finishes processing the submission. + * + * Tests use this to assert "the rule fired for this item" without coupling to + * the rule's internals — we only care that the resulting action execution + * landed in the warehouse. + */ +type ActionExecutionRow = { + action_id: string; + item_id: string | null; + item_type_id: string | null; + rules: string; +}; + +export async function waitForActionExecution( + deps: Pick, + opts: { + orgId: string; + actionId: string; + itemIdentifier: ItemIdentifier; + timeoutMs?: number; + }, +): Promise { + const { orgId, actionId, itemIdentifier } = opts; + return waitFor( + `action ${actionId} execution for item ${itemIdentifier.id} in ClickHouse ACTION_EXECUTIONS`, + async () => { + const rows = (await deps.DataWarehouse.query( + `SELECT action_id, item_id, item_type_id, rules + FROM analytics.ACTION_EXECUTIONS + WHERE org_id = ? + AND action_id = ? + AND item_id = ? + AND item_type_id = ? + LIMIT 1`, + deps.Tracer, + [orgId, actionId, itemIdentifier.id, itemIdentifier.typeId], + )) as readonly ActionExecutionRow[]; + return rows.length > 0 ? rows[0] : null; + }, + { timeoutMs: opts.timeoutMs }, + ); +} + +/** + * Asserts the *absence* of an action execution. Proving a negative requires + * waiting long enough that the worker would have written a row if it were + * going to — but a fixed sleep races against rule evaluation latency. + * + * Instead we gate on `analytics.CONTENT_API_REQUESTS`, which `ItemProcessingWorker` + * writes *after* `ruleEngine.runEnabledRules` returns (see the worker around + * Scylla → rules → contentApiLogger). Once that row is visible, the rule path + * for this submission has fully run; if no `ACTION_EXECUTIONS` row exists at + * that point, none will. + * + * Both writes go through `analytics.bulkWrite` with the same default + * batchTimeout, so seeing CONTENT_API_REQUESTS implies that batch interval + * has elapsed — long enough for an ACTION_EXECUTIONS batch from the same + * submission to have flushed. + */ +export async function assertNoActionExecution( + deps: Pick, + opts: { + orgId: string; + actionId: string; + itemIdentifier: ItemIdentifier; + timeoutMs?: number; + }, +) { + const { orgId, actionId, itemIdentifier, timeoutMs } = opts; + await waitForItemInClickHouse(deps, { orgId, itemIdentifier, timeoutMs }); + const rows = await deps.DataWarehouse.query( + `SELECT action_id, item_id + FROM analytics.ACTION_EXECUTIONS + WHERE org_id = ? + AND action_id = ? + AND item_id = ? + AND item_type_id = ? + LIMIT 1`, + deps.Tracer, + [orgId, actionId, itemIdentifier.id, itemIdentifier.typeId], + ); + if (rows.length > 0) { + throw new Error( + `Expected no action execution for action ${actionId} on item ${itemIdentifier.id}, but found one in ACTION_EXECUTIONS`, + ); + } +} + export type ReportingServiceReportRow = { request_id: string; reported_item_id: string;