=> ({
- messageId: 'test-message-id',
- textContent: 'test content',
- subject: 'Test Subject',
- content: 'test content',
- sendOn: 1234567890,
- from: 'sender@example.com',
- to: 'recipient@example.com',
- copy: [],
- incoming: true
- })
-
- it('should save message with attachments', async () => {
- const mockMessage = createMockGmailResponse()
- const mockChannel = { _id: 'test-channel-id' } as Channel
- mockWorkspace.getChannel = jest.fn().mockReturnValue(mockChannel)
- mockClient.findOne = jest.fn().mockResolvedValue(undefined)
- mockAttachmentHandler.getPartFiles = jest.fn().mockResolvedValue([])
-
- await messageManager.saveMessage(mockMessage, 'test@example.com')
-
- expect(mockClient.tx).toHaveBeenCalled()
- expect(mockAttachmentHandler.getPartFiles).toHaveBeenCalled()
- })
-
- it('should update existing message', async () => {
- const mockMessage = createMockGmailResponse()
- const mockChannel = { _id: 'test-channel-id' } as Channel
- const existingMessage = createMockMessage()
- mockWorkspace.getChannel = jest.fn().mockReturnValue(mockChannel)
- mockClient.findOne = jest.fn().mockResolvedValue(existingMessage)
- mockAttachmentHandler.getPartFiles = jest.fn().mockResolvedValue([])
-
- await messageManager.saveMessage(mockMessage, 'test@example.com')
-
- expect(mockClient.tx).toHaveBeenCalled()
- expect(mockAttachmentHandler.getPartFiles).toHaveBeenCalled()
- })
-
- it('should not save message when no channels found', async () => {
- const mockMessage = createMockGmailResponse()
- mockWorkspace.getChannel = jest.fn().mockReturnValue(undefined)
-
- await messageManager.saveMessage(mockMessage, 'test@example.com')
-
- expect(mockClient.tx).not.toHaveBeenCalled()
- expect(mockAttachmentHandler.getPartFiles).not.toHaveBeenCalled()
- })
-
- it('should handle incoming messages correctly', async () => {
- const mockMessage = createMockGmailResponse()
- const mockChannel = { _id: 'test-channel-id' } as Channel
- mockWorkspace.getChannel = jest.fn().mockReturnValue(mockChannel)
- mockClient.findOne = jest.fn().mockResolvedValue(undefined)
- mockAttachmentHandler.getPartFiles = jest.fn().mockResolvedValue([])
-
- await messageManager.saveMessage(mockMessage, 'recipient@example.com')
-
- expect(mockWorkspace.getChannel).toHaveBeenCalledWith('sender@example.com')
- })
-
- it('should handle outgoing messages correctly', async () => {
- const mockMessage = createMockGmailResponse()
- const mockChannel = { _id: 'test-channel-id' } as Channel
- mockWorkspace.getChannel = jest.fn().mockReturnValue(mockChannel)
- mockClient.findOne = jest.fn().mockResolvedValue(undefined)
- mockAttachmentHandler.getPartFiles = jest.fn().mockResolvedValue([])
-
- await messageManager.saveMessage(mockMessage, 'sender@example.com')
-
- expect(mockWorkspace.getChannel).toHaveBeenCalledWith('recipient@example.com')
- })
- })
-})
-
describe('sanitizeHtml', () => {
it('should remove all HTML tags', () => {
const html = 'This is bold and italic text
'
diff --git a/services/gmail/pod-gmail/src/__tests__/sync.test.ts b/services/gmail/pod-gmail/src/__tests__/sync.test.ts
index 81e47709ca..31ed66f8a6 100644
--- a/services/gmail/pod-gmail/src/__tests__/sync.test.ts
+++ b/services/gmail/pod-gmail/src/__tests__/sync.test.ts
@@ -13,12 +13,12 @@
// limitations under the License.
//
import { MeasureContext, PersonId } from '@hcengineering/core'
-import { SyncManager, SyncMutex } from '../message/sync'
-import { MessageManager } from '../message/message'
+import { SyncManager } from '../message/sync'
+import { MessageManagerV2 } from '../message/v2/message'
import { RateLimiter } from '../rateLimiter'
jest.mock('../config')
-jest.mock('../message/message')
+jest.mock('../message/adapter')
jest.mock('../utils', () => {
const originalModule = jest.requireActual('../utils')
return {
@@ -46,7 +46,7 @@ const mockKeyValueClientInstance = {
describe('SyncManager', () => {
// Mocked dependencies
let mockCtx: MeasureContext
- let mockMessageManager: jest.Mocked
+ let mockMessageManager: jest.Mocked
let mockGmail: any // Using any for easier mocking
let mockKeyValueClient: MockKeyValueClient
@@ -68,7 +68,7 @@ describe('SyncManager', () => {
mockMessageManager = {
saveMessage: jest.fn().mockResolvedValue(undefined)
- } as unknown as jest.Mocked
+ } as unknown as jest.Mocked
mockGmail = {
history: {
@@ -123,60 +123,6 @@ describe('SyncManager', () => {
jest.clearAllMocks()
})
- describe('getHistory', () => {
- it('should retrieve history from key-value store', async () => {
- // Arrange
- const mockHistory = {
- historyId,
- userId,
- workspace
- }
- mockKeyValueClient.getValue.mockResolvedValue(mockHistory)
-
- // Act
- const result = await (syncManager as any).getHistory(userId)
-
- // Assert
- expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`history:${workspace}:${userId}`)
- expect(result).toEqual(mockHistory)
- })
-
- it('should return null when history not found', async () => {
- // Arrange
- mockKeyValueClient.getValue.mockResolvedValue(null)
-
- // Act
- const result = await (syncManager as any).getHistory(userId)
-
- // Assert
- expect(result).toBeNull()
- })
- })
-
- describe('setHistoryId', () => {
- it('should store history in key-value store', async () => {
- // Act
- await (syncManager as any).setHistoryId(userId, historyId)
-
- // Assert
- expect(mockKeyValueClient.setValue).toHaveBeenCalledWith(`history:${workspace}:${userId}`, {
- historyId,
- userId,
- workspace
- })
- })
- })
-
- describe('clearHistory', () => {
- it('should delete history from key-value store', async () => {
- // Act
- await (syncManager as any).clearHistory(userId)
-
- // Assert
- expect(mockKeyValueClient.deleteKey).toHaveBeenCalledWith(`history:${workspace}:${userId}`)
- })
- })
-
describe('sync', () => {
it('should perform full sync when no history exists', async () => {
// Arrange
@@ -237,101 +183,3 @@ describe('SyncManager', () => {
})
})
})
-
-describe('SyncMutex', () => {
- let syncMutex: SyncMutex
-
- beforeEach(() => {
- syncMutex = new SyncMutex()
- })
-
- it('should allow sequential locking and unlocking', async () => {
- // Lock first time
- const release1 = await syncMutex.lock('test-key')
- release1()
-
- // Lock second time
- const release2 = await syncMutex.lock('test-key')
- release2()
-
- // If we got here without hanging, the test passes
- expect(true).toBe(true)
- })
-
- it('should queue up multiple requests for the same key', async () => {
- const results: number[] = []
-
- // Start 3 concurrent lock operations
- const promise1 = (async () => {
- const release = await syncMutex.lock('test-key')
- results.push(1)
- await new Promise((resolve) => setTimeout(resolve, 10))
- release()
- })()
-
- const promise2 = (async () => {
- const release = await syncMutex.lock('test-key')
- results.push(2)
- await new Promise((resolve) => setTimeout(resolve, 5))
- release()
- })()
-
- const promise3 = (async () => {
- const release = await syncMutex.lock('test-key')
- results.push(3)
- release()
- })()
-
- // Wait for all promises to resolve
- await Promise.all([promise1, promise2, promise3])
-
- // The operations should have happened in order
- expect(results).toEqual([1, 2, 3])
- })
-
- it('should allow concurrent operations on different keys', async () => {
- const results: string[] = []
-
- // Lock two different keys concurrently
- const promise1 = (async () => {
- const release = await syncMutex.lock('key1')
- results.push('key1-locked')
- await new Promise((resolve) => setTimeout(resolve, 20))
- results.push('key1-unlocked')
- release()
- })()
-
- const promise2 = (async () => {
- const release = await syncMutex.lock('key2')
- results.push('key2-locked')
- await new Promise((resolve) => setTimeout(resolve, 10))
- results.push('key2-unlocked')
- release()
- })()
-
- // Wait for both promises to resolve
- await Promise.all([promise1, promise2])
-
- // key2 operations should complete before key1 due to shorter timeout
- expect(results.indexOf('key2-locked')).toBeLessThan(results.indexOf('key1-unlocked'))
- })
-
- it('should release the lock properly even if an error occurs', async () => {
- try {
- const release = await syncMutex.lock('test-key')
- try {
- throw new Error('Test error')
- } finally {
- release()
- }
- } catch (error) {
- // Ignore the error
- }
-
- // Should be able to acquire the lock again
- const release = await syncMutex.lock('test-key')
- release()
-
- expect(true).toBe(true)
- })
-})
diff --git a/services/gmail/pod-gmail/src/__tests__/syncState.test.ts b/services/gmail/pod-gmail/src/__tests__/syncState.test.ts
new file mode 100644
index 0000000000..05a8c747e1
--- /dev/null
+++ b/services/gmail/pod-gmail/src/__tests__/syncState.test.ts
@@ -0,0 +1,253 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { PersonId } from '@hcengineering/core'
+import { KeyValueClient } from '@hcengineering/kvs-client'
+
+import { SyncStateManager } from '../message/syncState'
+import { IntegrationVersion } from '../types'
+import { History } from '../message/types'
+
+describe('SyncStateManager', () => {
+ const workspace = 'test-workspace'
+ const userId = 'test-user-id' as PersonId
+ const historyId = 'test-history-id'
+ const pageToken = 'test-page-token'
+
+ let mockKeyValueClient: jest.Mocked
+ let v1StateManager: SyncStateManager
+ let v2StateManager: SyncStateManager
+
+ beforeEach(() => {
+ // Reset mocks before each test
+ mockKeyValueClient = {
+ getValue: jest.fn(),
+ setValue: jest.fn().mockResolvedValue(undefined),
+ deleteKey: jest.fn().mockResolvedValue(undefined)
+ } as unknown as jest.Mocked
+
+ // Create state managers for both versions
+ v1StateManager = new SyncStateManager(mockKeyValueClient, workspace, IntegrationVersion.V1)
+ v2StateManager = new SyncStateManager(mockKeyValueClient, workspace, IntegrationVersion.V2)
+ })
+
+ describe('getHistory', () => {
+ it('should call getValue with correct key for V1', async () => {
+ const expectedHistory: History = {
+ historyId,
+ userId,
+ workspace
+ }
+
+ mockKeyValueClient.getValue.mockResolvedValue(expectedHistory)
+
+ const result = await v1StateManager.getHistory(userId)
+
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`history:${workspace}:${userId}`)
+ expect(result).toEqual(expectedHistory)
+ })
+
+ it('should call getValue with correct key for V2', async () => {
+ const expectedHistory: History = {
+ historyId,
+ userId,
+ workspace
+ }
+
+ mockKeyValueClient.getValue.mockResolvedValue(expectedHistory)
+
+ const result = await v2StateManager.getHistory(userId)
+
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`history-v2:${workspace}:${userId}`)
+ expect(result).toEqual(expectedHistory)
+ })
+
+ it('should return null when no history exists', async () => {
+ mockKeyValueClient.getValue.mockResolvedValue(null)
+
+ const result = await v1StateManager.getHistory(userId)
+
+ expect(result).toBeNull()
+ })
+
+ it('should propagate errors from KeyValueClient', async () => {
+ const error = new Error('Database error')
+ mockKeyValueClient.getValue.mockRejectedValue(error)
+
+ await expect(v1StateManager.getHistory(userId)).rejects.toThrow(error)
+ })
+ })
+
+ describe('clearHistory', () => {
+ it('should call deleteKey with correct key for V1', async () => {
+ await v1StateManager.clearHistory(userId)
+
+ expect(mockKeyValueClient.deleteKey).toHaveBeenCalledWith(`history:${workspace}:${userId}`)
+ })
+
+ it('should call deleteKey with correct key for V2', async () => {
+ await v2StateManager.clearHistory(userId)
+
+ expect(mockKeyValueClient.deleteKey).toHaveBeenCalledWith(`history-v2:${workspace}:${userId}`)
+ })
+
+ it('should propagate errors from KeyValueClient', async () => {
+ const error = new Error('Database error')
+ mockKeyValueClient.deleteKey.mockRejectedValue(error)
+
+ await expect(v1StateManager.clearHistory(userId)).rejects.toThrow(error)
+ })
+ })
+
+ describe('setHistoryId', () => {
+ it('should call setValue with correct key and value for V1', async () => {
+ await v1StateManager.setHistoryId(userId, historyId)
+
+ expect(mockKeyValueClient.setValue).toHaveBeenCalledWith(`history:${workspace}:${userId}`, {
+ historyId,
+ userId,
+ workspace
+ })
+ })
+
+ it('should call setValue with correct key and value for V2', async () => {
+ await v2StateManager.setHistoryId(userId, historyId)
+
+ expect(mockKeyValueClient.setValue).toHaveBeenCalledWith(`history-v2:${workspace}:${userId}`, {
+ historyId,
+ userId,
+ workspace
+ })
+ })
+
+ it('should propagate errors from KeyValueClient', async () => {
+ const error = new Error('Database error')
+ mockKeyValueClient.setValue.mockRejectedValue(error)
+
+ await expect(v1StateManager.setHistoryId(userId, historyId)).rejects.toThrow(error)
+ })
+ })
+
+ describe('getPageToken', () => {
+ it('should call getValue with correct key for V1', async () => {
+ mockKeyValueClient.getValue.mockResolvedValue(pageToken)
+
+ const result = await v1StateManager.getPageToken(userId)
+
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`page-token:${workspace}:${userId}`)
+ expect(result).toEqual(pageToken)
+ })
+
+ it('should call getValue with correct key for V2', async () => {
+ mockKeyValueClient.getValue.mockResolvedValue(pageToken)
+
+ const result = await v2StateManager.getPageToken(userId)
+
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`page-token-v2:${workspace}:${userId}`)
+ expect(result).toEqual(pageToken)
+ })
+
+ it('should return null when no page token exists', async () => {
+ mockKeyValueClient.getValue.mockResolvedValue(null)
+
+ const result = await v1StateManager.getPageToken(userId)
+
+ expect(result).toBeNull()
+ })
+
+ it('should propagate errors from KeyValueClient', async () => {
+ const error = new Error('Database error')
+ mockKeyValueClient.getValue.mockRejectedValue(error)
+
+ await expect(v1StateManager.getPageToken(userId)).rejects.toThrow(error)
+ })
+ })
+
+ describe('setPageToken', () => {
+ it('should call setValue with correct key and value for V1', async () => {
+ await v1StateManager.setPageToken(userId, pageToken)
+
+ expect(mockKeyValueClient.setValue).toHaveBeenCalledWith(`page-token:${workspace}:${userId}`, pageToken)
+ })
+
+ it('should call setValue with correct key and value for V2', async () => {
+ await v2StateManager.setPageToken(userId, pageToken)
+
+ expect(mockKeyValueClient.setValue).toHaveBeenCalledWith(`page-token-v2:${workspace}:${userId}`, pageToken)
+ })
+
+ it('should propagate errors from KeyValueClient', async () => {
+ const error = new Error('Database error')
+ mockKeyValueClient.setValue.mockRejectedValue(error)
+
+ await expect(v1StateManager.setPageToken(userId, pageToken)).rejects.toThrow(error)
+ })
+ })
+
+ describe('key generation', () => {
+ it('should generate different keys for V1 and V2', async () => {
+ // Test through the public methods to verify the keys are different
+ mockKeyValueClient.getValue.mockResolvedValue(null)
+
+ await v1StateManager.getHistory(userId)
+ await v2StateManager.getHistory(userId)
+
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`history:${workspace}:${userId}`)
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`history-v2:${workspace}:${userId}`)
+
+ mockKeyValueClient.getValue.mockClear()
+
+ await v1StateManager.getPageToken(userId)
+ await v2StateManager.getPageToken(userId)
+
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`page-token:${workspace}:${userId}`)
+ expect(mockKeyValueClient.getValue).toHaveBeenCalledWith(`page-token-v2:${workspace}:${userId}`)
+ })
+ })
+
+ describe('version migration scenario', () => {
+ it('should allow migrating from V1 to V2', async () => {
+ // Simulate V1 data existing
+ const v1History: History = { historyId: 'v1-history', userId, workspace }
+ mockKeyValueClient.getValue.mockImplementation((key: string) => {
+ if (key === `history:${workspace}:${userId}`) {
+ return Promise.resolve(v1History)
+ }
+ if (key === `history-v2:${workspace}:${userId}`) {
+ return Promise.resolve(null)
+ }
+ return Promise.resolve(null)
+ })
+
+ // V1 manager should find history
+ const historyFromV1 = await v1StateManager.getHistory(userId)
+ expect(historyFromV1).toEqual(v1History)
+
+ // V2 manager should not find history yet
+ const historyFromV2 = await v2StateManager.getHistory(userId)
+ expect(historyFromV2).toBeNull()
+
+ // Migrate by setting V2 history
+ await v2StateManager.setHistoryId(userId, 'v2-history')
+
+ // Verify the call to set V2 history
+ expect(mockKeyValueClient.setValue).toHaveBeenCalledWith(`history-v2:${workspace}:${userId}`, {
+ historyId: 'v2-history',
+ userId,
+ workspace
+ })
+ })
+ })
+})
diff --git a/services/gmail/pod-gmail/src/config.ts b/services/gmail/pod-gmail/src/config.ts
index df14ab4134..a565c50e51 100644
--- a/services/gmail/pod-gmail/src/config.ts
+++ b/services/gmail/pod-gmail/src/config.ts
@@ -13,20 +13,21 @@
// See the License for the specific language governing permissions and
// limitations under the License.
//
+import { BaseConfig } from '@hcengineering/mail-common'
import { config as dotenvConfig } from 'dotenv'
+import { IntegrationVersion } from './types'
dotenvConfig()
-interface Config {
+interface Config extends BaseConfig {
Port: number
- AccountsURL: string
ServiceID: string
Secret: string
Credentials: string
WATCH_TOPIC_NAME: string
FooterMessage: string
InitLimit: number
- KvsUrl: string
+ Version: IntegrationVersion
}
const envMap: { [key in keyof Config]: string } = {
@@ -38,12 +39,21 @@ const envMap: { [key in keyof Config]: string } = {
WATCH_TOPIC_NAME: 'WATCH_TOPIC_NAME',
FooterMessage: 'FOOTER_MESSAGE',
InitLimit: 'INIT_LIMIT',
- KvsUrl: 'KVS_URL'
+ KvsUrl: 'KVS_URL',
+ StorageConfig: 'STORAGE_CONFIG',
+ Version: 'VERSION'
}
const parseNumber = (str: string | undefined): number | undefined => (str !== undefined ? Number(str) : undefined)
const config: Config = (() => {
+ const versionStr = process.env[envMap.Version] ?? 'v1'
+ let version: IntegrationVersion
+ if (versionStr === IntegrationVersion.V1 || versionStr === IntegrationVersion.V2) {
+ version = versionStr as IntegrationVersion
+ } else {
+ throw new Error(`Invalid version: ${versionStr}. Must be 'v1' or 'v2'.`)
+ }
const params: Partial = {
Port: parseNumber(process.env[envMap.Port]) ?? 8087,
AccountsURL: process.env[envMap.AccountsURL],
@@ -53,7 +63,9 @@ const config: Config = (() => {
WATCH_TOPIC_NAME: process.env[envMap.WATCH_TOPIC_NAME],
InitLimit: parseNumber(process.env[envMap.InitLimit]) ?? 50,
FooterMessage: process.env[envMap.FooterMessage] ?? '
Sent via Huly
',
- KvsUrl: process.env[envMap.KvsUrl]
+ KvsUrl: process.env[envMap.KvsUrl],
+ StorageConfig: process.env[envMap.StorageConfig],
+ Version: version
}
const missingEnv = (Object.keys(params) as Array)
diff --git a/services/gmail/pod-gmail/src/gmail.ts b/services/gmail/pod-gmail/src/gmail.ts
index e013f547ce..933630aa6f 100644
--- a/services/gmail/pod-gmail/src/gmail.ts
+++ b/services/gmail/pod-gmail/src/gmail.ts
@@ -30,10 +30,11 @@ import { getOrCreateSocialId } from './accounts'
import { createIntegrationIfNotEsixts, disableIntegration, removeIntegration } from './integrations'
import { AttachmentHandler } from './message/attachments'
import { TokenStorage } from './tokens'
-import { MessageManager } from './message/message'
+import { createMessageManager } from './message/adapter'
import { SyncManager } from './message/sync'
import { getEmail } from './gmail/utils'
import { Integration } from '@hcengineering/account-client'
+import { IMessageManager } from './message/types'
const SCOPES = ['https://www.googleapis.com/auth/gmail.modify']
@@ -81,7 +82,7 @@ export class GmailClient {
private refreshTimer: NodeJS.Timeout | undefined = undefined
private readonly rateLimiter = new RateLimiter(1000, 200)
private readonly attachmentHandler: AttachmentHandler
- private readonly messageManager: MessageManager
+ private readonly messageManager: IMessageManager
private readonly syncManager: SyncManager
private readonly integrationToken: string
private integration: Integration | undefined = undefined
@@ -104,12 +105,13 @@ export class GmailClient {
this.client = new TxOperations(client, this.socialId._id)
this.account = this.user.userId
this.attachmentHandler = new AttachmentHandler(ctx, workspaceId, storageAdapter, this.gmail, this.client)
- this.messageManager = new MessageManager(
+ this.messageManager = createMessageManager(
ctx,
this.client,
this.attachmentHandler,
- this.socialId._id,
- this.workspace
+ this.workspace,
+ this.integrationToken,
+ socialId
)
const keyValueClient = getKvsClient(this.integrationToken)
this.syncManager = new SyncManager(
diff --git a/services/gmail/pod-gmail/src/message/adapter.ts b/services/gmail/pod-gmail/src/message/adapter.ts
new file mode 100644
index 0000000000..8aed40a5ac
--- /dev/null
+++ b/services/gmail/pod-gmail/src/message/adapter.ts
@@ -0,0 +1,38 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { MeasureContext, TxOperations, SocialId } from '@hcengineering/core'
+
+import config from '../config'
+import { AttachmentHandler } from './attachments'
+import { MessageManagerV2 } from './v2/message'
+import { MessageManagerV1 } from './v1/message'
+import { type IMessageManager } from './types'
+import { type Channel } from '../types'
+
+export function createMessageManager (
+ ctx: MeasureContext,
+ client: TxOperations,
+ attachmentHandler: AttachmentHandler,
+ workspace: { getChannel: (email: string) => Channel | undefined },
+ token: string,
+ socialId: SocialId
+): IMessageManager {
+ if (config.Version === 'v2') {
+ return new MessageManagerV2(ctx, attachmentHandler, token, socialId)
+ } else {
+ return new MessageManagerV1(ctx, client, attachmentHandler, socialId._id, workspace)
+ }
+}
diff --git a/services/gmail/pod-gmail/src/message/attachments.ts b/services/gmail/pod-gmail/src/message/attachments.ts
index 6cf3c325f4..0087a7768b 100644
--- a/services/gmail/pod-gmail/src/message/attachments.ts
+++ b/services/gmail/pod-gmail/src/message/attachments.ts
@@ -12,13 +12,14 @@
// See the License for the specific language governing permissions and
// limitations under the License.
//
+import { randomUUID } from 'crypto'
import attachment, { Attachment } from '@hcengineering/attachment'
import { AttachedData, Blob, MeasureContext, Ref, TxOperations, WorkspaceUuid } from '@hcengineering/core'
import { StorageAdapter } from '@hcengineering/server-core'
import { gmail_v1 } from 'googleapis'
import { v4 as uuid } from 'uuid'
import { encode64 } from '../base64'
-import type { AttachedFile } from '../types'
+import type { AttachedFile } from './types'
import { addFooter } from '../utils'
export class AttachmentHandler {
@@ -39,11 +40,11 @@ export class AttachmentHandler {
const data: AttachedData = {
name: file.name,
file: id as Ref,
- type: file.type ?? 'undefined',
- size: file.size ?? Buffer.from(file.file, 'base64').length,
+ type: file.data.toString('base64') ?? 'undefined',
+ size: file.size ?? file.data.length,
lastModified: file.lastModified
}
- await this.storageAdapter.put(this.ctx, this.workspaceId as any, id, file.file, data.type, data.size) // TODO: FIXME
+ await this.storageAdapter.put(this.ctx, this.workspaceId as any, id, file.data, data.type, data.size) // TODO: FIXME
await this.client.addCollection(
attachment.class.Attachment,
message.space,
@@ -65,9 +66,10 @@ export class AttachmentHandler {
if (attachment.data == null) return []
return [
{
- file: attachment.data,
+ id: randomUUID(),
name: part.filename,
- type: part.mimeType ?? undefined,
+ data: Buffer.from(attachment.data, 'base64'),
+ contentType: part.mimeType ?? 'application/octet-stream',
size: attachment.size ?? undefined,
lastModified: new Date().getTime()
}
@@ -76,9 +78,10 @@ export class AttachmentHandler {
if (part.body?.data == null) return []
return [
{
- file: part.body.data,
+ id: randomUUID(),
+ data: Buffer.from(part.body.data, 'base64'),
name: part.filename,
- type: part.mimeType ?? undefined,
+ contentType: part.mimeType ?? 'application/octet-stream',
size: part.body.size ?? undefined,
lastModified: new Date().getTime()
}
diff --git a/services/gmail/pod-gmail/src/message/sync.ts b/services/gmail/pod-gmail/src/message/sync.ts
index 28f1311fdc..cbd9759b98 100644
--- a/services/gmail/pod-gmail/src/message/sync.ts
+++ b/services/gmail/pod-gmail/src/message/sync.ts
@@ -17,85 +17,26 @@ import { gmail_v1 } from 'googleapis'
import { type MeasureContext, PersonId } from '@hcengineering/core'
import { type KeyValueClient } from '@hcengineering/kvs-client'
+import { SyncMutex } from '@hcengineering/mail-common'
import { RateLimiter } from '../rateLimiter'
-import { MessageManager } from './message'
-
-interface History {
- historyId: string
- userId: string
- workspace: string
-}
-
-export class SyncMutex {
- private readonly locks = new Map>()
-
- async lock (key: string): Promise<() => void> {
- // Wait for any existing lock to be released
- const currentLock = this.locks.get(key)
- if (currentLock != null) {
- await currentLock
- }
-
- // Create a new lock
- let releaseFn!: () => void
- const newLock = new Promise((resolve) => {
- releaseFn = resolve
- })
-
- // Store the lock
- this.locks.set(key, newLock)
-
- // Return the release function
- return () => {
- if (this.locks.get(key) === newLock) {
- this.locks.delete(key)
- }
- releaseFn()
- }
- }
-}
+import { IMessageManager } from './types'
+import { SyncStateManager } from './syncState'
+import config from '../config'
export class SyncManager {
private readonly syncMutex = new SyncMutex()
+ private readonly stateManager: SyncStateManager
constructor (
private readonly ctx: MeasureContext,
- private readonly messageManager: MessageManager,
+ private readonly messageManager: IMessageManager,
private readonly gmail: gmail_v1.Resource$Users,
private readonly workspace: string,
- private readonly keyValueClient: KeyValueClient,
+ keyValueClient: KeyValueClient,
private readonly rateLimiter: RateLimiter
- ) {}
-
- private async getHistory (userId: PersonId): Promise {
- const historyKey = this.getHistoryKey(userId)
- return await this.keyValueClient.getValue(historyKey)
- }
-
- private async clearHistory (userId: PersonId): Promise {
- const historyKey = this.getHistoryKey(userId)
- await this.keyValueClient.deleteKey(historyKey)
- }
-
- private async setHistoryId (userId: PersonId, historyId: string): Promise {
- const historyKey = this.getHistoryKey(userId)
- const history: History = {
- historyId,
- userId,
- workspace: this.workspace
- }
- await this.keyValueClient.setValue(historyKey, history)
- }
-
- private async getPageToken (userId: PersonId): Promise {
- const pageTokenKey = this.getPageTokenKey(userId)
- return await this.keyValueClient.getValue(pageTokenKey)
- }
-
- private async setPageToken (userId: PersonId, pageToken: string): Promise {
- const pageTokenKey = this.getPageTokenKey(userId)
- await this.keyValueClient.setValue(pageTokenKey, pageToken)
+ ) {
+ this.stateManager = new SyncStateManager(keyValueClient, workspace, config.Version)
}
private async partSync (userId: PersonId, userEmail: string | undefined, historyId: string): Promise {
@@ -115,7 +56,7 @@ export class SyncManager {
})
} catch (err: any) {
this.ctx.error('Part sync get history error', { workspaceUuid: this.workspace, userId, error: err.message })
- await this.clearHistory(userId)
+ await this.stateManager.clearHistory(userId)
void this.sync(userId)
return
}
@@ -141,7 +82,7 @@ export class SyncManager {
}
}
if (history.id != null) {
- await this.setHistoryId(userId, history.id)
+ await this.stateManager.setHistoryId(userId, history.id)
}
}
if (nextPageToken == null) {
@@ -162,8 +103,8 @@ export class SyncManager {
throw new Error('Cannot sync without user email')
}
- // Get saved page token if exists to resume sync
- let pageToken: string | undefined = (await this.getPageToken(userId)) ?? undefined
+ // Get saved page token to continue from
+ let pageToken: string | undefined = (await this.stateManager.getPageToken(userId)) ?? undefined
const query: gmail_v1.Params$Resource$Users$Messages$List = {
userId: 'me',
@@ -218,11 +159,11 @@ export class SyncManager {
// Update page token for the next iteration
pageToken = messages.data.nextPageToken
query.pageToken = pageToken
- await this.setPageToken(userId, pageToken)
+ await this.stateManager.setPageToken(userId, pageToken)
}
if (currentHistoryId != null) {
- await this.setHistoryId(userId, currentHistoryId)
+ await this.stateManager.setHistoryId(userId, currentHistoryId)
}
this.ctx.info('Full sync finished', { workspaceUuid: this.workspace, userId, userEmail })
} catch (err) {
@@ -246,7 +187,7 @@ export class SyncManager {
try {
this.ctx.info('Sync history', { workspaceUuid: this.workspace, userId, userEmail })
- const history = await this.getHistory(userId)
+ const history = await this.stateManager.getHistory(userId)
if (history?.historyId != null && history?.historyId !== '') {
this.ctx.info('Start part sync', { workspaceUuid: this.workspace, userId, historyId: history.historyId })
await this.partSync(userId, userEmail, history.historyId)
@@ -260,12 +201,4 @@ export class SyncManager {
releaseLock()
}
}
-
- private getHistoryKey (userId: PersonId): string {
- return `history:${this.workspace}:${userId}`
- }
-
- private getPageTokenKey (userId: PersonId): string {
- return `page-token:${this.workspace}:${userId}`
- }
}
diff --git a/services/gmail/pod-gmail/src/message/syncState.ts b/services/gmail/pod-gmail/src/message/syncState.ts
new file mode 100644
index 0000000000..efbb6d26bc
--- /dev/null
+++ b/services/gmail/pod-gmail/src/message/syncState.ts
@@ -0,0 +1,74 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { PersonId } from '@hcengineering/core'
+import { type KeyValueClient } from '@hcengineering/kvs-client'
+import { type History } from './types'
+import { IntegrationVersion } from '../types'
+
+/**
+ * Handles persistent storage for Gmail sync state
+ */
+export class SyncStateManager {
+ constructor (
+ private readonly keyValueClient: KeyValueClient,
+ private readonly workspace: string,
+ private readonly version: IntegrationVersion
+ ) {}
+
+ async getHistory (userId: PersonId): Promise {
+ const historyKey = this.getHistoryKey(userId)
+ return await this.keyValueClient.getValue(historyKey)
+ }
+
+ async clearHistory (userId: PersonId): Promise {
+ const historyKey = this.getHistoryKey(userId)
+ await this.keyValueClient.deleteKey(historyKey)
+ }
+
+ async setHistoryId (userId: PersonId, historyId: string): Promise {
+ const historyKey = this.getHistoryKey(userId)
+ const history: History = {
+ historyId,
+ userId,
+ workspace: this.workspace
+ }
+ await this.keyValueClient.setValue(historyKey, history)
+ }
+
+ async getPageToken (userId: PersonId): Promise {
+ const pageTokenKey = this.getPageTokenKey(userId)
+ return await this.keyValueClient.getValue(pageTokenKey)
+ }
+
+ async setPageToken (userId: PersonId, pageToken: string): Promise {
+ const pageTokenKey = this.getPageTokenKey(userId)
+ await this.keyValueClient.setValue(pageTokenKey, pageToken)
+ }
+
+ private getHistoryKey (userId: PersonId): string {
+ if (this.version === IntegrationVersion.V2) {
+ return `history-v2:${this.workspace}:${userId}`
+ }
+ return `history:${this.workspace}:${userId}`
+ }
+
+ private getPageTokenKey (userId: PersonId): string {
+ if (this.version === IntegrationVersion.V2) {
+ return `page-token-v2:${this.workspace}:${userId}`
+ }
+ return `page-token:${this.workspace}:${userId}`
+ }
+}
diff --git a/services/gmail/pod-gmail/src/message/types.ts b/services/gmail/pod-gmail/src/message/types.ts
new file mode 100644
index 0000000000..9c604f50c4
--- /dev/null
+++ b/services/gmail/pod-gmail/src/message/types.ts
@@ -0,0 +1,42 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the 'License');
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an 'AS IS' BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+import { type GaxiosResponse } from 'gaxios'
+import { gmail_v1 } from 'googleapis'
+
+export interface AttachedFile {
+ id: string
+ name: string
+ data: Buffer
+ contentType: string
+ size?: number
+ lastModified: number
+}
+
+export interface EmailContact {
+ email: string
+ firstName?: string
+ lastName?: string
+ photoUrl?: string | null
+}
+
+export interface History {
+ historyId: string
+ userId: string
+ workspace: string
+}
+
+export interface IMessageManager {
+ saveMessage: (message: GaxiosResponse, me: string) => Promise
+}
diff --git a/services/gmail/pod-gmail/src/message/message.ts b/services/gmail/pod-gmail/src/message/v1/message.ts
similarity index 96%
rename from services/gmail/pod-gmail/src/message/message.ts
rename to services/gmail/pod-gmail/src/message/v1/message.ts
index 87b7e7187a..ff9be48cf1 100644
--- a/services/gmail/pod-gmail/src/message/message.ts
+++ b/services/gmail/pod-gmail/src/message/v1/message.ts
@@ -32,15 +32,16 @@ import core from '@hcengineering/core'
import attachment, { Attachment } from '@hcengineering/attachment'
import sanitizeHtml from 'sanitize-html'
-import { type Channel } from '../types'
-import { AttachmentHandler } from './attachments'
-import { decode64 } from '../base64'
-import { diffAttributes } from '../utils'
+import { IMessageManager } from '../types'
+import { type Channel } from '../../types'
+import { AttachmentHandler } from '../attachments'
+import { decode64 } from '../../base64'
+import { diffAttributes } from '../../utils'
const EMAIL_REGEX =
/(([^<>()[\]\\.,;:\s@"]+(\.[^<>()[\]\\.,;:\s@"]+)*)|(".+"))@((\[[0-9]{1,3}\.[0-9]{1,3}\.[0-9]{1,3}\.[0-9]{1,3}])|(([a-zA-Z\-0-9]+\.)+[a-zA-Z]{2,}))/
-export class MessageManager {
+export class MessageManagerV1 implements IMessageManager {
constructor (
private readonly ctx: MeasureContext,
private readonly client: TxOperations,
diff --git a/services/gmail/pod-gmail/src/message/v2/message.ts b/services/gmail/pod-gmail/src/message/v2/message.ts
new file mode 100644
index 0000000000..cf16594d96
--- /dev/null
+++ b/services/gmail/pod-gmail/src/message/v2/message.ts
@@ -0,0 +1,100 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+import { type GaxiosResponse } from 'gaxios'
+import { gmail_v1 } from 'googleapis'
+import sanitizeHtml from 'sanitize-html'
+
+import { SocialId, type MeasureContext } from '@hcengineering/core'
+import { createMessages, parseEmailHeader, parseNameFromEmailHeader, EmailMessage } from '@hcengineering/mail-common'
+
+import { IMessageManager } from '../types'
+import config from '../../config'
+import { AttachmentHandler } from '../attachments'
+import { decode64 } from '../../base64'
+
+export class MessageManagerV2 implements IMessageManager {
+ constructor (
+ private readonly ctx: MeasureContext,
+ private readonly attachmentHandler: AttachmentHandler,
+ private readonly token: string,
+ private readonly socialId: SocialId
+ ) {}
+
+ async saveMessage (message: GaxiosResponse, me: string): Promise {
+ const res = convertMessage(message, me)
+ const attachments = await this.attachmentHandler.getPartFiles(message.data.payload, message.data.id ?? '')
+
+ await createMessages(config, this.ctx, this.token, res, attachments, me, this.socialId)
+ }
+}
+
+function getHeaderValue (payload: gmail_v1.Schema$MessagePart | undefined, name: string): string | undefined {
+ if (payload === undefined) return undefined
+ const headers = payload.headers
+
+ return headers?.find((header) => header.name?.toLowerCase() === name.toLowerCase())?.value ?? undefined
+}
+
+function getPartsMessage (parts: gmail_v1.Schema$MessagePart[] | undefined, mime: string): string {
+ let result = ''
+ if (parts !== undefined) {
+ const htmlPart = parts.find((part) => part.mimeType === mime)
+ const filtredParts = htmlPart !== undefined ? parts.filter((part) => part.mimeType === mime) : parts
+ for (const part of filtredParts ?? []) {
+ result += getPartMessage(part, mime)
+ }
+ }
+ return result
+}
+
+const sanitizeOptions: sanitizeHtml.IOptions = {
+ allowedTags: [],
+ allowedAttributes: {}
+}
+
+export function sanitizeText (input: string): string {
+ if (input == null) return ''
+ return sanitizeHtml(input, sanitizeOptions)
+}
+
+function getPartMessage (part: gmail_v1.Schema$MessagePart | undefined, mime: string): string {
+ if (part === undefined) return ''
+ if (part.body?.data != null) {
+ return decode64(part.body.data)
+ }
+ return getPartsMessage(part.parts, mime)
+}
+
+function convertMessage (message: GaxiosResponse, me: string): EmailMessage {
+ const date = message.data.internalDate != null ? new Date(Number.parseInt(message.data.internalDate)) : new Date()
+ const from = parseNameFromEmailHeader(getHeaderValue(message.data.payload, 'From') ?? '')
+ const to = parseEmailHeader(getHeaderValue(message.data.payload, 'To') ?? '')
+
+ const copy = parseEmailHeader(getHeaderValue(message.data.payload, 'Cc') ?? '')
+ const incoming = !from.email.includes(me)
+ return {
+ modifiedOn: date.getTime(),
+ mailId: getHeaderValue(message.data.payload, 'Message-ID') ?? '',
+ replyTo: getHeaderValue(message.data.payload, 'In-Reply-To'),
+ copy,
+ content: sanitizeHtml(getPartMessage(message.data.payload, 'text/html')),
+ textContent: sanitizeText(getPartMessage(message.data.payload, 'text/plain')),
+ from,
+ to,
+ incoming,
+ subject: getHeaderValue(message.data.payload, 'Subject') ?? '',
+ sendOn: date.getTime()
+ }
+}
diff --git a/services/gmail/pod-gmail/src/types.ts b/services/gmail/pod-gmail/src/types.ts
index 2e0b3b3184..1e4ddc0a74 100644
--- a/services/gmail/pod-gmail/src/types.ts
+++ b/services/gmail/pod-gmail/src/types.ts
@@ -37,14 +37,6 @@ export type State = User & {
redirectURL: string
}
-export interface AttachedFile {
- size?: number
- file: string
- type?: string
- lastModified: number
- name: string
-}
-
export type Channel = Pick
export type RequestType = 'get' | 'post'
@@ -75,3 +67,8 @@ export const GMAIL_INTEGRATION = 'gmail'
export enum SecretType {
TOKEN = 'token'
}
+
+export enum IntegrationVersion {
+ V1 = 'v1', // Save messages in legacy format using gmail.class.Message
+ V2 = 'v2' // Save messages as thread cards and communication messages
+}
diff --git a/services/mail/mail-common/.eslintrc.js b/services/mail/mail-common/.eslintrc.js
new file mode 100644
index 0000000000..72235dc283
--- /dev/null
+++ b/services/mail/mail-common/.eslintrc.js
@@ -0,0 +1,7 @@
+module.exports = {
+ extends: ['./node_modules/@hcengineering/platform-rig/profiles/default/eslint.config.json'],
+ parserOptions: {
+ tsconfigRootDir: __dirname,
+ project: './tsconfig.json'
+ }
+}
diff --git a/services/mail/mail-common/.gitignore b/services/mail/mail-common/.gitignore
new file mode 100644
index 0000000000..2eea525d88
--- /dev/null
+++ b/services/mail/mail-common/.gitignore
@@ -0,0 +1 @@
+.env
\ No newline at end of file
diff --git a/services/mail/mail-common/.npmignore b/services/mail/mail-common/.npmignore
new file mode 100644
index 0000000000..e3ec093c38
--- /dev/null
+++ b/services/mail/mail-common/.npmignore
@@ -0,0 +1,4 @@
+*
+!/lib/**
+!CHANGELOG.md
+/lib/**/__tests__/
diff --git a/services/mail/mail-common/config/rig.json b/services/mail/mail-common/config/rig.json
new file mode 100644
index 0000000000..0110930f55
--- /dev/null
+++ b/services/mail/mail-common/config/rig.json
@@ -0,0 +1,4 @@
+{
+ "$schema": "https://developer.microsoft.com/json-schemas/rig-package/rig.schema.json",
+ "rigPackageName": "@hcengineering/platform-rig"
+}
diff --git a/services/mail/mail-common/jest.config.js b/services/mail/mail-common/jest.config.js
new file mode 100644
index 0000000000..2cfd408b67
--- /dev/null
+++ b/services/mail/mail-common/jest.config.js
@@ -0,0 +1,7 @@
+module.exports = {
+ preset: 'ts-jest',
+ testEnvironment: 'node',
+ testMatch: ['**/?(*.)+(spec|test).[jt]s?(x)'],
+ roots: ["./src"],
+ coverageReporters: ["text-summary", "html"]
+}
diff --git a/services/mail/mail-common/package.json b/services/mail/mail-common/package.json
new file mode 100644
index 0000000000..d00d91d552
--- /dev/null
+++ b/services/mail/mail-common/package.json
@@ -0,0 +1,62 @@
+{
+ "name": "@hcengineering/mail-common",
+ "version": "0.6.0",
+ "main": "lib/index.js",
+ "svelte": "src/index.ts",
+ "types": "types/index.d.ts",
+ "files": [
+ "lib/**/*",
+ "types/**/*",
+ "tsconfig.json"
+ ],
+ "author": "Hardcore Engineering Inc.",
+ "scripts": {
+ "build": "compile",
+ "build:watch": "compile",
+ "format": "format src",
+ "test": "jest --passWithNoTests --silent",
+ "_phase:build": "compile transpile src",
+ "_phase:test": "jest --passWithNoTests --silent",
+ "_phase:format": "format src",
+ "_phase:validate": "compile validate"
+ },
+ "devDependencies": {
+ "@hcengineering/platform-rig": "^0.6.0",
+ "@tsconfig/node16": "^1.0.4",
+ "@types/express": "^4.17.13",
+ "@types/jest": "^29.5.5",
+ "@types/node": "~20.11.16",
+ "@types/turndown": "^5.0.5",
+ "@types/sanitize-html": "^2.15.0",
+ "@types/uuid": "^8.3.1",
+ "@typescript-eslint/eslint-plugin": "^6.11.0",
+ "@typescript-eslint/parser": "^6.11.0",
+ "esbuild": "^0.24.2",
+ "eslint": "^8.54.0",
+ "eslint-config-standard-with-typescript": "^40.0.0",
+ "eslint-plugin-import": "^2.26.0",
+ "eslint-plugin-n": "^15.4.0",
+ "eslint-plugin-node": "^11.1.0",
+ "eslint-plugin-promise": "^6.1.1",
+ "jest": "^29.7.0",
+ "prettier": "^3.1.0",
+ "ts-jest": "^29.1.1",
+ "ts-node": "^10.8.0",
+ "typescript": "^5.3.3"
+ },
+ "dependencies": {
+ "@hcengineering/account-client": "^0.6.0",
+ "@hcengineering/api-client": "^0.6.0",
+ "@hcengineering/card": "^0.6.0",
+ "@hcengineering/chat": "^0.6.0",
+ "@hcengineering/communication-rest-client": "0.1.178",
+ "@hcengineering/communication-types": "0.1.178",
+ "@hcengineering/contact": "^0.6.24",
+ "@hcengineering/core": "^0.6.32",
+ "@hcengineering/mail": "^0.6.0",
+ "@hcengineering/server-storage": "^0.6.0",
+ "sanitize-html": "^2.15.0",
+ "turndown": "^7.2.0",
+ "uuid": "^8.3.2"
+ }
+}
diff --git a/services/mail/mail-common/src/__tests__/channel.test.ts b/services/mail/mail-common/src/__tests__/channel.test.ts
new file mode 100644
index 0000000000..64454cb84a
--- /dev/null
+++ b/services/mail/mail-common/src/__tests__/channel.test.ts
@@ -0,0 +1,315 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { PersonId, Ref, WorkspaceUuid, MeasureContext, TxOperations, Doc, SocialId } from '@hcengineering/core'
+import { PersonSpace } from '@hcengineering/contact'
+import chat from '@hcengineering/chat'
+import mail from '@hcengineering/mail'
+import { ChannelCache, ChannelCacheFactory } from '../channel'
+
+/* eslint-disable @typescript-eslint/unbound-method */
+describe('ChannelCache', () => {
+ let mockCtx: jest.Mocked
+ let mockClient: jest.Mocked
+ let channelCache: ChannelCache
+
+ const workspace = 'test-workspace' as WorkspaceUuid
+ const spaceId = 'test-space-id' as Ref
+ const emailAccount = 'test@example.com'
+ const participants: PersonId[] = ['person1', 'person2'] as PersonId[]
+ const socialId: SocialId = { _id: 'social-id' as PersonId } as any
+
+ const mockChannel = {
+ _id: 'channel-id' as Ref,
+ title: emailAccount
+ }
+
+ const generatedId = 'generated-id' as Ref
+
+ beforeEach(() => {
+ jest.clearAllMocks()
+
+ mockCtx = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockClient = {
+ findOne: jest.fn(),
+ createDoc: jest.fn(),
+ createMixin: jest.fn()
+ } as unknown as jest.Mocked
+
+ channelCache = new ChannelCache(mockCtx, mockClient, workspace)
+ })
+
+ describe('getOrCreateChannel', () => {
+ it('should return cached channel if it exists', async () => {
+ // Set up a channel in the cache
+ const cache = (channelCache as any).cache
+ cache.set(`${spaceId}:${emailAccount}`, mockChannel._id)
+
+ const result = await channelCache.getOrCreateChannel(spaceId, participants, emailAccount, socialId)
+
+ expect(result).toBe(mockChannel._id)
+ expect(mockClient.findOne).not.toHaveBeenCalled()
+ })
+
+ it('should fetch existing channel if not in cache', async () => {
+ mockClient.findOne.mockResolvedValue(mockChannel as any)
+
+ const result = await channelCache.getOrCreateChannel(spaceId, participants, emailAccount, socialId)
+
+ expect(result).toBe(mockChannel._id)
+ expect(mockClient.findOne).toHaveBeenCalledWith(mail.tag.MailChannel, { title: emailAccount })
+ expect(mockClient.createDoc).not.toHaveBeenCalled()
+ expect(mockCtx.info).toHaveBeenCalledWith('Using existing channel', {
+ me: emailAccount,
+ space: spaceId,
+ channel: mockChannel._id
+ })
+ })
+
+ it('should create new channel if it does not exist', async () => {
+ // First findOne returns null (no existing channel)
+ // Second findOne (inside createNewChannel) also returns null
+ mockClient.findOne.mockResolvedValue(undefined)
+ mockClient.createDoc.mockResolvedValue(generatedId)
+ mockClient.createMixin.mockResolvedValue(undefined as any)
+
+ const result = await channelCache.getOrCreateChannel(spaceId, participants, emailAccount, socialId)
+
+ expect(result).toBe(generatedId)
+ expect(mockClient.findOne).toHaveBeenCalledTimes(2)
+ expect(mockClient.createDoc).toHaveBeenCalledWith(
+ chat.masterTag.Channel,
+ spaceId,
+ {
+ title: emailAccount,
+ private: true,
+ members: participants,
+ archived: false,
+ createdBy: socialId._id,
+ modifiedBy: socialId._id
+ },
+ expect.any(String),
+ expect.any(Number),
+ socialId._id
+ )
+ expect(mockClient.createMixin).toHaveBeenCalledWith(
+ expect.any(String),
+ chat.masterTag.Channel,
+ spaceId,
+ mail.tag.MailChannel,
+ {},
+ expect.any(Number),
+ socialId._id
+ )
+ })
+
+ it('should use existing channel if found after acquiring mutex lock', async () => {
+ // First findOne returns null (trigger createNewChannel)
+ // Second findOne inside createNewChannel returns the channel (simulate race condition handled)
+ mockClient.findOne.mockResolvedValueOnce(null as any).mockResolvedValueOnce(mockChannel as any)
+
+ const result = await channelCache.getOrCreateChannel(spaceId, participants, emailAccount, socialId)
+
+ expect(result).toBe(mockChannel._id)
+ expect(mockClient.findOne).toHaveBeenCalledTimes(2)
+ expect(mockClient.createDoc).not.toHaveBeenCalled()
+ expect(mockCtx.info).toHaveBeenCalledWith('Using existing channel (found after mutex lock)', {
+ me: emailAccount,
+ space: spaceId,
+ channel: mockChannel._id
+ })
+ })
+
+ it('should handle errors and remove failed lookup from cache', async () => {
+ const error = new Error('Database error')
+ mockClient.findOne.mockRejectedValue(error)
+
+ const result = await channelCache.getOrCreateChannel(spaceId, participants, emailAccount, socialId)
+
+ expect(result).toBeUndefined()
+ expect(mockCtx.error).toHaveBeenCalledWith('Failed to create channel', {
+ me: emailAccount,
+ space: spaceId,
+ workspace,
+ error: error.message
+ })
+
+ // Verify the cache doesn't contain the failed lookup
+ expect((channelCache as any).cache.has(`${spaceId}:${emailAccount}`)).toBe(false)
+ })
+ })
+
+ describe('clearCache', () => {
+ it('should clear cache for specific space and email', async () => {
+ // Set up some items in cache
+ const cache = (channelCache as any).cache
+ cache.set(`${spaceId}:${emailAccount}`, mockChannel._id)
+ cache.set(`${spaceId}:other@example.com`, 'another-id')
+
+ channelCache.clearCache(spaceId, emailAccount)
+
+ expect(cache.has(`${spaceId}:${emailAccount}`)).toBe(false)
+ expect(cache.has(`${spaceId}:other@example.com`)).toBe(true)
+ })
+ })
+
+ describe('clearAllCache', () => {
+ it('should clear all cached channels', async () => {
+ // Set up multiple items in cache
+ const cache = (channelCache as any).cache
+ cache.set(`${spaceId}:${emailAccount}`, mockChannel._id)
+ cache.set(`${spaceId}:other@example.com`, 'another-id')
+ cache.set(`other-space:${emailAccount}`, 'third-id')
+
+ channelCache.clearAllCache()
+
+ expect(cache.size).toBe(0)
+ })
+ })
+
+ describe('size', () => {
+ it('should return the number of cached channels', async () => {
+ const cache = (channelCache as any).cache
+
+ expect(channelCache.size).toBe(0)
+
+ cache.set(`${spaceId}:${emailAccount}`, mockChannel._id)
+ expect(channelCache.size).toBe(1)
+
+ cache.set(`${spaceId}:other@example.com`, 'another-id')
+ expect(channelCache.size).toBe(2)
+
+ channelCache.clearCache(spaceId, emailAccount)
+ expect(channelCache.size).toBe(1)
+
+ channelCache.clearAllCache()
+ expect(channelCache.size).toBe(0)
+ })
+ })
+})
+
+describe('ChannelCacheFactory', () => {
+ let mockCtx1: jest.Mocked
+ let mockClient1: jest.Mocked
+
+ const workspace1 = 'workspace-1' as WorkspaceUuid
+ const workspace2 = 'workspace-2' as WorkspaceUuid
+
+ beforeEach(() => {
+ jest.clearAllMocks()
+ ChannelCacheFactory.resetAllInstances()
+
+ mockCtx1 = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockClient1 = {} as unknown as jest.Mocked
+ })
+
+ describe('getInstance', () => {
+ it('should create a new instance for a workspace', () => {
+ expect(ChannelCacheFactory.instanceCount).toBe(0)
+
+ const cache = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+
+ expect(cache).toBeInstanceOf(ChannelCache)
+ expect(ChannelCacheFactory.instanceCount).toBe(1)
+ })
+
+ it('should return the same instance for the same workspace', () => {
+ const cache1 = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ const cache2 = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+
+ expect(cache1).toBe(cache2)
+ expect(ChannelCacheFactory.instanceCount).toBe(1)
+ })
+
+ it('should create different instances for different workspaces', () => {
+ const cache1 = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ const cache2 = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+
+ expect(cache1).not.toBe(cache2)
+ expect(ChannelCacheFactory.instanceCount).toBe(2)
+ })
+ })
+
+ describe('resetInstance', () => {
+ it('should remove the instance for a specific workspace', () => {
+ const cache1 = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ const cache2 = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+
+ expect(ChannelCacheFactory.instanceCount).toBe(2)
+
+ ChannelCacheFactory.resetInstance(workspace1)
+
+ expect(ChannelCacheFactory.instanceCount).toBe(1)
+
+ // Getting workspace1 again should create a new instance
+ const cache1New = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ expect(cache1New).not.toBe(cache1)
+
+ // Workspace2 instance should remain the same
+ const cache2Again = ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+ expect(cache2Again).toBe(cache2)
+ })
+
+ it('should do nothing when resetting a non-existent workspace', () => {
+ ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ expect(ChannelCacheFactory.instanceCount).toBe(1)
+
+ // Reset a workspace that doesn't have a cache
+ ChannelCacheFactory.resetInstance('non-existent-workspace' as WorkspaceUuid)
+
+ expect(ChannelCacheFactory.instanceCount).toBe(1)
+ })
+ })
+
+ describe('resetAllInstances', () => {
+ it('should remove all workspace instances', () => {
+ ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+ expect(ChannelCacheFactory.instanceCount).toBe(2)
+
+ ChannelCacheFactory.resetAllInstances()
+
+ expect(ChannelCacheFactory.instanceCount).toBe(0)
+ })
+ })
+
+ describe('instanceCount', () => {
+ it('should return the number of workspace instances', () => {
+ expect(ChannelCacheFactory.instanceCount).toBe(0)
+
+ ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ expect(ChannelCacheFactory.instanceCount).toBe(1)
+
+ ChannelCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+ expect(ChannelCacheFactory.instanceCount).toBe(2)
+
+ ChannelCacheFactory.resetInstance(workspace1)
+ expect(ChannelCacheFactory.instanceCount).toBe(1)
+
+ ChannelCacheFactory.resetAllInstances()
+ expect(ChannelCacheFactory.instanceCount).toBe(0)
+ })
+ })
+})
diff --git a/services/mail/mail-common/src/__tests__/mutex.test.ts b/services/mail/mail-common/src/__tests__/mutex.test.ts
new file mode 100644
index 0000000000..a7de310a97
--- /dev/null
+++ b/services/mail/mail-common/src/__tests__/mutex.test.ts
@@ -0,0 +1,179 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { SyncMutex } from '../mutex'
+
+describe('SyncMutex', () => {
+ let mutex: SyncMutex
+
+ beforeEach(() => {
+ mutex = new SyncMutex()
+ })
+
+ it('should allow a lock to be acquired and released', async () => {
+ const release = await mutex.lock('key1')
+ expect(typeof release).toBe('function')
+ release()
+ })
+
+ it('should allow different keys to be locked simultaneously', async () => {
+ const timestamps: number[] = []
+
+ void (async () => {
+ const release = await mutex.lock('key1')
+ timestamps.push(Date.now())
+ await new Promise((resolve) => setTimeout(resolve, 50))
+ release()
+ })()
+
+ void (async () => {
+ const release = await mutex.lock('key2')
+ timestamps.push(Date.now())
+ await new Promise((resolve) => setTimeout(resolve, 50))
+ release()
+ })()
+
+ // Give some time for both locks to be acquired
+ await new Promise((resolve) => setTimeout(resolve, 20))
+
+ // Both locks should be acquired within 20ms of each other
+ // since they use different keys and shouldn't block each other
+ expect(timestamps.length).toBe(2)
+ expect(Math.abs(timestamps[0] - timestamps[1])).toBeLessThan(20)
+ })
+
+ it('should make subsequent locks wait for release', async () => {
+ const executionOrder: string[] = []
+
+ // First lock
+ void (async () => {
+ const release = await mutex.lock('key1')
+ executionOrder.push('lock1 acquired')
+ await new Promise((resolve) => setTimeout(resolve, 50))
+ executionOrder.push('lock1 releasing')
+ release()
+ })()
+
+ // Small delay to ensure first lock is acquired first
+ await new Promise((resolve) => setTimeout(resolve, 10))
+
+ // Second lock on same key
+ void (async () => {
+ executionOrder.push('lock2 waiting')
+ const release = await mutex.lock('key1')
+ executionOrder.push('lock2 acquired')
+ release()
+ })()
+
+ // Wait for all operations to complete
+ await new Promise((resolve) => setTimeout(resolve, 100))
+
+ expect(executionOrder).toEqual(['lock1 acquired', 'lock2 waiting', 'lock1 releasing', 'lock2 acquired'])
+ })
+
+ it('should handle multiple sequential locks correctly', async () => {
+ const result: number[] = []
+
+ const lockAndExecute = async (value: number): Promise => {
+ const release = await mutex.lock('sequential')
+ result.push(value)
+ await new Promise((resolve) => setTimeout(resolve, 10)) // Simulate some work
+ release()
+ }
+
+ // Execute multiple locks in sequence
+ await Promise.all([lockAndExecute(1), lockAndExecute(2), lockAndExecute(3), lockAndExecute(4), lockAndExecute(5)])
+
+ // The values should be added to the result array in order
+ // since each lock waits for the previous one to complete
+ expect(result).toEqual([1, 2, 3, 4, 5])
+ })
+
+ it('should handle lock release correctly even if called multiple times', async () => {
+ const lockKey = 'multiple-release'
+ const release = await mutex.lock(lockKey)
+
+ // First release
+ release()
+
+ // Second release should not throw
+ release()
+
+ // We should be able to acquire the lock again
+ const newRelease = await mutex.lock(lockKey)
+ newRelease()
+ })
+
+ it('should handle errors within the locked code', async () => {
+ const lockKey = 'error-handling'
+ const executionOrder: string[] = []
+
+ // First lock with an error
+ try {
+ const release = await mutex.lock(lockKey)
+ executionOrder.push('lock1 acquired')
+ try {
+ throw new Error('Simulated error')
+ } finally {
+ executionOrder.push('lock1 releasing')
+ release()
+ }
+ } catch (err) {
+ executionOrder.push('error caught')
+ }
+
+ // Second lock should still work
+ const release2 = await mutex.lock(lockKey)
+ executionOrder.push('lock2 acquired')
+ release2()
+
+ expect(executionOrder).toEqual(['lock1 acquired', 'lock1 releasing', 'error caught', 'lock2 acquired'])
+ })
+
+ it('should work with complex async operations', async () => {
+ const results: string[] = []
+ const lockKey = 'async-operations'
+
+ // Helper function that performs complex async work under lock
+ const performWork = async (id: string, delay: number): Promise => {
+ const release = await mutex.lock(lockKey)
+ try {
+ results.push(`${id} started`)
+ await new Promise((resolve) => setTimeout(resolve, delay))
+ results.push(`${id} completed`)
+ } finally {
+ release()
+ }
+ }
+
+ // Start multiple work items with different delays
+ await Promise.all([performWork('task1', 30), performWork('task2', 10), performWork('task3', 20)])
+
+ // Tasks should be completed in the order they acquired the lock
+ expect(results.length).toBe(6)
+ expect(results.includes('task1 started')).toBe(true)
+ expect(results.includes('task1 completed')).toBe(true)
+ expect(results.includes('task2 started')).toBe(true)
+ expect(results.includes('task2 completed')).toBe(true)
+ expect(results.includes('task3 started')).toBe(true)
+ expect(results.includes('task3 completed')).toBe(true)
+
+ expect(results.indexOf('task1 started')).toBeLessThan(results.indexOf('task2 started'))
+ expect(results.indexOf('task2 started')).toBeLessThan(results.indexOf('task3 started'))
+
+ expect(results.indexOf('task1 completed')).toBeLessThan(results.indexOf('task2 completed'))
+ expect(results.indexOf('task2 completed')).toBeLessThan(results.indexOf('task3 completed'))
+ })
+})
diff --git a/services/mail/mail-common/src/__tests__/parseEmailHeader.test.ts b/services/mail/mail-common/src/__tests__/parseEmailHeader.test.ts
new file mode 100644
index 0000000000..cde80c4ea3
--- /dev/null
+++ b/services/mail/mail-common/src/__tests__/parseEmailHeader.test.ts
@@ -0,0 +1,237 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { parseEmailHeader } from '../utils'
+
+describe('parseEmailHeader', () => {
+ it('should handle undefined input', () => {
+ const result = parseEmailHeader(undefined)
+ expect(result).toEqual([])
+ })
+
+ it('should handle empty string input', () => {
+ const result = parseEmailHeader('')
+ expect(result).toEqual([])
+ })
+
+ it('should handle whitespace-only input', () => {
+ const result = parseEmailHeader(' \t\n ')
+ expect(result).toEqual([])
+ })
+
+ it('should parse a single plain email address', () => {
+ const result = parseEmailHeader('test@example.com')
+ expect(result).toEqual([
+ {
+ email: 'test@example.com',
+ firstName: 'test',
+ lastName: 'example.com'
+ }
+ ])
+ })
+
+ it('should parse a single email with name in angle brackets', () => {
+ const result = parseEmailHeader('John Doe ')
+ expect(result).toEqual([
+ {
+ email: 'john.doe@example.com',
+ firstName: 'John',
+ lastName: 'Doe'
+ }
+ ])
+ })
+
+ it('should parse a single email with quoted name in angle brackets', () => {
+ const result = parseEmailHeader('"John Doe" ')
+ expect(result).toEqual([
+ {
+ email: 'john.doe@example.com',
+ firstName: 'John',
+ lastName: 'Doe'
+ }
+ ])
+ })
+
+ it('should parse a single email with multi-word last name', () => {
+ const result = parseEmailHeader('John Doe Smith ')
+ expect(result).toEqual([
+ {
+ email: 'john.doe@example.com',
+ firstName: 'John',
+ lastName: 'Doe Smith'
+ }
+ ])
+ })
+
+ it('should parse multiple plain email addresses', () => {
+ const result = parseEmailHeader('test1@example.com, test2@example.com')
+ expect(result).toEqual([
+ {
+ email: 'test1@example.com',
+ firstName: 'test1',
+ lastName: 'example.com'
+ },
+ {
+ email: 'test2@example.com',
+ firstName: 'test2',
+ lastName: 'example.com'
+ }
+ ])
+ })
+
+ it('should parse multiple email addresses with names', () => {
+ const result = parseEmailHeader('John , Jane Doe ')
+ expect(result).toEqual([
+ {
+ email: 'john@example.com',
+ firstName: 'John',
+ lastName: 'example.com'
+ },
+ {
+ email: 'jane@example.com',
+ firstName: 'Jane',
+ lastName: 'Doe'
+ }
+ ])
+ })
+
+ it('should handle commas within quoted names', () => {
+ const result = parseEmailHeader('"Doe, John" , Jane Doe ')
+ expect(result).toEqual([
+ {
+ email: 'john@example.com',
+ firstName: 'Doe,',
+ lastName: 'John'
+ },
+ {
+ email: 'jane@example.com',
+ firstName: 'Jane',
+ lastName: 'Doe'
+ }
+ ])
+ })
+
+ it('should handle mixed format addresses', () => {
+ const result = parseEmailHeader('john@example.com, Jane Doe ')
+ expect(result).toEqual([
+ {
+ email: 'john@example.com',
+ firstName: 'john',
+ lastName: 'example.com'
+ },
+ {
+ email: 'jane@example.com',
+ firstName: 'Jane',
+ lastName: 'Doe'
+ }
+ ])
+ })
+
+ it('should handle email addresses with no name part', () => {
+ const result = parseEmailHeader(', ')
+ expect(result).toEqual([
+ {
+ email: 'john@example.com',
+ firstName: 'john',
+ lastName: 'example.com'
+ },
+ {
+ email: 'jane@example.com',
+ firstName: 'jane',
+ lastName: 'example.com'
+ }
+ ])
+ })
+
+ it('should handle extra whitespace between addresses', () => {
+ const result = parseEmailHeader('john@example.com , Jane Doe ')
+ expect(result).toEqual([
+ {
+ email: 'john@example.com',
+ firstName: 'john',
+ lastName: 'example.com'
+ },
+ {
+ email: 'jane@example.com',
+ firstName: 'Jane',
+ lastName: 'Doe'
+ }
+ ])
+ })
+
+ it('should handle the example from the prompt', () => {
+ const result = parseEmailHeader('example staff , personnel ')
+ expect(result).toEqual([
+ {
+ email: 'example-staff@example.com',
+ firstName: 'example',
+ lastName: 'staff'
+ },
+ {
+ email: 'personnel@example.com',
+ firstName: 'personnel',
+ lastName: 'example.com'
+ }
+ ])
+ })
+
+ it('should handle another example from the prompt', () => {
+ const result = parseEmailHeader('abc@test.com, 123@test.com')
+ expect(result).toEqual([
+ {
+ email: 'abc@test.com',
+ firstName: 'abc',
+ lastName: 'test.com'
+ },
+ {
+ email: '123@test.com',
+ firstName: '123',
+ lastName: 'test.com'
+ }
+ ])
+ })
+
+ it('should handle addresses with trailing comma', () => {
+ const result = parseEmailHeader('John , Jane ')
+ expect(result).toEqual([
+ {
+ email: 'john@example.com',
+ firstName: 'John',
+ lastName: 'example.com'
+ },
+ {
+ email: 'jane@example.com',
+ firstName: 'Jane',
+ lastName: 'example.com'
+ }
+ ])
+ })
+
+ it('should handle empty addresses between commas', () => {
+ const result = parseEmailHeader('John , , Jane ')
+ expect(result).toEqual([
+ {
+ email: 'john@example.com',
+ firstName: 'John',
+ lastName: 'example.com'
+ },
+ {
+ email: 'jane@example.com',
+ firstName: 'Jane',
+ lastName: 'example.com'
+ }
+ ])
+ })
+})
diff --git a/services/mail/mail-common/src/__tests__/parseNameFromEmailHeader.test.ts b/services/mail/mail-common/src/__tests__/parseNameFromEmailHeader.test.ts
new file mode 100644
index 0000000000..594a5285f1
--- /dev/null
+++ b/services/mail/mail-common/src/__tests__/parseNameFromEmailHeader.test.ts
@@ -0,0 +1,113 @@
+import { parseNameFromEmailHeader } from '../utils'
+import { EmailContact } from '../types'
+
+describe('parseNameFromEmailHeader', () => {
+ it('should parse email with name in double quotes', () => {
+ const input = '"John Doe" '
+ const expected: EmailContact = {
+ email: 'john.doe@example.com',
+ firstName: 'John',
+ lastName: 'Doe'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should parse email with name without quotes', () => {
+ const input = 'Jane Smith '
+ const expected: EmailContact = {
+ email: 'jane.smith@example.com',
+ firstName: 'Jane',
+ lastName: 'Smith'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should parse email without name', () => {
+ const input = 'no-reply@example.com'
+ const expected: EmailContact = {
+ email: 'no-reply@example.com',
+ firstName: 'no-reply',
+ lastName: 'example.com'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should parse email with angle brackets only', () => {
+ const input = ''
+ const expected: EmailContact = {
+ email: 'support@example.com',
+ firstName: 'support',
+ lastName: 'example.com'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should parse email with multi-word last name', () => {
+ const input = 'Maria Van Der Berg '
+ const expected: EmailContact = {
+ email: 'maria@example.com',
+ firstName: 'Maria',
+ lastName: 'Van Der Berg'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should handle undefined input', () => {
+ const expected: EmailContact = {
+ email: '',
+ firstName: '',
+ lastName: ''
+ }
+
+ expect(parseNameFromEmailHeader(undefined)).toEqual(expected)
+ })
+
+ it('should handle empty string input', () => {
+ const input = ''
+ const expected: EmailContact = {
+ email: '',
+ firstName: '',
+ lastName: ''
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should handle malformed email formats', () => {
+ const input = 'John Doe john.doe@example.com'
+ const expected: EmailContact = {
+ email: 'John Doe john.doe@example.com',
+ firstName: 'John Doe john.doe',
+ lastName: 'example.com'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should parse single-name format', () => {
+ const input = 'Support '
+ const expected: EmailContact = {
+ email: 'help@example.com',
+ firstName: 'Support',
+ lastName: 'example.com'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+
+ it('should handle name with special characters', () => {
+ const input = '"O\'Neill, James" '
+ const expected: EmailContact = {
+ email: 'james.oneill@example.com',
+ firstName: "O'Neill,",
+ lastName: 'James'
+ }
+
+ expect(parseNameFromEmailHeader(input)).toEqual(expected)
+ })
+})
diff --git a/services/mail/mail-common/src/__tests__/person.test.ts b/services/mail/mail-common/src/__tests__/person.test.ts
new file mode 100644
index 0000000000..f70dade268
--- /dev/null
+++ b/services/mail/mail-common/src/__tests__/person.test.ts
@@ -0,0 +1,181 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { MeasureContext, PersonId, PersonUuid, SocialIdType } from '@hcengineering/core'
+import { RestClient } from '@hcengineering/api-client'
+import { PersonCache, CachedPerson } from '../person'
+import { EmailContact } from '../types'
+
+describe('PersonCache', () => {
+ let personCache: PersonCache
+ let mockCtx: MeasureContext
+ let mockRestClient: jest.Mocked
+
+ const mockContact: EmailContact = {
+ email: 'test@example.com',
+ firstName: 'Test',
+ lastName: 'User'
+ }
+
+ const mockPerson: CachedPerson = {
+ socialId: 'test-social-id' as PersonId,
+ uuid: 'test-uuid' as PersonUuid,
+ localPerson: 'test-local-person'
+ }
+
+ beforeEach(() => {
+ // Set up mocks
+ mockCtx = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockRestClient = {
+ ensurePerson: jest.fn().mockResolvedValue(mockPerson)
+ } as unknown as jest.Mocked
+
+ // Create the cache instance
+ personCache = new PersonCache(mockCtx, mockRestClient)
+ })
+
+ afterEach(() => {
+ jest.resetAllMocks()
+ })
+
+ describe('ensurePerson', () => {
+ it('should call ensurePerson on the RestClient and return the result', async () => {
+ const result = await personCache.ensurePerson(mockContact)
+
+ expect(result).toEqual(mockPerson)
+ expect(mockRestClient.ensurePerson).toHaveBeenCalledWith(
+ SocialIdType.EMAIL,
+ mockContact.email,
+ mockContact.firstName,
+ mockContact.lastName
+ )
+ })
+
+ it('should normalize email addresses by trimming and lowercasing', async () => {
+ const uppercaseContact: EmailContact = {
+ email: ' TEST@EXAMPLE.COM ',
+ firstName: 'Test',
+ lastName: 'User'
+ }
+
+ await personCache.ensurePerson(uppercaseContact)
+
+ expect(mockRestClient.ensurePerson).toHaveBeenCalledWith(
+ SocialIdType.EMAIL,
+ 'test@example.com',
+ uppercaseContact.firstName,
+ uppercaseContact.lastName
+ )
+ })
+
+ it('should return cached result for subsequent calls with the same email', async () => {
+ await personCache.ensurePerson(mockContact)
+ await personCache.ensurePerson(mockContact)
+
+ expect(mockRestClient.ensurePerson).toHaveBeenCalledTimes(1)
+ })
+
+ it('should handle errors from ensurePerson and remove failed lookup from cache', async () => {
+ const error = new Error('Network error')
+ mockRestClient.ensurePerson.mockRejectedValueOnce(error)
+
+ await expect(personCache.ensurePerson(mockContact)).rejects.toThrow(error)
+
+ expect(mockCtx.error).toHaveBeenCalledWith('Error ensuring person exists', {
+ email: mockContact.email,
+ firstName: mockContact.firstName,
+ lastName: mockContact.lastName,
+ error: error.message
+ })
+
+ // The failed lookup should be removed from cache,
+ // so a subsequent call should retry the API call
+ mockRestClient.ensurePerson.mockResolvedValueOnce(mockPerson)
+ await personCache.ensurePerson(mockContact)
+ expect(mockRestClient.ensurePerson).toHaveBeenCalledTimes(2)
+ })
+
+ it('should cache different contacts separately', async () => {
+ const contact1: EmailContact = {
+ email: 'test1@example.com',
+ firstName: 'Test1',
+ lastName: 'User1'
+ }
+
+ const contact2: EmailContact = {
+ email: 'test2@example.com',
+ firstName: 'Test2',
+ lastName: 'User2'
+ }
+
+ const mockPerson1: CachedPerson = {
+ socialId: 'test-social-id-1' as PersonId,
+ uuid: 'test-uuid-1' as PersonUuid,
+ localPerson: 'test-local-person-1'
+ }
+
+ const mockPerson2: CachedPerson = {
+ socialId: 'test-social-id-2' as PersonId,
+ uuid: 'test-uuid-2' as PersonUuid,
+ localPerson: 'test-local-person-2'
+ }
+
+ mockRestClient.ensurePerson.mockResolvedValueOnce(mockPerson1).mockResolvedValueOnce(mockPerson2)
+
+ const result1 = await personCache.ensurePerson(contact1)
+ const result2 = await personCache.ensurePerson(contact2)
+
+ expect(result1).toEqual(mockPerson1)
+ expect(result2).toEqual(mockPerson2)
+ expect(mockRestClient.ensurePerson).toHaveBeenCalledTimes(2)
+ })
+ })
+
+ describe('clearCache', () => {
+ it('should clear the cache', async () => {
+ await personCache.ensurePerson(mockContact)
+ expect(personCache.size()).toBe(1)
+
+ personCache.clearCache()
+ expect(personCache.size()).toBe(0)
+
+ // Subsequent call should call the API again
+ await personCache.ensurePerson(mockContact)
+ expect(mockRestClient.ensurePerson).toHaveBeenCalledTimes(2)
+ })
+ })
+
+ describe('size', () => {
+ it('should return the number of cached entries', async () => {
+ expect(personCache.size()).toBe(0)
+
+ await personCache.ensurePerson(mockContact)
+ expect(personCache.size()).toBe(1)
+
+ const anotherContact: EmailContact = {
+ email: 'another@example.com',
+ firstName: 'Another',
+ lastName: 'User'
+ }
+ await personCache.ensurePerson(anotherContact)
+ expect(personCache.size()).toBe(2)
+ })
+ })
+})
diff --git a/services/mail/mail-common/src/__tests__/personFactory.test.ts b/services/mail/mail-common/src/__tests__/personFactory.test.ts
new file mode 100644
index 0000000000..69e0393991
--- /dev/null
+++ b/services/mail/mail-common/src/__tests__/personFactory.test.ts
@@ -0,0 +1,225 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { MeasureContext, PersonUuid, PersonId, WorkspaceUuid } from '@hcengineering/core'
+import { RestClient } from '@hcengineering/api-client'
+import { PersonCache, PersonCacheFactory, CachedPerson } from '../person'
+import { EmailContact } from '../types'
+
+describe('PersonCacheFactory', () => {
+ let mockCtx1: jest.Mocked
+ let mockCtx2: jest.Mocked
+ let mockRestClient1: jest.Mocked
+ let mockRestClient2: jest.Mocked
+
+ const workspace1 = 'workspace-1' as WorkspaceUuid
+ const workspace2 = 'workspace-2' as WorkspaceUuid
+
+ const mockContact: EmailContact = {
+ email: 'test@example.com',
+ firstName: 'Test',
+ lastName: 'User'
+ }
+
+ const mockPerson: CachedPerson = {
+ socialId: 'test-social-id' as PersonId,
+ uuid: 'test-uuid' as PersonUuid,
+ localPerson: 'test-local-id'
+ }
+
+ beforeEach(() => {
+ // Reset all instances before each test
+ PersonCacheFactory.resetAllInstances()
+
+ // Set up mocks
+ mockCtx1 = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockCtx2 = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockRestClient1 = {
+ ensurePerson: jest.fn().mockResolvedValue(mockPerson)
+ } as unknown as jest.Mocked
+
+ mockRestClient2 = {
+ ensurePerson: jest.fn().mockResolvedValue({
+ ...mockPerson,
+ uuid: 'different-uuid' as PersonUuid
+ })
+ } as unknown as jest.Mocked
+ })
+
+ describe('getInstance', () => {
+ it('should create a new instance for a workspace', () => {
+ expect(PersonCacheFactory.instanceCount).toBe(0)
+
+ const cache = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+
+ expect(cache).toBeInstanceOf(PersonCache)
+ expect(PersonCacheFactory.instanceCount).toBe(1)
+ })
+
+ it('should return the same instance for the same workspace', () => {
+ const cache1 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ const cache2 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+
+ expect(cache1).toBe(cache2)
+ expect(PersonCacheFactory.instanceCount).toBe(1)
+ })
+
+ it('should create different instances for different workspaces', () => {
+ const cache1 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ const cache2 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace2)
+
+ expect(cache1).not.toBe(cache2)
+ expect(PersonCacheFactory.instanceCount).toBe(2)
+ })
+
+ it('should use the first context and client for subsequent calls with the same workspace', async () => {
+ // First instance with first context and client
+ const cache1 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+
+ // Try to get instance with different context and client, but same workspace
+ const cache2 = PersonCacheFactory.getInstance(mockCtx2, mockRestClient2, workspace1)
+
+ // Should be the same instance
+ expect(cache1).toBe(cache2)
+
+ // Test that it uses the first client, not the second
+ await cache2.ensurePerson(mockContact)
+ expect(mockRestClient1.ensurePerson).toHaveBeenCalled()
+ expect(mockRestClient2.ensurePerson).not.toHaveBeenCalled()
+
+ // Now get a different workspace with second client
+ const cache3 = PersonCacheFactory.getInstance(mockCtx2, mockRestClient2, workspace2)
+ await cache3.ensurePerson(mockContact)
+
+ // Second client should be used for second workspace
+ expect(mockRestClient2.ensurePerson).toHaveBeenCalled()
+ })
+ })
+
+ describe('resetInstance', () => {
+ it('should remove the instance for a specific workspace', () => {
+ const cache1 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ const cache2 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace2)
+
+ expect(PersonCacheFactory.instanceCount).toBe(2)
+
+ PersonCacheFactory.resetInstance(workspace1)
+
+ expect(PersonCacheFactory.instanceCount).toBe(1)
+
+ // Getting workspace1 again should create a new instance
+ const cache1New = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ expect(cache1New).not.toBe(cache1)
+
+ // Workspace2 instance should remain the same
+ const cache2Again = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace2)
+ expect(cache2Again).toBe(cache2)
+ })
+
+ it('should do nothing when resetting a non-existent workspace', () => {
+ PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ expect(PersonCacheFactory.instanceCount).toBe(1)
+
+ // Reset a workspace that doesn't have a cache
+ PersonCacheFactory.resetInstance('non-existent-workspace' as WorkspaceUuid)
+
+ expect(PersonCacheFactory.instanceCount).toBe(1)
+ })
+ })
+
+ describe('resetAllInstances', () => {
+ it('should remove all workspace instances', () => {
+ PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace2)
+ expect(PersonCacheFactory.instanceCount).toBe(2)
+
+ PersonCacheFactory.resetAllInstances()
+
+ expect(PersonCacheFactory.instanceCount).toBe(0)
+ })
+ })
+
+ describe('instanceCount', () => {
+ it('should return the number of workspace instances', () => {
+ expect(PersonCacheFactory.instanceCount).toBe(0)
+
+ PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ expect(PersonCacheFactory.instanceCount).toBe(1)
+
+ PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace2)
+ expect(PersonCacheFactory.instanceCount).toBe(2)
+
+ PersonCacheFactory.resetInstance(workspace1)
+ expect(PersonCacheFactory.instanceCount).toBe(1)
+
+ PersonCacheFactory.resetAllInstances()
+ expect(PersonCacheFactory.instanceCount).toBe(0)
+ })
+ })
+
+ describe('integration with PersonCache', () => {
+ it('should maintain separate caches for each workspace', async () => {
+ const cache1 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ const cache2 = PersonCacheFactory.getInstance(mockCtx2, mockRestClient2, workspace2)
+
+ // Populate cache for workspace1
+ await cache1.ensurePerson(mockContact)
+ expect(mockRestClient1.ensurePerson).toHaveBeenCalledTimes(1)
+ mockRestClient1.ensurePerson.mockClear()
+
+ // Accessing the same contact in workspace1 should use cache
+ await cache1.ensurePerson(mockContact)
+ expect(mockRestClient1.ensurePerson).toHaveBeenCalledTimes(0)
+
+ // Accessing the same contact in workspace2 should make a new API call
+ // since it has a separate cache
+ await cache2.ensurePerson(mockContact)
+ expect(mockRestClient2.ensurePerson).toHaveBeenCalledTimes(1)
+ })
+
+ it('should allow clearing cache for a specific workspace', async () => {
+ const cache1 = PersonCacheFactory.getInstance(mockCtx1, mockRestClient1, workspace1)
+ const cache2 = PersonCacheFactory.getInstance(mockCtx2, mockRestClient2, workspace2)
+
+ // Populate caches
+ await cache1.ensurePerson(mockContact)
+ await cache2.ensurePerson(mockContact)
+
+ mockRestClient1.ensurePerson.mockClear()
+ mockRestClient2.ensurePerson.mockClear()
+
+ // Clear cache for workspace1
+ cache1.clearCache()
+
+ // Accessing contact in workspace1 should make a new API call
+ await cache1.ensurePerson(mockContact)
+ expect(mockRestClient1.ensurePerson).toHaveBeenCalledTimes(1)
+
+ // Accessing contact in workspace2 should still use cache
+ await cache2.ensurePerson(mockContact)
+ expect(mockRestClient2.ensurePerson).toHaveBeenCalledTimes(0)
+ })
+ })
+})
diff --git a/services/mail/mail-common/src/__tests__/personSpaces.test.ts b/services/mail/mail-common/src/__tests__/personSpaces.test.ts
new file mode 100644
index 0000000000..ff12f8b187
--- /dev/null
+++ b/services/mail/mail-common/src/__tests__/personSpaces.test.ts
@@ -0,0 +1,321 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import {
+ MeasureContext,
+ PersonUuid,
+ TxOperations,
+ WorkspaceUuid,
+ Ref,
+ Doc,
+ Space,
+ toFindResult
+} from '@hcengineering/core'
+import contact, { PersonSpace } from '@hcengineering/contact'
+import { PersonSpacesCache, PersonSpacesCacheFactory } from '../personSpaces'
+
+/* eslint-disable @typescript-eslint/unbound-method */
+describe('PersonSpacesCache', () => {
+ let mockCtx: jest.Mocked
+ let mockClient: jest.Mocked
+ let personSpacesCache: PersonSpacesCache
+
+ const workspace = 'test-workspace' as WorkspaceUuid
+ const mailId = 'test-mail-id'
+ const personUuid = 'test-person-uuid' as PersonUuid
+ const email = 'test@example.com'
+
+ const mockPerson: Doc = { _id: 'person1' as Ref } as any
+ const mockPersonSpaces: PersonSpace[] = [
+ { _id: 'space1', person: mockPerson._id } as unknown as PersonSpace,
+ { _id: 'space2', person: mockPerson._id } as unknown as PersonSpace
+ ]
+
+ beforeEach(() => {
+ jest.clearAllMocks()
+
+ mockCtx = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockClient = {
+ findAll: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockClient.findAll.mockImplementation((_class, _query, _options) => {
+ if (_class === contact.class.Person) {
+ return Promise.resolve(toFindResult([mockPerson]))
+ }
+ if (_class === contact.class.PersonSpace) {
+ return Promise.resolve(toFindResult(mockPersonSpaces))
+ }
+ return Promise.resolve(toFindResult([]))
+ })
+
+ personSpacesCache = new PersonSpacesCache(mockCtx, mockClient, workspace)
+ })
+
+ describe('getPersonSpaces', () => {
+ it('should fetch person spaces from the database on first call', async () => {
+ const spaces = await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+
+ expect(spaces).toEqual(mockPersonSpaces)
+ expect(mockClient.findAll).toHaveBeenCalledTimes(2)
+ expect(mockClient.findAll).toHaveBeenCalledWith(contact.class.Person, { personUuid }, { projection: { _id: 1 } })
+ expect(mockClient.findAll).toHaveBeenCalledWith(contact.class.PersonSpace, { person: { $in: [mockPerson._id] } })
+ })
+
+ it('should return cached spaces on subsequent calls', async () => {
+ await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+ mockClient.findAll.mockClear()
+
+ const spaces = await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+
+ expect(spaces).toEqual(mockPersonSpaces)
+ expect(mockClient.findAll).not.toHaveBeenCalled()
+ })
+
+ it('should warn when no spaces are found', async () => {
+ mockClient.findAll.mockImplementation((clazz) => {
+ if (clazz === contact.class.Person) {
+ return Promise.resolve(toFindResult([mockPerson]))
+ }
+ if (clazz === contact.class.PersonSpace) {
+ return Promise.resolve(toFindResult([]))
+ }
+ return Promise.resolve(toFindResult([]))
+ })
+
+ const spaces = await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+
+ expect(spaces.length).toEqual(0)
+ expect(mockCtx.warn).toHaveBeenCalledWith('No personal space found, skip', {
+ mailId,
+ personUuid,
+ email,
+ workspace
+ })
+ })
+
+ it('should handle errors and remove failed lookups from cache', async () => {
+ const error = new Error('Database error')
+ mockClient.findAll.mockRejectedValueOnce(error)
+
+ await expect(personSpacesCache.getPersonSpaces(mailId, personUuid, email)).rejects.toThrow(error)
+
+ expect(mockCtx.error).toHaveBeenCalledWith('Error fetching person spaces', {
+ mailId,
+ personUuid,
+ email,
+ workspace,
+ error: error.message
+ })
+
+ // Verify the cache entry was removed
+ mockClient.findAll.mockImplementation((clazz) => {
+ if (clazz === contact.class.Person) {
+ return Promise.resolve(toFindResult([mockPerson]))
+ }
+ if (clazz === contact.class.PersonSpace) {
+ return Promise.resolve(toFindResult(mockPersonSpaces))
+ }
+ return Promise.resolve(toFindResult([]))
+ })
+
+ // Next call should try again
+ await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+ expect(mockClient.findAll).toHaveBeenCalledTimes(3)
+ })
+ })
+
+ describe('clearPersonCache', () => {
+ it('should clear the cache for a specific person', async () => {
+ // Populate cache
+ await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+ mockClient.findAll.mockClear()
+
+ // Clear cache for this person
+ personSpacesCache.clearPersonCache(personUuid)
+
+ // Next call should fetch from database again
+ await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+ expect(mockClient.findAll).toHaveBeenCalledTimes(2)
+ })
+ })
+
+ describe('clearCache', () => {
+ it('should clear the entire cache', async () => {
+ // Populate cache
+ await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+ mockClient.findAll.mockClear()
+
+ // Clear all cache
+ personSpacesCache.clearCache()
+
+ // Next call should fetch from database again
+ await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+ expect(mockClient.findAll).toHaveBeenCalledTimes(2)
+ })
+ })
+
+ describe('size', () => {
+ it('should return the number of cached entries', async () => {
+ expect(personSpacesCache.size).toBe(0)
+
+ await personSpacesCache.getPersonSpaces(mailId, personUuid, email)
+ expect(personSpacesCache.size).toBe(1)
+
+ const anotherPersonUuid = 'another-person-uuid' as PersonUuid
+ await personSpacesCache.getPersonSpaces(mailId, anotherPersonUuid, 'another@example.com')
+ expect(personSpacesCache.size).toBe(2)
+
+ personSpacesCache.clearPersonCache(personUuid)
+ expect(personSpacesCache.size).toBe(1)
+
+ personSpacesCache.clearCache()
+ expect(personSpacesCache.size).toBe(0)
+ })
+ })
+})
+
+describe('PersonSpacesCacheFactory', () => {
+ let mockCtx1: jest.Mocked
+ let mockCtx2: jest.Mocked
+ let mockClient1: jest.Mocked
+ let mockClient2: jest.Mocked
+
+ const workspace1 = 'workspace-1' as WorkspaceUuid
+ const workspace2 = 'workspace-2' as WorkspaceUuid
+
+ beforeEach(() => {
+ PersonSpacesCacheFactory.resetAllInstances()
+
+ mockCtx1 = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockCtx2 = {
+ info: jest.fn(),
+ warn: jest.fn(),
+ error: jest.fn()
+ } as unknown as jest.Mocked
+
+ mockClient1 = {
+ findAll: jest.fn().mockResolvedValue([])
+ } as unknown as jest.Mocked
+
+ mockClient2 = {
+ findAll: jest.fn().mockResolvedValue([])
+ } as unknown as jest.Mocked
+ })
+
+ describe('getInstance', () => {
+ it('should create a new instance for a workspace', () => {
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(0)
+
+ const cache = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+
+ expect(cache).toBeInstanceOf(PersonSpacesCache)
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(1)
+ })
+
+ it('should return the same instance for the same workspace', () => {
+ const cache1 = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ const cache2 = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+
+ expect(cache1).toBe(cache2)
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(1)
+ })
+
+ it('should create different instances for different workspaces', () => {
+ const cache1 = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ const cache2 = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+
+ expect(cache1).not.toBe(cache2)
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(2)
+ })
+
+ it('should use the first context and client for subsequent calls with the same workspace', async () => {
+ const mailId = 'test-mail-id'
+ const personUuid = 'test-person-uuid' as PersonUuid
+ const email = 'test@example.com'
+
+ // First instance with first context and client
+ const cache1 = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+
+ // Try to get instance with different context and client, but same workspace
+ const cache2 = PersonSpacesCacheFactory.getInstance(mockCtx2, mockClient2, workspace1)
+
+ // Should be the same instance
+ expect(cache1).toBe(cache2)
+
+ // Test that it uses the first client, not the second
+ await cache2.getPersonSpaces(mailId, personUuid, email)
+ expect(mockClient1.findAll).toHaveBeenCalled()
+ expect(mockClient2.findAll).not.toHaveBeenCalled()
+ })
+ })
+
+ describe('resetInstance', () => {
+ it('should remove the instance for a specific workspace', () => {
+ const cache1 = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(2)
+
+ PersonSpacesCacheFactory.resetInstance(workspace1)
+
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(1)
+
+ // Getting workspace1 again should create a new instance
+ const cache1New = PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ expect(cache1New).not.toBe(cache1)
+ })
+ })
+
+ describe('resetAllInstances', () => {
+ it('should remove all workspace instances', () => {
+ PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(2)
+
+ PersonSpacesCacheFactory.resetAllInstances()
+
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(0)
+ })
+ })
+
+ describe('instanceCount', () => {
+ it('should return the number of workspace instances', () => {
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(0)
+
+ PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace1)
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(1)
+
+ PersonSpacesCacheFactory.getInstance(mockCtx1, mockClient1, workspace2)
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(2)
+
+ PersonSpacesCacheFactory.resetInstance(workspace1)
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(1)
+
+ PersonSpacesCacheFactory.resetAllInstances()
+ expect(PersonSpacesCacheFactory.instanceCount).toBe(0)
+ })
+ })
+})
diff --git a/services/mail/mail-common/src/channel.ts b/services/mail/mail-common/src/channel.ts
new file mode 100644
index 0000000000..4fa46ef7ad
--- /dev/null
+++ b/services/mail/mail-common/src/channel.ts
@@ -0,0 +1,198 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import {
+ MeasureContext,
+ PersonId,
+ Ref,
+ TxOperations,
+ Doc,
+ WorkspaceUuid,
+ generateId,
+ SocialId
+} from '@hcengineering/core'
+import chat from '@hcengineering/chat'
+import mail from '@hcengineering/mail'
+import { PersonSpace } from '@hcengineering/contact'
+import { SyncMutex } from './mutex'
+
+const createMutex = new SyncMutex()
+
+/**
+ * Caches channel references to reduce calls to create mail channels
+ */
+export class ChannelCache {
+ // Key is `${spaceId}:${emailAccount}`
+ private readonly cache = new Map>()
+
+ constructor (
+ private readonly ctx: MeasureContext,
+ private readonly client: TxOperations,
+ private readonly workspace: WorkspaceUuid
+ ) {}
+
+ /**
+ * Gets or creates a mail channel with caching
+ */
+ async getOrCreateChannel (
+ spaceId: Ref,
+ participants: PersonId[],
+ emailAccount: string,
+ socialId: SocialId
+ ): Promise[ | undefined> {
+ const cacheKey = `${spaceId}:${emailAccount}`
+
+ let channel = this.cache.get(cacheKey)
+ if (channel != null) {
+ return channel
+ }
+
+ channel = await this.fetchOrCreateChannel(spaceId, participants, emailAccount, socialId)
+ if (channel != null) {
+ this.cache.set(cacheKey, channel)
+ }
+
+ return channel
+ }
+
+ clearCache (spaceId: Ref, emailAccount: string): void {
+ this.cache.delete(`${spaceId}:${emailAccount}`)
+ }
+
+ clearAllCache (): void {
+ this.cache.clear()
+ }
+
+ get size (): number {
+ return this.cache.size
+ }
+
+ private async fetchOrCreateChannel (
+ space: Ref,
+ participants: PersonId[],
+ emailAccount: string,
+ socialId: SocialId
+ ): Promise][ | undefined> {
+ try {
+ // First try to find existing channel
+ const channel = await this.client.findOne(mail.tag.MailChannel, { title: emailAccount })
+
+ if (channel != null) {
+ this.ctx.info('Using existing channel', { me: emailAccount, space, channel: channel._id })
+ return channel._id
+ }
+
+ return await this.createNewChannel(space, participants, emailAccount, socialId)
+ } catch (err) {
+ this.ctx.error('Failed to create channel', {
+ me: emailAccount,
+ space,
+ workspace: this.workspace,
+ error: err instanceof Error ? err.message : String(err)
+ })
+
+ // Remove failed lookup from cache
+ this.cache.delete(`${space}:${emailAccount}`)
+
+ return undefined
+ }
+ }
+
+ private async createNewChannel (
+ space: Ref,
+ participants: PersonId[],
+ emailAccount: string,
+ socialId: SocialId
+ ): Promise][ | undefined> {
+ const mutexKey = `channel:${this.workspace}:${space}:${emailAccount}`
+ const releaseLock = await createMutex.lock(mutexKey)
+
+ try {
+ // Double-check that channel doesn't exist after acquiring lock
+ const existingChannel = await this.client.findOne(mail.tag.MailChannel, { title: emailAccount })
+ if (existingChannel != null) {
+ this.ctx.info('Using existing channel (found after mutex lock)', {
+ me: emailAccount,
+ space,
+ channel: existingChannel._id
+ })
+ return existingChannel._id
+ }
+
+ // Create new channel if it doesn't exist
+ this.ctx.info('Creating new channel', { me: emailAccount, space, personId: socialId._id })
+ const channelId = await this.client.createDoc(
+ chat.masterTag.Channel,
+ space,
+ {
+ title: emailAccount,
+ private: true,
+ members: participants,
+ archived: false,
+ createdBy: socialId._id,
+ modifiedBy: socialId._id
+ },
+ generateId(),
+ Date.now(),
+ socialId._id
+ )
+
+ this.ctx.info('Creating mixin', { me: emailAccount, space, personId: socialId._id, channelId })
+ await this.client.createMixin(
+ channelId,
+ chat.masterTag.Channel,
+ space,
+ mail.tag.MailChannel,
+ {},
+ Date.now(),
+ socialId._id
+ )
+
+ return channelId
+ } finally {
+ releaseLock()
+ }
+ }
+}
+
+/**
+ * Factory for creating ChannelCache instances per workspace
+ */
+export const ChannelCacheFactory = {
+ instances: new Map(),
+
+ getInstance (ctx: MeasureContext, client: TxOperations, workspace: WorkspaceUuid): ChannelCache {
+ let instance = ChannelCacheFactory.instances.get(workspace)
+
+ if (instance === undefined) {
+ instance = new ChannelCache(ctx, client, workspace)
+ ChannelCacheFactory.instances.set(workspace, instance)
+ }
+
+ return instance
+ },
+
+ resetInstance (workspace: WorkspaceUuid): void {
+ ChannelCacheFactory.instances.delete(workspace)
+ },
+
+ resetAllInstances (): void {
+ ChannelCacheFactory.instances.clear()
+ },
+
+ get instanceCount (): number {
+ return ChannelCacheFactory.instances.size
+ }
+}
diff --git a/services/mail/mail-common/src/index.ts b/services/mail/mail-common/src/index.ts
new file mode 100644
index 0000000000..a169b586a6
--- /dev/null
+++ b/services/mail/mail-common/src/index.ts
@@ -0,0 +1,19 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the 'License');
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an 'AS IS' BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+export * from './message'
+export * from './types'
+export * from './utils'
+export * from './mutex'
diff --git a/services/mail/mail-common/src/message.ts b/services/mail/mail-common/src/message.ts
new file mode 100644
index 0000000000..07e1898ee0
--- /dev/null
+++ b/services/mail/mail-common/src/message.ts
@@ -0,0 +1,290 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the 'License');
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an 'AS IS' BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+import { getClient as getAccountClient, isWorkspaceLoginInfo } from '@hcengineering/account-client'
+import { createRestTxOperations, createRestClient } from '@hcengineering/api-client'
+import { type Card } from '@hcengineering/card'
+import {
+ type RestClient as CommunicationClient,
+ createRestClient as getCommunicationClient
+} from '@hcengineering/communication-rest-client'
+import { MessageType } from '@hcengineering/communication-types'
+import chat from '@hcengineering/chat'
+import { PersonSpace } from '@hcengineering/contact'
+import {
+ type Blob,
+ type MeasureContext,
+ type PersonId,
+ type Ref,
+ type TxOperations,
+ generateId,
+ PersonUuid,
+ RateLimiter,
+ SocialId
+} from '@hcengineering/core'
+import mail from '@hcengineering/mail'
+
+import { buildStorageFromConfig, storageConfigFromEnv } from '@hcengineering/server-storage'
+
+import { BaseConfig, type Attachment } from './types'
+import { EmailMessage } from './types'
+import { getMdContent } from './utils'
+import { PersonCacheFactory } from './person'
+import { PersonSpacesCacheFactory } from './personSpaces'
+import { ChannelCache, ChannelCacheFactory } from './channel'
+
+export async function createMessages (
+ config: BaseConfig,
+ ctx: MeasureContext,
+ token: string,
+ message: EmailMessage,
+ attachments: Attachment[],
+ me: string,
+ socialId: SocialId
+): Promise {
+ const { mailId, from, subject, replyTo } = message
+ const tos = [...(message.to ?? []), ...(message.copy ?? [])]
+ ctx.info('Sending message', { mailId, from, to: tos.join(',') })
+
+ const accountClient = getAccountClient(config.AccountsURL, token)
+ const wsInfo = await accountClient.getLoginInfoByToken()
+
+ if (!isWorkspaceLoginInfo(wsInfo)) {
+ ctx.error('Unable to get workspace info', { mailId, from, tos })
+ return
+ }
+
+ const transactorUrl = wsInfo.endpoint.replace('ws://', 'http://').replace('wss://', 'https://')
+ const txClient = await createRestTxOperations(transactorUrl, wsInfo.workspace, wsInfo.token)
+ const msgClient = getCommunicationClient(wsInfo.endpoint, wsInfo.workspace, wsInfo.token)
+ const restClient = createRestClient(transactorUrl, wsInfo.workspace, wsInfo.token)
+ const personCache = PersonCacheFactory.getInstance(ctx, restClient, wsInfo.workspace)
+ const personSpacesCache = PersonSpacesCacheFactory.getInstance(ctx, txClient, wsInfo.workspace)
+ const channelCache = ChannelCacheFactory.getInstance(ctx, txClient, wsInfo.workspace)
+
+ const fromPerson = await personCache.ensurePerson(from)
+
+ const toPersons: { address: string, uuid: PersonUuid, socialId: PersonId }[] = []
+ for (const to of tos) {
+ const toPerson = await personCache.ensurePerson(to)
+ if (toPerson === undefined) {
+ continue
+ }
+ toPersons.push({ address: to.email, ...toPerson })
+ }
+ if (toPersons.length === 0) {
+ ctx.error('Unable to create message without a proper TO', { mailId, from })
+ return
+ }
+
+ const modifiedBy = fromPerson.socialId
+ const participants = [fromPerson.socialId, ...toPersons.map((p) => p.socialId)]
+ const content = getMdContent(ctx, message)
+
+ const attachedBlobs: Attachment[] = []
+ if (config.StorageConfig !== undefined) {
+ const storageConfig = storageConfigFromEnv(config.StorageConfig)
+ const storageAdapter = buildStorageFromConfig(storageConfig)
+ try {
+ for (const a of attachments ?? []) {
+ try {
+ await storageAdapter.put(
+ ctx,
+ {
+ uuid: wsInfo.workspace,
+ url: wsInfo.workspaceUrl,
+ dataId: wsInfo.workspaceDataId
+ },
+ a.id,
+ a.data,
+ a.contentType
+ )
+ attachedBlobs.push(a)
+ ctx.info('Uploaded attachment', { mailId, blobId: a.id, name: a.name, contentType: a.contentType })
+ } catch (error) {
+ ctx.error('Failed to upload attachment', { name: a.name, error, mailId })
+ }
+ }
+ } finally {
+ await storageAdapter.close()
+ }
+ }
+
+ try {
+ const spaces = await personSpacesCache.getPersonSpaces(mailId, fromPerson.uuid, from.email)
+ if (spaces.length > 0) {
+ await saveMessageToSpaces(
+ ctx,
+ txClient,
+ msgClient,
+ mailId,
+ spaces,
+ participants,
+ modifiedBy,
+ subject,
+ content,
+ attachedBlobs,
+ me,
+ socialId,
+ message.sendOn,
+ channelCache,
+ replyTo
+ )
+ }
+ } catch (error) {
+ ctx.error('Failed to save message to personal spaces', {
+ error,
+ mailId,
+ personUuid: fromPerson.uuid,
+ email: from
+ })
+ }
+
+ for (const to of toPersons) {
+ try {
+ const spaces = await personSpacesCache.getPersonSpaces(mailId, to.uuid, to.address)
+ if (spaces.length > 0) {
+ await saveMessageToSpaces(
+ ctx,
+ txClient,
+ msgClient,
+ mailId,
+ spaces,
+ participants,
+ modifiedBy,
+ subject,
+ content,
+ attachedBlobs,
+ me,
+ socialId,
+ message.sendOn,
+ channelCache,
+ replyTo
+ )
+ }
+ } catch (error) {
+ ctx.error('Failed to save message spaces', { error, mailId, personUuid: to.uuid, email: to.address })
+ }
+ }
+}
+
+async function saveMessageToSpaces (
+ ctx: MeasureContext,
+ client: TxOperations,
+ msgClient: CommunicationClient,
+ mailId: string,
+ spaces: PersonSpace[],
+ participants: PersonId[],
+ modifiedBy: PersonId,
+ subject: string,
+ content: string,
+ attachments: Attachment[],
+ me: string,
+ socialId: SocialId,
+ createdDate: number,
+ channelCache: ChannelCache,
+ inReplyTo?: string
+): Promise {
+ const rateLimiter = new RateLimiter(10)
+ for (const space of spaces) {
+ const spaceId = space._id
+ await rateLimiter.add(async () => {
+ ctx.info('Saving message to space', { mailId, space: spaceId })
+
+ const route = await client.findOne(mail.class.MailRoute, { mailId, space: spaceId })
+ if (route !== undefined) {
+ ctx.info('Message is already in the thread, skip', { mailId, threadId: route.threadId, spaceId })
+ return
+ }
+
+ let threadId: Ref | undefined
+ if (inReplyTo !== undefined) {
+ const route = await client.findOne(mail.class.MailRoute, { mailId: inReplyTo, space: spaceId })
+ if (route !== undefined) {
+ threadId = route.threadId as Ref
+ ctx.info('Found existing thread', { mailId, threadId, spaceId })
+ }
+ }
+ if (threadId === undefined) {
+ const channel = await channelCache.getOrCreateChannel(spaceId, participants, me, socialId)
+ const newThreadId = await client.createDoc(
+ chat.masterTag.Thread,
+ space._id,
+ {
+ title: subject,
+ description: content,
+ private: true,
+ members: participants,
+ archived: false,
+ createdBy: modifiedBy,
+ modifiedBy,
+ parent: channel
+ },
+ generateId(),
+ undefined,
+ modifiedBy
+ )
+ await client.createMixin(
+ newThreadId,
+ chat.masterTag.Thread,
+ space._id,
+ mail.tag.MailThread,
+ {},
+ Date.now(),
+ socialId._id
+ )
+ threadId = newThreadId as Ref
+ ctx.info('Created new thread', { mailId, threadId, spaceId })
+ }
+
+ const { id: messageId, created: messageCreated } = await msgClient.createMessage(
+ threadId,
+ chat.masterTag.Thread,
+ content,
+ modifiedBy,
+ MessageType.Message,
+ {
+ created: createdDate
+ }
+ )
+ ctx.info('Created message', { mailId, messageId, threadId, content })
+
+ for (const a of attachments) {
+ await msgClient.createFile(
+ threadId,
+ messageId,
+ messageCreated,
+ a.id as Ref,
+ a.contentType,
+ a.name,
+ a.data.length,
+ modifiedBy
+ )
+ }
+
+ await client.createDoc(
+ mail.class.MailRoute,
+ space._id,
+ {
+ mailId,
+ threadId
+ },
+ generateId(),
+ undefined,
+ modifiedBy
+ )
+ })
+ }
+ await rateLimiter.waitProcessing()
+}
diff --git a/services/mail/mail-common/src/mutex.ts b/services/mail/mail-common/src/mutex.ts
new file mode 100644
index 0000000000..29fbf4ab85
--- /dev/null
+++ b/services/mail/mail-common/src/mutex.ts
@@ -0,0 +1,42 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+export class SyncMutex {
+ private readonly locks = new Map>()
+
+ async lock (key: string): Promise<() => void> {
+ // Wait for any existing lock to be released
+ const currentLock = this.locks.get(key)
+ if (currentLock != null) {
+ await currentLock
+ }
+
+ // Create a new lock
+ let releaseFn!: () => void
+ const newLock = new Promise((resolve) => {
+ releaseFn = resolve
+ })
+
+ // Store the lock
+ this.locks.set(key, newLock)
+
+ // Return the release function
+ return () => {
+ if (this.locks.get(key) === newLock) {
+ this.locks.delete(key)
+ }
+ releaseFn()
+ }
+ }
+}
diff --git a/services/mail/mail-common/src/person.ts b/services/mail/mail-common/src/person.ts
new file mode 100644
index 0000000000..8d97b22851
--- /dev/null
+++ b/services/mail/mail-common/src/person.ts
@@ -0,0 +1,122 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { type MeasureContext, PersonId, PersonUuid, SocialIdType, WorkspaceUuid } from '@hcengineering/core'
+import { type RestClient } from '@hcengineering/api-client'
+import { EmailContact } from './types'
+
+export interface CachedPerson {
+ socialId: PersonId
+ uuid: PersonUuid
+ localPerson: string
+}
+
+/**
+ * Caches persons to reduce ensure person API calls
+ */
+export class PersonCache {
+ private readonly cache = new Map>()
+
+ constructor (
+ private readonly ctx: MeasureContext,
+ private readonly restClient: RestClient
+ ) {}
+
+ /**
+ * Gets or creates a person by email address with caching
+ */
+ async ensurePerson (contact: EmailContact): Promise {
+ const email = contact.email.toLowerCase().trim()
+
+ let personPromise = this.cache.get(email)
+
+ if (personPromise === undefined) {
+ personPromise = this.fetchAndCachePerson(email, contact.firstName, contact.lastName)
+ this.cache.set(email, personPromise)
+ }
+
+ const result = await personPromise
+ if (result === undefined) {
+ throw new Error(`Failed to ensure person exists for email: ${email}`)
+ }
+ return result
+ }
+
+ size (): number {
+ return this.cache.size
+ }
+
+ clearCache (): void {
+ this.cache.clear()
+ }
+
+ private async fetchAndCachePerson (
+ email: string,
+ firstName: string,
+ lastName: string
+ ): Promise {
+ try {
+ const result = await this.restClient.ensurePerson(SocialIdType.EMAIL, email, firstName, lastName)
+
+ if (result === undefined) {
+ this.ctx.warn('Failed to ensure person exists', { email })
+ return undefined
+ }
+
+ return result
+ } catch (err) {
+ this.ctx.error('Error ensuring person exists', {
+ email,
+ firstName,
+ lastName,
+ error: err instanceof Error ? err.message : String(err)
+ })
+
+ this.cache.delete(email)
+
+ throw err
+ }
+ }
+}
+
+/**
+ * Factory for creating and managing PersonCache instances per workspace
+ */
+export const PersonCacheFactory = {
+ instances: new Map(),
+
+ getInstance (ctx: MeasureContext, restClient: RestClient, workspace: WorkspaceUuid): PersonCache {
+ let instance = PersonCacheFactory.instances.get(workspace)
+
+ if (instance === undefined) {
+ instance = new PersonCache(ctx, restClient)
+ PersonCacheFactory.instances.set(workspace, instance)
+ }
+
+ return instance
+ },
+
+ resetInstance (workspace: WorkspaceUuid): void {
+ PersonCacheFactory.instances.delete(workspace)
+ },
+
+ resetAllInstances (): void {
+ PersonCacheFactory.instances.clear()
+ },
+
+ get instanceCount (): number {
+ return PersonCacheFactory.instances.size
+ }
+}
diff --git a/services/mail/mail-common/src/personSpaces.ts b/services/mail/mail-common/src/personSpaces.ts
new file mode 100644
index 0000000000..a54c427809
--- /dev/null
+++ b/services/mail/mail-common/src/personSpaces.ts
@@ -0,0 +1,112 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the "License");
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+
+import { MeasureContext, PersonUuid, TxOperations, WorkspaceUuid } from '@hcengineering/core'
+import contact, { PersonSpace } from '@hcengineering/contact'
+
+/**
+ * Cache for person spaces
+ */
+export class PersonSpacesCache {
+ private readonly cache = new Map>()
+
+ constructor (
+ private readonly ctx: MeasureContext,
+ private readonly client: TxOperations,
+ private readonly workspace: WorkspaceUuid
+ ) {}
+
+ async getPersonSpaces (mailId: string, personUuid: PersonUuid, email: string): Promise {
+ let spacesPromise = this.cache.get(personUuid)
+
+ if (spacesPromise === undefined) {
+ spacesPromise = this.fetchPersonSpaces(mailId, personUuid, email)
+ this.cache.set(personUuid, spacesPromise)
+ }
+
+ return await spacesPromise
+ }
+
+ clearPersonCache (personUuid: PersonUuid): void {
+ this.cache.delete(personUuid)
+ }
+
+ clearCache (): void {
+ this.cache.clear()
+ }
+
+ get size (): number {
+ return this.cache.size
+ }
+
+ private async fetchPersonSpaces (mailId: string, personUuid: PersonUuid, email: string): Promise {
+ try {
+ const persons = await this.client.findAll(contact.class.Person, { personUuid }, { projection: { _id: 1 } })
+
+ const personRefs = persons.map((p) => p._id)
+ const spaces = await this.client.findAll(contact.class.PersonSpace, { person: { $in: personRefs } })
+
+ if (spaces.length === 0) {
+ this.ctx.warn('No personal space found, skip', { mailId, personUuid, email, workspace: this.workspace })
+ }
+
+ return spaces
+ } catch (err) {
+ this.ctx.error('Error fetching person spaces', {
+ mailId,
+ personUuid,
+ email,
+ workspace: this.workspace,
+ error: err instanceof Error ? err.message : String(err)
+ })
+
+ // Remove failed lookup from cache to allow retry
+ this.cache.delete(personUuid)
+
+ // Re-throw to allow handling upstream
+ throw err
+ }
+ }
+}
+
+/**
+ * Factory for creating and managing PersonSpacesCache instances per workspace
+ */
+export const PersonSpacesCacheFactory = {
+ instances: new Map(),
+
+ getInstance (ctx: MeasureContext, client: TxOperations, workspace: WorkspaceUuid): PersonSpacesCache {
+ let instance = PersonSpacesCacheFactory.instances.get(workspace)
+
+ if (instance === undefined) {
+ instance = new PersonSpacesCache(ctx, client, workspace)
+ PersonSpacesCacheFactory.instances.set(workspace, instance)
+ }
+
+ return instance
+ },
+
+ resetInstance (workspace: WorkspaceUuid): void {
+ PersonSpacesCacheFactory.instances.delete(workspace)
+ },
+
+ resetAllInstances (): void {
+ PersonSpacesCacheFactory.instances.clear()
+ },
+
+ get instanceCount (): number {
+ return PersonSpacesCacheFactory.instances.size
+ }
+}
diff --git a/services/mail/mail-common/src/types.ts b/services/mail/mail-common/src/types.ts
new file mode 100644
index 0000000000..3c1cfe4c9e
--- /dev/null
+++ b/services/mail/mail-common/src/types.ts
@@ -0,0 +1,48 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the 'License');
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an 'AS IS' BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+export interface Attachment {
+ id: string
+ name: string
+ data: Buffer
+ contentType: string
+ size?: number
+ lastModified: number
+}
+
+export interface EmailContact {
+ email: string
+ firstName: string
+ lastName: string
+}
+
+export interface EmailMessage {
+ modifiedOn: number
+ mailId: string
+ replyTo?: string
+ copy?: EmailContact[]
+ content: string
+ textContent: string
+ from: EmailContact
+ to: EmailContact[]
+ incoming: boolean
+ subject: string
+ sendOn: number
+}
+
+export interface BaseConfig {
+ AccountsURL: string
+ KvsUrl: string
+ StorageConfig: string
+}
diff --git a/services/mail/mail-common/src/utils.ts b/services/mail/mail-common/src/utils.ts
new file mode 100644
index 0000000000..6639c8973c
--- /dev/null
+++ b/services/mail/mail-common/src/utils.ts
@@ -0,0 +1,112 @@
+//
+// Copyright © 2025 Hardcore Engineering Inc.
+//
+// Licensed under the Eclipse Public License, Version 2.0 (the 'License');
+// you may not use this file except in compliance with the License. You may
+// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an 'AS IS' BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+//
+// See the License for the specific language governing permissions and
+// limitations under the License.
+//
+import TurndownService from 'turndown'
+import sanitizeHtml from 'sanitize-html'
+
+import { MeasureContext } from '@hcengineering/core'
+import { EmailContact, EmailMessage } from './types'
+
+export function getMdContent (ctx: MeasureContext, email: EmailMessage): string {
+ if (email.content !== undefined) {
+ try {
+ const html = sanitizeHtml(email.content)
+ const tds = new TurndownService()
+ return tds.turndown(html)
+ } catch (error) {
+ ctx.warn('Failed to parse html content', { error })
+ }
+ }
+ return email.textContent
+}
+
+/**
+ * Parse email header into EmailContact objects
+ * Supports both single and multiple addresses in formats like:
+ * - "Name"
+ * - Name
+ * - email@example.com
+ * - Multiple comma-separated addresses in any of the above formats
+ *
+ * @param headerValue Email header value to parse
+ * @returns Array of EmailContact objects
+ */
+export function parseEmailHeader (headerValue: string | undefined): EmailContact[] {
+ if (headerValue == null || headerValue.trim() === '') {
+ return []
+ }
+
+ // Split the header by commas, but ignore commas inside quotes
+ const regex = /,(?=(?:[^"]*"[^"]*")*[^"]*$)/
+ const addresses = headerValue
+ .split(regex)
+ .map((addr) => addr.trim())
+ .filter((addr) => addr !== '')
+
+ return addresses.map((address) => parseNameFromEmailHeader(address))
+}
+
+export function parseNameFromEmailHeader (headerValue: string | undefined): EmailContact {
+ if (headerValue == null || headerValue.trim() === '') {
+ return {
+ email: '',
+ firstName: '',
+ lastName: ''
+ }
+ }
+
+ // Match pattern like: "Name" or Name
+ const nameEmailPattern = /^(?:"?([^"<]+)"?\s*)?<([^>]+)>$/
+ const match = headerValue.trim().match(nameEmailPattern)
+
+ if (match == null) {
+ const address = headerValue.trim()
+ const parts = address.split('@')
+ return {
+ email: address,
+ firstName: parts[0],
+ lastName: parts[1]
+ }
+ }
+
+ const displayName = match[1]?.trim()
+ const email = match[2].trim()
+
+ if (displayName == null || displayName === '') {
+ const parts = email.split('@')
+ return {
+ email,
+ firstName: parts[0],
+ lastName: parts[1]
+ }
+ }
+
+ const nameParts = displayName.split(/\s+/)
+ let firstName: string | undefined
+ let lastName: string | undefined
+
+ if (nameParts.length === 1) {
+ firstName = nameParts[0]
+ } else if (nameParts.length > 1) {
+ firstName = nameParts[0]
+ lastName = nameParts.slice(1).join(' ')
+ }
+
+ const parts = email.split('@')
+ return {
+ email,
+ firstName: firstName ?? parts[0],
+ lastName: lastName ?? parts[1]
+ }
+}
diff --git a/services/mail/mail-common/tsconfig.json b/services/mail/mail-common/tsconfig.json
new file mode 100644
index 0000000000..59e4fd4297
--- /dev/null
+++ b/services/mail/mail-common/tsconfig.json
@@ -0,0 +1,10 @@
+{
+ "extends": "./node_modules/@hcengineering/platform-rig/profiles/default/tsconfig.json",
+
+ "compilerOptions": {
+ "rootDir": "./src",
+ "outDir": "./lib",
+ "declarationDir": "./types",
+ "tsBuildInfoFile": ".build/build.tsbuildinfo"
+ }
+}
\ No newline at end of file
]