diff --git a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service.ts b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service.ts index 678efa303ca..71303870400 100644 --- a/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service.ts +++ b/packages/twenty-server/src/engine/metadata-modules/ai/ai-chat/services/agent-chat.service.ts @@ -76,6 +76,7 @@ export class AgentChatService { type: 'created', entityName: 'agentChatThread', recordId: savedThread.id, + recipientUserWorkspaceIds: [userWorkspaceId], properties: { after: serializeThreadForBroadcast(savedThread), }, @@ -339,6 +340,7 @@ export class AgentChatService { type: 'updated', entityName: 'agentChatThread', recordId: threadId, + recipientUserWorkspaceIds: [thread.userWorkspaceId], properties: { updatedFields: ['title'], after: serializeThreadForBroadcast({ ...thread, title }), diff --git a/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/types/workspace-broadcast-event.type.ts b/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/types/workspace-broadcast-event.type.ts index c03a4ba76a2..bc73d038ceb 100644 --- a/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/types/workspace-broadcast-event.type.ts +++ b/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/types/workspace-broadcast-event.type.ts @@ -8,4 +8,9 @@ export type WorkspaceBroadcastEvent = { after?: Record; diff?: Record; }; + // Restricts delivery to streams whose authContext.userWorkspaceId is in this + // list. Omit for workspace-wide events (shared metadata like views, objects, + // fields). Set for user-scoped entities (e.g. agentChatThread) so other users + // in the same workspace don't receive them. + recipientUserWorkspaceIds?: string[]; }; diff --git a/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/workspace-event-broadcaster.service.ts b/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/workspace-event-broadcaster.service.ts index 119a512fed6..36f6ca5dae7 100644 --- a/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/workspace-event-broadcaster.service.ts +++ b/packages/twenty-server/src/engine/subscriptions/workspace-event-broadcaster/workspace-event-broadcaster.service.ts @@ -41,23 +41,45 @@ export class WorkspaceEventBroadcaster { const streamIdsToRemove: string[] = []; - const payload: EventStreamPayload = { - objectRecordEventsWithQueryIds: [], - metadataEvents: events.map((event) => ({ - metadataName: event.entityName, - type: event.type, - recordId: event.recordId, - properties: event.properties, - updatedCollectionHash, - })), - }; - for (const [streamChannelId, streamData] of streamsData) { if (!isDefined(streamData)) { streamIdsToRemove.push(streamChannelId); continue; } + const streamUserWorkspaceId = streamData.authContext.userWorkspaceId; + + const metadataEventsForStream = events + .filter((event) => { + // Events without recipientUserWorkspaceIds are workspace-wide; delivered + // to every stream. Events with the field are user-scoped; only delivered + // to streams whose authContext.userWorkspaceId is in the list. + if (!isDefined(event.recipientUserWorkspaceIds)) { + return true; + } + + return ( + isDefined(streamUserWorkspaceId) && + event.recipientUserWorkspaceIds.includes(streamUserWorkspaceId) + ); + }) + .map((event) => ({ + metadataName: event.entityName, + type: event.type, + recordId: event.recordId, + properties: event.properties, + updatedCollectionHash, + })); + + if (metadataEventsForStream.length === 0) { + continue; + } + + const payload: EventStreamPayload = { + objectRecordEventsWithQueryIds: [], + metadataEvents: metadataEventsForStream, + }; + await this.subscriptionService.publishToEventStream({ workspaceId, eventStreamChannelId: streamChannelId,