diff --git a/common/config/rush/pnpm-lock.yaml b/common/config/rush/pnpm-lock.yaml index 5195fe978a..eda1faf365 100644 --- a/common/config/rush/pnpm-lock.yaml +++ b/common/config/rush/pnpm-lock.yaml @@ -4241,7 +4241,7 @@ packages: version: 0.0.0 '@rush-temp/model-card@file:projects/model-card.tgz': - resolution: {integrity: sha512-7yhnpluYR/Gypgi3kKPtEVe1NpcG8YfE7S3ACF56GIqXideFnlXYWJ8PjUDQ9Ub5BHUulGYTXKsFUqcHRxQDZA==, tarball: file:projects/model-card.tgz} + resolution: {integrity: sha512-diGAGkQiQNXD27/ln6i6al/sr7X2dOUdiTx6114JOahvZngprIhzdJLdXuW44BnQHTugr4j0Sp1zOwQXKcTBVA==, tarball: file:projects/model-card.tgz} version: 0.0.0 '@rush-temp/model-chunter@file:projects/model-chunter.tgz': @@ -4353,7 +4353,7 @@ packages: version: 0.0.0 '@rush-temp/model-server-card@file:projects/model-server-card.tgz': - resolution: {integrity: sha512-SP381dOkAOymE4k9VaEP2o28yv7oJxJEol5LeYjcqHXkyN/al/v8UetEGmS4/gvuxk0Np0KZPpLnczV9O3ApUg==, tarball: file:projects/model-server-card.tgz} + resolution: {integrity: sha512-yMw0GV5uufsqDDev756PLadZksaAwUcL1eq52hijUSI4yYC1cr5wvJLW97/g/znWX6kZuHC7lAGvrBcgHf2axw==, tarball: file:projects/model-server-card.tgz} version: 0.0.0 '@rush-temp/model-server-chunter@file:projects/model-server-chunter.tgz': @@ -4589,7 +4589,7 @@ packages: version: 0.0.0 '@rush-temp/pod-calendar@file:projects/pod-calendar.tgz': - resolution: {integrity: sha512-gUk3jshnHJ0CUucxZ1Stk8oQ2DTdSVfJzFHIl/EBUISbRv/YySoXesQCnqfn1kVpqOBZdIvI/N4NRAL/yREbWw==, tarball: file:projects/pod-calendar.tgz} + resolution: {integrity: sha512-BvN38ScXSeB5Mui24Sawekl81QZMFZ47EK6KdIqRJJQwM+Qj0owWul3y/r7ZtNAfKnD9bV3AeEWaqkZ8EXYjwg==, tarball: file:projects/pod-calendar.tgz} version: 0.0.0 '@rush-temp/pod-collaborator@file:projects/pod-collaborator.tgz': @@ -4621,7 +4621,7 @@ packages: version: 0.0.0 '@rush-temp/pod-server@file:projects/pod-server.tgz': - resolution: {integrity: sha512-U7wNYw+h8a6hEV6fwXYJFStTk5noW4LAsL1G4JCJjyQ7PVzlmXZqGncFe+bq5EqgYgbOewZ2tR+T8rB8HWi9Cw==, tarball: file:projects/pod-server.tgz} + resolution: {integrity: sha512-MJ5TQobQy4kv67qy7qpQgz056fhaZw8Zs/NaXniFafG/xjstgWeIJCXNHDkKvdBm+wT6EdHgNLIN3BaRpMSDsA==, tarball: file:projects/pod-server.tgz} version: 0.0.0 '@rush-temp/pod-ses@file:projects/pod-ses.tgz': @@ -4809,7 +4809,7 @@ packages: version: 0.0.0 '@rush-temp/server-calendar-resources@file:projects/server-calendar-resources.tgz': - resolution: {integrity: sha512-Rpdy7Nx56rXrJ1apUsvlancAVa5USFSwXboZ8iexTJ4Jaebv1yf32j5LUo/fLKwpZyL9KaUHT+A4N+nJKoRVWg==, tarball: file:projects/server-calendar-resources.tgz} + resolution: {integrity: sha512-pFGwg9Q/vogm2In88l39LD/B3Da/8sVTD3BVw3Or0VJ9peXU/YJHrqkgd4+ZU3RDRf4amDVz/rIpqMqjD0rS6Q==, tarball: file:projects/server-calendar-resources.tgz} version: 0.0.0 '@rush-temp/server-calendar@file:projects/server-calendar.tgz': @@ -4817,7 +4817,7 @@ packages: version: 0.0.0 '@rush-temp/server-card-resources@file:projects/server-card-resources.tgz': - resolution: {integrity: sha512-lp3RSXAHV78ux6p3I0iSbVkfY3aB0fXDfxMhxQZfjU9sSXLOsKolgDxMaUNXVm+/WIpS6dwLk4jWHgGopH8Nmw==, tarball: file:projects/server-card-resources.tgz} + resolution: {integrity: sha512-sJdtOVlhTFawBXtSTzcx+Vs8NALsiqBfwAz9cqUBj45vobiSF97D987LPBgjhCARk2+qx8lZyUPPTUJkJdaLOQ==, tarball: file:projects/server-card-resources.tgz} version: 0.0.0 '@rush-temp/server-card@file:projects/server-card.tgz': @@ -5193,7 +5193,7 @@ packages: version: 0.0.0 '@rush-temp/text-editor-resources@file:projects/text-editor-resources.tgz': - resolution: {integrity: sha512-JAu7yul5yHwjxWkBN9p1LW7HR8vAaRlf8Ioiyc7FYAKaQ1EBV05sCCpYsVayVpmU0vU92XzQffryB8TH/dZ6+A==, tarball: file:projects/text-editor-resources.tgz} + resolution: {integrity: sha512-PkvG582CT0XfnKC+4uNE1ZnwUgT9ywmfDmVr5vgcXti04Ntx1rB9VzTqnNDFI/A7h5dFAeEbR8s08gzgO6fdrw==, tarball: file:projects/text-editor-resources.tgz} version: 0.0.0 '@rush-temp/text-editor@file:projects/text-editor.tgz': @@ -20586,6 +20586,7 @@ snapshots: '@typescript-eslint/eslint-plugin': 6.21.0(@typescript-eslint/parser@6.21.0(eslint@8.56.0)(typescript@5.3.3))(eslint@8.56.0)(typescript@5.3.3) '@typescript-eslint/parser': 6.21.0(eslint@8.56.0)(typescript@5.3.3) cors: 2.8.5 + cross-env: 7.0.3 dotenv: 16.0.3 esbuild: 0.24.2 eslint: 8.56.0 diff --git a/packages/core/src/classes.ts b/packages/core/src/classes.ts index 766994966d..70693b1f3c 100644 --- a/packages/core/src/classes.ts +++ b/packages/core/src/classes.ts @@ -791,3 +791,13 @@ export interface BaseWorkspaceInfo { backupInfo?: BackupStatus } + +/** + * @public + */ +export type ClientWorkspaceInfo = Omit & { + lastProcessingTime?: number + attempts?: number + message?: string + workspaceId: string +} diff --git a/pods/server/package.json b/pods/server/package.json index ca61179794..7fa99af653 100644 --- a/pods/server/package.json +++ b/pods/server/package.json @@ -70,6 +70,7 @@ "@hcengineering/analytics-service": "^0.6.0", "@hcengineering/contact": "^0.6.24", "@hcengineering/notification": "^0.6.23", + "@hcengineering/server-calendar": "^0.6.0", "@hcengineering/server-notification": "^0.6.1", "@hcengineering/server-telegram": "^0.6.0", "@hcengineering/pod-telegram-bot": "^0.6.0", diff --git a/pods/server/src/__start.ts b/pods/server/src/__start.ts index 2d049cf40a..cc026823a4 100644 --- a/pods/server/src/__start.ts +++ b/pods/server/src/__start.ts @@ -10,6 +10,7 @@ import { MeasureMetricsContext, newMetrics, setOperationLogProfiling } from '@hc import { setMetadata } from '@hcengineering/platform' import { serverConfigFromEnv } from '@hcengineering/server' import serverAiBot from '@hcengineering/server-ai-bot' +import serverCalendar from '@hcengineering/server-calendar' import serverCore, { type ConnectionSocket, type Session, @@ -75,6 +76,7 @@ setMetadata(serverNotification.metadata.SesAuthToken, config.sesAuthToken) setMetadata(serverTelegram.metadata.BotUrl, process.env.TELEGRAM_BOT_URL) setMetadata(serverAiBot.metadata.SupportWorkspaceId, process.env.SUPPORT_WORKSPACE) setMetadata(serverAiBot.metadata.EndpointURL, process.env.AI_BOT_URL) +setMetadata(serverCalendar.metadata.EndpointURL, process.env.CALENDAR_URL) const { shutdown, sessionManager } = start(metricsContext, config.dbUrl, { fulltextUrl: config.fulltextUrl, diff --git a/server-plugins/calendar-resources/package.json b/server-plugins/calendar-resources/package.json index dc1718ba3c..c62140575a 100644 --- a/server-plugins/calendar-resources/package.json +++ b/server-plugins/calendar-resources/package.json @@ -39,6 +39,8 @@ "dependencies": { "@hcengineering/core": "^0.6.32", "@hcengineering/platform": "^0.6.11", + "@hcengineering/server-calendar": "^0.6.0", + "@hcengineering/server-token": "^0.6.11", "@hcengineering/calendar": "^0.6.24", "@hcengineering/contact": "^0.6.24", "@hcengineering/server-core": "^0.6.1", diff --git a/server-plugins/calendar-resources/src/index.ts b/server-plugins/calendar-resources/src/index.ts index edfc41c11c..0afb07f16f 100644 --- a/server-plugins/calendar-resources/src/index.ts +++ b/server-plugins/calendar-resources/src/index.ts @@ -17,6 +17,7 @@ import calendar, { Calendar, Event, ExternalCalendar } from '@hcengineering/cale import contactPlugin, { Contact, Person, PersonAccount } from '@hcengineering/contact' import core, { Class, + concatLink, Data, Doc, DocumentQuery, @@ -24,6 +25,7 @@ import core, { FindResult, Hierarchy, Ref, + systemAccountEmail, Tx, TxCreateDoc, TxCUD, @@ -31,9 +33,11 @@ import core, { TxRemoveDoc, TxUpdateDoc } from '@hcengineering/core' -import { getResource } from '@hcengineering/platform' +import serverCalendar from '@hcengineering/server-calendar' +import { getMetadata, getResource } from '@hcengineering/platform' import { TriggerControl } from '@hcengineering/server-core' import { getHTMLPresenter, getTextPresenter } from '@hcengineering/server-notification-resources' +import { generateToken } from '@hcengineering/server-token' /** * @public @@ -145,6 +149,9 @@ async function onEventUpdate (ctx: TxUpdateDoc, control: TriggerControl): if (Object.keys(otherOps).length === 0) return [] const event = (await control.findAll(control.ctx, calendar.class.Event, { _id: ctx.objectId }, { limit: 1 }))[0] if (event === undefined) return [] + if (ctx.modifiedBy !== core.account.System) { + void sendEventToService(event, 'update', control) + } if (event.access !== 'owner') return [] const events = await control.findAll(control.ctx, calendar.class.Event, { eventId: event.eventId }) const res: Tx[] = [] @@ -222,8 +229,43 @@ async function eventForNewParticipants ( return res } +async function sendEventToService ( + event: Event, + type: 'create' | 'update' | 'delete', + control: TriggerControl +): Promise { + const url = getMetadata(serverCalendar.metadata.EndpointURL) ?? '' + + if (url === '') { + return + } + + const workspace = control.workspace.name + + try { + await fetch(concatLink(url, '/event'), { + method: 'POST', + keepalive: true, + headers: { + Authorization: 'Bearer ' + generateToken(systemAccountEmail, control.workspace), + 'Content-Type': 'application/json' + }, + body: JSON.stringify({ + event, + workspace, + type + }) + }) + } catch (err) { + control.ctx.error('Could not send calendar event to service', { err }) + } +} + async function onEventCreate (ctx: TxCreateDoc, control: TriggerControl): Promise { const event = TxProcessor.createDoc2Doc(ctx) + if (ctx.modifiedBy !== core.account.System) { + void sendEventToService(event, 'create', control) + } if (event.access !== 'owner') return [] const res: Tx[] = [] const { _class, space, attachedTo, attachedToClass, collection, ...attr } = event @@ -265,6 +307,9 @@ async function onRemoveEvent (ctx: TxRemoveDoc, control: TriggerControl): const removed = control.removedMap.get(ctx.objectId) as Event const res: Tx[] = [] if (removed !== undefined) { + if (ctx.modifiedBy !== core.account.System) { + void sendEventToService(removed, 'delete', control) + } if (removed.access !== 'owner') return [] const current = await control.findAll(control.ctx, calendar.class.Event, { eventId: removed.eventId }) for (const cur of current) { diff --git a/server-plugins/calendar/src/index.ts b/server-plugins/calendar/src/index.ts index 804ec16bea..a7298e1b49 100644 --- a/server-plugins/calendar/src/index.ts +++ b/server-plugins/calendar/src/index.ts @@ -14,7 +14,7 @@ // limitations under the License. // -import type { Plugin, Resource } from '@hcengineering/platform' +import type { Metadata, Plugin, Resource } from '@hcengineering/platform' import { plugin } from '@hcengineering/platform' import type { ObjectDDParticipantFunc, TriggerFunc } from '@hcengineering/server-core' import { Presenter } from '@hcengineering/server-notification' @@ -28,6 +28,9 @@ export const serverCalendarId = 'server-calendar' as Plugin * @public */ export default plugin(serverCalendarId, { + metadata: { + EndpointURL: '' as Metadata + }, function: { ReminderHTMLPresenter: '' as Resource, ReminderTextPresenter: '' as Resource, diff --git a/server/account/src/operations.ts b/server/account/src/operations.ts index c265fac33f..772c009dd3 100644 --- a/server/account/src/operations.ts +++ b/server/account/src/operations.ts @@ -26,6 +26,7 @@ import contact, { import core, { AccountRole, Client, + ClientWorkspaceInfo, concatLink, Data, generateId, @@ -68,7 +69,6 @@ import type { Account, AccountDB, AccountInfo, - ClientWorkspaceInfo, Invite, LoginInfo, ObjectId, @@ -1822,6 +1822,34 @@ export async function getWorkspaceInfo ( return clientWs } +/** + * @public + */ +export async function getWorkspacesInfo ( + ctx: MeasureContext, + db: AccountDB, + branding: Branding | null, + token: string, + ids: string[] +): Promise { + const { email } = decodeToken(ctx, token) + + if (email !== systemAccountEmail) { + ctx.error('getWorkspaceInfos with wrong email', { email, token }) + throw new PlatformError(new Status(Severity.ERROR, platform.status.Forbidden, {})) + } + const query: Query = { + workspace: { $in: ids } + } + const workspaces = await ctx.with( + 'get-workspace', + {}, + async () => await db.workspace.find(query, { lastVisit: 'descending' }) + ) + + return workspaces.map(mapToClientWorkspace) +} + async function getUpgradeStatistics (db: AccountDB, region: string): Promise { return ( (await db.upgrade.findOne({ @@ -2937,6 +2965,7 @@ export function getMethods (hasSignUp: boolean = true): Record & { workspaceId: string } - /** * @public */ diff --git a/server/client/src/account.ts b/server/client/src/account.ts index b05d86f0c4..f866a2fc6f 100644 --- a/server/client/src/account.ts +++ b/server/client/src/account.ts @@ -16,6 +16,7 @@ import { AccountRole, BackupStatus, + ClientWorkspaceInfo, Doc, Ref, type BaseWorkspaceInfo, @@ -262,6 +263,25 @@ export async function getWorkspaceInfo ( return workspaceInfo.result as BaseWorkspaceInfo | undefined } +export async function getWorkspacesInfo (token: string, workspaces: string[]): Promise { + const accountsUrl = getAccoutsUrlOrFail() + const workspaceInfo = await ( + await fetch(accountsUrl, { + method: 'POST', + headers: { + Authorization: 'Bearer ' + token, + 'Content-Type': 'application/json' + }, + body: JSON.stringify({ + method: 'getWorkspacesInfo', + params: [workspaces] + }) + }) + ).json() + + return workspaceInfo.result as ClientWorkspaceInfo[] +} + export async function login (user: string, password: string, workspace: string): Promise { const accountsUrl = getAccoutsUrlOrFail() const response = await fetch(accountsUrl, { diff --git a/server/collaborator/src/account.ts b/server/collaborator/src/account.ts index 40c456b070..297eb7acb0 100644 --- a/server/collaborator/src/account.ts +++ b/server/collaborator/src/account.ts @@ -13,7 +13,7 @@ // limitations under the License. // -import { ClientWorkspaceInfo } from '@hcengineering/account' +import { ClientWorkspaceInfo } from '@hcengineering/core' import config from './config' export async function getWorkspaceInfo (token: string): Promise { diff --git a/services/calendar/pod-calendar/src/calendar.ts b/services/calendar/pod-calendar/src/calendar.ts index 9e5627cd80..3b0af5d449 100644 --- a/services/calendar/pod-calendar/src/calendar.ts +++ b/services/calendar/pod-calendar/src/calendar.ts @@ -23,7 +23,6 @@ import calendar, { } from '@hcengineering/calendar' import { Contact } from '@hcengineering/contact' import core, { - Account, AttachedData, Client, Data, @@ -32,134 +31,114 @@ import core, { DocumentUpdate, Mixin, Ref, - TxOperations, - TxUpdateDoc, - generateId + TxOperations } from '@hcengineering/core' import setting from '@hcengineering/setting' import { htmlToMarkup, markupToHTML } from '@hcengineering/text' import { deepEqual } from 'fast-equals' -import type { Credentials, OAuth2Client } from 'google-auth-library' -import { calendar_v3, google } from 'googleapis' import type { Collection, Db } from 'mongodb' -import { encode64 } from './base64' import { CalendarController } from './calendarController' -import config from './config' -import { RateLimiter } from './rateLimiter' -import type { CalendarHistory, EventHistory, EventWatch, ProjectCredentials, State, Token, User, Watch } from './types' +import type { CalendarHistory, DummyWatch, EventHistory, Token, User } from './types' import { encodeReccuring, isToken, parseRecurrenceStrings } from './utils' import type { WorkspaceClient } from './workspaceClient' - -const SCOPES = [ - 'https://www.googleapis.com/auth/calendar.calendars.readonly', - 'https://www.googleapis.com/auth/calendar.calendarlist.readonly', - 'https://www.googleapis.com/auth/calendar.events', - 'https://www.googleapis.com/auth/userinfo.email' -] -const DUMMY_RESOURCE = 'Dummy' +import { GoogleClient } from './googleClient' +import { calendar_v3 } from 'googleapis' +import { WatchController } from './watch' export class CalendarClient { - private readonly oAuth2Client: OAuth2Client private readonly calendar: calendar_v3.Calendar - private readonly tokens: Collection private readonly calendarHistories: Collection private readonly histories: Collection private readonly client: TxOperations - private me: string | undefined = undefined - private readonly watches: EventWatch[] = [] - private calendarWatch: Watch | undefined = undefined - private refreshTimer: NodeJS.Timeout | undefined = undefined + private readonly systemTxOp: TxOperations private readonly activeSync: Record = {} - private readonly rateLimiter = new RateLimiter(1000, 500) + private readonly dummyWatches: DummyWatch[] = [] + // to do< find!!!! + private readonly googleClient + + private inactiveTimer: NodeJS.Timeout isClosed: boolean = false private constructor ( - credentials: ProjectCredentials, private readonly user: User, - mongo: Db, + private readonly mongo: Db, client: Client, private readonly workspace: WorkspaceClient ) { - const { client_secret, client_id, redirect_uris } = credentials.web // eslint-disable-line - this.oAuth2Client = new google.auth.OAuth2(client_id, client_secret, redirect_uris[0]) // eslint-disable-line - this.calendar = google.calendar({ version: 'v3', auth: this.oAuth2Client }) - this.tokens = mongo.collection('tokens') + this.client = new TxOperations(client, this.user.userId) + this.systemTxOp = new TxOperations(client, core.account.System) + this.googleClient = new GoogleClient(user, mongo, this) + this.calendar = this.googleClient.calendar this.histories = mongo.collection('histories') this.calendarHistories = mongo.collection('calendarHistories') - this.client = new TxOperations(client, this.user.userId) + this.inactiveTimer = setTimeout(() => { + this.closeByTimer() + }, 60 * 1000) + } + + async cleanIntegration (): Promise { + const integration = await this.client.findOne(setting.class.Integration, { + createdBy: this.user.userId, + type: calendar.integrationType.Calendar, + value: this.user.email + }) + if (integration !== undefined) { + await this.client.update(integration, { disabled: true }) + } + this.workspace.removeClient(this.user.email) + } + + private updateTimer (): void { + clearTimeout(this.inactiveTimer) + this.inactiveTimer = setTimeout(() => { + this.closeByTimer() + }, 60 * 1000) } static async create ( - credentials: ProjectCredentials, user: User | Token, mongo: Db, client: Client, workspace: WorkspaceClient ): Promise { - const calendarClient = new CalendarClient(credentials, user, mongo, client, workspace) + const calendarClient = new CalendarClient(user, mongo, client, workspace) if (isToken(user)) { - await calendarClient.setToken(user) - await calendarClient.refreshToken() + await calendarClient.googleClient.init(user) await calendarClient.addClient() } return calendarClient } - static getAutUrl (redirectURL: string, workspace: string, userId: Ref, token: string): string { - const credentials = JSON.parse(config.Credentials) - const { client_secret, client_id, redirect_uris } = credentials.web // eslint-disable-line - const oAuth2Client = new google.auth.OAuth2(client_id, client_secret, redirect_uris[0]) // eslint-disable-line - const state: State = { - token, - redirectURL, - workspace, - userId, - email: '' - } - const authUrl = oAuth2Client.generateAuthUrl({ - access_type: 'offline', - scope: SCOPES, - state: encode64(JSON.stringify(state)) - }) - return authUrl - } - async authorize (code: string): Promise { - const token = await this.oAuth2Client.getToken(code) - await this.setToken(token.tokens) - const me = await this.getMe() - const providedScopes = token.tokens.scope?.split(' ') ?? [] - for (const scope of SCOPES) { - if (providedScopes.findIndex((p) => p === scope) === -1) { - const integrations = await this.client.findAll(setting.class.Integration, { - createdBy: this.user.userId, - type: calendar.integrationType.Calendar - }) - for (const integration of integrations.filter((p) => p.value === '')) { - await this.client.remove(integration) - } - - const updated = integrations.find((p) => p.disabled && p.value === me) - if (updated !== undefined) { - await this.client.update(updated, { - disabled: true, - error: calendar.string.NotAllPermissions - }) - } else { - await this.client.createDoc(setting.class.Integration, core.space.Workspace, { - type: calendar.integrationType.Calendar, - disabled: true, - error: calendar.string.NotAllPermissions, - value: me - }) - } - throw new Error( - `Not all scopes provided, provided: ${providedScopes.join(', ')} required: ${SCOPES.join(', ')}` - ) + this.updateTimer() + const me = await this.googleClient.authorize(code) + if (me === undefined) { + const integrations = await this.client.findAll(setting.class.Integration, { + createdBy: this.user.userId, + type: calendar.integrationType.Calendar + }) + for (const integration of integrations.filter((p) => p.value === '')) { + await this.client.remove(integration) } + + const updated = integrations.find((p) => p.disabled && p.value === me) + if (updated !== undefined) { + await this.client.update(updated, { + disabled: true, + error: calendar.string.NotAllPermissions + }) + } else { + const value = await this.googleClient.getMe() + await this.client.createDoc(setting.class.Integration, core.space.Workspace, { + type: calendar.integrationType.Calendar, + disabled: true, + error: calendar.string.NotAllPermissions, + value + }) + } + throw new Error('Not all scopes provided') } - await this.refreshToken() await this.addClient() const integrations = await this.client.findAll(setting.class.Integration, { @@ -185,32 +164,31 @@ export class CalendarClient { }) } - await this.startSync(me) - void this.syncOurEvents() + void this.syncOurEvents().then(async () => { + await this.startSync(me) + }) return me } - async signout (byError: boolean = false): Promise { + async signout (): Promise { + this.updateTimer() try { - await this.close() - await this.oAuth2Client.revokeCredentials() + this.close() + if (isToken(this.user)) { + const watch = WatchController.get(this.mongo) + await watch.unsubscribe(this.user) + } + await this.googleClient.signout() } catch {} - await this.tokens.deleteOne({ - userId: this.user.userId, - workspace: this.user.workspace - }) + const integration = await this.client.findOne(setting.class.Integration, { createdBy: this.user.userId, type: calendar.integrationType.Calendar, value: this.user.email }) if (integration !== undefined) { - if (byError) { - await this.client.update(integration, { disabled: true }) - } else { - await this.client.remove(integration) - } + await this.client.remove(integration) } this.workspace.removeClient(this.user.email) } @@ -218,194 +196,62 @@ export class CalendarClient { async startSync (me?: string): Promise { try { if (me === undefined) { - me = await this.getMe() + me = await this.googleClient.getMe() } await this.syncCalendars(me) const calendars = this.workspace.getMyCalendars(me) for (const calendar of calendars) { if (calendar.externalId !== undefined) { - void this.sync(calendar.externalId, me) + await this.sync(calendar.externalId, me) } } } catch (err) { - console.log('Start sync error', this.user.workspace, this.user.userId, err) + console.error('Start sync error', this.user.workspace, this.user.userId, err) } } async startSyncCalendar (calendar: ExternalCalendar): Promise { - const me = await this.getMe() + const me = await this.googleClient.getMe() void this.sync(calendar.externalId, me) } - async close (): Promise { - if (this.refreshTimer !== undefined) clearTimeout(this.refreshTimer) - for (const watch of this.watches) { + private closeByTimer (): void { + this.close() + this.workspace.removeClient(this.user.email) + } + + close (): void { + this.googleClient.close() + for (const watch of this.dummyWatches) { clearTimeout(watch.timer) - try { - if (watch.resourceId !== DUMMY_RESOURCE) { - await this.rateLimiter.take(1) - await this.calendar.channels.stop({ requestBody: { id: watch.channelId, resourceId: watch.resourceId } }) - } - } catch (err) { - console.log('close error', err) - } - } - if (this.calendarWatch !== undefined) { - clearTimeout(this.calendarWatch.timer) - try { - await this.rateLimiter.take(1) - await this.calendar.channels.stop({ - requestBody: { id: this.calendarWatch.channelId, resourceId: this.calendarWatch.resourceId } - }) - } catch (err) { - console.log('close error', err) - } } this.isClosed = true } - private async getMe (): Promise { - if (this.me !== undefined) { - return this.me - } - - const info = await google.oauth2({ version: 'v2', auth: this.oAuth2Client }).userinfo.get() - this.me = info.data.email ?? '' - return this.me - } - - // #region Token - - private async getCurrentToken (): Promise { - return await this.tokens.findOne({ - userId: this.user.userId, - workspace: this.user.workspace, - email: this.me ?? this.user.email - }) - } - - private async updateCurrentToken (token: Credentials): Promise { - await this.tokens.updateOne( - { - userId: this.user.userId, - workspace: this.user.workspace, - email: this.me ?? this.user.email - }, - { - $set: { - ...token - } - } - ) - } - private async addClient (): Promise { try { - const me = await this.getMe() + const me = await this.googleClient.getMe() const controller = CalendarController.getCalendarController() controller.addClient(me, this) + this.updateTimer() } catch (err) { - console.log('Add client error', this.user.workspace, this.user.userId, err) + console.error('Add client error', this.user.workspace, this.user.userId, err) } } - private async setToken (token: Credentials): Promise { - try { - this.oAuth2Client.setCredentials(token) - } catch (err: any) { - console.log('Set token error', this.user.workspace, this.user.userId, err) - await this.checkError(err) - throw err - } - } - - private async updateToken (token: Credentials): Promise { - try { - const currentToken = await this.getCurrentToken() - if (currentToken != null) { - await this.updateCurrentToken(token) - } else { - await this.tokens.insertOne({ - userId: this.user.userId, - email: this.me ?? this.user.email, - workspace: this.user.workspace, - token: this.user.token, - ...token - }) - } - } catch (err) { - console.log('update token error', this.user.workspace, this.user.userId, err) - } - } - - private async refreshToken (): Promise { - try { - const res = await this.oAuth2Client.refreshAccessToken() - await this.updateToken(res.credentials) - this.refreshTimer = setTimeout( - () => { - void this.refreshToken() - }, - 30 * 60 * 1000 - ) - } catch (err: any) { - console.log("Couldn't refresh token, error:", err) - if (err?.response?.data?.error === 'invalid_grant' || err.message === 'No refresh token is set.') { - await this.signout(true) - } else { - this.refreshTimer = setTimeout( - () => { - void this.refreshToken() - }, - 15 * 60 * 1000 - ) - } - throw err - } - } - - // #endregion - // #region Calendars - private async watchCalendar (): Promise { - try { - const current = this.calendarWatch - if (current !== undefined) { - clearTimeout(current.timer) - await this.rateLimiter.take(1) - await this.calendar.channels.stop({ requestBody: { id: current.channelId, resourceId: current.resourceId } }) - } - const channelId = generateId() - const me = await this.getMe() - const body = { id: channelId, address: config.WATCH_URL, type: 'webhook', token: `user=${me}&mode=calendar` } - await this.rateLimiter.take(1) - const res = await this.calendar.calendarList.watch({ requestBody: body }) - if (res.data.expiration != null && res.data.resourceId !== null) { - const time = Number(res.data.expiration) - new Date().getTime() - // eslint-disable-next-line - const timer = setTimeout(() => void this.watchCalendar(), time) - this.calendarWatch = { - channelId, - resourceId: res.data.resourceId ?? '', - timer - } - } - } catch (err) { - console.log('Calendar watch error', err) - } - } - async syncCalendars (me: string): Promise { const history = await this.getCalendarHistory(me) await this.calendarSync(history?.historyId) - await this.watchCalendar() + await this.googleClient.watchCalendar() } private async calendarSync (syncToken?: string, pageToken?: string): Promise { try { - await this.rateLimiter.take(1) - const res = await this.calendar.calendarList.list({ + this.updateTimer() + await this.googleClient.rateLimiter.take(1) + const res = await this.googleClient.calendar.calendarList.list({ syncToken, pageToken }) @@ -418,7 +264,7 @@ export class CalendarClient { try { await this.syncCalendar(calendar) } catch (err) { - console.log('save calendar error', JSON.stringify(event), err) + console.error('save calendar error', JSON.stringify(calendar), err) } } if (nextPageToken != null) { @@ -432,15 +278,16 @@ export class CalendarClient { await this.calendarSync() return } - console.log('Calendar sync error', this.user.workspace, this.user.userId, err) + console.error('Calendar sync error', this.user.workspace, this.user.userId, err) } } private async syncCalendar (val: calendar_v3.Schema$CalendarListEntry): Promise { if (val.id != null) { + const me = await this.googleClient.getMe() const exists = await this.client.findOne(calendar.class.ExternalCalendar, { externalId: val.id, - externalUser: this.me ?? '' + externalUser: me }) if (exists === undefined) { const data: Data = { @@ -448,7 +295,7 @@ export class CalendarClient { visibility: 'freeBusy', hidden: false, externalId: val.id, - externalUser: this.me ?? '', + externalUser: me, default: false } if (val.primary === true) { @@ -485,7 +332,7 @@ export class CalendarClient { } private async setCalendarHistoryId (historyId: string): Promise { - const me = await this.getMe() + const me = await this.googleClient.getMe() await this.calendarHistories.updateOne( { userId: this.user.userId, @@ -506,90 +353,22 @@ export class CalendarClient { // #region Events // #region Incoming - - async stopWatch (calendar: ExternalCalendar): Promise { - for (const watch of this.watches) { - if (watch.calendarId === calendar.externalId) { - clearTimeout(watch.timer) - try { - if (watch.resourceId !== DUMMY_RESOURCE) { - await this.rateLimiter.take(1) - await this.calendar.channels.stop({ requestBody: { id: watch.channelId, resourceId: watch.resourceId } }) - } - } catch (err) { - console.log('close error', err) - } - } - } - } - private async watch (calendarId: string): Promise { - try { - const index = this.watches.findIndex((p) => p.calendarId === calendarId) - if (index !== -1) { - const current = this.watches[index] - if (current !== undefined) { - clearTimeout(current.timer) - if (current.resourceId !== DUMMY_RESOURCE) { - await this.rateLimiter.take(1) - await this.calendar.channels.stop({ - requestBody: { id: current.channelId, resourceId: current.resourceId } - }) - } - } - this.watches.splice(index, 1) - } - const channelId = generateId() - const me = await this.getMe() - const body = { - id: channelId, - address: config.WATCH_URL, - type: 'webhook', - token: `user=${me}&mode=events&calendarId=${calendarId}` - } - await this.rateLimiter.take(1) - const res = await this.calendar.events.watch({ calendarId, requestBody: body }) - if (res.data.expiration != null && res.data.resourceId != null) { - const time = Number(res.data.expiration) - new Date().getTime() - // eslint-disable-next-line - const timer = setTimeout(() => void this.watch(calendarId), time) - this.watches.push({ - calendarId, - channelId, - resourceId: res.data.resourceId ?? '', - timer - }) - } - } catch (err: any) { - if (err?.errors?.[0]?.reason === 'pushNotSupportedForRequestedResource') { - await this.dummyWatch(calendarId) - } else { - console.log('Watch error', err) - await this.checkError(err) - } + if (!(await this.googleClient.watch(calendarId))) { + await this.dummyWatch(calendarId) } } - private async checkError (err: any): Promise { - if (err?.response?.data?.error === 'invalid_grant') { - await this.signout(true) - return true - } - return false - } - private async dummyWatch (calendarId: string): Promise { - const me = await this.getMe() + const me = await this.googleClient.getMe() const timer = setTimeout( () => { void this.sync(calendarId, me) }, 6 * 60 * 60 * 1000 ) - this.watches.push({ + this.dummyWatches.push({ calendarId, - channelId: DUMMY_RESOURCE, - resourceId: DUMMY_RESOURCE, timer }) } @@ -613,7 +392,7 @@ export class CalendarClient { } private async setEventHistoryId (calendarId: string, historyId: string): Promise { - const me = await this.getMe() + const me = await this.googleClient.getMe() await this.histories.updateOne( { calendarId, @@ -637,7 +416,7 @@ export class CalendarClient { private async eventsSync (calendarId: string, syncToken?: string, pageToken?: string): Promise { try { - await this.rateLimiter.take(1) + await this.googleClient.rateLimiter.take(1) const res = await this.calendar.events.list({ calendarId, syncToken, @@ -653,7 +432,7 @@ export class CalendarClient { try { await this.syncEvent(calendarId, event, res.data.accessRole ?? 'reader') } catch (err) { - console.log('save event error', JSON.stringify(event), err) + console.error('save event error', JSON.stringify(event), err) } } if (nextPageToken != null) { @@ -667,14 +446,15 @@ export class CalendarClient { await this.eventsSync(calendarId) return } - await this.checkError(err) - console.log('Event sync error', this.user.workspace, this.user.userId, err) + await this.googleClient.checkError(err) + console.error('Event sync error', this.user.workspace, this.user.userId, err) } } private async syncEvent (calendarId: string, event: calendar_v3.Schema$Event, accessRole: string): Promise { + this.updateTimer() if (event.id != null) { - const me = await this.getMe() + const me = await this.googleClient.getMe() const calendars = this.workspace.getMyCalendars(me) const _calendar = calendars.find((p) => p.externalId === event.organizer?.email) ?? @@ -695,6 +475,7 @@ export class CalendarClient { } private async updateExtEvent (event: calendar_v3.Schema$Event, current: Event): Promise { + this.updateTimer() if (event.status === 'cancelled' && current._class !== calendar.class.ReccuringInstance) { await this.client.remove(current) return @@ -747,7 +528,7 @@ export class CalendarClient { if (this.client.getHierarchy().hasMixin(current, mixin as Ref>)) { const diff = this.getDiff(attr, this.client.getHierarchy().as(current, mixin as Ref>)) if (Object.keys(diff).length > 0) { - await this.client.updateMixin( + await this.systemTxOp.updateMixin( current._id, current._class, calendar.space.Calendar, @@ -756,7 +537,7 @@ export class CalendarClient { ) } } else { - await this.client.createMixin( + await this.systemTxOp.createMixin( current._id, current._class, calendar.space.Calendar, @@ -782,7 +563,7 @@ export class CalendarClient { for (const mixin in mixins) { const attr = mixins[mixin] if (typeof attr === 'object' && Object.keys(attr).length > 0) { - await this.client.createMixin( + await this.systemTxOp.createMixin( _id, calendar.class.Event, calendar.space.Calendar, @@ -799,10 +580,11 @@ export class CalendarClient { accessRole: string, _calendar: ExternalCalendar ): Promise { + this.updateTimer() const data: AttachedData = await this.parseData(event, accessRole, _calendar._id) if (event.recurringEventId != null) { const parseRule = parseRecurrenceStrings(event.recurrence ?? []) - const id = await this.client.addCollection( + const id = await this.systemTxOp.addCollection( calendar.class.ReccuringInstance, calendar.space.Calendar, calendar.ids.NoAttached, @@ -823,7 +605,7 @@ export class CalendarClient { } else if (event.status !== 'cancelled') { if (event.recurrence != null) { const parseRule = parseRecurrenceStrings(event.recurrence) - const id = await this.client.addCollection( + const id = await this.systemTxOp.addCollection( calendar.class.ReccuringEvent, calendar.space.Calendar, calendar.ids.NoAttached, @@ -840,7 +622,7 @@ export class CalendarClient { ) await this.saveMixins(event, id) } else { - const id = await this.client.addCollection( + const id = await this.systemTxOp.addCollection( calendar.class.Event, calendar.space.Calendar, calendar.ids.NoAttached, @@ -1004,12 +786,14 @@ export class CalendarClient { } private async createRecInstance (calendarId: string, event: ReccuringInstance): Promise { - const body = this.convertBody(event) + this.updateTimer() + const me = await this.googleClient.getMe() + const body = this.convertBody(event, me) const req: calendar_v3.Params$Resource$Events$Instances = { calendarId, eventId: event.recurringEventId } - await this.rateLimiter.take(1) + await this.googleClient.rateLimiter.take(1) const instancesResp = await this.calendar.events.instances(req) const items = instancesResp.data.items const target = items?.find( @@ -1020,7 +804,7 @@ export class CalendarClient { ) if (target?.id != null) { body.id = target.id - await this.rateLimiter.take(1) + await this.googleClient.rateLimiter.take(1) await this.calendar.events.update({ calendarId, eventId: target.id, @@ -1031,14 +815,15 @@ export class CalendarClient { } async createEvent (event: Event): Promise { + const me = await this.googleClient.getMe() try { const _calendar = this.workspace.calendars.byId.get(event.calendar as Ref) if (_calendar !== undefined) { if (event._class === calendar.class.ReccuringInstance) { await this.createRecInstance(_calendar.externalId, event as ReccuringInstance) } else { - const body = this.convertBody(event) - await this.rateLimiter.take(1) + const body = this.convertBody(event, me) + await this.googleClient.rateLimiter.take(1) await this.calendar.events.insert({ calendarId: _calendar.externalId, requestBody: body @@ -1046,23 +831,24 @@ export class CalendarClient { } } } catch (err: any) { - await this.checkError(err) + await this.googleClient.checkError(err) // eslint-disable-next-line throw new Error(`Create event error, ${this.user.workspace}, ${this.user.userId}, ${event._id}, ${err?.message}`) } } - async updateEvent (event: Event, tx: TxUpdateDoc): Promise { + async updateEvent (event: Event): Promise { + const me = await this.googleClient.getMe() const _calendar = this.workspace.calendars.byId.get(event.calendar as Ref) const calendarId = _calendar?.externalId if (calendarId !== undefined) { try { - await this.rateLimiter.take(1) + await this.googleClient.rateLimiter.take(1) const current = await this.calendar.events.get({ calendarId, eventId: event.eventId }) if (current?.data !== undefined) { if (current.data.organizer?.self === true) { - const ev = this.applyUpdate(current.data, event) - await this.rateLimiter.take(1) + const ev = this.applyUpdate(current.data, event, me) + await this.googleClient.rateLimiter.take(1) await this.calendar.events.update({ calendarId, eventId: event.eventId, @@ -1074,18 +860,19 @@ export class CalendarClient { if (err.code === 404) { await this.createEvent(event) } else { - console.log('Update event error', this.user.workspace, this.user.userId, err) - await this.checkError(err) + console.error('Update event error', this.user.workspace, this.user.userId, err) + await this.googleClient.checkError(err) } } } } async remove (eventId: string, calendarId: string): Promise { + this.updateTimer() const current = await this.calendar.events.get({ calendarId, eventId }) if (current?.data !== undefined) { if (current.data.organizer?.self === true) { - await this.rateLimiter.take(1) + await this.googleClient.rateLimiter.take(1) await this.calendar.events.delete({ eventId, calendarId @@ -1101,12 +888,13 @@ export class CalendarClient { await this.remove(event.eventId, _calendar.externalId) } } catch (err) { - console.log('Remove event error', this.user.workspace, this.user.userId, err) + console.error('Remove event error', this.user.workspace, this.user.userId, err) } } async syncOurEvents (): Promise { - const me = await this.getMe() + this.updateTimer() + const me = await this.googleClient.getMe() const events = await this.client.findAll(calendar.class.Event, { access: 'owner', createdBy: this.user.userId, @@ -1119,25 +907,29 @@ export class CalendarClient { } async syncMyEvent (event: Event): Promise { + const me = await this.googleClient.getMe() if (event.access === 'owner' || event.access === 'writer') { try { const space = this.workspace.calendars.byId.get(event.calendar as Ref) - if (space !== undefined && space.externalUser === this.me) { + if (space !== undefined && space.externalUser === me) { + this.updateTimer() if (!(await this.update(event, space))) { await this.create(event, space) } } } catch (err: any) { - console.log('Sync event error', this.user.workspace, this.user.userId, event._id, err.message) + console.error('Sync event error', this.user.workspace, this.user.userId, event._id, err.message) } } } private async create (event: Event, space: ExternalCalendar): Promise { - const body = this.convertBody(event) + this.updateTimer() + const me = await this.googleClient.getMe() + const body = this.convertBody(event, me) const calendarId = space?.externalId if (calendarId !== undefined) { - await this.rateLimiter.take(1) + await this.googleClient.rateLimiter.take(1) await this.calendar.events.insert({ calendarId, requestBody: body @@ -1146,14 +938,16 @@ export class CalendarClient { } private async update (event: Event, space: ExternalCalendar): Promise { + this.updateTimer() + const me = await this.googleClient.getMe() const calendarId = space?.externalId if (calendarId !== undefined) { try { - await this.rateLimiter.take(1) + await this.googleClient.rateLimiter.take(1) const current = await this.calendar.events.get({ calendarId, eventId: event.eventId }) if (current !== undefined) { - const ev = this.applyUpdate(current.data, event) - await this.rateLimiter.take(1) + const ev = this.applyUpdate(current.data, event, me) + await this.googleClient.rateLimiter.take(1) await this.calendar.events.update({ calendarId, eventId: event.eventId, @@ -1190,7 +984,7 @@ export class CalendarClient { return res } - private convertBody (event: Event): calendar_v3.Schema$Event { + private convertBody (event: Event, me: string): calendar_v3.Schema$Event { const res: calendar_v3.Schema$Event = { start: convertDate(event.date, event.allDay, getTimezone(event)), end: convertDate(event.dueDate, event.allDay, getTimezone(event)), @@ -1229,10 +1023,10 @@ export class CalendarClient { }) } } - const attendees = this.getAttendees(event) + const attendees = this.getAttendees(event, me) if (attendees.length > 0) { res.attendees = attendees.map((p) => { - if (p === this.me) { + if (p === me) { return { email: p, responseStatus: 'accepted', self: true } } return { email: p } @@ -1251,7 +1045,7 @@ export class CalendarClient { return res } - private applyUpdate (event: calendar_v3.Schema$Event, current: Event): calendar_v3.Schema$Event { + private applyUpdate (event: calendar_v3.Schema$Event, current: Event, me: string): calendar_v3.Schema$Event { if (current.title !== event.summary) { event.summary = current.title } @@ -1276,7 +1070,7 @@ export class CalendarClient { if (current.location !== event.location) { event.location = current.location } - const attendees = this.getAttendees(current) + const attendees = this.getAttendees(current, me) if (attendees.length > 0 && event.attendees !== undefined) { for (const attendee of attendees) { if (event.attendees.findIndex((p) => p.email === attendee) === -1) { @@ -1293,11 +1087,11 @@ export class CalendarClient { return event } - private getAttendees (event: Event): string[] { + private getAttendees (event: Event, me: string): string[] { const res = new Set() for (const participant of event.participants) { const integrations = this.workspace.integrations.byContact.get(participant) ?? [] - const integration = integrations.find((p) => p === this.me) ?? integrations[0] + const integration = integrations.find((p) => p === me) ?? integrations[0] if (integration !== undefined && integration !== '') { res.add(integration) } else { diff --git a/services/calendar/pod-calendar/src/calendarController.ts b/services/calendar/pod-calendar/src/calendarController.ts index ad6783d4c9..6ef093c324 100644 --- a/services/calendar/pod-calendar/src/calendarController.ts +++ b/services/calendar/pod-calendar/src/calendarController.ts @@ -13,27 +13,35 @@ // limitations under the License. // -import { Account, isActiveMode, RateLimiter, Ref, systemAccountEmail } from '@hcengineering/core' -import { type Db } from 'mongodb' +import { Account, isActiveMode, isDeletingMode, RateLimiter, Ref, systemAccountEmail } from '@hcengineering/core' +import { Event } from '@hcengineering/calendar' +import { Collection, type Db } from 'mongodb' import { type CalendarClient } from './calendar' import config from './config' -import { type ProjectCredentials, type Token, type User } from './types' +import { type Token, type User } from './types' import { WorkspaceClient } from './workspaceClient' -import { getWorkspaceInfo } from '@hcengineering/server-client' +import { getWorkspacesInfo } from '@hcengineering/server-client' import { generateToken } from '@hcengineering/server-token' export class CalendarController { - private readonly workspaces: Map = new Map() + private readonly workspaces: Map> = new Map< + string, + WorkspaceClient | Promise + >() - private readonly credentials: ProjectCredentials + private readonly tokens: Collection private readonly clients: Map = new Map() - private readonly initLimitter = new RateLimiter(config.InitLimit) protected static _instance: CalendarController private constructor (private readonly mongo: Db) { - this.credentials = JSON.parse(config.Credentials) + this.tokens = mongo.collection('tokens') CalendarController._instance = this + setInterval(() => { + if (this.workspaces.size > 0) { + console.log('active workspaces', this.workspaces.size) + } + }, 60000) } static getCalendarController (mongo?: Db): CalendarController { @@ -45,7 +53,7 @@ export class CalendarController { } async startAll (): Promise { - const tokens = await this.mongo.collection('tokens').find().toArray() + const tokens = await this.tokens.find().toArray() const groups = new Map() console.log('start calendar service', tokens.length) for (const token of tokens) { @@ -59,69 +67,86 @@ export class CalendarController { } const limiter = new RateLimiter(config.InitLimit) - - for (const [workspace, tokens] of groups) { + const token = generateToken(systemAccountEmail, { name: '' }) + const ids = [...groups.keys()] + console.log('start workspaces', ids) + const infos = await getWorkspacesInfo(token, ids) + console.log('infos', infos) + for (const info of infos) { + const tokens = groups.get(info.workspaceId) + if (tokens === undefined) { + console.log('no tokens for workspace', info.workspaceId) + continue + } + if (isDeletingMode(info.mode)) { + if (tokens !== undefined) { + for (const token of tokens) { + await this.tokens.deleteOne({ userId: token.userId, workspace: token.workspace }) + } + } + continue + } + if (!isActiveMode(info.mode)) { + continue + } await limiter.add(async () => { - const wstok = generateToken(systemAccountEmail, { name: workspace }) - const info = await getWorkspaceInfo(wstok) - if (info === undefined) { - console.log('workspace not found', workspace) - return - } - if (!isActiveMode(info.mode)) { - console.log('workspace is not active', workspace) - return - } - const startPromise = this.startWorkspace(workspace, tokens) - const timeoutPromise = new Promise((resolve) => { - setTimeout(() => { - resolve() - }, 60000) - }) - await Promise.race([startPromise, timeoutPromise]) + console.log('start workspace', info.workspaceId) + const workspace = await this.startWorkspace(info.workspaceId, tokens) + await workspace.sync() }) } - - await limiter.waitProcessing() - console.log('Calendar service started') } - async startWorkspace (workspace: string, tokens: Token[]): Promise { + async startWorkspace (workspace: string, tokens: Token[]): Promise { const workspaceClient = await this.getWorkspaceClient(workspace) - const clients: CalendarClient[] = [] for (const token of tokens) { try { const timeout = setTimeout(() => { - console.log('init client hang', token.workspace, token.userId) + console.warn('init client hang', token.workspace, token.userId) }, 60000) - const client = await workspaceClient.createCalendarClient(token) + console.log('init client', token.workspace, token.userId) + await workspaceClient.createCalendarClient(token) clearTimeout(timeout) - clients.push(client) } catch (err) { console.error(`Couldn't create client for ${workspace} ${token.userId} ${token.email}`) } } - for (const client of clients) { - void this.initLimitter.add(async () => { - await client.startSync() - }) - } - void workspaceClient.sync() - console.log('Workspace started', workspace) + return workspaceClient } - push (email: string, mode: 'events' | 'calendar', calendarId?: string): void { - const clients = this.clients.get(email) - for (const client of clients ?? []) { + async push (email: string, mode: 'events' | 'calendar', calendarId?: string): Promise { + const tokens = await this.tokens.find({ email, access_token: { $exists: true } }).toArray() + const token = generateToken(systemAccountEmail, { name: '' }) + const workspaces = [...new Set(tokens.map((p) => p.workspace))] + const infos = await getWorkspacesInfo(token, workspaces) + for (const token of tokens) { + const info = infos.find((p) => p.workspace === token.workspace) + if (info === undefined) { + continue + } + if (isDeletingMode(info.mode)) { + await this.tokens.deleteOne({ userId: token.userId, workspace: token.workspace }) + continue + } + if (!isActiveMode(info.mode)) { + continue + } + const workspace = await this.getWorkspaceClient(token.workspace) + const calendarClient = await workspace.createCalendarClient(token) if (mode === 'calendar') { - void client.syncCalendars(email) + await calendarClient.syncCalendars(email) } if (mode === 'events' && calendarId !== undefined) { - void client.sync(calendarId, email) + await calendarClient.sync(calendarId, email) } } } + async pushEvent (workspace: string, event: Event, type: 'create' | 'update' | 'delete'): Promise { + const workspaceController = await this.getWorkspaceClient(workspace) + await workspaceController.pushEvent(event, type) + } + addClient (email: string, client: CalendarClient): void { const clients = this.clients.get(email) if (clients === undefined) { @@ -135,10 +160,12 @@ export class CalendarController { removeClient (email: string): void { const clients = this.clients.get(email) if (clients !== undefined) { - this.clients.set( - email, - clients.filter((p) => !p.isClosed) - ) + const filtered = clients.filter((p) => !p.isClosed) + if (filtered.length === 0) { + this.clients.delete(email) + } else { + this.clients.set(email, filtered) + } } } @@ -151,7 +178,7 @@ export class CalendarController { const workspaceClient = await this.getWorkspaceClient(workspace) const clients = await workspaceClient.signout(value) if (clients === 0) { - this.workspaces.delete(workspace) + this.removeWorkspace(workspace) } } @@ -160,7 +187,10 @@ export class CalendarController { } async close (): Promise { - for (const workspace of this.workspaces.values()) { + for (let workspace of this.workspaces.values()) { + if (workspace instanceof Promise) { + workspace = await workspace + } await workspace.close() } this.workspaces.clear() @@ -179,16 +209,22 @@ export class CalendarController { } private async getWorkspaceClient (workspace: string): Promise { - let res = this.workspaces.get(workspace) - if (res === undefined) { - try { - res = await WorkspaceClient.create(this.credentials, this.mongo, workspace, this) - this.workspaces.set(workspace, res) - } catch (err) { - console.error(`Couldn't create workspace worker for ${workspace}, reason: ${JSON.stringify(err)}`) - throw err + const res = this.workspaces.get(workspace) + if (res !== undefined) { + if (res instanceof Promise) { + return await res } + return res + } + try { + const client = WorkspaceClient.create(this.mongo, workspace, this) + this.workspaces.set(workspace, client) + const res = await client + this.workspaces.set(workspace, res) + return res + } catch (err) { + console.error(`Couldn't create workspace worker for ${workspace}, reason: ${JSON.stringify(err)}`) + throw err } - return res } } diff --git a/services/calendar/pod-calendar/src/googleClient.ts b/services/calendar/pod-calendar/src/googleClient.ts new file mode 100644 index 0000000000..7f4c6be317 --- /dev/null +++ b/services/calendar/pod-calendar/src/googleClient.ts @@ -0,0 +1,313 @@ +// +// 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 { Credentials, OAuth2Client } from 'google-auth-library' +import { calendar_v3, google } from 'googleapis' +import { ProjectCredentials, State, Token, User, Watch, WatchBase } from './types' +import config from './config' +import { encode64 } from './base64' +import { Account, generateId, Ref } from '@hcengineering/core' +import { Collection, Db } from 'mongodb' +import { RateLimiter } from './rateLimiter' +import { CalendarClient } from './calendar' + +export const DUMMY_RESOURCE = 'Dummy' + +const SCOPES = [ + 'https://www.googleapis.com/auth/calendar.calendars.readonly', + 'https://www.googleapis.com/auth/calendar.calendarlist.readonly', + 'https://www.googleapis.com/auth/calendar.events', + 'https://www.googleapis.com/auth/userinfo.email' +] + +export class GoogleClient { + private me: string | undefined = undefined + private readonly credentials: ProjectCredentials + private readonly oAuth2Client: OAuth2Client + readonly calendar: calendar_v3.Calendar + private readonly tokens: Collection + private readonly watches: Collection + + private refreshTimer: NodeJS.Timeout | undefined = undefined + + readonly rateLimiter = new RateLimiter(1000, 500) + + constructor ( + private readonly user: User, + mongo: Db, + private readonly calendarClient: CalendarClient + ) { + this.tokens = mongo.collection('tokens') + this.credentials = JSON.parse(config.Credentials) + const { client_secret, client_id, redirect_uris } = this.credentials.web // eslint-disable-line + this.oAuth2Client = new google.auth.OAuth2(client_id, client_secret, redirect_uris[0]) // eslint-disable-line + this.calendar = google.calendar({ version: 'v3', auth: this.oAuth2Client }) + this.watches = mongo.collection('watch') + } + + static getAutUrl (redirectURL: string, workspace: string, userId: Ref, token: string): string { + const credentials = JSON.parse(config.Credentials) + const { client_secret, client_id, redirect_uris } = credentials.web // eslint-disable-line + const oAuth2Client = new google.auth.OAuth2(client_id, client_secret, redirect_uris[0]) // eslint-disable-line + const state: State = { + token, + redirectURL, + workspace, + userId, + email: '' + } + const authUrl = oAuth2Client.generateAuthUrl({ + access_type: 'offline', + scope: SCOPES, + state: encode64(JSON.stringify(state)) + }) + return authUrl + } + + async signout (): Promise { + // get watch controller and unsubscibe + await this.oAuth2Client.revokeCredentials() + await this.tokens.deleteOne({ + userId: this.user.userId, + workspace: this.user.workspace + }) + } + + async init (token: Token): Promise { + await this.setToken(token) + await this.refreshToken() + } + + async authorize (code: string): Promise { + const token = await this.oAuth2Client.getToken(code) + await this.setToken(token.tokens) + const me = await this.getMe() + const providedScopes = token.tokens.scope?.split(' ') ?? [] + for (const scope of SCOPES) { + if (providedScopes.findIndex((p) => p === scope) === -1) { + console.error(`Not all scopes provided, provided: ${providedScopes.join(', ')} required: ${SCOPES.join(', ')}`) + return undefined + } + } + await this.refreshToken() + + return me + } + + close (): void { + if (this.refreshTimer !== undefined) clearTimeout(this.refreshTimer) + } + + async getMe (): Promise { + if (this.me !== undefined) { + return this.me + } + + const info = await google.oauth2({ version: 'v2', auth: this.oAuth2Client }).userinfo.get() + this.me = info.data.email ?? '' + return this.me + } + + private async setToken (token: Credentials): Promise { + try { + this.oAuth2Client.setCredentials(token) + } catch (err: any) { + console.error('Set token error', this.user.workspace, this.user.userId, err) + await this.checkError(err) + throw err + } + } + + async checkError (err: any): Promise { + if (err?.response?.data?.error === 'invalid_grant') { + await this.calendarClient.cleanIntegration() + return true + } + return false + } + + private async updateToken (token: Credentials): Promise { + try { + const currentToken = await this.getCurrentToken() + if (currentToken != null) { + await this.updateCurrentToken(token) + } else { + await this.tokens.insertOne({ + userId: this.user.userId, + email: this.me ?? this.user.email, + workspace: this.user.workspace, + token: this.user.token, + ...token + }) + } + } catch (err) { + console.error('update token error', this.user.workspace, this.user.userId, err) + } + } + + private async refreshToken (): Promise { + try { + const res = await this.oAuth2Client.refreshAccessToken() + await this.updateToken(res.credentials) + this.refreshTimer = setTimeout( + () => { + void this.refreshToken() + }, + 30 * 60 * 1000 + ) + } catch (err: any) { + console.error("Couldn't refresh token, error:", err) + if (err?.response?.data?.error === 'invalid_grant' || err.message === 'No refresh token is set.') { + await this.calendarClient.cleanIntegration() + } else { + this.refreshTimer = setTimeout( + () => { + void this.refreshToken() + }, + 15 * 60 * 1000 + ) + } + throw err + } + } + + private async getCurrentToken (): Promise { + return await this.tokens.findOne({ + userId: this.user.userId, + workspace: this.user.workspace, + email: this.me ?? this.user.email + }) + } + + private async updateCurrentToken (token: Credentials): Promise { + await this.tokens.updateOne( + { + userId: this.user.userId, + workspace: this.user.workspace, + email: this.me ?? this.user.email + }, + { + $set: { + ...token + } + } + ) + } + + async watchCalendar (): Promise { + try { + const current = await this.watches.findOne({ + userId: this.user.userId, + workspace: this.user.workspace, + calendarId: null + }) + if (current != null) { + await this.rateLimiter.take(1) + await this.calendar.channels.stop({ requestBody: { id: current.channelId, resourceId: current.resourceId } }) + } + const channelId = generateId() + const me = await this.getMe() + const body = { id: channelId, address: config.WATCH_URL, type: 'webhook', token: `user=${me}&mode=calendar` } + await this.rateLimiter.take(1) + const res = await this.calendar.calendarList.watch({ requestBody: body }) + if (res.data.expiration != null && res.data.resourceId !== null) { + if (current != null) { + await this.watches.updateOne( + { + userId: this.user.userId, + workspace: this.user.workspace, + calendarId: null + }, + { + channelId, + expired: Number.parseInt(res.data.expiration), + resourceId: res.data.resourceId ?? '' + } + ) + } else { + await this.watches.insertOne({ + calendarId: null, + channelId, + expired: Number.parseInt(res.data.expiration), + resourceId: res.data.resourceId ?? '', + userId: this.user.userId, + workspace: this.user.workspace + }) + } + } + } catch (err) { + console.error('Calendar watch error', err) + } + } + + async watch (calendarId: string): Promise { + try { + const current = await this.watches.findOne({ + userId: this.user.userId, + workspace: this.user.workspace, + calendarId + }) + if (current != null) { + await this.rateLimiter.take(1) + await this.calendar.channels.stop({ + requestBody: { id: current.channelId, resourceId: current.resourceId } + }) + } + const channelId = generateId() + const me = await this.getMe() + const body = { + id: channelId, + address: config.WATCH_URL, + type: 'webhook', + token: `user=${me}&mode=events&calendarId=${calendarId}` + } + await this.rateLimiter.take(1) + const res = await this.calendar.events.watch({ calendarId, requestBody: body }) + if (res.data.expiration != null && res.data.resourceId != null) { + if (current != null) { + await this.watches.updateOne( + { + userId: this.user.userId, + workspace: this.user.workspace, + calendarId + }, + { + channelId, + expired: Number.parseInt(res.data.expiration), + resourceId: res.data.resourceId ?? '' + } + ) + } else { + await this.watches.insertOne({ + calendarId, + channelId, + expired: Number.parseInt(res.data.expiration), + resourceId: res.data.resourceId ?? '', + userId: this.user.userId, + workspace: this.user.workspace + }) + } + } + return true + } catch (err: any) { + if (err?.errors?.[0]?.reason === 'pushNotSupportedForRequestedResource') { + return false + } else { + console.error('Watch error', err) + await this.checkError(err) + return false + } + } + } +} diff --git a/services/calendar/pod-calendar/src/main.ts b/services/calendar/pod-calendar/src/main.ts index 324b2a516b..2367745f37 100644 --- a/services/calendar/pod-calendar/src/main.ts +++ b/services/calendar/pod-calendar/src/main.ts @@ -15,7 +15,6 @@ import { type IncomingHttpHeaders } from 'http' import { decode64 } from './base64' -import { CalendarClient } from './calendar' import { CalendarController } from './calendarController' import config from './config' import { createServer, listen } from './server' @@ -24,6 +23,8 @@ import { type Endpoint, type State } from './types' import { setMetadata } from '@hcengineering/platform' import serverClient from '@hcengineering/server-client' import serverToken, { decodeToken } from '@hcengineering/server-token' +import { GoogleClient } from './googleClient' +import { WatchController } from './watch' const extractToken = (header: IncomingHttpHeaders): any => { try { @@ -41,6 +42,8 @@ export const main = async (): Promise => { const db = await getDB() const calendarController = CalendarController.getCalendarController(db) await calendarController.startAll() + const watchController = WatchController.get(db) + watchController.startCheck() const endpoints: Endpoint[] = [ { endpoint: '/signin', @@ -57,10 +60,10 @@ export const main = async (): Promise => { const { email, workspace } = decodeToken(token) const userId = await calendarController.getUserId(email, workspace.name) - const url = CalendarClient.getAutUrl(redirectURL, workspace.name, userId, token) + const url = GoogleClient.getAutUrl(redirectURL, workspace.name, userId, token) res.send(url) } catch (err) { - console.log('signin error', err) + console.error('signin error', err) res.status(500).send() } } @@ -75,7 +78,7 @@ export const main = async (): Promise => { await calendarController.newClient(state, code) res.redirect(state.redirectURL) } catch (err) { - console.log(err) + console.error(err) res.redirect(state.redirectURL) } } @@ -97,7 +100,7 @@ export const main = async (): Promise => { const { workspace } = decodeToken(token) await calendarController.signout(workspace.name, value) } catch (err) { - console.log('signout error', err) + console.error('signout error', err) } res.send() @@ -122,9 +125,23 @@ export const main = async (): Promise => { res.status(400).send({ err: "'data' is missing" }) return } - calendarController.push(data.user, data.mode as 'events' | 'calendar', data.calendarId) + void calendarController.push(data.user, data.mode as 'events' | 'calendar', data.calendarId) } + res.send() + } + }, + { + endpoint: '/event', + type: 'post', + handler: async (req, res) => { + const { event, workspace, type } = req.body + + if (event === undefined || workspace === undefined || type === undefined) { + res.status(400).send({ err: "'event' or 'workspace' or 'type' is missing" }) + return + } + void calendarController.pushEvent(workspace, event, type) res.send() } } @@ -134,6 +151,7 @@ export const main = async (): Promise => { const shutdown = (): void => { server.close(() => { + watchController.stop() void calendarController .close() .then(async () => { diff --git a/services/calendar/pod-calendar/src/types.ts b/services/calendar/pod-calendar/src/types.ts index 45b9ede9c4..3ec0e10e6a 100644 --- a/services/calendar/pod-calendar/src/types.ts +++ b/services/calendar/pod-calendar/src/types.ts @@ -18,17 +18,28 @@ import type { Account, Ref, Timestamp } from '@hcengineering/core' import type { NextFunction, Request, Response } from 'express' import type { Credentials } from 'google-auth-library' -export interface Watch { - timer: NodeJS.Timeout +export interface WatchBase { + userId: Ref + workspace: string + expired: Timestamp channelId: string resourceId: string + calendarId: string | null } -export interface EventWatch { +export interface CalendarsWatch extends WatchBase { + calendarId: null +} + +export interface EventWatch extends WatchBase { + calendarId: string +} + +export type Watch = CalendarsWatch | EventWatch + +export interface DummyWatch { timer: NodeJS.Timeout calendarId: string - channelId: string - resourceId: string } export type Token = User & Credentials diff --git a/services/calendar/pod-calendar/src/watch.ts b/services/calendar/pod-calendar/src/watch.ts new file mode 100644 index 0000000000..940343e552 --- /dev/null +++ b/services/calendar/pod-calendar/src/watch.ts @@ -0,0 +1,230 @@ +import { Collection, Db } from 'mongodb' +import { EventWatch, Token, Watch, WatchBase } from './types' +import { generateId, isActiveMode, systemAccountEmail } from '@hcengineering/core' +import { getWorkspacesInfo } from '@hcengineering/server-client' +import { generateToken } from '@hcengineering/server-token' +import config from './config' +import { Credentials, OAuth2Client } from 'google-auth-library' +import { calendar_v3, google } from 'googleapis' +import { RateLimiter } from './rateLimiter' + +export class WatchClient { + private readonly watches: Collection + private readonly oAuth2Client: OAuth2Client + private readonly calendar: calendar_v3.Calendar + private readonly user: Token + private me: string = '' + readonly rateLimiter = new RateLimiter(1000, 500) + + private constructor (mongo: Db, token: Token) { + this.user = token + this.watches = mongo.collection('watch') + const credentials = JSON.parse(config.Credentials) + const { client_secret, client_id, redirect_uris } = credentials.web // eslint-disable-line + this.oAuth2Client = new google.auth.OAuth2(client_id, client_secret, redirect_uris[0]) // eslint-disable-line + this.calendar = google.calendar({ version: 'v3', auth: this.oAuth2Client }) + } + + static async Create (mongo: Db, token: Token): Promise { + const watchClient = new WatchClient(mongo, token) + await watchClient.init(token) + return watchClient + } + + private async setToken (token: Credentials): Promise { + try { + this.oAuth2Client.setCredentials(token) + const info = await google.oauth2({ version: 'v2', auth: this.oAuth2Client }).userinfo.get() + this.me = info.data.email ?? '' + } catch (err: any) { + console.error('Set token error', this.user.workspace, this.user.userId, err) + await this.checkError(err) + throw err + } + } + + async checkError (err: any): Promise { + if (err?.response?.data?.error === 'invalid_grant') { + await this.watches.deleteMany({ userId: this.user.userId, workspace: this.user.workspace }) + } + } + + private async init (token: Token): Promise { + await this.setToken(token) + } + + async subscribe (watches: Watch[]): Promise { + for (const watch of watches) { + if (watch.calendarId == null) { + await this.watchCalendars(watch) + } else { + await this.watchCalendar(watch) + } + } + } + + async unsubscribe (watches: Watch[]): Promise { + for (const watch of watches) { + await this.unsubscribeWatch(watch) + } + } + + private async unsubscribeWatch (current: Watch): Promise { + await this.rateLimiter.take(1) + await this.calendar.channels.stop({ requestBody: { id: current.channelId, resourceId: current.resourceId } }) + } + + private async watchCalendars (current: Watch): Promise { + try { + await this.unsubscribeWatch(current) + const channelId = generateId() + const body = { id: channelId, address: config.WATCH_URL, type: 'webhook', token: `user=${this.me}&mode=calendar` } + await this.rateLimiter.take(1) + const res = await this.calendar.calendarList.watch({ requestBody: body }) + if (res.data.expiration != null && res.data.resourceId !== null) { + // eslint-disable-next-line + this.watches.updateOne( + { + userId: current.userId, + workspace: current.workspace, + calendarId: null + }, + { + $set: { + channelId, + expired: Number.parseInt(res.data.expiration), + resourceId: res.data.resourceId ?? '' + } + } + ) + } + } catch (err) { + console.error('Calendar watch error', err) + } + } + + private async watchCalendar (current: EventWatch): Promise { + try { + await this.unsubscribeWatch(current) + const channelId = generateId() + const body = { + id: channelId, + address: config.WATCH_URL, + type: 'webhook', + token: `user=${this.me}&mode=events&calendarId=${current.calendarId}` + } + await this.rateLimiter.take(1) + const res = await this.calendar.events.watch({ calendarId: current.calendarId, requestBody: body }) + if (res.data.expiration != null && res.data.resourceId != null) { + // eslint-disable-next-line + this.watches.updateOne( + { + userId: current.userId, + workspace: current.workspace, + calendarId: current.calendarId + }, + { + $set: { + channelId, + expired: Number.parseInt(res.data.expiration), + resourceId: res.data.resourceId ?? '' + } + } + ) + } + } catch (err: any) { + await this.checkError(err) + } + } +} + +// we have to refresh channels approx each week +export class WatchController { + private readonly watches: Collection + private readonly tokens: Collection + + private timer: NodeJS.Timeout | undefined = undefined + protected static _instance: WatchController + + private constructor (private readonly mongo: Db) { + this.watches = mongo.collection('watch') + this.tokens = mongo.collection('tokens') + console.log('watch started') + } + + static get (mongo: Db): WatchController { + if (WatchController._instance !== undefined) { + return WatchController._instance + } + return new WatchController(mongo) + } + + async unsubscribe (user: Token): Promise { + const allWatches = await this.watches.find({ userId: user.userId, workspae: user.workspace }).toArray() + await this.watches.deleteMany({ userId: user.userId, workspae: user.workspace }) + const token = this.tokens.findOne({ user: user.userId, workspace: user.workspace }) + if (token == null) return + const watchClient = await WatchClient.Create(this.mongo, user) + await watchClient.unsubscribe(allWatches) + } + + stop (): void { + if (this.timer !== undefined) { + clearInterval(this.timer) + } + } + + startCheck (): void { + this.timer = setInterval( + () => { + void this.checkAll() + }, + 1000 * 60 * 60 * 24 + ) + void this.checkAll() + } + + async checkAll (): Promise { + const expired = Date.now() + 24 * 60 * 60 * 1000 + const watches = await this.watches + .find({ + expired: { $lt: expired } + }) + .toArray() + console.log('watch, found for update', watches.length) + const groups = new Map() + const workspaces = new Set() + for (const watch of watches) { + workspaces.add(watch.workspace) + const key = `${watch.userId}:${watch.workspace}` + const group = groups.get(key) + if (group !== undefined) { + group.push(watch) + } else { + groups.set(key, [watch]) + } + } + const token = generateToken(systemAccountEmail, { name: '' }) + const infos = await getWorkspacesInfo(token, [...groups.keys()]) + const tokens = await this.tokens.find({ workspace: { $in: [...workspaces] } }).toArray() + for (const group of groups.values()) { + try { + const userId = group[0].userId + const workspace = group[0].workspace + const token = tokens.find((p) => p.workspace === workspace && p.userId === userId) + if (token === undefined) { + await this.watches.deleteMany({ userId, workspace }) + continue + } + const info = infos.find((p) => p.workspace === workspace) + if (info === undefined || isActiveMode(info.mode)) { + await this.watches.deleteMany({ userId, workspace }) + continue + } + const watchClient = await WatchClient.Create(this.mongo, token) + await watchClient.subscribe(group) + } catch {} + } + console.log('watch check done') + } +} diff --git a/services/calendar/pod-calendar/src/workspaceClient.ts b/services/calendar/pod-calendar/src/workspaceClient.ts index 189e39ecbb..c7531da5e9 100644 --- a/services/calendar/pod-calendar/src/workspaceClient.ts +++ b/services/calendar/pod-calendar/src/workspaceClient.ts @@ -16,6 +16,7 @@ import calendar, { Event, ExternalCalendar } from '@hcengineering/calendar' import contact, { Channel, Contact, type Employee, type PersonAccount } from '@hcengineering/contact' import core, { + RateLimiter, TxOperations, TxProcessor, systemAccountEmail, @@ -35,14 +36,20 @@ import { Collection, type Db } from 'mongodb' import { CalendarClient } from './calendar' import { CalendarController } from './calendarController' import { getClient } from './client' -import { SyncHistory, type ProjectCredentials, type User } from './types' +import { SyncHistory, Token, type User } from './types' +import config from './config' export class WorkspaceClient { private readonly txHandlers: ((...tx: Tx[]) => Promise)[] = [] - private client!: Client - private readonly clients: Map = new Map() + client!: Client + private readonly clients: Map> = new Map< + string, + CalendarClient | Promise + >() + private readonly syncHistory: Collection + private readonly tokens: Collection private channels = new Map, Channel>() private readonly calendarsByEmail = new Map() readonly calendars = { @@ -62,39 +69,45 @@ export class WorkspaceClient { } private constructor ( - private readonly credentials: ProjectCredentials, private readonly mongo: Db, private readonly workspace: string, private readonly serviceController: CalendarController ) { + this.tokens = mongo.collection('tokens') this.syncHistory = mongo.collection('syncHistories') } - static async create ( - credentials: ProjectCredentials, - mongo: Db, - workspace: string, - serviceController: CalendarController - ): Promise { - const instance = new WorkspaceClient(credentials, mongo, workspace, serviceController) + static async getSystemClient (workspace: string): Promise { + const token = generateToken(systemAccountEmail, { name: workspace }) + return await getClient(token) + } + + static async create (mongo: Db, workspace: string, serviceController: CalendarController): Promise { + const instance = new WorkspaceClient(mongo, workspace, serviceController) await instance.initClient(workspace) return instance } async createCalendarClient (user: User): Promise { const current = this.getCalendarClient(user.email) - if (current !== undefined) return current - const newClient = await CalendarClient.create(this.credentials, user, this.mongo, this.client, this) + if (current !== undefined) { + if (current instanceof Promise) { + return await current + } + return current + } + const newClient = CalendarClient.create(user, this.mongo, this.client, this) this.clients.set(user.email, newClient) - console.log('create new client', user.email, this.workspace) - return newClient + const res = await newClient + this.clients.set(user.email, res) + return res } async newCalendarClient (user: User, code: string): Promise { - const newClient = await CalendarClient.create(this.credentials, user, this.mongo, this.client, this) + const newClient = await CalendarClient.create(user, this.mongo, this.client, this) const email = await newClient.authorize(code) if (this.clients.has(email)) { - await newClient.close() + newClient.close() throw new Error('Client already exist') } this.clients.set(email, newClient) @@ -102,8 +115,11 @@ export class WorkspaceClient { } async close (): Promise { - for (const client of this.clients.values()) { - await client.close() + for (let client of this.clients.values()) { + if (client instanceof Promise) { + client = await client + } + client.close() } this.clients.clear() await this.client?.close() @@ -118,9 +134,12 @@ export class WorkspaceClient { } async signout (value: string, byError: boolean = false): Promise { - const client = this.clients.get(value) + let client = this.clients.get(value) if (client !== undefined) { - await client.signout(byError) + if (client instanceof Promise) { + client = await client + } + await client.signout() } else { const integration = await this.client.findOne(setting.class.Integration, { type: calendar.integrationType.Calendar, @@ -143,19 +162,38 @@ export class WorkspaceClient { this.clients.delete(email) this.serviceController.removeClient(email) if (this.clients.size > 0) return + void this.close() this.serviceController.removeWorkspace(this.workspace) } - private getCalendarClient (email: string): CalendarClient | undefined { + private getCalendarClient (email: string): CalendarClient | Promise | undefined { return this.clients.get(email) } - private getCalendarClientByCalendar (id: Ref): CalendarClient | undefined { + private async getCalendarClientByCalendar ( + id: Ref, + create: boolean = false + ): Promise { const calendar = this.calendars.byId.get(id) if (calendar === undefined) { - console.log("couldn't find calendar by id", id) + console.warn("couldn't find calendar by id", id) + return } - return calendar != null ? this.clients.get(calendar.externalUser) : undefined + const client = this.clients.get(calendar.externalUser) + if (client instanceof Promise) { + return await client + } + if (client === undefined && create) { + const user = await this.tokens.findOne({ + workspace: this.workspace, + access_token: { $exists: true }, + email: calendar.externalUser + }) + if (user != null) { + return await this.createCalendarClient(user) + } + } + return client } private async initClient (workspace: string): Promise { @@ -187,6 +225,15 @@ export class WorkspaceClient { async sync (): Promise { await this.getNewEvents() + const limiter = new RateLimiter(config.InitLimit) + for (let client of this.clients.values()) { + void limiter.add(async () => { + if (client instanceof Promise) { + client = await client + } + await client.startSync() + }) + } } // #region Events @@ -213,6 +260,20 @@ export class WorkspaceClient { ) } + async pushEvent (event: Event, type: 'create' | 'update' | 'delete'): Promise { + const client = await this.getCalendarClientByCalendar(event.calendar as Ref, true) + if (client === undefined) { + console.warn('Client not found', event.calendar, this.workspace) + return + } + if (type === 'delete') { + await client.removeEvent(event) + } else { + await client.syncMyEvent(event) + } + await this.updateSyncTime() + } + async getNewEvents (): Promise { const lastSync = await this.getSyncTime() const query = lastSync !== undefined ? { modifiedOn: { $gt: lastSync } } : {} @@ -220,17 +281,16 @@ export class WorkspaceClient { this.txHandlers.push(async (...tx: Tx[]) => { await this.txEventHandler(...tx) }) - console.log('receive new events', this.workspace, newEvents.length) for (const newEvent of newEvents) { - const client = this.getCalendarClientByCalendar(newEvent.calendar as Ref) + const client = await this.getCalendarClientByCalendar(newEvent.calendar as Ref) if (client === undefined) { - console.log('Client not found', newEvent.calendar, this.workspace) + console.warn('Client not found', newEvent.calendar, this.workspace) return } await client.syncMyEvent(newEvent) await this.updateSyncTime() } - console.log('all messages synced', this.workspace) + console.log('all outcoming messages synced', this.workspace) } private async txEventHandler (...txes: Tx[]): Promise { @@ -256,7 +316,7 @@ export class WorkspaceClient { if (hierarhy.isDerived(tx.objectClass, calendar.class.Event)) { const doc = TxProcessor.createDoc2Doc(tx as TxCreateDoc) if (doc.access !== 'owner') return - const client = this.getCalendarClientByCalendar(doc.calendar as Ref) + const client = await this.getCalendarClientByCalendar(doc.calendar as Ref) if (client === undefined) { return } @@ -264,7 +324,7 @@ export class WorkspaceClient { await client.createEvent(doc) await this.updateSyncTime() } catch (err) { - console.log(err) + console.error(err) } } } @@ -281,7 +341,7 @@ export class WorkspaceClient { const extracted = txes.filter((p) => p._id !== tx._id) const ev = TxProcessor.buildDoc2Doc(extracted) if (ev !== undefined) { - const oldClient = this.getCalendarClientByCalendar(ev.calendar as Ref) + const oldClient = await this.getCalendarClientByCalendar(ev.calendar as Ref) if (oldClient !== undefined) { const oldCalendar = this.calendars.byId.get(ev.calendar as Ref) if (oldCalendar !== undefined) { @@ -290,16 +350,16 @@ export class WorkspaceClient { } } } catch (err) { - console.log('Error on remove event', err) + console.error('Error on remove event', err) } try { - const client = this.getCalendarClientByCalendar(event.calendar as Ref) + const client = await this.getCalendarClientByCalendar(event.calendar as Ref) if (client !== undefined) { await client.syncMyEvent(event) } await this.updateSyncTime() } catch (err) { - console.log('Error on move event', err) + console.error('Error on move event', err) } } @@ -315,15 +375,15 @@ export class WorkspaceClient { return } if (event.access !== 'owner' && event.access !== 'writer') return - const client = this.getCalendarClientByCalendar(event.calendar as Ref) + const client = await this.getCalendarClientByCalendar(event.calendar as Ref) if (client === undefined) { return } try { - await client.updateEvent(event, tx) + await client.updateEvent(event) await this.updateSyncTime() } catch (err) { - console.log(err) + console.error(err) } } } @@ -337,7 +397,7 @@ export class WorkspaceClient { const ev = TxProcessor.buildDoc2Doc(txes) if (ev === undefined) return if (ev.access !== 'owner' && ev.access !== 'writer') return - const client = this.getCalendarClientByCalendar(ev?.calendar as Ref) + const client = await this.getCalendarClientByCalendar(ev?.calendar as Ref) if (client === undefined) { return } diff --git a/services/github/pod-github/src/account.ts b/services/github/pod-github/src/account.ts index cf013c2a71..092d0743e6 100644 --- a/services/github/pod-github/src/account.ts +++ b/services/github/pod-github/src/account.ts @@ -1,4 +1,4 @@ -import { ClientWorkspaceInfo } from '@hcengineering/account' +import { ClientWorkspaceInfo } from '@hcengineering/core' import config from './config' /** diff --git a/services/github/pod-github/src/platform.ts b/services/github/pod-github/src/platform.ts index 4de5d9befe..a9e4333c99 100644 --- a/services/github/pod-github/src/platform.ts +++ b/services/github/pod-github/src/platform.ts @@ -9,6 +9,7 @@ import core, { BrandingMap, Client, ClientConnectEvent, + ClientWorkspaceInfo, DocumentUpdate, isActiveMode, isDeletingMode, @@ -29,7 +30,6 @@ import { Installation, type InstallationCreatedEvent, type InstallationUnsuspend import { Collection } from 'mongodb' import { App, Octokit } from 'octokit' -import { ClientWorkspaceInfo } from '@hcengineering/account' import { Analytics } from '@hcengineering/analytics' import { SplitLogger } from '@hcengineering/analytics-service' import contact, { Person, PersonAccount } from '@hcengineering/contact'