diff --git a/client/src/graphql/generated.ts b/client/src/graphql/generated.ts index 1066cd7..3d6b33e 100644 --- a/client/src/graphql/generated.ts +++ b/client/src/graphql/generated.ts @@ -3443,6 +3443,7 @@ export type GQLQuery = { readonly itemWithHistory: GQLItemHistoryResponse; readonly itemsWithId: ReadonlyArray; readonly latestItemSubmissions: ReadonlyArray; + readonly latestItemsByIpAddress: ReadonlyArray; readonly latestItemsCreatedBy: ReadonlyArray; readonly latestItemsCreatedByWithThread: ReadonlyArray; readonly locationBank?: Maybe; @@ -3602,6 +3603,13 @@ export type GQLQueryLatestItemSubmissionsArgs = { itemIdentifiers: ReadonlyArray; }; +export type GQLQueryLatestItemsByIpAddressArgs = { + earliestReturnedSubmissionDate?: InputMaybe; + ipAddress: Scalars['String']['input']; + limit?: InputMaybe; + oldestReturnedSubmissionDate?: InputMaybe; +}; + export type GQLQueryLatestItemsCreatedByArgs = { earliestReturnedSubmissionDate?: InputMaybe; itemIdentifier: GQLItemIdentifierInput; @@ -7310,6 +7318,55 @@ export type GQLInvestigationItemsQuery = { | { readonly __typename: 'NotFoundError'; readonly title: string }; }; +export type GQLGetItemsByIpAddressQueryVariables = Exact<{ + ipAddress: Scalars['String']['input']; + limit?: InputMaybe; +}>; + +export type GQLGetItemsByIpAddressQuery = { + readonly __typename: 'Query'; + readonly latestItemsByIpAddress: ReadonlyArray<{ + readonly __typename: 'ItemSubmissions'; + readonly latest: + | { + readonly __typename: 'ContentItem'; + readonly id: string; + readonly submissionId: string; + readonly submissionTime?: Date | string | null; + readonly type: { + readonly __typename: 'ContentItemType'; + readonly id: string; + readonly name: string; + readonly version: string; + }; + } + | { + readonly __typename: 'ThreadItem'; + readonly id: string; + readonly submissionId: string; + readonly submissionTime?: Date | string | null; + readonly type: { + readonly __typename: 'ThreadItemType'; + readonly id: string; + readonly name: string; + readonly version: string; + }; + } + | { + readonly __typename: 'UserItem'; + readonly id: string; + readonly submissionId: string; + readonly submissionTime?: Date | string | null; + readonly type: { + readonly __typename: 'UserItemType'; + readonly id: string; + readonly name: string; + readonly version: string; + }; + }; + }>; +}; + export type GQLGetAuthorInfoQueryVariables = Exact<{ userIdentifiers: | ReadonlyArray @@ -30688,6 +30745,123 @@ export type GQLInvestigationItemsQueryResult = Apollo.QueryResult< GQLInvestigationItemsQuery, GQLInvestigationItemsQueryVariables >; +export const GQLGetItemsByIpAddressDocument = gql` + query GetItemsByIpAddress($ipAddress: String!, $limit: Int) { + latestItemsByIpAddress(ipAddress: $ipAddress, limit: $limit) { + latest { + ... on ItemBase { + id + submissionId + submissionTime + type { + ... on ItemTypeBase { + id + name + version + } + } + } + } + } + } +`; + +/** + * __useGQLGetItemsByIpAddressQuery__ + * + * To run a query within a React component, call `useGQLGetItemsByIpAddressQuery` and pass it any options that fit your needs. + * When your component renders, `useGQLGetItemsByIpAddressQuery` returns an object from Apollo Client that contains loading, error, and data properties + * you can use to render your UI. + * + * @param baseOptions options that will be passed into the query, supported options are listed on: https://www.apollographql.com/docs/react/api/react-hooks/#options; + * + * @example + * const { data, loading, error } = useGQLGetItemsByIpAddressQuery({ + * variables: { + * ipAddress: // value for 'ipAddress' + * limit: // value for 'limit' + * }, + * }); + */ +export function useGQLGetItemsByIpAddressQuery( + baseOptions: Apollo.QueryHookOptions< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + > & + ( + | { variables: GQLGetItemsByIpAddressQueryVariables; skip?: boolean } + | { skip: boolean } + ), +) { + const options = { ...defaultOptions, ...baseOptions }; + return Apollo.useQuery< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + >(GQLGetItemsByIpAddressDocument, options); +} +export function useGQLGetItemsByIpAddressLazyQuery( + baseOptions?: Apollo.LazyQueryHookOptions< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + >, +) { + const options = { ...defaultOptions, ...baseOptions }; + return Apollo.useLazyQuery< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + >(GQLGetItemsByIpAddressDocument, options); +} +// @ts-ignore +export function useGQLGetItemsByIpAddressSuspenseQuery( + baseOptions?: Apollo.SuspenseQueryHookOptions< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + >, +): Apollo.UseSuspenseQueryResult< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables +>; +export function useGQLGetItemsByIpAddressSuspenseQuery( + baseOptions?: + | Apollo.SkipToken + | Apollo.SuspenseQueryHookOptions< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + >, +): Apollo.UseSuspenseQueryResult< + GQLGetItemsByIpAddressQuery | undefined, + GQLGetItemsByIpAddressQueryVariables +>; +export function useGQLGetItemsByIpAddressSuspenseQuery( + baseOptions?: + | Apollo.SkipToken + | Apollo.SuspenseQueryHookOptions< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + >, +) { + const options = + baseOptions === Apollo.skipToken + ? baseOptions + : { ...defaultOptions, ...baseOptions }; + return Apollo.useSuspenseQuery< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables + >(GQLGetItemsByIpAddressDocument, options); +} +export type GQLGetItemsByIpAddressQueryHookResult = ReturnType< + typeof useGQLGetItemsByIpAddressQuery +>; +export type GQLGetItemsByIpAddressLazyQueryHookResult = ReturnType< + typeof useGQLGetItemsByIpAddressLazyQuery +>; +export type GQLGetItemsByIpAddressSuspenseQueryHookResult = ReturnType< + typeof useGQLGetItemsByIpAddressSuspenseQuery +>; +export type GQLGetItemsByIpAddressQueryResult = Apollo.QueryResult< + GQLGetItemsByIpAddressQuery, + GQLGetItemsByIpAddressQueryVariables +>; export const GQLGetAuthorInfoDocument = gql` query GetAuthorInfo($userIdentifiers: [ItemIdentifierInput!]!) { latestItemSubmissions(itemIdentifiers: $userIdentifiers) { @@ -44831,6 +45005,7 @@ export const namedOperations = { GetOrgData: 'GetOrgData', GetItemsWithId: 'GetItemsWithId', InvestigationItems: 'InvestigationItems', + GetItemsByIpAddress: 'GetItemsByIpAddress', GetAuthorInfo: 'GetAuthorInfo', ItemType: 'ItemType', ItemTypeFormOrg: 'ItemTypeFormOrg', diff --git a/client/src/webpages/dashboard/investigation/InvestigationDashboard.tsx b/client/src/webpages/dashboard/investigation/InvestigationDashboard.tsx index dadc231..bfdf0e9 100644 --- a/client/src/webpages/dashboard/investigation/InvestigationDashboard.tsx +++ b/client/src/webpages/dashboard/investigation/InvestigationDashboard.tsx @@ -20,10 +20,11 @@ gql` export default function InvestigationDashboard() { const [searchParams] = useSearchParams(); - const [id, typeId, submissionTime] = [ + const [id, typeId, submissionTime, ip] = [ searchParams.get('id') ?? undefined, searchParams.get('typeId') ?? undefined, searchParams.get('submissionTime') ?? undefined, + searchParams.get('ip') ?? undefined, ]; return ( @@ -37,9 +38,11 @@ export default function InvestigationDashboard() { />
); diff --git a/client/src/webpages/dashboard/investigation/ItemInvestigation.tsx b/client/src/webpages/dashboard/investigation/ItemInvestigation.tsx index 01ec249..37f223a 100644 --- a/client/src/webpages/dashboard/investigation/ItemInvestigation.tsx +++ b/client/src/webpages/dashboard/investigation/ItemInvestigation.tsx @@ -1,7 +1,7 @@ import { gql } from '@apollo/client'; import { ItemIdentifier } from '@roostorg/coop-types'; import { Alert, Input } from 'antd'; -import { useEffect, useState } from 'react'; +import { useEffect, useState, type ReactNode } from 'react'; import { useNavigate } from 'react-router-dom'; import ComponentLoading from '../../../components/common/ComponentLoading'; @@ -17,11 +17,13 @@ import { useGQLGetOrgDataQuery, } from '../../../graphql/generated'; import { filterNullOrUndefined } from '../../../utils/collections'; +import { getFieldValueForRole } from '../../../utils/itemUtils'; import { __throw } from '../../../utils/misc'; import ManualReviewJobPrimaryUserComponent from '../mrt/manual_review_job/v2/user/ManualReviewJobPrimaryUserComponent'; import { ITEM_TYPE_FRAGMENT } from '../rules/rule_form/RuleForm'; import ItemInvestigationRuleResults from './ItemInvestigationRuleResults'; import ItemInvestigationSummary from './ItemInvestigationSummary'; +import ItemsByIpAddress from './ItemsByIpAddress'; import ThreadInvestigation from './ThreadInvestigation'; export type RuleExecutionHistory = GQLItemHistoryResult['executions'][0]; @@ -215,11 +217,13 @@ export default function ItemInvestigation(props: { itemId: string | undefined; itemTypeId: string | undefined; submissionTime: string | undefined; + ipAddress?: string | undefined; }) { const { itemId: initialItemId, itemTypeId: initialItemTypeId, submissionTime: initialSubmissionTime, + ipAddress, } = props; if (initialItemTypeId && !initialItemId) { @@ -388,9 +392,39 @@ export default function ItemInvestigation(props: { allowMultiplePoliciesPerAction = false, } = orgData?.myOrg ?? {}; + // The item type's `ipAddress` field role (when configured) tells us which + // field holds the IP, so we can offer a reverse lookup of other items that + // share it. Named distinctly from the `ipAddress` prop (URL `?ip=`) to avoid + // shadowing it. + const derivedIpAddress = (() => { + try { + return getFieldValueForRole( + { + // eslint-disable-next-line custom-rules/no-casting-in-getFieldValueForRole + type: item.type as Parameters< + typeof getFieldValueForRole + >[0]['type'], + data: item.data, + }, + 'ipAddress', + ); + } catch { + return undefined; + } + })(); + + const ipPanel = derivedIpAddress ? ( + + ) : null; + + let investigationContent: ReactNode = null; switch (item.__typename) { case 'ContentItem': - return ( + investigationContent = (
); + break; case 'UserItem': - return ( + investigationContent = (
{syntheticBanner}
); + break; case 'ThreadItem': - return ( + investigationContent = (
); + break; } + + return ( + <> + {investigationContent} + {ipPanel} + + ); })(); return ( @@ -493,6 +537,9 @@ export default function ItemInvestigation(props: { {selectedItem ? : null} + {!selectedItem && ipAddress ? ( + + ) : null} {results} ); diff --git a/client/src/webpages/dashboard/investigation/ItemsByIpAddress.tsx b/client/src/webpages/dashboard/investigation/ItemsByIpAddress.tsx new file mode 100644 index 0000000..57a4a16 --- /dev/null +++ b/client/src/webpages/dashboard/investigation/ItemsByIpAddress.tsx @@ -0,0 +1,105 @@ +import { gql } from '@apollo/client'; +import { Empty } from 'antd'; +import { Link } from 'react-router-dom'; + +import ComponentLoading from '../../../components/common/ComponentLoading'; +import FormSectionHeader from '../components/FormSectionHeader'; + +import { useGQLGetItemsByIpAddressQuery } from '../../../graphql/generated'; +import { filterNullOrUndefined } from '../../../utils/collections'; + +gql` + query GetItemsByIpAddress($ipAddress: String!, $limit: Int) { + latestItemsByIpAddress(ipAddress: $ipAddress, limit: $limit) { + latest { + ... on ItemBase { + id + submissionId + submissionTime + type { + ... on ItemTypeBase { + id + name + version + } + } + } + } + } + } +`; + +/** + * Lists other items that share the same IP address as the item currently being + * investigated, so a moderator can pivot from one item to everything else + * associated with that IP (e.g. to spot ban evasion or coordinated abuse). + * + * Each row links into the standard item investigation view for that item. + */ +export default function ItemsByIpAddress(props: { + ipAddress: string; + currentItemId?: string; + currentItemTypeId?: string; +}) { + const { ipAddress, currentItemId, currentItemTypeId } = props; + + const { data, loading, error } = useGQLGetItemsByIpAddressQuery({ + variables: { ipAddress, limit: 50 }, + }); + + const items = filterNullOrUndefined(data?.latestItemsByIpAddress ?? []) + .map((it) => it.latest) + .filter( + (it) => !(it.id === currentItemId && it.type.id === currentItemTypeId), + ); + + return ( +
+
+ + {(() => { + if (loading) { + return ; + } + if (error) { + return ( +
+ Error loading items for this IP: {error.message} +
+ ); + } + if (items.length === 0) { + return ( + + ); + } + return ( +
+ {items.map((item) => ( + +
+
{item.type.name}
+
{item.id}
+
+ {item.submissionTime ? ( +
+ {new Date(item.submissionTime).toLocaleString()} +
+ ) : null} + + ))} +
+ ); + })()} +
+ ); +} diff --git a/client/src/webpages/dashboard/mrt/manual_review_job/v2/ManualReviewJobFieldsComponent.tsx b/client/src/webpages/dashboard/mrt/manual_review_job/v2/ManualReviewJobFieldsComponent.tsx index dfd7027..0639f9f 100644 --- a/client/src/webpages/dashboard/mrt/manual_review_job/v2/ManualReviewJobFieldsComponent.tsx +++ b/client/src/webpages/dashboard/mrt/manual_review_job/v2/ManualReviewJobFieldsComponent.tsx @@ -195,7 +195,6 @@ function TableRowComponent(props: { case 'ID': case 'NUMBER': case 'POLICY_ID': - case 'IP_ADDRESS': case 'STRING': { return (
@@ -208,6 +207,29 @@ function TableRowComponent(props: {
); } + case 'IP_ADDRESS': { + // Make the IP clickable so a moderator can pivot to every other item + // associated with the same IP (ban evasion, coordinated abuse, etc.). + return ( +
+ {label ? ( +
+ {label} +
+ ) : null} + + {String(value)} + +
+ ); + } case 'USER_ID': { return (
diff --git a/db/src/scripts/clickhouse/2026.06.09T22.57.32.add_item_ip_address_to_content_api_requests.sql b/db/src/scripts/clickhouse/2026.06.09T22.57.32.add_item_ip_address_to_content_api_requests.sql new file mode 100644 index 0000000..b759cdb --- /dev/null +++ b/db/src/scripts/clickhouse/2026.06.09T22.57.32.add_item_ip_address_to_content_api_requests.sql @@ -0,0 +1,18 @@ +-- Adds `item_ip_address` to `analytics.CONTENT_API_REQUESTS` to support +-- investigating by IP address over the analytics lookback window (longer than +-- Scylla's 30-day TTL). +-- +-- The value is denormalized at write time by the server (extracted from the +-- item's `ipAddress` schema field role) rather than parsed out of `item_data` +-- JSON at query time, because the field holding the IP is org/item-type +-- specific. It defaults to '' so existing rows remain queryable without a +-- backfill and older code paths that don't send it keep working. +-- +-- The bloom_filter data-skipping index keeps high-cardinality IP equality +-- lookups from scanning every granule. + +ALTER TABLE analytics.CONTENT_API_REQUESTS + ADD COLUMN IF NOT EXISTS item_ip_address String DEFAULT ''; + +ALTER TABLE analytics.CONTENT_API_REQUESTS + ADD INDEX IF NOT EXISTS item_ip_address_idx item_ip_address TYPE bloom_filter GRANULARITY 4; diff --git a/db/src/scripts/scylla/2026.06.09T22.57.31.add_ip_address_index_to_item_submissions.cjs b/db/src/scripts/scylla/2026.06.09T22.57.31.add_ip_address_index_to_item_submissions.cjs new file mode 100644 index 0000000..53828eb --- /dev/null +++ b/db/src/scripts/scylla/2026.06.09T22.57.31.add_ip_address_index_to_item_submissions.cjs @@ -0,0 +1,47 @@ +'use strict'; + +// NB: The Datastax Cassandra/Scylla Driver only supports 1 SQL statement +// per API call. So Scylla migrations should use sequential, raw `query` calls. + +// Adds reverse-lookup support for investigating by IP address. +// +// `item_ip_address` is denormalized onto each item submission at write time by +// the server (extracted from the item's `ipAddress` schema field role). The +// `item_submission_by_ip` materialized view then makes "find every item +// associated with this IP" a partition lookup, mirroring the existing +// `item_submission_by_creator` view that powers creator-based investigation. +// +// Rows without an IP (the common case) are excluded from the view by the +// `item_ip_address IS NOT NULL` filter, exactly like the creator view. The view +// inherits the base table's 30-day TTL. + +/** + * @param {{ context: import("cassandra-driver").Client }} context + */ +exports.up = async function ({ context }) { + const query = context.execute.bind(context); + + await query( + 'ALTER TABLE item_submission_by_thread ADD item_ip_address text;', + ); + + await query(`CREATE MATERIALIZED VIEW IF NOT EXISTS item_submission_by_ip AS + SELECT * FROM item_submission_by_thread + WHERE org_id IS NOT NULL AND item_ip_address IS NOT NULL + AND item_synthetic_created_at IS NOT NULL AND item_identifier IS NOT NULL + AND synthetic_thread_id IS NOT NULL AND parent_identifier IS NOT NULL + AND submission_id IS NOT NULL + PRIMARY KEY((org_id, item_ip_address), item_synthetic_created_at, item_identifier, synthetic_thread_id, parent_identifier, submission_id) + WITH CLUSTERING ORDER BY (item_synthetic_created_at DESC) + AND compression = { 'sstable_compression': 'LZ4Compressor', 'chunk_length_in_kb': 128 };`); +}; + +/** + * @param {{ context: import("cassandra-driver").Client }} context + */ +exports.down = async function ({ context }) { + const query = context.execute.bind(context); + + await query('DROP MATERIALIZED VIEW IF EXISTS item_submission_by_ip;'); + await query('ALTER TABLE item_submission_by_thread DROP item_ip_address;'); +}; diff --git a/server/graphql/generated.ts b/server/graphql/generated.ts index 7315137..7b48083 100644 --- a/server/graphql/generated.ts +++ b/server/graphql/generated.ts @@ -3511,6 +3511,7 @@ export type GQLQuery = { readonly itemWithHistory: GQLItemHistoryResponse; readonly itemsWithId: ReadonlyArray; readonly latestItemSubmissions: ReadonlyArray; + readonly latestItemsByIpAddress: ReadonlyArray; readonly latestItemsCreatedBy: ReadonlyArray; readonly latestItemsCreatedByWithThread: ReadonlyArray; readonly locationBank?: Maybe; @@ -3670,6 +3671,13 @@ export type GQLQueryLatestItemSubmissionsArgs = { itemIdentifiers: ReadonlyArray; }; +export type GQLQueryLatestItemsByIpAddressArgs = { + earliestReturnedSubmissionDate?: InputMaybe; + ipAddress: Scalars['String']['input']; + limit?: InputMaybe; + oldestReturnedSubmissionDate?: InputMaybe; +}; + export type GQLQueryLatestItemsCreatedByArgs = { earliestReturnedSubmissionDate?: InputMaybe; itemIdentifier: GQLItemIdentifierInput; @@ -12538,6 +12546,12 @@ export type GQLQueryResolvers< ContextType, RequireFields >; + latestItemsByIpAddress?: Resolver< + ReadonlyArray, + ParentType, + ContextType, + RequireFields + >; latestItemsCreatedBy?: Resolver< ReadonlyArray, ParentType, diff --git a/server/graphql/modules/investigation.ts b/server/graphql/modules/investigation.ts index 4e4dd90..8b5c97e 100644 --- a/server/graphql/modules/investigation.ts +++ b/server/graphql/modules/investigation.ts @@ -1,4 +1,5 @@ /* eslint-disable max-lines */ +import { isIP } from 'node:net'; import { type DateString } from '@roostorg/coop-types'; import _ from 'lodash'; @@ -27,6 +28,10 @@ import { formatItemSubmissionForGQL } from '../types.js'; import { unauthenticatedError } from '../utils/errors.js'; import { gqlErrorResult, gqlSuccessResult } from '../utils/gqlResult.js'; +// Upper bound on how many items an IP reverse-lookup can return in one call, so +// a caller-supplied `limit` can't trigger an unbounded index scan. +const MAX_ITEMS_BY_IP_LIMIT = 200; + const typeDefs = /* GraphQL */ ` type Query { itemSubmissions( @@ -42,6 +47,15 @@ const typeDefs = /* GraphQL */ ` itemIdentifier: ItemIdentifierInput! ): [ThreadWithMessages!]! + # Reverse lookup of every item associated with a given IP address. Powers + # the "other items from this IP" investigation view. + latestItemsByIpAddress( + ipAddress: String! + limit: Int + oldestReturnedSubmissionDate: DateTime + earliestReturnedSubmissionDate: DateTime + ): [ItemSubmissions!]! + userHistory(itemIdentifier: ItemIdentifierInput!): UserHistoryResponse! # This query enables the caller to get all item submissions for a given @@ -514,6 +528,58 @@ const Query: GQLQueryResolvers = { }); }, + async latestItemsByIpAddress( + _, + { + ipAddress, + limit, + oldestReturnedSubmissionDate, + earliestReturnedSubmissionDate, + }, + context, + ) { + const user = context.getUser(); + if (user == null) { + throw unauthenticatedError('Unauthenticated User'); + } + + // Validate the IP server-side so we never run an index scan on arbitrary + // attacker-controlled input. Return an empty result for anything that + // isn't a well-formed IPv4/IPv6 address. + const normalizedIp = ipAddress.trim(); + if (!isIP(normalizedIp)) { + return []; + } + + // Clamp the caller-supplied limit to a sane, bounded range. + const validatedLimit = + limit == null + ? undefined + : Math.min(Math.max(Math.trunc(limit), 1), MAX_ITEMS_BY_IP_LIMIT); + + const items = await asyncIterableToArray( + context.services.ItemInvestigationService.getItemSubmissionsByIpAddress({ + orgId: user.orgId, + ipAddress: normalizedIp, + limit: validatedLimit, + oldestReturnedSubmissionDate: oldestReturnedSubmissionDate + ? new Date(oldestReturnedSubmissionDate) + : undefined, + earliestReturnedSubmissionDate: earliestReturnedSubmissionDate + ? new Date(earliestReturnedSubmissionDate) + : undefined, + }), + ); + + return items.map((contentItems) => { + const { latestSubmission, priorSubmissions = [] } = contentItems; + return { + latest: formatItemSubmissionForGQL(latestSubmission), + prior: priorSubmissions.map(formatItemSubmissionForGQL), + }; + }); + }, + async latestItemsCreatedByWithThread(__, { itemIdentifier }, context) { const user = context.getUser(); if (user == null) { diff --git a/server/plugins/warehouse/queries/ClickhouseContentApiRequestsAdapter.ts b/server/plugins/warehouse/queries/ClickhouseContentApiRequestsAdapter.ts index 17c828a..064c800 100644 --- a/server/plugins/warehouse/queries/ClickhouseContentApiRequestsAdapter.ts +++ b/server/plugins/warehouse/queries/ClickhouseContentApiRequestsAdapter.ts @@ -6,6 +6,7 @@ import { SIX_MONTHS_MS } from '../../../utils/time.js'; import { formatClickhouseQuery } from '../utils/clickhouseSql.js'; import { type ContentApiImageCountRecord, + type ContentApiRequestByIpRecord, type ContentApiRequestCountRecord, type ContentApiRequestQueryOptions, type ContentApiRequestRecord, @@ -24,6 +25,11 @@ interface ClickhouseContentApiRow { item_type_schema_variant: string; } +interface ClickhouseContentApiByIpRow extends ClickhouseContentApiRow { + item_id: string; + item_type_id: string; +} + interface CountRow { date: string; count: number; @@ -89,6 +95,65 @@ export class ClickhouseContentApiRequestsAdapter implements IContentApiRequestsA })); } + async getSuccessfulRequestsByIpAddress( + orgId: string, + ipAddress: string, + options?: ContentApiRequestQueryOptions, + ): Promise> { + const { lookbackWindowMs = SIX_MONTHS_MS, limit } = options ?? {}; + + const lookbackStart = new Date(Date.now() - Math.max(1, lookbackWindowMs)); + const lookbackStartDate = lookbackStart.toISOString().slice(0, 10); + + // Bound the result set so a high-volume IP can't pull an unbounded number of + // rows into memory. Only apply when the caller asks for a positive cap. + const hasLimit = + typeof limit === 'number' && Number.isFinite(limit) && limit > 0; + + const sql = ` + SELECT + item_id, + item_type_id, + item_data, + submission_id, + ts, + item_creator_id, + item_creator_type_id, + item_type_version, + item_type_schema_variant + FROM analytics.CONTENT_API_REQUESTS + WHERE org_id = ? + AND event = 'REQUEST_SUCCEEDED' + AND item_ip_address = ? + AND item_ip_address != '' + AND ds >= toDate(?) + ORDER BY ts DESC + ${hasLimit ? 'LIMIT ?' : ''} + `; + + const params: unknown[] = [orgId, ipAddress, lookbackStartDate]; + if (hasLimit) { + params.push(Math.floor(limit)); + } + + const rows = (await this.query( + sql, + params, + )) as ClickhouseContentApiByIpRow[]; + + return rows.map((row) => ({ + itemId: row.item_id, + itemTypeId: row.item_type_id, + submissionId: row.submission_id, + itemData: row.item_data, + itemTypeVersion: row.item_type_version, + itemTypeSchemaVariant: row.item_type_schema_variant, + itemCreatorId: row.item_creator_id, + itemCreatorTypeId: row.item_creator_type_id, + occurredAt: new Date(row.ts), + })); + } + async getSuccessfulRequestCountsByDay( orgId: string, start: Date, diff --git a/server/plugins/warehouse/queries/IContentApiRequestsAdapter.ts b/server/plugins/warehouse/queries/IContentApiRequestsAdapter.ts index 9a0e84e..90082b1 100644 --- a/server/plugins/warehouse/queries/IContentApiRequestsAdapter.ts +++ b/server/plugins/warehouse/queries/IContentApiRequestsAdapter.ts @@ -10,9 +10,20 @@ export interface ContentApiRequestRecord { occurredAt: Date; } +/** + * Like {@link ContentApiRequestRecord} but also carries the item's identity, + * since IP-based lookups can return submissions across many different items + * (and item types). + */ +export interface ContentApiRequestByIpRecord extends ContentApiRequestRecord { + itemId: string; + itemTypeId: string; +} + export interface ContentApiRequestQueryOptions { latestOnly?: boolean; lookbackWindowMs?: number; + limit?: number; } export interface ContentApiRequestCountRecord { @@ -43,6 +54,17 @@ export interface IContentApiRequestsAdapter { options?: ContentApiRequestQueryOptions, ): Promise>; + /** + * Returns successful submissions whose denormalized `item_ip_address` matches + * the given IP, ordered most-recent first. Used by investigation to find every + * item associated with an IP beyond Scylla's TTL window. + */ + getSuccessfulRequestsByIpAddress( + orgId: string, + ipAddress: string, + options?: ContentApiRequestQueryOptions, + ): Promise>; + getSuccessfulRequestCountsByDay( orgId: string, start: Date, diff --git a/server/routes/reporting/ReportingRoutes.test.ts b/server/routes/reporting/ReportingRoutes.test.ts index bfa96fb..805d214 100644 --- a/server/routes/reporting/ReportingRoutes.test.ts +++ b/server/routes/reporting/ReportingRoutes.test.ts @@ -308,6 +308,75 @@ describe('POST Report', () => { ]); }); + test('Should accept non-Content (user) items in additional items', async () => { + const payload = { + reporter: { kind: 'user', id: '5123521', typeId: contentTypeId }, + reportedAt: new Date().toISOString(), + reportedForReason: { policyId: '1231241254', reason: 'Some Reason' }, + reportedItem: { + id: '21342135', + typeId: contentTypeId, + data: { name: 'Some name' }, + }, + additionalItems: [ + { + id: '12345123', + typeId: userTypeId, + data: { name: 'Some name' }, + }, + ], + }; + + // Spy (calling through) so we can assert what gets forwarded to MRT. + const enqueueSpy = jest.spyOn(deps.ManualReviewToolService, 'enqueue'); + + try { + await request + .post('/api/v1/report') + .set('x-api-key', apiKey) + .send(payload) + .expect(201); + + await new Promise((resolve) => setTimeout(resolve, 2000)); + + // The user item is still indexed/recorded on the report row... + expect(getBulkWriteMock().mock.calls[0]).toMatchObject([ + 'REPORTING_SERVICE.REPORTS', + [ + { + org_id: orgId, + reported_item_id: '21342135', + reported_item_type_kind: 'CONTENT', + additional_items: [ + { + id: '12345123', + typeIdentifier: { + id: userTypeId, + version: expect.any(String), + schemaVariant: 'original', + }, + data: { name: 'Some name' }, + }, + ], + }, + ], + ]); + + // ...but it must NOT leak into MRT's Content-only `additionalContentItems`. + expect(enqueueSpy).toHaveBeenCalled(); + const enqueueArg = enqueueSpy.mock.calls[0]?.[0] as + | { + payload?: { + additionalContentItems?: ReadonlyArray<{ id: string }>; + }; + } + | undefined; + expect(enqueueArg?.payload?.additionalContentItems ?? []).toEqual([]); + } finally { + enqueueSpy.mockRestore(); + } + }); + test('Should return the expected response for content report and additional items', async () => { const payload = { reporter: { kind: 'user', id: '5123521', typeId: contentTypeId }, diff --git a/server/routes/reporting/submitReport.ts b/server/routes/reporting/submitReport.ts index ecd5eb6..3d0cb1c 100644 --- a/server/routes/reporting/submitReport.ts +++ b/server/routes/reporting/submitReport.ts @@ -220,6 +220,15 @@ export default function submitReport({ ); }; + const isAllValidItems = ( + maybeItemSubmissions: Awaited>[], + ): maybeItemSubmissions is { + itemSubmission: ItemSubmission; + error: undefined; + }[] => { + return maybeItemSubmissions.every((it) => !it.error); + }; + // We disable this lint rule here and below because using `??` here would // match the intended semantics less well/be less clear, but casting these // expressions to strict booleans with Boolean() confuses TS control flow @@ -228,7 +237,7 @@ export default function submitReport({ (reportedThreadSubmission && !isAllValidContentItems(reportedThreadSubmission)) || (additionalItemSubmissions && - !isAllValidContentItems(additionalItemSubmissions)); + !isAllValidItems(additionalItemSubmissions)); const isInvalidReportedAtDate = !isValidDate( new Date(req.body.reportedAt), @@ -259,7 +268,7 @@ export default function submitReport({ threadOrAdditionalItemsHadInvalidOrIllegalItems ? [ makeBadRequestError( - `Invalid report containing a thread or additional items containing items that aren't entirely Content Types`, + `Invalid report: thread messages must all be valid Content Types, and additional items must all be valid item submissions`, { shouldErrorSpan: true }, ), ] @@ -430,12 +439,19 @@ export default function submitReport({ : {}), ...(additionalItemSubmissions ? { - additionalContentItems: additionalItemSubmissions.map( - (it) => + // MRT's `additionalContentItems` is Content-only, so + // drop non-Content items (e.g. USER); they're still + // indexed above. + additionalContentItems: additionalItemSubmissions + .filter( + (it) => + it.itemSubmission.itemType.kind === 'CONTENT', + ) + .map((it) => itemSubmissionToItemSubmissionWithTypeIdentifier( it.itemSubmission, ), - ), + ), } : {}), ...(req.body.reportedItemsInThread diff --git a/server/services/analyticsLoggers/ContentApiLogger.ts b/server/services/analyticsLoggers/ContentApiLogger.ts index 64d6495..9b8fee9 100644 --- a/server/services/analyticsLoggers/ContentApiLogger.ts +++ b/server/services/analyticsLoggers/ContentApiLogger.ts @@ -1,6 +1,7 @@ import { type Dependencies } from '../../iocContainer/index.js'; import { inject } from '../../iocContainer/utils.js'; import { + getFieldValueForRole, type ItemSubmission, type NormalizedItemData, type RawItemData, @@ -51,6 +52,23 @@ class ContentApiLogger { const { failureReason, itemSubmission } = data; const { itemType } = itemSubmission; const now = new Date(); + + // Denormalize the IP address (if the item type maps an `ipAddress` field + // role) into its own column so investigation can look up items by IP without + // parsing `item_data` JSON at query time. Only attempt this for successful + // requests, since failed requests may carry un-normalized item data. + const rawIpAddress = + failureReason == null + ? getFieldValueForRole( + itemType.schema, + itemType.schemaFieldRoles, + 'ipAddress', + itemSubmission.data as NormalizedItemData, + ) + : undefined; + const ipAddress = + typeof rawIpAddress === 'string' ? rawIpAddress.trim() : undefined; + await this.analytics.bulkWrite( 'CONTENT_API_REQUESTS', [ @@ -59,6 +77,7 @@ class ContentApiLogger { ts: now.valueOf(), item_id: itemSubmission.itemId, item_data: jsonStringifyUnstable(itemSubmission.data), + item_ip_address: ipAddress ?? '', ...(itemSubmission.creator !== undefined ? { item_creator_id: itemSubmission.creator.id, diff --git a/server/services/itemInvestigationService/dbTypes.ts b/server/services/itemInvestigationService/dbTypes.ts index a78cd1e..fb9cb52 100644 --- a/server/services/itemInvestigationService/dbTypes.ts +++ b/server/services/itemInvestigationService/dbTypes.ts @@ -55,6 +55,7 @@ export type ScyllaItemSubmissionsRow = { item_type_schema_field_roles: JsonOf; item_type_schema: JsonOf; item_type_schema_variant: 'original' | 'partial'; + item_ip_address: string | null; }; export type ScyllaTables = { @@ -65,6 +66,7 @@ export type ScyllaViews = { item_submission_by_item_id: ScyllaItemSubmissionsRow; item_submission_by_thread_and_time: ScyllaItemSubmissionsRow; item_submission_by_creator: ScyllaItemSubmissionsRow; + item_submission_by_ip: ScyllaItemSubmissionsRow; }; export type ScyllaRelations = ScyllaTables & ScyllaViews; diff --git a/server/services/itemInvestigationService/itemInvestigationService.ts b/server/services/itemInvestigationService/itemInvestigationService.ts index 344a4eb..3b8f917 100644 --- a/server/services/itemInvestigationService/itemInvestigationService.ts +++ b/server/services/itemInvestigationService/itemInvestigationService.ts @@ -5,6 +5,7 @@ import _ from 'lodash'; import { type Dependencies } from '../../iocContainer/index.js'; import { type IActionExecutionsAdapter } from '../../plugins/warehouse/queries/IActionExecutionsAdapter.js'; import { + type ContentApiRequestByIpRecord, type ContentApiRequestRecord, type IContentApiRequestsAdapter, } from '../../plugins/warehouse/queries/IContentApiRequestsAdapter.js'; @@ -207,6 +208,20 @@ export class ItemInvestigationService { ] : []; + // Denormalize the IP for the `item_submission_by_ip` MV. Trim to null for + // blanks: it's part of the partition key, so empty strings would hot-spot + // and whitespace would miss trimmed lookups. + const rawIpAddress = getFieldValueForRole( + itemType.schema, + itemType.schemaFieldRoles, + 'ipAddress', + item.data, + ); + const ipAddress = + typeof rawIpAddress === 'string' && rawIpAddress.trim() !== '' + ? rawIpAddress.trim() + : null; + const itemIdentifier = { id: item.itemId, typeId: itemType.id }; const syntheticThreadId = getSyntheticThreadId(itemIdentifier, threadId); @@ -237,6 +252,7 @@ export class ItemInvestigationService { item_type_schema_field_roles: jsonStringify(itemType.schemaFieldRoles), item_type_schema: jsonStringify(itemType.schema), item_type_schema_variant: itemType.schemaVariant, + item_ip_address: ipAddress, }, }); @@ -946,6 +962,155 @@ export class ItemInvestigationService { } } + /** + * Reverse lookup of items by IP, newest first. Reads the Scylla + * `item_submission_by_ip` MV first, then falls back to ClickHouse for older + * history, de-duplicating items already seen in Scylla. + */ + async *getItemSubmissionsByIpAddress(opts: { + orgId: string; + ipAddress: string; + limit?: number; + oldestReturnedSubmissionDate?: Date; + earliestReturnedSubmissionDate?: Date; + latestSubmissionsOnly?: boolean; + }): AsyncIterable { + const { + orgId, + ipAddress, + limit = 100, + // Bounds apply to `item_synthetic_created_at` (item `createdAt`, not write + // time), so search the full range (TTL bounds Scylla); + oldestReturnedSubmissionDate = new Date(0), + earliestReturnedSubmissionDate = new Date(Date.now() + DAY_MS), + latestSubmissionsOnly = true, + } = opts; + + const seenItemKeys = new Set(); + let returnedItemCount = 0; + + // Tier 1: Scylla hot storage (bounded by the table's ~30-day write TTL). + let emptyQueryResult = true; + let searchStartDate = earliestReturnedSubmissionDate; + const now = Date.now(); + while (returnedItemCount < limit) { + const stream = this.selectStream({ + from: 'item_submission_by_ip', + select: '*', + where: [ + ['org_id', '=', orgId], + ['item_ip_address', '=', ipAddress], + ['item_synthetic_created_at', '<', searchStartDate], + ['item_synthetic_created_at', '>', oldestReturnedSubmissionDate], + ], + limit: Math.floor( + (limit - returnedItemCount) * ARTIFICIAL_LIMIT_MULTIPLIER, + ), + sortOrder: [['item_synthetic_created_at', 'DESC']], + }); + + const groupedStream = chunkAsyncIterableByKey( + stream, + (it: ScyllaItemSubmissionsRow) => jsonStringify(it.item_identifier), + ); + + for await (const itemGroup of groupedStream) { + const { latestSubmission, priorSubmissions } = + partitionLatestAndPriorSubmissions(itemGroup); + const difference = Math.abs( + now - latestSubmission.item_submission_time.getTime(), + ); + this.meter.scyllaRecordAgeHistogram.record( + Math.floor(difference / DAY_MS), + ); + + seenItemKeys.add(jsonStringify(latestSubmission.item_identifier)); + yield { + latestSubmission: + dbRowToItemSubmissionWithItemTypeIdentifier(latestSubmission), + priorSubmissions: latestSubmissionsOnly + ? undefined + : priorSubmissions.map(dbRowToItemSubmissionWithItemTypeIdentifier), + }; + returnedItemCount++; + emptyQueryResult = false; + if (returnedItemCount >= limit) { + break; + } + searchStartDate = latestSubmission.item_synthetic_created_at; + } + if (emptyQueryResult) { + break; + } + emptyQueryResult = true; + } + + if (returnedItemCount >= limit) { + return; + } + + // Tier 2: ClickHouse analytics warehouse (history beyond Scylla's TTL). + const remaining = limit - returnedItemCount; + const records = await this.contentApiRequestsAdapter + .getSuccessfulRequestsByIpAddress(orgId, ipAddress, { + lookbackWindowMs: 6 * MONTH_MS, + limit: Math.floor(remaining * ARTIFICIAL_LIMIT_MULTIPLIER), + }) + .catch(() => [] as ContentApiRequestByIpRecord[]); + + // Rows arrive ordered most-recent first, so the first record per item is + // its latest submission. + const recordsByItem = new Map(); + for (const record of records) { + const key = jsonStringify({ + id: record.itemId, + type_id: record.itemTypeId, + }); + if (seenItemKeys.has(key)) { + continue; + } + const group = recordsByItem.get(key) ?? []; + group.push(record); + recordsByItem.set(key, group); + } + + const toSubmission = (record: ContentApiRequestByIpRecord) => + dbRowToItemSubmissionWithItemTypeIdentifier({ + submission_id: record.submissionId as SubmissionId, + item_identifier: { + id: tryParseNonEmptyString(record.itemId), + type_id: tryParseNonEmptyString(record.itemTypeId), + }, + item_type_version: record.itemTypeVersion, + item_creator_identifier: + record.itemCreatorId && record.itemCreatorTypeId + ? { + id: tryParseNonEmptyString(record.itemCreatorId), + type_id: tryParseNonEmptyString(record.itemCreatorTypeId), + } + : ({ id: '', type_id: '' } as const), + item_data: record.itemData as JsonOf, + item_submission_time: record.occurredAt, + item_type_schema_variant: record.itemTypeSchemaVariant as + | 'original' + | 'partial', + }); + + for (const group of recordsByItem.values()) { + if (returnedItemCount >= limit) { + break; + } + const [latestRecord, ...priorRecords] = group; + yield { + latestSubmission: toSubmission(latestRecord), + priorSubmissions: latestSubmissionsOnly + ? undefined + : priorRecords.map(toSubmission), + }; + returnedItemCount++; + } + } + async #updateItemSubmissionTTL(_opts: { org_id: string; itemID: ItemIdentifier; diff --git a/server/services/itemInvestigationService/itemInvestigationServiceAdapter.ts b/server/services/itemInvestigationService/itemInvestigationServiceAdapter.ts index 1acf6e7..7844ac7 100644 --- a/server/services/itemInvestigationService/itemInvestigationServiceAdapter.ts +++ b/server/services/itemInvestigationService/itemInvestigationServiceAdapter.ts @@ -153,6 +153,18 @@ export class ItemInvestigationServiceAdapter { return this.#adaptInternalStreamToItemSubmissionsForItem(opts.orgId, raw); } + getItemSubmissionsByIpAddress(opts: { + orgId: string; + ipAddress: string; + limit?: number; + oldestReturnedSubmissionDate?: Date; + earliestReturnedSubmissionDate?: Date; + latestSubmissionsOnly?: boolean; + }): AdaptedReturnType<'getItemSubmissionsByIpAddress'> { + const raw = this.service.getItemSubmissionsByIpAddress(opts); + return this.#adaptInternalStreamToItemSubmissionsForItem(opts.orgId, raw); + } + getThreadSubmissionsByTime(opts: { orgId: string; threadId: ItemIdentifier; diff --git a/server/storage/dataWarehouse/IDataWarehouseAnalytics.ts b/server/storage/dataWarehouse/IDataWarehouseAnalytics.ts index 520a064..b5bde58 100644 --- a/server/storage/dataWarehouse/IDataWarehouseAnalytics.ts +++ b/server/storage/dataWarehouse/IDataWarehouseAnalytics.ts @@ -107,6 +107,7 @@ export type AnalyticsSchema = { item_type_schema: string; item_type_schema_variant: string; item_type_schema_field_roles?: Record; + item_ip_address: string; item_creator_id?: string; item_creator_type_id?: string; request_id: string; diff --git a/server/storage/dataWarehouse/warehouseSchema.ts b/server/storage/dataWarehouse/warehouseSchema.ts index 0dd8197..edc88b1 100644 --- a/server/storage/dataWarehouse/warehouseSchema.ts +++ b/server/storage/dataWarehouse/warehouseSchema.ts @@ -151,6 +151,7 @@ export type ItemSubmissionsRow = { ITEM_TYPE_SCHEMA_FIELD_ROLES: SchemaFieldRoles; ITEM_TYPE_ID: string; ITEM_TYPE_SCHEMA: JsonOf; + ITEM_IP_ADDRESS: string; TS: ColumnType; DS: ColumnType; } & (