mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-28 04:25:03 +02:00
Improve calendar sync (#9400)
Signed-off-by: Denis Bykhov <bykhov.denis@gmail.com>
This commit is contained in:
@@ -59,30 +59,38 @@ export class CalendarController {
|
||||
groups.set(int.workspaceUuid, group)
|
||||
}
|
||||
}
|
||||
|
||||
const ids = [...groups.keys()]
|
||||
if (ids.length === 0) return
|
||||
const limiter = new RateLimiter(config.InitLimit)
|
||||
const infos = await this.accountClient.getWorkspacesInfo(ids)
|
||||
for (const info of infos) {
|
||||
const integrations = groups.get(info.uuid) ?? []
|
||||
if (await this.checkWorkspace(info, integrations)) {
|
||||
await limiter.add(async () => {
|
||||
try {
|
||||
this.ctx.info('start workspace', { workspace: info.uuid })
|
||||
await WorkspaceClient.run(this.ctx, this.accountClient, info.uuid)
|
||||
} catch (err) {
|
||||
this.ctx.error('Failed to start workspace', { workspace: info.uuid, error: err })
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
await limiter.waitProcessing()
|
||||
void this.runAll(groups)
|
||||
} catch (err: any) {
|
||||
this.ctx.error('Failed to start existing integrations', err)
|
||||
}
|
||||
}
|
||||
|
||||
private async runAll (groups: Map<WorkspaceUuid, Integration[]>): Promise<void> {
|
||||
const ids = [...groups.keys()]
|
||||
if (ids.length === 0) return
|
||||
const limiter = new RateLimiter(config.InitLimit)
|
||||
const infos = await this.accountClient.getWorkspacesInfo(ids)
|
||||
for (let index = 0; index < infos.length; index++) {
|
||||
const info = infos[index]
|
||||
const integrations = groups.get(info.uuid) ?? []
|
||||
if (await this.checkWorkspace(info, integrations)) {
|
||||
await limiter.add(async () => {
|
||||
try {
|
||||
this.ctx.info('start workspace', { workspace: info.uuid })
|
||||
await WorkspaceClient.run(this.ctx, this.accountClient, info.uuid)
|
||||
} catch (err) {
|
||||
this.ctx.error('Failed to start workspace', { workspace: info.uuid, error: err })
|
||||
}
|
||||
})
|
||||
}
|
||||
if (index % 10 === 0) {
|
||||
this.ctx.info('starting progress', { value: index + 1, total: infos.length })
|
||||
}
|
||||
}
|
||||
await limiter.waitProcessing()
|
||||
this.ctx.info('Started all workspaces', { count: infos.length })
|
||||
}
|
||||
|
||||
private async checkWorkspace (info: WorkspaceInfoWithStatus, integrations: Integration[]): Promise<boolean> {
|
||||
if (isDeletingMode(info.mode)) {
|
||||
if (integrations !== undefined) {
|
||||
|
||||
@@ -66,6 +66,7 @@ export const main = async (): Promise<void> => {
|
||||
|
||||
const calendarController = CalendarController.getCalendarController(ctx, accountClient)
|
||||
await calendarController.startAll()
|
||||
ctx.info('Calendar controller started')
|
||||
watchController.startCheck()
|
||||
const endpoints: Endpoint[] = [
|
||||
{
|
||||
|
||||
@@ -35,6 +35,7 @@ import {
|
||||
removeIntegrationSecret,
|
||||
setCredentials
|
||||
} from './utils'
|
||||
import { synced } from './sync'
|
||||
|
||||
export class OutcomingClient {
|
||||
private readonly calendar: calendar_v3.Calendar
|
||||
@@ -65,7 +66,9 @@ export class OutcomingClient {
|
||||
} else if (type === 'create') {
|
||||
await this.create(event, calendar)
|
||||
}
|
||||
await setSyncHistory(this.workspace, Date.now())
|
||||
if (synced.has(this.user.workspace)) {
|
||||
await setSyncHistory(this.workspace, Date.now())
|
||||
}
|
||||
}
|
||||
|
||||
private async convertBody (event: Event): Promise<calendar_v3.Schema$Event> {
|
||||
|
||||
@@ -35,7 +35,8 @@ import core, {
|
||||
Ref,
|
||||
SocialIdType,
|
||||
TxOperations,
|
||||
TxProcessor
|
||||
TxProcessor,
|
||||
WorkspaceUuid
|
||||
} from '@hcengineering/core'
|
||||
import setting from '@hcengineering/setting'
|
||||
import { htmlToMarkup } from '@hcengineering/text'
|
||||
@@ -60,6 +61,8 @@ import {
|
||||
} from './utils'
|
||||
import { WatchController } from './watch'
|
||||
|
||||
export const synced = new Set<WorkspaceUuid>()
|
||||
|
||||
const locks = new Map<string, Promise<void>>()
|
||||
|
||||
export async function lock (key: string): Promise<() => void> {
|
||||
@@ -269,6 +272,12 @@ export class IncomingSyncManager {
|
||||
showDeleted: syncToken != null
|
||||
})
|
||||
if (res.status === 410) {
|
||||
this.ctx.warn('Sync token is no longer valid, resyncing calendar', {
|
||||
workspace: this.user.workspace,
|
||||
user: this.user.userId,
|
||||
email: this.email,
|
||||
calendarId
|
||||
})
|
||||
await this.eventsSync(calendarId)
|
||||
return
|
||||
}
|
||||
@@ -290,6 +299,12 @@ export class IncomingSyncManager {
|
||||
} catch (err: any) {
|
||||
if (err?.response?.status === 410) {
|
||||
await this.eventsSync(calendarId)
|
||||
this.ctx.warn('Sync token is no longer valid, resyncing calendar', {
|
||||
workspace: this.user.workspace,
|
||||
user: this.user.userId,
|
||||
email: this.email,
|
||||
calendarId
|
||||
})
|
||||
return
|
||||
}
|
||||
this.ctx.error('Event sync error', { workspace: this.user.workspace, user: this.user.userId, err })
|
||||
|
||||
@@ -28,7 +28,7 @@ import { CalendarClient } from './calendar'
|
||||
import { getClient } from './client'
|
||||
import config from './config'
|
||||
import { addUserByEmail, getSyncHistory, setSyncHistory } from './kvsUtils'
|
||||
import { IncomingSyncManager } from './sync'
|
||||
import { IncomingSyncManager, synced } from './sync'
|
||||
import { getWorkspaceTokens } from './tokens'
|
||||
import { GoogleEmail, Token } from './types'
|
||||
import { getWorkspaceToken } from './utils'
|
||||
@@ -136,7 +136,10 @@ export class WorkspaceClient {
|
||||
private async getNewEvents (): Promise<void> {
|
||||
const lastSync = await getSyncHistory(this.workspace)
|
||||
this.lastSync = lastSync ?? 0
|
||||
const query = lastSync !== undefined ? { modifiedOn: { $gt: lastSync } } : {}
|
||||
const baseQuery = {
|
||||
calendar: { $in: Array.from(this.calendarsById.keys()) }
|
||||
}
|
||||
const query = lastSync !== undefined ? { modifiedOn: { $gt: lastSync }, ...baseQuery } : baseQuery
|
||||
const newEvents = await this.client.findAll(calendar.class.Event, query, { sort: { modifiedOn: 1 } })
|
||||
const interval = setInterval(() => {
|
||||
void this.updateSyncHistory()
|
||||
@@ -159,6 +162,8 @@ export class WorkspaceClient {
|
||||
}
|
||||
}
|
||||
clearInterval(interval)
|
||||
await setSyncHistory(this.workspace, Date.now())
|
||||
synced.add(this.workspace)
|
||||
this.ctx.info('all outcoming messages synced', { workspace: this.workspace })
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user