From 7cf2fe20412e2f941a6d34ea4b5aefb92058706d Mon Sep 17 00:00:00 2001 From: Kristina Date: Sun, 29 Jun 2025 15:24:20 +0400 Subject: [PATCH] Try to fix missing communication broadcast events (#9384) --- communication | 2 +- packages/core/src/server.ts | 1 + .../activity-resources/src/newActivity.ts | 8 ++- server-plugins/card-resources/src/index.ts | 6 +-- server/core/src/types.ts | 9 ++-- server/middleware/src/queue.ts | 1 + server/middleware/src/triggers.ts | 4 ++ server/server-pipeline/src/communication.ts | 54 ++++++++++++------- server/server/src/client.ts | 13 +++-- 9 files changed, 65 insertions(+), 33 deletions(-) diff --git a/communication b/communication index 2e04e5c091..74ae6fd435 160000 --- a/communication +++ b/communication @@ -1 +1 @@ -Subproject commit 2e04e5c0914ec3ce99e896d8f0eb91622570e873 +Subproject commit 74ae6fd435311b225fa4241965d48657d6e87dae diff --git a/packages/core/src/server.ts b/packages/core/src/server.ts index dbde07d201..31b269fd6a 100644 --- a/packages/core/src/server.ts +++ b/packages/core/src/server.ts @@ -63,6 +63,7 @@ export interface SessionData { sessionId: string admin?: boolean isTriggerCtx?: boolean + hasDomainBroadcast?: boolean workspace: WorkspaceIds socialStringsToUsers: Map< PersonId, diff --git a/server-plugins/activity-resources/src/newActivity.ts b/server-plugins/activity-resources/src/newActivity.ts index af6295e005..5098fce6b9 100644 --- a/server-plugins/activity-resources/src/newActivity.ts +++ b/server-plugins/activity-resources/src/newActivity.ts @@ -75,6 +75,7 @@ async function createMessages (tx: TxCUD, control: TriggerControl, card: C result.push({ action }) } + const events: CreateMessageEvent[] = [] for (const data of result) { const event: CreateMessageEvent = { type: MessageEventType.CreateMessage, @@ -86,9 +87,12 @@ async function createMessages (tx: TxCUD, control: TriggerControl, card: C extra: data, date: new Date(tx.modifiedOn) } - - void control.domainRequest(control.ctx, 'communication' as OperationDomain, { event }) + events.push(event) } + + await Promise.all( + events.map((event) => control.domainRequest(control.ctx, 'communication' as OperationDomain, { event })) + ) } function getActivityAction (control: ActivityControl, tx: TxCUD): 'create' | 'remove' | 'update' { diff --git a/server-plugins/card-resources/src/index.ts b/server-plugins/card-resources/src/index.ts index a869f9c2c2..ae56a891de 100644 --- a/server-plugins/card-resources/src/index.ts +++ b/server-plugins/card-resources/src/index.ts @@ -294,7 +294,7 @@ async function OnCardRemove (ctx: TxRemoveDoc[], control: TriggerControl): socialId: removedCard.modifiedBy } - void control.domainRequest(control.ctx, 'communication' as OperationDomain, { + await control.domainRequest(control.ctx, 'communication' as OperationDomain, { event }) @@ -366,7 +366,7 @@ async function OnCardUpdate (ctx: TxUpdateDoc[], control: TriggerControl): socialId: updateTx.createdBy ?? updateTx.modifiedBy, date: new Date(updateTx.createdOn ?? updateTx.modifiedOn) } - void control.domainRequest(control.ctx, 'communication' as OperationDomain, { + await control.domainRequest(control.ctx, 'communication' as OperationDomain, { event }) } @@ -464,7 +464,7 @@ async function updateCollaborators (control: TriggerControl, ctx: TxCreateDoc) => void } -export interface BroadcastSessionsFunc { - broadcast: (ctx: MeasureContext, sessionIds: string[], result: any) => void - enqueue: (ctx: MeasureContext, result: any) => void +export interface CommunicationCallbacks { + registerAsyncRequest: (ctx: MeasureContext, promise: (ctx: MeasureContext) => Promise) => void + broadcast: (ctx: MeasureContext, sessionIds: Record) => void + enqueue: (ctx: MeasureContext, result: any[]) => void } /** @@ -253,7 +254,7 @@ export type PipelineFactory = ( export type CommunicationApiFactory = ( ctx: MeasureContext, ws: WorkspaceIds, - broadcastSessions: BroadcastSessionsFunc + callbacks: CommunicationCallbacks ) => Promise /** diff --git a/server/middleware/src/queue.ts b/server/middleware/src/queue.ts index 5386ac2809..6e4e175acc 100644 --- a/server/middleware/src/queue.ts +++ b/server/middleware/src/queue.ts @@ -48,6 +48,7 @@ export class QueueMiddleware extends BaseMiddleware { await this.connected this.connected = undefined } + await Promise.all([ this.provideBroadcast(ctx), this.txProducer.send( diff --git a/server/middleware/src/triggers.ts b/server/middleware/src/triggers.ts index 3c7de2c1a5..f06f772a25 100644 --- a/server/middleware/src/triggers.ts +++ b/server/middleware/src/triggers.ts @@ -253,6 +253,10 @@ export class TriggersMiddleware extends BaseMiddleware implements Middleware { // We need to send all to recipients await this.context.head?.handleBroadcast(ctx) }) + } else if (ctx.contextData.hasDomainBroadcast === true) { + await ctx.with('broadcast-async-domain-request', {}, async (ctx) => { + await this.context.head?.handleBroadcast(ctx) + }) } } diff --git a/server/server-pipeline/src/communication.ts b/server/server-pipeline/src/communication.ts index c09a88bfa8..902bdfb364 100644 --- a/server/server-pipeline/src/communication.ts +++ b/server/server-pipeline/src/communication.ts @@ -23,13 +23,13 @@ import core, { type SessionData, type TxDomainEvent } from '@hcengineering/core' -import type { - CommunicationApiFactory, - Middleware, - MiddlewareCreator, - PipelineContext +import { + type CommunicationApiFactory, + type Middleware, + type MiddlewareCreator, + type PipelineContext, + BaseMiddleware } from '@hcengineering/server-core' -import { BaseMiddleware } from '@hcengineering/server-core' export const COMMUNICATION_DOMAIN = 'communication' as OperationDomain @@ -49,34 +49,47 @@ export class CommunicationMiddleware extends BaseMiddleware implements Middlewar static create (communicationApiFactory: CommunicationApiFactory): MiddlewareCreator { return async (ctx, context, next): Promise => { const communicationApi = await communicationApiFactory(ctx, context.workspace, { - broadcast: (ctx, sessions, result: Event) => { - const { contextData, evt } = CommunicationMiddleware.wrapEvent(ctx, result) - for (const s of sessions) { - contextData.broadcast.sessions[s] = (contextData.broadcast.sessions[s] ?? []).concat(evt) + registerAsyncRequest: (ctx, promise) => { + const contextData = ctx.contextData as SessionData + contextData.asyncRequests = [ + ...(contextData.asyncRequests ?? []), + async (_ctx) => { + await promise(_ctx) + } + ] + }, + broadcast: (ctx, result) => { + const contextData = ctx.contextData as SessionData + contextData.hasDomainBroadcast = true + for (const [sessionId, events] of Object.entries(result)) { + const txEvents = CommunicationMiddleware.wrapEvents(contextData, events) + contextData.broadcast.sessions[sessionId] = (contextData.broadcast.sessions[sessionId] ?? []).concat( + txEvents + ) } }, - enqueue: (ctx, result: Event) => { - const { contextData, evt } = CommunicationMiddleware.wrapEvent(ctx, result) - contextData.broadcast.queue.push(evt) + enqueue: (ctx, result: Event[]) => { + const contextData = ctx.contextData as SessionData + const txEvents = CommunicationMiddleware.wrapEvents(contextData, result) + contextData.hasDomainBroadcast = true + contextData.broadcast.queue.push(...txEvents) } }) return new CommunicationMiddleware(ctx, context, next, communicationApi) } } - private static wrapEvent (ctx: MeasureContext, result: Event): { contextData: SessionData, evt: TxDomainEvent } { - const contextData = ctx.contextData as SessionData - const evt: TxDomainEvent = { + private static wrapEvents (ctx: SessionData, result: Event[]): TxDomainEvent[] { + return result.map((it) => ({ _id: generateId(), space: core.space.Tx, objectSpace: core.space.Domain, _class: core.class.TxDomainEvent, domain: COMMUNICATION_DOMAIN, - event: result, - modifiedBy: contextData.account.primarySocialId, + event: it, + modifiedBy: ctx.account.primarySocialId, modifiedOn: Date.now() - } - return { contextData, evt } + })) } async domainRequest (ctx: MeasureContext, domain: OperationDomain, params: DomainParams): Promise { @@ -132,6 +145,7 @@ export class CommunicationMiddleware extends BaseMiddleware implements Middlewar return { ...ctx, sessionId: ctx.contextData.sessionId, + asyncData: [], derived: ctx.contextData.isTriggerCtx === true, // TODO: We should decide what to do with communications package and remove this workaround account: ctx.contextData.account diff --git a/server/server/src/client.ts b/server/server/src/client.ts index 12378f01d5..571ab9c378 100644 --- a/server/server/src/client.ts +++ b/server/server/src/client.ts @@ -388,6 +388,12 @@ export class ClientSession implements Session { // We need to broadcast all collected transactions const broadcastPromise = ctx.pipeline.handleBroadcast(ctx.ctx) + await broadcastPromise + + ctx.ctx.contextData.broadcast.queue = [] + ctx.ctx.contextData.broadcast.txes = [] + ctx.ctx.contextData.broadcast.sessions = {} + // ok we could perform async requests if any const asyncs = (ctx.ctx.contextData as SessionData).asyncRequests ?? [] let asyncsPromise: Promise | undefined @@ -403,13 +409,14 @@ export class ClientSession implements Session { } asyncsPromise = handleAyncs() } - await broadcastPromise + if (asyncsPromise !== undefined) { await asyncsPromise + await ctx.pipeline?.handleBroadcast(ctx.ctx) } } catch (err) { - await ctx.sendError(ctx.requestId, 'Failed to findAll', unknownError(err)) - ctx.ctx.error('failed to findAll', { err }) + await ctx.sendError(ctx.requestId, 'Failed to domainRequest', unknownError(err)) + ctx.ctx.error('failed to domainRequest', { err }) } }