Try to fix missing communication broadcast events (#9384)

This commit is contained in:
Kristina
2025-06-29 18:24:20 +07:00
committed by GitHub
parent 47f7e107f2
commit 7cf2fe2041
9 changed files with 65 additions and 33 deletions
+1
View File
@@ -63,6 +63,7 @@ export interface SessionData {
sessionId: string
admin?: boolean
isTriggerCtx?: boolean
hasDomainBroadcast?: boolean
workspace: WorkspaceIds
socialStringsToUsers: Map<
PersonId,
@@ -75,6 +75,7 @@ async function createMessages (tx: TxCUD<Card>, 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<Card>, 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<Doc>): 'create' | 'remove' | 'update' {
+3 -3
View File
@@ -294,7 +294,7 @@ async function OnCardRemove (ctx: TxRemoveDoc<Card>[], 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<Card>[], 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<Ca
socialId: tx.createdBy ?? tx.modifiedBy,
date: new Date((tx.createdOn ?? tx.modifiedOn) + 1)
}
void control.domainRequest(control.ctx, 'communication' as OperationDomain, {
await control.domainRequest(control.ctx, 'communication' as OperationDomain, {
event
})
}
+5 -4
View File
@@ -132,9 +132,10 @@ export interface BroadcastOps {
broadcastSessions: (measure: MeasureContext, sessionIds: Record<string, Tx[]>) => 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>) => void
broadcast: (ctx: MeasureContext, sessionIds: Record<string, any[]>) => 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<CommunicationApi>
/**
+1
View File
@@ -48,6 +48,7 @@ export class QueueMiddleware extends BaseMiddleware {
await this.connected
this.connected = undefined
}
await Promise.all([
this.provideBroadcast(ctx),
this.txProducer.send(
+4
View File
@@ -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)
})
}
}
+34 -20
View File
@@ -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<Middleware> => {
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<any>, 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<DomainResult> {
@@ -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
+10 -3
View File
@@ -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<void> | 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 })
}
}