Calendar duplicate watch fix (#9383)

Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>
This commit is contained in:
Denis Bykhov
2025-06-25 23:44:13 +05:00
committed by GitHub
parent 4a876d6bb8
commit cde79b05ce
4 changed files with 84 additions and 83 deletions
+1 -1
View File
@@ -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<void> {
+7 -5
View File
@@ -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<string, Ref<Person>>()
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<Event> = 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,
@@ -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
+76 -75
View File
@@ -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<WatchClient> {
const watchClient = new WatchClient(token)
static async Create (ctx: MeasureContext, token: Token, accountClient: AccountClient): Promise<WatchClient> {
const watchClient = new WatchClient(ctx, token, accountClient)
await watchClient.setToken(token)
return watchClient
}
private async getWatches (): Promise<Record<string, Watch>> {
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<Watch>(key)
return watches ?? {}
}
@@ -48,10 +52,22 @@ export class WatchClient {
async checkError (err: any): Promise<void> {
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<void> {
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<void> {
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<void> {
async function watchCalendars (email: GoogleEmail, googleClient: calendar_v3.Calendar): Promise<void> {
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<Watch>(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<Watch>(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<Record<string, Watch>> {
private async getUserWatches (email: GoogleEmail): Promise<Record<string, Watch>> {
const client = getKvsClient()
const key = `${CALENDAR_INTEGRATION}:watch:${workspace}:${userId}`
const key = `${CALENDAR_INTEGRATION}:watch:${email}`
return (await client.listKeys<Watch>(key)) ?? {}
}
async unsubscribe (user: Token): Promise<void> {
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<string, WatchBase[]>()
const workspaces = new Set<WorkspaceUuid>()
const groups = new Map<GoogleEmail, Watch[]>()
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<void> {
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<Watch>(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', {