mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(api): batch the message lookups behind message search (#2640)
This commit is contained in:
@@ -365,10 +365,16 @@ FLUXER_DISCOVERY_ENABLED=true
|
|||||||
#FLUXER_ERLANG_SCHEDULERS_MAX=16
|
#FLUXER_ERLANG_SCHEDULERS_MAX=16
|
||||||
|
|
||||||
# In-flight request ceiling for the four services Compose forwards it to: the
|
# In-flight request ceiling for the four services Compose forwards it to: the
|
||||||
# users and messages routers and their shards. The Rust built-in defaults are 192
|
# users and messages routers and their shards. Leave it unset and each service
|
||||||
# for messages, 320 for snowflakes and 64 elsewhere, and they govern every service
|
# uses its own built-in default, which is what the numbers below describe. Set it
|
||||||
# Compose does not forward this to.
|
# and the one value replaces the built-in default on all four, so size it for the
|
||||||
#FLUXER_SVC_MAX_CONCURRENT_REQUESTS=20
|
# busiest of them rather than for the smallest. The built-in defaults are 192 for
|
||||||
|
# messages, 320 for snowflakes and 64 elsewhere, and they govern every service
|
||||||
|
# Compose does not forward this to. A router holds a slot for the whole round
|
||||||
|
# trip to its shard, so this is a ceiling on requests in flight at once and not a
|
||||||
|
# rate: too low a value does not slow requests down, it rejects them, and the api
|
||||||
|
# turns that rejection into a 503.
|
||||||
|
#FLUXER_SVC_MAX_CONCURRENT_REQUESTS=192
|
||||||
|
|
||||||
# The api and the Rust services name their fixed Postgres statement shapes so the
|
# The api and the Rust services name their fixed Postgres statement shapes so the
|
||||||
# server can reuse their plans. Named prepared statements require a session that
|
# server can reuse their plans. Named prepared statements require a session that
|
||||||
|
|||||||
@@ -565,7 +565,7 @@ services:
|
|||||||
<<: *fluxer-env
|
<<: *fluxer-env
|
||||||
FLUXER_SVC_NAME: users
|
FLUXER_SVC_NAME: users
|
||||||
FLUXER_SVC_MODE: router
|
FLUXER_SVC_MODE: router
|
||||||
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
|
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
|
||||||
healthcheck: *fluxer-svc-healthcheck
|
healthcheck: *fluxer-svc-healthcheck
|
||||||
depends_on:
|
depends_on:
|
||||||
nats: {condition: service_healthy}
|
nats: {condition: service_healthy}
|
||||||
@@ -583,7 +583,7 @@ services:
|
|||||||
FLUXER_SVC_MODE: shard
|
FLUXER_SVC_MODE: shard
|
||||||
FLUXER_SVC_SHARD_ID: "0"
|
FLUXER_SVC_SHARD_ID: "0"
|
||||||
FLUXER_POSTGRES_MAX_CONNECTIONS: "20"
|
FLUXER_POSTGRES_MAX_CONNECTIONS: "20"
|
||||||
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
|
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
|
||||||
healthcheck: *fluxer-svc-healthcheck
|
healthcheck: *fluxer-svc-healthcheck
|
||||||
depends_on:
|
depends_on:
|
||||||
nats: {condition: service_healthy}
|
nats: {condition: service_healthy}
|
||||||
@@ -633,7 +633,7 @@ services:
|
|||||||
<<: *fluxer-env
|
<<: *fluxer-env
|
||||||
FLUXER_SVC_NAME: messages
|
FLUXER_SVC_NAME: messages
|
||||||
FLUXER_SVC_MODE: router
|
FLUXER_SVC_MODE: router
|
||||||
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
|
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
|
||||||
healthcheck: *fluxer-svc-healthcheck
|
healthcheck: *fluxer-svc-healthcheck
|
||||||
depends_on:
|
depends_on:
|
||||||
nats: {condition: service_healthy}
|
nats: {condition: service_healthy}
|
||||||
@@ -651,7 +651,7 @@ services:
|
|||||||
FLUXER_SVC_MODE: shard
|
FLUXER_SVC_MODE: shard
|
||||||
FLUXER_SVC_SHARD_ID: "0"
|
FLUXER_SVC_SHARD_ID: "0"
|
||||||
FLUXER_POSTGRES_MAX_CONNECTIONS: "20"
|
FLUXER_POSTGRES_MAX_CONNECTIONS: "20"
|
||||||
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-20}"
|
FLUXER_SVC_MAX_CONCURRENT_REQUESTS: "${FLUXER_SVC_MAX_CONCURRENT_REQUESTS:-}"
|
||||||
healthcheck: *fluxer-svc-healthcheck
|
healthcheck: *fluxer-svc-healthcheck
|
||||||
depends_on:
|
depends_on:
|
||||||
nats: {condition: service_healthy}
|
nats: {condition: service_healthy}
|
||||||
|
|||||||
@@ -229,14 +229,11 @@ export class AdminMessageService {
|
|||||||
hitsPerPage: limit,
|
hitsPerPage: limit,
|
||||||
page: 1,
|
page: 1,
|
||||||
});
|
});
|
||||||
const messageEntries = result.hits.map((hit) => ({
|
const messageResponses = await createMessageResponseDataService().buildMessages({
|
||||||
channelId: createChannelID(BigInt(hit.channelId)),
|
userId: createUserID(0n),
|
||||||
messageId: createMessageID(BigInt(hit.id)),
|
messages: result.messages,
|
||||||
}));
|
access: await this.getMessageResponseAccessForAdmin(channelId),
|
||||||
const resolvedMessages = await Promise.all(
|
});
|
||||||
messageEntries.map(({channelId, messageId}) => this.getMessageResponseForAdmin(channelId, messageId)),
|
|
||||||
);
|
|
||||||
const messageResponses = resolvedMessages.filter((message): message is MessageResponse => message !== null);
|
|
||||||
const attachmentStatuses = await this.getAttachmentStatusesForMessages(messageResponses);
|
const attachmentStatuses = await this.getAttachmentStatusesForMessages(messageResponses);
|
||||||
const priorReports = await this.getPriorReportsForMessages(messageResponses);
|
const priorReports = await this.getPriorReportsForMessages(messageResponses);
|
||||||
const adminMessages = messageResponses.map((message) =>
|
const adminMessages = messageResponses.map((message) =>
|
||||||
@@ -282,19 +279,6 @@ export class AdminMessageService {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
private async getMessageResponseForAdmin(
|
|
||||||
channelId: ChannelID,
|
|
||||||
messageId: MessageID,
|
|
||||||
): Promise<MessageResponse | null> {
|
|
||||||
const access = await this.getMessageResponseAccessForAdmin(channelId);
|
|
||||||
return createMessageResponseDataService().getMessage({
|
|
||||||
userId: createUserID(0n),
|
|
||||||
channelId,
|
|
||||||
messageId,
|
|
||||||
access,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
private async getPriorReportsForMessages(messages: Array<MessageResponse>): Promise<Map<string, Array<string>>> {
|
private async getPriorReportsForMessages(messages: Array<MessageResponse>): Promise<Map<string, Array<string>>> {
|
||||||
const authorIds = messages.map((message) => createUserID(BigInt(message.author.id)));
|
const authorIds = messages.map((message) => createUserID(BigInt(message.author.id)));
|
||||||
return this.deps.ncmecSubmissionService.getUserPriorReportIds(authorIds);
|
return this.deps.ncmecSubmissionService.getUserPriorReportIds(authorIds);
|
||||||
|
|||||||
@@ -4,11 +4,10 @@ import {ChannelTypes, Permissions} from '@fluxer/constants/src/ChannelConstants'
|
|||||||
import {UnknownMessageError} from '@fluxer/errors/src/domains/channel/UnknownMessageError';
|
import {UnknownMessageError} from '@fluxer/errors/src/domains/channel/UnknownMessageError';
|
||||||
import {FeatureTemporarilyDisabledError} from '@fluxer/errors/src/domains/core/FeatureTemporarilyDisabledError';
|
import {FeatureTemporarilyDisabledError} from '@fluxer/errors/src/domains/core/FeatureTemporarilyDisabledError';
|
||||||
import type {MessageSearchRequest} from '@fluxer/schema/src/domains/message/MessageRequestSchemas';
|
import type {MessageSearchRequest} from '@fluxer/schema/src/domains/message/MessageRequestSchemas';
|
||||||
import type {MessageResponse, MessageSearchResponse} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
|
import type {MessageSearchResponse} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
|
||||||
import {snowflakeToDate} from '@fluxer/snowflake/src/Snowflake';
|
import {snowflakeToDate} from '@fluxer/snowflake/src/Snowflake';
|
||||||
import {AttachmentDecayService} from '../../../attachment/AttachmentDecayService';
|
import {AttachmentDecayService} from '../../../attachment/AttachmentDecayService';
|
||||||
import type {AttachmentID, ChannelID, MessageID, UserID} from '../../../BrandedTypes';
|
import type {AttachmentID, ChannelID, MessageID, UserID} from '../../../BrandedTypes';
|
||||||
import {createChannelID, createMessageID} from '../../../BrandedTypes';
|
|
||||||
import type {UserCacheService} from '../../../infrastructure/UserCacheService';
|
import type {UserCacheService} from '../../../infrastructure/UserCacheService';
|
||||||
import type {RequestCache} from '../../../middleware/RequestCacheMiddleware';
|
import type {RequestCache} from '../../../middleware/RequestCacheMiddleware';
|
||||||
import type {Channel} from '../../../models/Channel';
|
import type {Channel} from '../../../models/Channel';
|
||||||
@@ -192,29 +191,19 @@ export class MessageRetrievalService {
|
|||||||
hitsPerPage,
|
hitsPerPage,
|
||||||
page,
|
page,
|
||||||
});
|
});
|
||||||
const messageEntries = result.hits.map((hit) => ({
|
|
||||||
channelId: createChannelID(BigInt(hit.channelId)),
|
|
||||||
messageId: createMessageID(BigInt(hit.id)),
|
|
||||||
}));
|
|
||||||
const access = {
|
const access = {
|
||||||
sourceGuildId: channel.guildId,
|
sourceGuildId: channel.guildId,
|
||||||
messageHistoryCutoff: !hasReadHistory ? (authChannel.guild?.message_history_cutoff ?? null) : null,
|
messageHistoryCutoff: !hasReadHistory ? (authChannel.guild?.message_history_cutoff ?? null) : null,
|
||||||
canReadMessageHistory: hasReadHistory,
|
canReadMessageHistory: hasReadHistory,
|
||||||
};
|
};
|
||||||
const responseDataService = createMessageResponseDataService();
|
const builtMessages = await createMessageResponseDataService().buildMessages({
|
||||||
const foundMessages = await Promise.all(
|
|
||||||
messageEntries.map(({channelId, messageId}) =>
|
|
||||||
responseDataService.getMessage({
|
|
||||||
userId,
|
userId,
|
||||||
channelId,
|
messages: result.messages,
|
||||||
messageId,
|
|
||||||
access,
|
access,
|
||||||
}),
|
});
|
||||||
),
|
const messageResponses = builtMessages.map(
|
||||||
|
({referenced_message: _referencedMessage, ...searchMessage}) => searchMessage,
|
||||||
);
|
);
|
||||||
const messageResponses = foundMessages
|
|
||||||
.filter((message): message is MessageResponse => message !== null)
|
|
||||||
.map(({referenced_message: _referencedMessage, ...searchMessage}) => searchMessage);
|
|
||||||
return {
|
return {
|
||||||
channels: messageResponses.length > 0 ? [await this.mapSearchChannelResponse(channel, userId, requestCache)] : [],
|
channels: messageResponses.length > 0 ? [await this.mapSearchChannelResponse(channel, userId, requestCache)] : [],
|
||||||
messages: messageResponses,
|
messages: messageResponses,
|
||||||
|
|||||||
@@ -156,7 +156,7 @@ export class GuildSearchService {
|
|||||||
page,
|
page,
|
||||||
cursor,
|
cursor,
|
||||||
});
|
});
|
||||||
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result, userId, requestCache);
|
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result.messages, userId, requestCache);
|
||||||
return {
|
return {
|
||||||
messages: mappedResponses.messages,
|
messages: mappedResponses.messages,
|
||||||
channels: mappedResponses.channels,
|
channels: mappedResponses.channels,
|
||||||
@@ -238,7 +238,7 @@ export class GuildSearchService {
|
|||||||
page,
|
page,
|
||||||
cursor,
|
cursor,
|
||||||
});
|
});
|
||||||
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result, userId, requestCache);
|
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result.messages, userId, requestCache);
|
||||||
return {
|
return {
|
||||||
messages: mappedResponses.messages,
|
messages: mappedResponses.messages,
|
||||||
channels: mappedResponses.channels,
|
channels: mappedResponses.channels,
|
||||||
|
|||||||
@@ -253,7 +253,7 @@ export class GlobalSearchService {
|
|||||||
page,
|
page,
|
||||||
cursor,
|
cursor,
|
||||||
});
|
});
|
||||||
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result, userId, requestCache);
|
const mappedResponses = await this.responseMapper.mapSearchResultToResponses(result.messages, userId, requestCache);
|
||||||
return {
|
return {
|
||||||
messages: mappedResponses.messages,
|
messages: mappedResponses.messages,
|
||||||
channels: mappedResponses.channels,
|
channels: mappedResponses.channels,
|
||||||
|
|||||||
@@ -1,20 +1,18 @@
|
|||||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||||
|
|
||||||
import type {
|
import type {MessageSearchResultsResponse} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
|
||||||
MessageResponse,
|
|
||||||
MessageSearchResultsResponse,
|
|
||||||
} from '@fluxer/schema/src/domains/message/MessageResponseSchemas';
|
|
||||||
import type {UserID} from '../BrandedTypes';
|
import type {UserID} from '../BrandedTypes';
|
||||||
import {createChannelID, createMessageID} from '../BrandedTypes';
|
import {createChannelID} from '../BrandedTypes';
|
||||||
import {mapChannelToResponse} from '../channel/ChannelMappers';
|
import {mapChannelToResponse} from '../channel/ChannelMappers';
|
||||||
import type {IChannelRepository} from '../channel/IChannelRepository';
|
import type {IChannelRepository} from '../channel/IChannelRepository';
|
||||||
import {
|
import {createMessageResponseDataService} from '../channel/services/message/MessageResponseDataService';
|
||||||
createMessageResponseDataService,
|
|
||||||
messageResponseAccessForChannel,
|
|
||||||
} from '../channel/services/message/MessageResponseDataService';
|
|
||||||
import type {UserCacheService} from '../infrastructure/UserCacheService';
|
import type {UserCacheService} from '../infrastructure/UserCacheService';
|
||||||
import type {RequestCache} from '../middleware/RequestCacheMiddleware';
|
import type {RequestCache} from '../middleware/RequestCacheMiddleware';
|
||||||
import type {Channel} from '../models/Channel';
|
import type {Channel} from '../models/Channel';
|
||||||
|
import type {Message} from '../models/Message';
|
||||||
|
import {mapWithConcurrency} from '../utils/ConcurrencyUtils';
|
||||||
|
|
||||||
|
const CHANNEL_LOOKUP_CONCURRENCY = 16;
|
||||||
|
|
||||||
export class MessageSearchResponseMapper {
|
export class MessageSearchResponseMapper {
|
||||||
constructor(
|
constructor(
|
||||||
@@ -23,71 +21,46 @@ export class MessageSearchResponseMapper {
|
|||||||
) {}
|
) {}
|
||||||
|
|
||||||
async mapSearchResultToResponses(
|
async mapSearchResultToResponses(
|
||||||
result: {
|
messages: Array<Message>,
|
||||||
hits: Array<{
|
|
||||||
channelId: string;
|
|
||||||
id: string;
|
|
||||||
}>;
|
|
||||||
total: number;
|
|
||||||
},
|
|
||||||
userId: UserID,
|
userId: UserID,
|
||||||
requestCache: RequestCache,
|
requestCache: RequestCache,
|
||||||
): Promise<{
|
): Promise<{
|
||||||
messages: Array<MessageSearchResultsResponse['messages'][number]>;
|
messages: Array<MessageSearchResultsResponse['messages'][number]>;
|
||||||
channels: Array<MessageSearchResultsResponse['channels'][number]>;
|
channels: Array<MessageSearchResultsResponse['channels'][number]>;
|
||||||
}> {
|
}> {
|
||||||
const messageEntries = result.hits.map((hit) => ({
|
const orderedChannelIds = Array.from(new Set(messages.map((message) => message.channelId.toString())));
|
||||||
channelId: createChannelID(BigInt(hit.channelId)),
|
const channels = await mapWithConcurrency(orderedChannelIds, CHANNEL_LOOKUP_CONCURRENCY, (channelId) =>
|
||||||
messageId: createMessageID(BigInt(hit.id)),
|
|
||||||
}));
|
|
||||||
const orderedChannelIds = new Set<string>();
|
|
||||||
for (const entry of messageEntries) {
|
|
||||||
orderedChannelIds.add(entry.channelId.toString());
|
|
||||||
}
|
|
||||||
const channels = await Promise.all(
|
|
||||||
Array.from(orderedChannelIds).map((channelId) =>
|
|
||||||
this.channelRepository.findUnique(createChannelID(BigInt(channelId))),
|
this.channelRepository.findUnique(createChannelID(BigInt(channelId))),
|
||||||
),
|
|
||||||
);
|
);
|
||||||
const validChannels = channels.filter((channel): channel is Channel => channel !== null);
|
const channelById = new Map(
|
||||||
const channelById = new Map(validChannels.map((channel) => [channel.id.toString(), channel] as const));
|
channels
|
||||||
const responseDataService = createMessageResponseDataService();
|
.filter((channel): channel is Channel => channel !== null)
|
||||||
const messageResponsesWithEntries = await Promise.all(
|
.map((channel) => [channel.id.toString(), channel] as const),
|
||||||
messageEntries.map(async (entry) => {
|
);
|
||||||
const channel = channelById.get(entry.channelId.toString());
|
const renderableMessages = messages.filter((message) => channelById.has(message.channelId.toString()));
|
||||||
if (!channel) return null;
|
const messageResponses = await createMessageResponseDataService().buildMessagesForChannels({
|
||||||
const message = await responseDataService.getMessage({
|
|
||||||
userId,
|
userId,
|
||||||
channelId: entry.channelId,
|
messages: renderableMessages,
|
||||||
messageId: entry.messageId,
|
channelById,
|
||||||
access: messageResponseAccessForChannel(channel),
|
|
||||||
});
|
});
|
||||||
return message ? {message, channelId: entry.channelId.toString()} : null;
|
const searchMessages = messageResponses.map(
|
||||||
}),
|
({referenced_message: _referencedMessage, ...searchMessage}) => searchMessage,
|
||||||
);
|
);
|
||||||
const validMessageResponseEntries = messageResponsesWithEntries.filter(
|
const respondedChannelIds = new Set(searchMessages.map((message) => message.channel_id));
|
||||||
(entry): entry is {message: MessageResponse; channelId: string} => entry !== null,
|
const orderedChannels = orderedChannelIds
|
||||||
);
|
.filter((channelId) => respondedChannelIds.has(channelId))
|
||||||
const messageResponses = validMessageResponseEntries.map((entry) => {
|
|
||||||
const {referenced_message: _referencedMessage, ...searchMessage} = entry.message;
|
|
||||||
return searchMessage;
|
|
||||||
});
|
|
||||||
const orderedResponseChannelIds = new Set(validMessageResponseEntries.map((entry) => entry.channelId));
|
|
||||||
const orderedChannels = Array.from(orderedResponseChannelIds)
|
|
||||||
.map((channelId) => channelById.get(channelId))
|
.map((channelId) => channelById.get(channelId))
|
||||||
.filter((channel): channel is Channel => channel !== undefined);
|
.filter((channel): channel is Channel => channel !== undefined);
|
||||||
const channelResponses = await Promise.all(
|
const channelResponses = await mapWithConcurrency(orderedChannels, CHANNEL_LOOKUP_CONCURRENCY, (channel) =>
|
||||||
orderedChannels.map((channel) =>
|
|
||||||
mapChannelToResponse({
|
mapChannelToResponse({
|
||||||
channel,
|
channel,
|
||||||
currentUserId: userId,
|
currentUserId: userId,
|
||||||
userCacheService: this.userCacheService,
|
userCacheService: this.userCacheService,
|
||||||
requestCache,
|
requestCache,
|
||||||
}),
|
}),
|
||||||
),
|
|
||||||
);
|
);
|
||||||
return {
|
return {
|
||||||
messages: messageResponses,
|
messages: searchMessages,
|
||||||
channels: channelResponses,
|
channels: channelResponses,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import {type ChannelID, createChannelID, createMessageID, type MessageID} from '
|
|||||||
import type {IMessageRepository} from '../channel/repositories/IMessageRepository';
|
import type {IMessageRepository} from '../channel/repositories/IMessageRepository';
|
||||||
import {Logger} from '../Logger';
|
import {Logger} from '../Logger';
|
||||||
import type {Message} from '../models/Message';
|
import type {Message} from '../models/Message';
|
||||||
|
import {mapWithConcurrency} from '../utils/ConcurrencyUtils';
|
||||||
import type {IMessageSearchService} from './IMessageSearchService';
|
import type {IMessageSearchService} from './IMessageSearchService';
|
||||||
import {deleteMessageSearchDocuments} from './MessageSearchIndexCleanup';
|
import {deleteMessageSearchDocuments} from './MessageSearchIndexCleanup';
|
||||||
|
|
||||||
@@ -13,11 +14,16 @@ const RECONCILE_BATCH_SIZE = 250;
|
|||||||
const MAX_RECONCILE_PAGES = 40;
|
const MAX_RECONCILE_PAGES = 40;
|
||||||
const MAX_STALE_DELETE_ABSOLUTE = 250;
|
const MAX_STALE_DELETE_ABSOLUTE = 250;
|
||||||
const MAX_STALE_DELETE_RATIO = 0.5;
|
const MAX_STALE_DELETE_RATIO = 0.5;
|
||||||
|
const HIT_LOOKUP_CONCURRENCY = 32;
|
||||||
|
|
||||||
interface MessageLookupRepository {
|
interface MessageLookupRepository {
|
||||||
readonly messages: Pick<IMessageRepository, 'getMessage'>;
|
readonly messages: Pick<IMessageRepository, 'getMessage'>;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export interface MessageSearchLookupResult extends SearchResult<SearchableMessage> {
|
||||||
|
messages: Array<Message>;
|
||||||
|
}
|
||||||
|
|
||||||
interface SearchExistingMessagesParams {
|
interface SearchExistingMessagesParams {
|
||||||
searchService: IMessageSearchService;
|
searchService: IMessageSearchService;
|
||||||
messageRepository: MessageLookupRepository;
|
messageRepository: MessageLookupRepository;
|
||||||
@@ -28,12 +34,21 @@ interface SearchExistingMessagesParams {
|
|||||||
cursor?: Array<string>;
|
cursor?: Array<string>;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
interface ValidatedHit {
|
||||||
|
hit: SearchableMessage;
|
||||||
|
message: Message | null;
|
||||||
|
}
|
||||||
|
|
||||||
interface ValidatedHits {
|
interface ValidatedHits {
|
||||||
validHits: Array<SearchableMessage>;
|
validHits: Array<ValidatedHit>;
|
||||||
staleMessageIds: Array<MessageID>;
|
staleMessageIds: Array<MessageID>;
|
||||||
lookupErrorCount: number;
|
lookupErrorCount: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
function resolvedMessages(validHits: Array<ValidatedHit>): Array<Message> {
|
||||||
|
return validHits.map((entry) => entry.message).filter((message): message is Message => message !== null);
|
||||||
|
}
|
||||||
|
|
||||||
export async function searchExistingMessages({
|
export async function searchExistingMessages({
|
||||||
searchService,
|
searchService,
|
||||||
messageRepository,
|
messageRepository,
|
||||||
@@ -42,7 +57,7 @@ export async function searchExistingMessages({
|
|||||||
hitsPerPage,
|
hitsPerPage,
|
||||||
page,
|
page,
|
||||||
cursor,
|
cursor,
|
||||||
}: SearchExistingMessagesParams): Promise<SearchResult<SearchableMessage>> {
|
}: SearchExistingMessagesParams): Promise<MessageSearchLookupResult> {
|
||||||
const result = await searchService.searchMessages(query, filters, {
|
const result = await searchService.searchMessages(query, filters, {
|
||||||
hitsPerPage,
|
hitsPerPage,
|
||||||
page: cursor?.length ? undefined : page,
|
page: cursor?.length ? undefined : page,
|
||||||
@@ -50,7 +65,7 @@ export async function searchExistingMessages({
|
|||||||
});
|
});
|
||||||
const validated = await validateSearchHits(messageRepository, result.hits);
|
const validated = await validateSearchHits(messageRepository, result.hits);
|
||||||
if (validated.staleMessageIds.length === 0) {
|
if (validated.staleMessageIds.length === 0) {
|
||||||
return result;
|
return {...result, messages: resolvedMessages(validated.validHits)};
|
||||||
}
|
}
|
||||||
if (cursor?.length) {
|
if (cursor?.length) {
|
||||||
if (validated.lookupErrorCount === 0) {
|
if (validated.lookupErrorCount === 0) {
|
||||||
@@ -58,8 +73,9 @@ export async function searchExistingMessages({
|
|||||||
}
|
}
|
||||||
return {
|
return {
|
||||||
...result,
|
...result,
|
||||||
hits: validated.validHits,
|
hits: validated.validHits.map((entry) => entry.hit),
|
||||||
total: Math.max(validated.validHits.length, result.total - validated.staleMessageIds.length),
|
total: Math.max(validated.validHits.length, result.total - validated.staleMessageIds.length),
|
||||||
|
messages: resolvedMessages(validated.validHits),
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
return reconcileOffsetSearchResult({
|
return reconcileOffsetSearchResult({
|
||||||
@@ -79,9 +95,9 @@ async function reconcileOffsetSearchResult({
|
|||||||
filters,
|
filters,
|
||||||
hitsPerPage,
|
hitsPerPage,
|
||||||
page,
|
page,
|
||||||
}: Omit<SearchExistingMessagesParams, 'cursor'>): Promise<SearchResult<SearchableMessage>> {
|
}: Omit<SearchExistingMessagesParams, 'cursor'>): Promise<MessageSearchLookupResult> {
|
||||||
const requestedOffset = (page - 1) * hitsPerPage;
|
const requestedOffset = (page - 1) * hitsPerPage;
|
||||||
const pageHits: Array<SearchableMessage> = [];
|
const pageHits: Array<ValidatedHit> = [];
|
||||||
const staleMessageIds: Array<MessageID> = [];
|
const staleMessageIds: Array<MessageID> = [];
|
||||||
let lookupErrorCount = 0;
|
let lookupErrorCount = 0;
|
||||||
let examinedCount = 0;
|
let examinedCount = 0;
|
||||||
@@ -102,9 +118,9 @@ async function reconcileOffsetSearchResult({
|
|||||||
lookupErrorCount += validated.lookupErrorCount;
|
lookupErrorCount += validated.lookupErrorCount;
|
||||||
examinedCount += result.hits.length;
|
examinedCount += result.hits.length;
|
||||||
staleMessageIds.push(...validated.staleMessageIds);
|
staleMessageIds.push(...validated.staleMessageIds);
|
||||||
for (const hit of validated.validHits) {
|
for (const entry of validated.validHits) {
|
||||||
if (validTotal >= requestedOffset && pageHits.length < hitsPerPage) {
|
if (validTotal >= requestedOffset && pageHits.length < hitsPerPage) {
|
||||||
pageHits.push(hit);
|
pageHits.push(entry);
|
||||||
}
|
}
|
||||||
validTotal += 1;
|
validTotal += 1;
|
||||||
}
|
}
|
||||||
@@ -121,8 +137,9 @@ async function reconcileOffsetSearchResult({
|
|||||||
await deleteStaleSearchDocuments(searchService, staleMessageIds, examinedCount);
|
await deleteStaleSearchDocuments(searchService, staleMessageIds, examinedCount);
|
||||||
}
|
}
|
||||||
return {
|
return {
|
||||||
hits: pageHits,
|
hits: pageHits.map((entry) => entry.hit),
|
||||||
total: Math.max(pageHits.length, corpusTotal - staleMessageIds.length),
|
total: Math.max(pageHits.length, corpusTotal - staleMessageIds.length),
|
||||||
|
messages: resolvedMessages(pageHits),
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -130,15 +147,14 @@ async function validateSearchHits(
|
|||||||
messageRepository: MessageLookupRepository,
|
messageRepository: MessageLookupRepository,
|
||||||
hits: Array<SearchableMessage>,
|
hits: Array<SearchableMessage>,
|
||||||
): Promise<ValidatedHits> {
|
): Promise<ValidatedHits> {
|
||||||
const checked = await Promise.all(
|
const checked = await mapWithConcurrency(hits, HIT_LOOKUP_CONCURRENCY, async (hit) => {
|
||||||
hits.map(async (hit) => {
|
|
||||||
let channelId: ChannelID;
|
let channelId: ChannelID;
|
||||||
let messageId: MessageID;
|
let messageId: MessageID;
|
||||||
try {
|
try {
|
||||||
channelId = createChannelID(BigInt(hit.channelId));
|
channelId = createChannelID(BigInt(hit.channelId));
|
||||||
messageId = createMessageID(BigInt(hit.id));
|
messageId = createMessageID(BigInt(hit.id));
|
||||||
} catch (_invalidId) {
|
} catch (_invalidId) {
|
||||||
return {hit: null, staleMessageId: null, lookupError: false};
|
return {entry: null, staleMessageId: null, lookupError: false};
|
||||||
}
|
}
|
||||||
let message: Message | null;
|
let message: Message | null;
|
||||||
try {
|
try {
|
||||||
@@ -148,20 +164,19 @@ async function validateSearchHits(
|
|||||||
{error, messageId: hit.id, channelId: hit.channelId},
|
{error, messageId: hit.id, channelId: hit.channelId},
|
||||||
'Search read repair lookup failed; keeping document',
|
'Search read repair lookup failed; keeping document',
|
||||||
);
|
);
|
||||||
return {hit, staleMessageId: null, lookupError: true};
|
return {entry: {hit, message: null}, staleMessageId: null, lookupError: true};
|
||||||
}
|
}
|
||||||
if (message && message.channelId.toString() === hit.channelId) {
|
if (message && message.channelId.toString() === hit.channelId) {
|
||||||
return {hit, staleMessageId: null, lookupError: false};
|
return {entry: {hit, message}, staleMessageId: null, lookupError: false};
|
||||||
}
|
}
|
||||||
return {hit: null, staleMessageId: messageId, lookupError: false};
|
return {entry: null, staleMessageId: messageId, lookupError: false};
|
||||||
}),
|
});
|
||||||
);
|
const validHits: Array<ValidatedHit> = [];
|
||||||
const validHits: Array<SearchableMessage> = [];
|
|
||||||
const staleMessageIds: Array<MessageID> = [];
|
const staleMessageIds: Array<MessageID> = [];
|
||||||
let lookupErrorCount = 0;
|
let lookupErrorCount = 0;
|
||||||
for (const item of checked) {
|
for (const item of checked) {
|
||||||
if (item.hit) {
|
if (item.entry) {
|
||||||
validHits.push(item.hit);
|
validHits.push(item.entry);
|
||||||
}
|
}
|
||||||
if (item.staleMessageId) {
|
if (item.staleMessageId) {
|
||||||
staleMessageIds.push(item.staleMessageId);
|
staleMessageIds.push(item.staleMessageId);
|
||||||
|
|||||||
Reference in New Issue
Block a user