/* * Copyright (c) 2022. * Author Peter Placzek (tada5hi) * For the full copyright and license information, * view the LICENSE file that was distributed with this source code. */ import {consumeMessageQueue, handleMessageQueueChannel} from "../modules/message-queue"; import {MemorySocketEventMessage, publishEventForSocketServer} from "../modules/memory-event-bus/socket"; import {MemoryMessagePayload} from "../modules/memory-event-bus/message"; import {AuthContact} from "@chamesu/common/domains/auth/contact"; import {AuthContactData} from "@chamesu/common/domains/auth/contact-data"; import {useContactEventCache} from "../app/auth/socket/entities/contact/cache"; import {QueueMessage} from "../modules/message-queue/type"; function createContactAggregatorHandlers({namespace} : {namespace: string}) { const publishMessage = (event: string, payload: MemoryMessagePayload, room?: string) => { const message : MemorySocketEventMessage = { room, event, namespace, payload } publishEventForSocketServer(message); } const getSocketContactRoom = (contactId: string) => { return 'userContact' + ':' + contactId; } const getSocketUserRoom = (userId: number) => { return 'user' + ':' + userId; } return { contactCreated: async(message: QueueMessage) => { const contact = message.data; // Socket publishMessage(message.type, { type: 'contact', id: message.data.id, data: message.data }, getSocketUserRoom(contact.user_one_id)); publishMessage(message.type, { type: 'contact', id: message.data.id, data: message.data }, getSocketUserRoom(contact.user_two_id)); // Cache const contactCache = useContactEventCache(); contactCache .setContact(contact) .then(r => r); contactCache .addUserContactId(contact.user_one_id, contact.id) .then(r => r); contactCache .addUserContactId(contact.user_two_id, contact.id) .then(r => r); }, contactUpdated: async(message: QueueMessage) => { const contact = message.data; // Socket publishMessage(message.type, { type: 'contact', id: message.data.id, data: message.data }, getSocketUserRoom(contact.user_one_id)); publishMessage(message.type, { type: 'contact', id: message.data.id, data: message.data }, getSocketUserRoom(contact.user_two_id)); // Cache useContactEventCache() .setContact(contact) .then(r => r); }, contactDeleted: async(message: QueueMessage) => { const contact = message.data; publishMessage(message.type, { type: 'contact', id: message.data.id, data: message.data }, getSocketUserRoom(contact.user_one_id)); // Cache const contactCache = useContactEventCache(); contactCache .dropContact(contact.id) .then(r => r); contactCache .dropUserContactId(contact.user_one_id, contact.id) .then(r => r); contactCache .dropUserContactId(contact.user_two_id, contact.id) .then(r => r); }, contactDataCreated: async(message: QueueMessage) => { const shareData = message.data; publishMessage(message.type, { type: 'contactData', id: message.data.id, data: message.data }, getSocketUserRoom(shareData.to_user_id)); }, contactDataUpdated: async(message: QueueMessage) => { const shareData = message.data; publishMessage(message.type, { type: 'contactData', id: message.data.id, data: message.data }, getSocketUserRoom(shareData.to_user_id)) }, contactDataDeleted: async(message: QueueMessage) => { const shareData = message.data; publishMessage(message.type, { type: 'contactData', id: message.data.id, data: message.data }, getSocketUserRoom(shareData.to_user_id)); } } } export function buildContactAggregator() { const handlers = createContactAggregatorHandlers({ namespace: '/' }); function start() { return consumeMessageQueue('contact.event', ((async (channel, msg) => { try { await handleMessageQueueChannel(channel, handlers, msg); await channel.ack(msg); } catch (e) { console.log(e); await channel.reject(msg, false); } }))); } return { start } }