mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-28 04:25:03 +02:00
Fix communication indexing after new communication api (#9931)
* Support full reindexing cards communication Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> * Remove unnecessary groups fetch Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> * Fix indexing of attachment patches Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> * Remove BlobPatchEvent processing from indexer Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com> --------- Signed-off-by: Nikolay Marchuk <nikolay.marchuk@hardcoreeng.com>
This commit is contained in:
@@ -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<T extends Event> = Omit<T, 'date'> & { date: string }
|
||||
type IndexableCommunicationEvent =
|
||||
| QueueSourced<CreateMessageEvent>
|
||||
| QueueSourced<UpdatePatchEvent>
|
||||
| QueueSourced<BlobPatchEvent>
|
||||
| QueueSourced<AttachmentPatchEvent> // TODO: handle
|
||||
| QueueSourced<AttachmentPatchEvent>
|
||||
| QueueSourced<RemovePatchEvent>
|
||||
| QueueSourced<UpdateCardTypeEvent>
|
||||
| QueueSourced<RemoveCardEvent>
|
||||
@@ -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<number> {
|
||||
const communicationApi = this.communicationApi
|
||||
if (communicationApi === undefined) {
|
||||
return 0
|
||||
}
|
||||
let processed = 0
|
||||
const cardsInfo = new Map<CardID, { space: Ref<Space>, _class: Ref<Class<Doc>> }>()
|
||||
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<EventType> = [
|
||||
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<Doc>,
|
||||
_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)
|
||||
|
||||
Reference in New Issue
Block a user