From cde79b05cefc9e413f357f863844e57420869d09 Mon Sep 17 00:00:00 2001 From: Denis Bykhov Date: Wed, 25 Jun 2025 23:44:13 +0500 Subject: [PATCH] Calendar duplicate watch fix (#9383) Signed-off-by: Denis Bykhov --- services/calendar/pod-calendar/src/auth.ts | 2 +- services/calendar/pod-calendar/src/sync.ts | 12 +- services/calendar/pod-calendar/src/types.ts | 2 - services/calendar/pod-calendar/src/watch.ts | 151 ++++++++++---------- 4 files changed, 84 insertions(+), 83 deletions(-) diff --git a/services/calendar/pod-calendar/src/auth.ts b/services/calendar/pod-calendar/src/auth.ts index 64dd7d8c2c..dbfe4a9034 100644 --- a/services/calendar/pod-calendar/src/auth.ts +++ b/services/calendar/pod-calendar/src/auth.ts @@ -131,9 +131,9 @@ export class AuthController { const secret = await this.accountClient.getIntegrationSecret(data) if (secret == null) return const token = JSON.parse(secret.secret) + await removeIntegrationSecret(this.ctx, this.accountClient, data) const watchController = WatchController.get(this.ctx, this.accountClient) await watchController.unsubscribe(token) - await removeIntegrationSecret(this.ctx, this.accountClient, data) } private async process (code: string): Promise { diff --git a/services/calendar/pod-calendar/src/sync.ts b/services/calendar/pod-calendar/src/sync.ts index 4ac98b4517..60298d3b9c 100644 --- a/services/calendar/pod-calendar/src/sync.ts +++ b/services/calendar/pod-calendar/src/sync.ts @@ -13,7 +13,7 @@ // limitations under the License. // -import { +import core, { AttachedData, Data, Doc, @@ -78,6 +78,7 @@ export class IncomingSyncManager { private readonly rateLimiter: RateLimiter private calendars: ExternalCalendar[] = [] private readonly participants = new Map>() + private readonly systemClient: TxOperations private constructor ( private readonly ctx: MeasureContext, private readonly accountClient: AccountClient, @@ -87,6 +88,7 @@ export class IncomingSyncManager { private readonly googleClient: calendar_v3.Calendar ) { this.rateLimiter = getRateLimitter(this.email) + this.systemClient = new TxOperations(client.client, core.account.System) } static async sync ( @@ -479,7 +481,7 @@ export class IncomingSyncManager { 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.systemClient.addCollection( calendar.class.ReccuringInstance, calendar.space.Calendar, calendar.ids.NoAttached, @@ -500,7 +502,7 @@ export class IncomingSyncManager { } else if (event.status !== 'cancelled') { if (event.recurrence != null) { const parseRule = parseRecurrenceStrings(event.recurrence) - const id = await this.client.addCollection( + const id = await this.systemClient.addCollection( calendar.class.ReccuringEvent, calendar.space.Calendar, calendar.ids.NoAttached, @@ -517,7 +519,7 @@ export class IncomingSyncManager { ) await this.saveMixins(event, id) } else { - const id = await this.client.addCollection( + const id = await this.systemClient.addCollection( calendar.class.Event, calendar.space.Calendar, calendar.ids.NoAttached, @@ -536,7 +538,7 @@ export class IncomingSyncManager { for (const mixin in mixins) { const attr = mixins[mixin] if (typeof attr === 'object' && Object.keys(attr).length > 0) { - await this.client.createMixin( + await this.systemClient.createMixin( _id, calendar.class.Event, calendar.space.Calendar, diff --git a/services/calendar/pod-calendar/src/types.ts b/services/calendar/pod-calendar/src/types.ts index 3582f131e0..164f8b647e 100644 --- a/services/calendar/pod-calendar/src/types.ts +++ b/services/calendar/pod-calendar/src/types.ts @@ -24,8 +24,6 @@ export type GoogleEmail = string & { googleEmail: true } export interface WatchBase { email: GoogleEmail - userId: PersonId - workspace: WorkspaceUuid expired: Timestamp channelId: string resourceId: string diff --git a/services/calendar/pod-calendar/src/watch.ts b/services/calendar/pod-calendar/src/watch.ts index 96a88c9362..41cd05db24 100644 --- a/services/calendar/pod-calendar/src/watch.ts +++ b/services/calendar/pod-calendar/src/watch.ts @@ -1,11 +1,11 @@ import { AccountClient } from '@hcengineering/account-client' -import { generateId, isActiveMode, MeasureContext, PersonId, WorkspaceUuid } from '@hcengineering/core' +import { generateId, MeasureContext } from '@hcengineering/core' import { Credentials, OAuth2Client } from 'google-auth-library' import { calendar_v3 } from 'googleapis' import config from './config' import { getKvsClient } from './kvsUtils' import { getRateLimitter, RateLimiter } from './rateLimiter' -import { CALENDAR_INTEGRATION, EventWatch, GoogleEmail, Token, User, Watch, WatchBase } from './types' +import { CALENDAR_INTEGRATION, EventWatch, GoogleEmail, Token, User, Watch } from './types' import { getGoogleClient } from './utils' export class WatchClient { @@ -14,7 +14,11 @@ export class WatchClient { private readonly user: Token readonly rateLimiter: RateLimiter - private constructor (token: Token) { + private constructor ( + private readonly ctx: MeasureContext, + token: Token, + private readonly accountClient: AccountClient + ) { this.user = token this.rateLimiter = getRateLimitter(this.user.email) @@ -23,15 +27,15 @@ export class WatchClient { this.oAuth2Client = res.auth } - static async Create (token: Token): Promise { - const watchClient = new WatchClient(token) + static async Create (ctx: MeasureContext, token: Token, accountClient: AccountClient): Promise { + const watchClient = new WatchClient(ctx, token, accountClient) await watchClient.setToken(token) return watchClient } private async getWatches (): Promise> { const client = getKvsClient() - const key = `${CALENDAR_INTEGRATION}:watch:${this.user.workspace}:${this.user.userId}:${this.user.email}` + const key = `${CALENDAR_INTEGRATION}:watch:${this.user.email}` const watches = await client.listKeys(key) return watches ?? {} } @@ -48,10 +52,22 @@ export class WatchClient { async checkError (err: any): Promise { if (err?.response?.data?.error === 'invalid_grant') { - const watches = await this.getWatches() - const client = getKvsClient() - for (const key in watches) { - await client.deleteKey(key) + await this.accountClient.deleteIntegrationSecret({ + socialId: this.user.userId, + kind: CALENDAR_INTEGRATION, + workspaceUuid: this.user.workspace, + key: this.user.email + }) + const active = await this.accountClient.listIntegrationsSecrets({ + kind: CALENDAR_INTEGRATION, + key: this.user.email + }) + if (active.length === 0) { + const watches = await this.getWatches() + const client = getKvsClient() + for (const key in watches) { + await client.deleteKey(key) + } } } } @@ -82,32 +98,38 @@ export class WatchClient { private async watchCalendars (current: Watch): Promise { try { await this.unsubscribeWatch(current) - await watchCalendars(this.user, this.user.email, this.calendar) + } catch (err: any) { + this.ctx.error('Calendars watch unsubscribe error', { message: err.message }) + } + try { + await watchCalendars(this.user.email, this.calendar) } catch (err) { - console.error('Calendar watch error', err) + this.ctx.error('Calendars watch error', { message: (err as any).message }) } } private async watchCalendar (current: EventWatch): Promise { try { await this.unsubscribeWatch(current) - await watchCalendar(this.user, this.user.email, current.calendarId, this.calendar) } catch (err: any) { - await this.checkError(err) + this.ctx.error('Calendar watch unsubscribe error', { message: err.message }) + } + try { + await watchCalendar(this.user.email, current.calendarId, this.calendar) + } catch (err: any) { + this.ctx.error('Calendar watch error', { message: err.message, calendarId: current.calendarId }) } } } -async function watchCalendars (user: User, email: GoogleEmail, googleClient: calendar_v3.Calendar): Promise { +async function watchCalendars (email: GoogleEmail, googleClient: calendar_v3.Calendar): Promise { const channelId = generateId() const body = { id: channelId, address: config.WATCH_URL, type: 'webhook', token: `user=${email}&mode=calendar` } const res = await googleClient.calendarList.watch({ requestBody: body }) if (res.data.expiration != null && res.data.resourceId !== null) { const client = getKvsClient() - const key = `${CALENDAR_INTEGRATION}:watch:${user.workspace}:${user.userId}:${email}:null` + const key = `${CALENDAR_INTEGRATION}:watch:${email}:null` await client.setValue(key, { - userId: user.userId, - workspace: user.workspace, email, calendarId: null, channelId, @@ -118,7 +140,6 @@ async function watchCalendars (user: User, email: GoogleEmail, googleClient: cal } async function watchCalendar ( - user: User, email: GoogleEmail, calendarId: string, googleClient: calendar_v3.Calendar @@ -133,10 +154,8 @@ async function watchCalendar ( const res = await googleClient.events.watch({ calendarId, requestBody: body }) if (res.data.expiration != null && res.data.resourceId != null) { const client = getKvsClient() - const key = `${CALENDAR_INTEGRATION}:watch:${user.workspace}:${user.userId}:${email}:${calendarId}` + const key = `${CALENDAR_INTEGRATION}:watch:${email}:${calendarId}` await client.setValue(key, { - userId: user.userId, - workspace: user.workspace, email, calendarId, channelId, @@ -165,27 +184,27 @@ export class WatchController { return WatchController._instance } - private async getUserWatches (userId: PersonId, workspace: WorkspaceUuid): Promise> { + private async getUserWatches (email: GoogleEmail): Promise> { const client = getKvsClient() - const key = `${CALENDAR_INTEGRATION}:watch:${workspace}:${userId}` + const key = `${CALENDAR_INTEGRATION}:watch:${email}` return (await client.listKeys(key)) ?? {} } async unsubscribe (user: Token): Promise { - const client = getKvsClient() - const watches = await this.getUserWatches(user.userId, user.workspace) - for (const key in watches) { - await client.deleteKey(key) - } - const token = await this.accountClient.getIntegrationSecret({ - socialId: user.userId, + const active = await this.accountClient.listIntegrationsSecrets({ kind: CALENDAR_INTEGRATION, - workspaceUuid: user.workspace, key: user.email }) - if (token == null) return - const watchClient = await WatchClient.Create(user) - await watchClient.unsubscribe(Object.values(watches)) + if (active.length === 0) { + const client = getKvsClient() + const watches = await this.getUserWatches(user.email) + for (const key in watches) { + await client.deleteKey(key) + } + + const watchClient = await WatchClient.Create(this.ctx, user, this.accountClient) + await watchClient.unsubscribe(Object.values(watches)) + } } stop (): void { @@ -218,45 +237,27 @@ export class WatchController { } this.ctx.info('watch, found for update', { count: toRefresh.length }) if (toRefresh.length === 0) return - const groups = new Map() - const workspaces = new Set() + const groups = new Map() for (const watch of toRefresh) { - 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 curr = groups.get(watch.email) ?? [] + curr.push(watch) + groups.set(watch.email, curr) } - const ids = [...workspaces] - if (ids.length === 0) return - const infos = await this.accountClient.getWorkspacesInfo(ids) - const tokens = await this.accountClient.listIntegrationsSecrets({ kind: CALENDAR_INTEGRATION }) - for (const group of groups.values()) { - try { - const userId = group[0].userId - const workspace = group[0].workspace - const token = tokens.find((p) => p.workspaceUuid === workspace && p.socialId === userId) - if (token === undefined) { - const toRemove = await this.getUserWatches(userId, workspace) - for (const key in toRemove) { - await client.deleteKey(key) - } - continue - } - const info = infos.find((p) => p.uuid === workspace) - if (info === undefined || !isActiveMode(info.mode)) { - const toRemove = await this.getUserWatches(userId, workspace) - for (const key in toRemove) { - await client.deleteKey(key) - } - continue - } - const watchClient = await WatchClient.Create(JSON.parse(token.secret)) - await watchClient.subscribe(group) - } catch {} + for (const group of groups) { + const secrets = await this.accountClient.listIntegrationsSecrets({ + kind: CALENDAR_INTEGRATION, + key: group[0] + }) + for (const val of group[1]) { + await client.deleteKey(`${CALENDAR_INTEGRATION}:watch:${group[0]}:${val.calendarId ?? 'null'}`) + } + for (const secret of secrets) { + try { + const watchClient = await WatchClient.Create(this.ctx, JSON.parse(secret.secret), this.accountClient) + await watchClient.subscribe(group[1]) + break + } catch {} + } } this.ctx.info('watch check done') } @@ -268,16 +269,16 @@ export class WatchController { googleClient: calendar_v3.Calendar ): Promise { const client = getKvsClient() - const key = `${CALENDAR_INTEGRATION}:watch:${user.workspace}:${user.userId}:${email}:${calendarId ?? 'null'}` + const key = `${CALENDAR_INTEGRATION}:watch:${email}:${calendarId ?? 'null'}` const exists = await client.getValue(key) if (exists != null) { return } try { if (calendarId != null) { - await watchCalendar(user, email, calendarId, googleClient) + await watchCalendar(email, calendarId, googleClient) } else { - await watchCalendars(user, email, googleClient) + await watchCalendars(email, googleClient) } } catch (err: any) { this.ctx.error('Watch add error', {