diff --git a/server/indexer/src/indexer/indexer.ts b/server/indexer/src/indexer/indexer.ts index f2a1394462..e4aa7f57e7 100644 --- a/server/indexer/src/indexer/indexer.ts +++ b/server/indexer/src/indexer/indexer.ts @@ -66,7 +66,6 @@ import { type FullTextPipeline } from './types' import { blobPseudoClass, createIndexedDoc, createIndexedDocFromMessage, getContent, messagePseudoClass } from './utils' import { type AttachmentPatchEvent, - type BlobPatchEvent, CardEventType, type CreateMessageEvent, type Event, @@ -89,6 +88,7 @@ import { } from '@hcengineering/communication-types' import { isBlobAttachment, + isBlobAttachmentType, isLinkPreviewAttachment, loadMessages, loadMessagesGroups @@ -110,8 +110,7 @@ export type QueueSourced = Omit & { date: string } type IndexableCommunicationEvent = | QueueSourced | QueueSourced - | QueueSourced - | QueueSourced // TODO: handle + | QueueSourced | QueueSourced | QueueSourced | QueueSourced @@ -286,7 +285,6 @@ export class FullTextIndexPipeline implements FullTextPipeline { let processed = 0 let processedCommunication = 0 - let hasCards = false await ctx.with( 'reindex domain', { domain }, @@ -309,12 +307,22 @@ export class FullTextIndexPipeline implements FullTextPipeline { // Skip non indexable classes continue } - if (!hasCards && this.hierarchy.isDerived(v, card.class.Card)) { - hasCards = true - } await this.indexDocuments(ctx, v, values, pushQueue) await control?.heartbeat() + + if (this.hierarchy.isDerived(v, card.class.Card)) { + for (const card of values) { + processedCommunication += await this.indexCommunication( + ctx, + control, + pushQueue, + card as Card, + processedCommunication + ) + } + } + await control?.heartbeat() } processed += docs.length @@ -340,22 +348,6 @@ export class FullTextIndexPipeline implements FullTextPipeline { } finally { await allDocs.close() } - if (hasCards) { - await ctx.with( - 'reindex-communication', - {}, - async (ctx) => { - try { - const pushQueue = new ElasticPushQueue(this.fulltextAdapter, this.workspace, ctx, control) - processedCommunication = await this.indexCommunication(ctx, control, pushQueue) - await pushQueue.waitProcessing() - } catch (err: any) { - ctx.error('failed to restore index state', { err }) - } - }, - { workspace: this.workspace.uuid } - ) - } }, { domain, @@ -577,131 +569,42 @@ export class FullTextIndexPipeline implements FullTextPipeline { async indexCommunication ( ctx: MeasureContext, control: ConsumerControl | undefined, - pushQueue: ElasticPushQueue + pushQueue: ElasticPushQueue, + card: Card, + processedCommunication: number ): Promise { - const communicationApi = this.communicationApi - if (communicationApi === undefined) { - return 0 - } - let processed = 0 - const cardsInfo = new Map, _class: Ref> }>() + let processed = processedCommunication const rateLimit = new RateLimiter(10) let lastPrint = platformNow() - await ctx.with('process-message-groups', {}, async (ctx) => { - // let groups = await communicationApi.findMessagesGroups(this.communicationSession, { - // limit: messageGroupsLimit, - // order: SortingOrder.Ascending - // }) - let groups = [] as any[] - while (groups.length > 0) { - if (this.cancelling) { - return processed - } - for (const group of groups) { - if (control !== undefined) { - await control.heartbeat() - } - try { - let cardInfo = cardsInfo.get(group.cardId) - if (cardInfo === undefined) { - const cardDoc = await this.storage.findAll(ctx, card.class.Card, { _id: group.cardId }, { limit: 1 }) - if (cardDoc.length !== 1) { - continue - } - cardInfo = { space: cardDoc[0].space, _class: cardDoc[0]._class } - cardsInfo.set(group.cardId, cardInfo) - } - // const blob = await this.storageAdapter.read(ctx, this.workspace, group.blobId) - // const messagesFile = Buffer.concat(blob as any).toString() - // const messagesParsedFile = parseYaml(messagesFile) - // const messages = messagesParsedFile.messages - const messages = [] as Message[] - - for (const message of messages) { - await rateLimit.add(async () => { - await this.processCommunicationMessage( - ctx, - pushQueue, - group.cardId, - cardInfo.space, - cardInfo._class, - message - ) - }) - processed += 1 - const now = platformNow() - if (now - lastPrint > printThresholdMs) { - ctx.info('processed', { - processedCommunication: processed, - elapsed: Math.round(now - lastPrint), - workspace: this.workspace.uuid - }) - lastPrint = now - } - } - } catch (err: any) { - ctx.error('Failed to process message group', { - cardId: group.cardId, - blobId: group.blobId, - error: err - }) - Analytics.handleError(err) - } - } - if (this.cancelling) { - return processed - } - // groups = await communicationApi.findMessagesGroups(this.communicationSession, { - // limit: messageGroupsLimit, - // order: SortingOrder.Ascending, - // fromDate: { - // greater: groups[groups.length - 1].toDate - // } - // }) - groups = [] + let messagesGroups = [] + try { + messagesGroups = await loadMessagesGroups(this.hulylake, card._id) + } catch (err: any) { + ctx.error('Failed to get message groups', { + cardId: card._id, + error: err + }) + Analytics.handleError(err) + return 0 + } + for (const groupInfo of messagesGroups) { + if (this.cancelling) { + return processed } - }) - await ctx.with('process-messages', {}, async (ctx) => { - // let messages = await communicationApi.findMessages(this.communicationSession, { - // limit: messagesLimit, - // order: SortingOrder.Ascending - // }) - let messages = [] as any[] - while (messages.length > 0) { + if (control !== undefined) { + await control.heartbeat() + } + try { + const messages = await loadMessages( + this.hulylake, + groupInfo.blobId, + { cardId: card._id }, + { attachments: true } + ) for (const message of messages) { - if (control !== undefined) { - await control.heartbeat() - } - try { - let cardInfo = cardsInfo.get(message.cardId) - if (cardInfo === undefined) { - const cardDoc = await this.storage.findAll(ctx, card.class.Card, { _id: message.cardId }, { limit: 1 }) - if (cardDoc.length !== 1) { - continue - } - cardInfo = { space: cardDoc[0].space, _class: cardDoc[0]._class } - cardsInfo.set(message.cardId, cardInfo) - } - if (this.cancelling) { - return processed - } - await rateLimit.add(async () => { - await this.processCommunicationMessage( - ctx, - pushQueue, - message.cardId, - cardInfo.space, - cardInfo._class, - message - ) - }) - } catch (err: any) { - ctx.error('Failed to processed message', { - cardId: message.cardId, - id: message.id, - error: err - }) - } + await rateLimit.add(async () => { + await this.processCommunicationMessage(ctx, pushQueue, card._id, card.space, card._class, message) + }) processed += 1 const now = platformNow() if (now - lastPrint > printThresholdMs) { @@ -713,16 +616,15 @@ export class FullTextIndexPipeline implements FullTextPipeline { lastPrint = now } } - // messages = await communicationApi.findMessages(this.communicationSession, { - // limit: messagesLimit, - // order: SortingOrder.Ascending, - // created: { - // greater: messages[messages.length - 1].created - // } - // }) - messages = [] + } catch (err: any) { + ctx.error('Failed to process message group', { + cardId: groupInfo.cardId, + blobId: groupInfo.blobId, + error: err + }) + Analytics.handleError(err) } - }) + } await rateLimit.waitProcessing() return processed } @@ -738,7 +640,6 @@ export class FullTextIndexPipeline implements FullTextPipeline { const indexableCommunicationEventTypes: Array = [ MessageEventType.CreateMessage, MessageEventType.UpdatePatch, - MessageEventType.BlobPatch, MessageEventType.AttachmentPatch, MessageEventType.RemovePatch, CardEventType.UpdateCardType, @@ -843,16 +744,10 @@ export class FullTextIndexPipeline implements FullTextPipeline { if (meta === undefined) { return undefined } - const messagesGroups = await loadMessagesGroups(this.hulylake, cardId) - const group = messagesGroups.find((it) => it.blobId === meta.blobId) - - if (group === undefined) { - return undefined - } return ( await loadMessages( this.hulylake, - group.blobId, + meta.blobId, { cardId, id: msgId @@ -883,18 +778,22 @@ export class FullTextIndexPipeline implements FullTextPipeline { } await this.processCommunicationMessage(ctx, pushQueue, cardDoc._id, cardDoc.space, cardDoc._class, message) messagesUpdated.add(event.messageId) - } else if (tx.event.type === MessageEventType.BlobPatch) { + } else if (tx.event.type === MessageEventType.AttachmentPatch) { const event = tx.event if (messagesUpdated.has(event.messageId)) { continue } for (const operation of event.operations) { - if (operation.opcode === 'attach' || operation.opcode === 'set' || operation.opcode === 'update') { - for (const blobData of operation.blobs) { + if (operation.opcode === 'add' || operation.opcode === 'set') { + for (const blobData of operation.attachments) { + if (!isBlobAttachmentType(blobData.mimeType)) { + continue + } + const params = blobData.params as BlobParams const blobAttachment: BlobAttachment = { - id: blobData.blobId as any as AttachmentID, + id: params.blobId as any as AttachmentID, mimeType: blobData.mimeType ?? '', - params: blobData as BlobParams, + params, creator: event.socialId, created: new Date(Date.parse(event.date)) } @@ -911,8 +810,37 @@ export class FullTextIndexPipeline implements FullTextPipeline { blobAttachment ) } - } else if (operation.opcode === 'detach') { - for (const blobId of operation.blobIds) { + } else if (operation.opcode === 'update') { + if (messagesUpdated.has(event.messageId)) { + continue + } + const message = await getMessage(cardId, event.messageId) + if (message === undefined) { + continue + } + const blobIds = new Set(operation.attachments.map((d) => d.id)) + for (const attachment of message.attachments) { + if (!blobIds.has(attachment.id)) { + continue + } + if (!isBlobAttachmentType(attachment.mimeType)) { + continue + } + const blobAttachment = attachment as BlobAttachment + await this.processCommunicationBlob( + ctx, + pushQueue, + { + id: `${event.messageId}@${cardDoc._id}` as any, + _class: [messagePseudoClass], + space: cardDoc.space, + attachedTo: cardDoc._id + }, + blobAttachment + ) + } + } else if (operation.opcode === 'remove') { + for (const blobId of operation.ids) { toRemove.push({ _id: `${blobId}@${cardDoc._id}` as Ref, _class: blobPseudoClass @@ -920,8 +848,6 @@ export class FullTextIndexPipeline implements FullTextPipeline { } } } - } else if (tx.event.type === MessageEventType.AttachmentPatch) { - // TODO: implement } else if (tx.event.type === MessageEventType.RemovePatch) { const event = tx.event messagesUpdated.add(event.messageId)