From 0c9c985ab81d2004bfd3abf56f1ca3da7ed51dbd Mon Sep 17 00:00:00 2001 From: Andrey Sobolev Date: Wed, 16 Apr 2025 20:50:58 +0700 Subject: [PATCH] UBERF-9521: Refactor session manager (#8560) Signed-off-by: Andrey Sobolev --- .vscode/launch.json | 18 +- common/config/rush/pnpm-lock.yaml | 11 +- dev/tool/src/workspace.ts | 2 +- packages/account-client/src/client.ts | 41 +- packages/account-client/src/types.ts | 30 +- packages/core/src/server.ts | 1 - pods/server/package.json | 1 + pods/server/src/__start.ts | 35 +- pods/server/src/__tests__/server.test.ts | 4 +- pods/server/src/server.ts | 13 +- pods/server/src/server_http.ts | 40 +- server/account/src/__tests__/utils.test.ts | 19 +- server/account/src/collections/mongo.ts | 65 +- server/account/src/collections/postgres.ts | 11 +- server/account/src/operations.ts | 122 ++- server/account/src/types.ts | 46 +- server/account/src/utils.ts | 104 ++- server/backup/src/service.ts | 4 +- server/core/src/stats.ts | 4 +- server/core/src/types.ts | 64 +- server/core/src/utils.ts | 3 - server/indexer/src/indexer/indexer.ts | 1 - server/middleware/src/triggers.ts | 1 - server/server-pipeline/src/pipeline.ts | 6 +- server/server/package.json | 3 +- server/server/src/client.ts | 49 +- server/server/src/sessionManager.ts | 800 ++++++++---------- server/server/src/stats.ts | 18 +- server/server/src/workspace.ts | 130 +++ server/workspace-service/src/ws-operations.ts | 1 - 30 files changed, 897 insertions(+), 750 deletions(-) create mode 100644 server/server/src/workspace.ts diff --git a/.vscode/launch.json b/.vscode/launch.json index 453b5bf8a7..3fd00f7070 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -97,7 +97,7 @@ "MODEL_JSON": "${workspaceRoot}/models/all/bundle/model.json", // "SERVER_PROVIDER":"uweb" "SERVER_PROVIDER": "ws", - "MODEL_VERSION": "0.7.48", + "MODEL_VERSION": "0.7.75", // "VERSION": "0.6.289", "ELASTIC_INDEX_NAME": "local_storage_index", "UPLOAD_URL": "/files", @@ -120,9 +120,9 @@ "args": ["src/__start.ts"], "env": { "FULLTEXT_URL": "http://localhost:4710", - "DB_URL": "mongodb://localhost:27018", + // "DB_URL": "mongodb://localhost:27018", // "DB_URL": "postgresql://postgres:example@localhost:5432", - // "DB_URL": "postgresql://root@huly.local:26258/defaultdb?sslmode=disable", + "DB_URL": "postgresql://root@huly.local:26258/defaultdb?sslmode=disable", // "GREEN_URL": "http://huly.local:6767?token=secret", "SERVER_PORT": "3334", "METRICS_CONSOLE": "false", @@ -135,8 +135,9 @@ "FRONT_URL": "http://localhost:8083", "ACCOUNTS_URL": "http://localhost:3003", "MODEL_JSON": "${workspaceRoot}/models/all/bundle/model.json", - "MODEL_VERSION": "0.7.1", - "STATS_URL": "http://huly.local:4901" + "MODEL_VERSION": "0.7.75", + "STATS_URL": "http://huly.local:4901", + "QUEUE_CONFIG": "localhost:19093" }, "runtimeArgs": ["--nolazy", "-r", "ts-node/register"], "runtimeVersion": "20", @@ -213,12 +214,13 @@ "request": "launch", "args": ["src/__start.ts"], "env": { - "MONGO_URL": "mongodb://localhost:27018", - "DB_URL": "mongodb://localhost:27018", + // "MONGO_URL": "mongodb://localhost:27018", + // "DB_URL": "mongodb://localhost:27018", // "DB_URL": "postgresql://postgres:example@localhost:5432", + "DB_URL": "postgresql://root@huly.local:26258/defaultdb?sslmode=disable", "SERVER_SECRET": "secret", "REGION_INFO": "|Mongo;pg|Postgres;cockroach|CockroachDB", - "TRANSACTOR_URL": "ws://huly.local:3334;;,ws://huly.local:3335;;europe", + "TRANSACTOR_URL": "ws://transactor:3334;ws://localhost:3334", "ACCOUNTS_URL": "http://localhost:3003", "ACCOUNT_PORT": "3003", "FRONT_URL": "http://localhost:8083", diff --git a/common/config/rush/pnpm-lock.yaml b/common/config/rush/pnpm-lock.yaml index d395f5843a..7479d88a11 100644 --- a/common/config/rush/pnpm-lock.yaml +++ b/common/config/rush/pnpm-lock.yaml @@ -774,7 +774,7 @@ importers: version: file:projects/pod-print.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))(bufferutil@4.0.8)(utf-8-validate@6.0.4) '@rush-temp/pod-server': specifier: file:./projects/pod-server.tgz - version: file:projects/pod-server.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))(utf-8-validate@6.0.4) + version: file:projects/pod-server.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9)) '@rush-temp/pod-ses': specifier: file:./projects/pod-ses.tgz version: file:projects/pod-ses.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9)) @@ -4806,7 +4806,7 @@ packages: version: 0.0.0 '@rush-temp/pod-server@file:projects/pod-server.tgz': - resolution: {integrity: sha512-nlXF3FUjnhvEM9Jz2qnS0pYDDJC4Akt9EZEmku7H+TFP6znCb/pBVelUlY5W+CSRBUFpRkmNv2e/KS/knh4b7w==, tarball: file:projects/pod-server.tgz} + resolution: {integrity: sha512-9Ff7F+voylF+V6WMOoCiftEMWTcPA4odM1CleTnf6cfA8Kg8vUz4Di2kzNcH/yxmiywvp47OGSsfuly+djtnkw==, tarball: file:projects/pod-server.tgz} version: 0.0.0 '@rush-temp/pod-ses@file:projects/pod-ses.tgz': @@ -5274,7 +5274,7 @@ packages: version: 0.0.0 '@rush-temp/server@file:projects/server.tgz': - resolution: {integrity: sha512-zYTPrh4adUZ36sz8IfoCcJ2yGCSbXgLaeYBv4++J+g/njPytlCfEN9L1x5aYT5zu91cZ/Gt4DkHGFgHsPwNjAw==, tarball: file:projects/server.tgz} + resolution: {integrity: sha512-ZeFqWmcBAIlHDDCSms07YvcEuy5G2LJ/Y4wAw743ykDDgqxQ7ptxJpQ0DATsSdOnAF916DE8TYktKbV9Hu1cug==, tarball: file:projects/server.tgz} version: 0.0.0 '@rush-temp/setting-assets@file:projects/setting-assets.tgz': @@ -22126,7 +22126,7 @@ snapshots: - supports-color - utf-8-validate - '@rush-temp/pod-server@file:projects/pod-server.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))(utf-8-validate@6.0.4)': + '@rush-temp/pod-server@file:projects/pod-server.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))': dependencies: '@hcengineering/communication-server': 0.1.170(typescript@5.3.3) '@types/body-parser': 1.19.5 @@ -22160,6 +22160,7 @@ snapshots: ts-jest: 29.1.2(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))(esbuild@0.24.2)(jest@29.7.0(@types/node@20.11.19)(ts-node@10.9.2(@types/node@20.11.19)(typescript@5.3.3)))(typescript@5.3.3) ts-node: 10.9.2(@types/node@20.11.19)(typescript@5.3.3) typescript: 5.3.3 + utf-8-validate: 6.0.4 ws: 8.18.0(bufferutil@4.0.8)(utf-8-validate@6.0.4) transitivePeerDependencies: - '@babel/core' @@ -22170,7 +22171,6 @@ snapshots: - babel-plugin-macros - node-notifier - supports-color - - utf-8-validate '@rush-temp/pod-ses@file:projects/pod-ses.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))': dependencies: @@ -25474,6 +25474,7 @@ snapshots: prettier: 3.2.5 ts-jest: 29.1.2(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))(esbuild@0.24.2)(jest@29.7.0(@types/node@20.11.19)(ts-node@10.9.2(@types/node@20.11.19)(typescript@5.3.3)))(typescript@5.3.3) typescript: 5.3.3 + utf-8-validate: 6.0.4 transitivePeerDependencies: - '@babel/core' - '@jest/types' diff --git a/dev/tool/src/workspace.ts b/dev/tool/src/workspace.ts index b568f72b24..4cfda47138 100644 --- a/dev/tool/src/workspace.ts +++ b/dev/tool/src/workspace.ts @@ -164,7 +164,7 @@ export async function backupRestore ( recheck: false, storageAdapter: workspaceStorage, getConnection: async () => { - return wrapPipeline(ctx, await pipelineFactory(ctx, wsUrl, true, () => {}, null, null), wsUrl) + return wrapPipeline(ctx, await pipelineFactory(ctx, wsUrl, () => {}, null, null), wsUrl) } }) ) diff --git a/packages/account-client/src/client.ts b/packages/account-client/src/client.ts index 027c71e5dd..d787b6173c 100644 --- a/packages/account-client/src/client.ts +++ b/packages/account-client/src/client.ts @@ -13,38 +13,39 @@ // limitations under the License. // import { - type AccountRole, type AccountInfo, + type AccountRole, + type AccountUuid, BackupStatus, + concatLink, Data, type Person, - type PersonUuid, + type PersonId, type PersonInfo, + type PersonUuid, + type SocialIdType, Version, type WorkspaceInfoWithStatus, type WorkspaceMemberInfo, WorkspaceMode, - concatLink, type WorkspaceUserOperation, - type WorkspaceUuid, - type PersonId, - type SocialIdType, - type AccountUuid + type WorkspaceUuid } from '@hcengineering/core' import platform, { PlatformError, Severity, Status } from '@hcengineering/platform' import type { - LoginInfo, - MailboxOptions, - OtpInfo, - WorkspaceLoginInfo, - RegionInfo, - WorkspaceOperation, - MailboxInfo, Integration, IntegrationKey, IntegrationSecret, IntegrationSecretKey, - SocialId + LoginInfo, + LoginInfoWithWorkspaces, + MailboxInfo, + MailboxOptions, + OtpInfo, + RegionInfo, + SocialId, + WorkspaceLoginInfo, + WorkspaceOperation } from './types' import { getClientTimezone } from './utils' @@ -63,6 +64,7 @@ export interface AccountClient { validateOtp: (email: string, code: string) => Promise loginOtp: (email: string) => Promise getLoginInfoByToken: () => Promise + getLoginWithWorkspaceInfo: () => Promise restorePassword: (password: string) => Promise confirm: () => Promise requestPasswordReset: (email: string) => Promise @@ -308,6 +310,15 @@ class AccountClientImpl implements AccountClient { return await this.rpc(request) } + async getLoginWithWorkspaceInfo (): Promise { + const request = { + method: 'getLoginWithWorkspaceInfo' as const, + params: {} + } + + return await this.rpc(request) + } + async restorePassword (password: string): Promise { const request = { method: 'restorePassword' as const, diff --git a/packages/account-client/src/types.ts b/packages/account-client/src/types.ts index 65b0d945da..e1848bf990 100644 --- a/packages/account-client/src/types.ts +++ b/packages/account-client/src/types.ts @@ -6,7 +6,8 @@ import { type AccountRole, type Timestamp, type SocialId as SocialIdBase, - PersonUuid + PersonUuid, + type WorkspaceMode } from '@hcengineering/core' export interface LoginInfo { @@ -16,6 +17,33 @@ export interface LoginInfo { token?: string } +export interface EndpointInfo { + internalUrl: string + externalUrl: string + region: string +} +export interface WorkspaceVersion { + versionMajor: number + versionMinor: number + versionPatch: number +} + +export interface LoginInfoWorkspace { + url: string + dataId?: WorkspaceDataId + mode: WorkspaceMode + version: WorkspaceVersion + endpoint: EndpointInfo + role: AccountRole | null + progress?: number +} + +export interface LoginInfoWithWorkspaces extends LoginInfo { + // Information necessary to handle user <--> transactor connectivity. + workspaces: Record + socialIds: SocialId[] +} + /** * @public */ diff --git a/packages/core/src/server.ts b/packages/core/src/server.ts index be3bcb5f49..96fd04ef13 100644 --- a/packages/core/src/server.ts +++ b/packages/core/src/server.ts @@ -50,7 +50,6 @@ export interface SessionData { admin?: boolean isTriggerCtx?: boolean workspace: WorkspaceIds - branding: Branding | null socialStringsToUsers: Map asyncRequests?: (() => Promise)[] diff --git a/pods/server/package.json b/pods/server/package.json index 555b04964a..3204af99f3 100644 --- a/pods/server/package.json +++ b/pods/server/package.json @@ -82,6 +82,7 @@ "@hcengineering/server-storage": "^0.6.0", "@hcengineering/server-telegram": "^0.6.0", "@hcengineering/server-token": "^0.6.11", + "utf-8-validate": "^6.0.4", "bufferutil": "^4.0.8", "msgpackr": "^1.11.2", "msgpackr-extract": "^3.0.3", diff --git a/pods/server/src/__start.ts b/pods/server/src/__start.ts index 183f211068..d4e2ea13c0 100644 --- a/pods/server/src/__start.ts +++ b/pods/server/src/__start.ts @@ -16,11 +16,7 @@ import serverCalendar from '@hcengineering/server-calendar' import serverCore, { initStatisticsContext, loadBrandingMap, - type ConnectionSocket, - type Session, type StorageConfiguration, - type UserStatistics, - type Workspace, type WorkspaceStatistics } from '@hcengineering/server-core' import serverNotification from '@hcengineering/server-notification' @@ -33,7 +29,7 @@ import { profileStart, profileStop } from './inspector' configureAnalytics(process.env.SENTRY_DSN, {}) Analytics.setTag('application', 'transactor') -let getUsers: () => WorkspaceStatistics[] = () => { +let getStats: () => WorkspaceStatistics[] = () => { return [] } @@ -50,8 +46,8 @@ void queue.createTopics(10).catch((err) => { // Force create server metrics context with proper logging const metricsContext = initStatisticsContext('transactor', { - getUsers: (): WorkspaceStatistics[] => { - return getUsers() + getStats: (): WorkspaceStatistics[] => { + return getStats() }, factory: () => new MeasureMetricsContext( @@ -103,29 +99,8 @@ const { shutdown, sessionManager } = start(metricsContext, config.dbUrl, { queue }) -const entryToUserStats = (session: Session, socket: ConnectionSocket): UserStatistics => { - return { - current: session.current, - mins5: session.mins5, - userId: session.getUser(), - sessionId: socket.id, - total: session.total, - data: socket.data - } -} - -const workspaceToWorkspaceStats = (ws: Workspace): WorkspaceStatistics => { - return { - clientsTotal: new Set(Array.from(ws.sessions.values()).map((it) => it.session.getUser())).size, - sessionsTotal: ws.sessions.size, - workspaceName: ws.workspaceName, - wsId: ws.workspaceUuid, - sessions: Array.from(ws.sessions.values()).map((it) => entryToUserStats(it.session, it.socket)) - } -} - -getUsers = () => { - return Array.from(sessionManager.workspaces.values()).map((it) => workspaceToWorkspaceStats(it)) +getStats = (): WorkspaceStatistics[] => { + return sessionManager.getStatistics() } const close = (): void => { diff --git a/pods/server/src/__tests__/server.test.ts b/pods/server/src/__tests__/server.test.ts index 908cc34e30..4271a1a375 100644 --- a/pods/server/src/__tests__/server.test.ts +++ b/pods/server/src/__tests__/server.test.ts @@ -40,7 +40,7 @@ import { type TxResult, type WorkspaceUuid } from '@hcengineering/core' -import { ClientSession, startSessionManager, type SessionManagerOptions } from '@hcengineering/server' +import { startSessionManager, type SessionManagerOptions } from '@hcengineering/server' import { createDummyQueue, createDummyStorageAdapter } from '@hcengineering/server-core' import { startHttpServer } from '../server_http' import { genMinModel } from './minmodel' @@ -100,7 +100,6 @@ describe('server', () => { communicationApiFactory: async () => { return {} as any }, - sessionFactory: (token, workspace, account) => new ClientSession(token, workspace, account, true), brandingMap: {}, accountsUrl: '', queue: createDummyQueue() @@ -216,7 +215,6 @@ describe('server', () => { communicationApiFactory: async () => { return {} as any }, - sessionFactory: (token, workspace, account) => new ClientSession(token, workspace, account, true), brandingMap: {}, accountsUrl: '', queue: createDummyQueue() diff --git a/pods/server/src/server.ts b/pods/server/src/server.ts index 7066490811..c50d3d2e68 100644 --- a/pods/server/src/server.ts +++ b/pods/server/src/server.ts @@ -14,19 +14,16 @@ // limitations under the License. // -import { type Account, type BrandingMap, type MeasureContext, type Tx } from '@hcengineering/core' +import { type BrandingMap, type MeasureContext, type Tx } from '@hcengineering/core' import { buildStorageFromConfig } from '@hcengineering/server-storage' -import { ClientSession, startSessionManager } from '@hcengineering/server' +import { startSessionManager } from '@hcengineering/server' import { type CommunicationApiFactory, type PlatformQueue, - type Session, type SessionManager, - type StorageConfiguration, - type Workspace + type StorageConfiguration } from '@hcengineering/server-core' -import { type Token } from '@hcengineering/server-token' import { Api as CommunicationApi } from '@hcengineering/communication-server' import { @@ -120,9 +117,6 @@ export function start ( { ...opt, externalStorage, adapterSecurity: isAdapterSecurity(dbUrl), queue: opt.queue }, {} ) - const sessionFactory = (token: Token, workspace: Workspace, account: Account): Session => { - return new ClientSession(token, workspace, account, token.extra?.mode === 'backup') - } const communicationApiFactory: CommunicationApiFactory = async (ctx, workspace, broadcastSessions) => { if (dbUrl.startsWith('mongodb')) { return { @@ -150,7 +144,6 @@ export function start ( const sessionManager = startSessionManager(metrics, { pipelineFactory, - sessionFactory, communicationApiFactory, brandingMap: opt.brandingMap, enableCompression: opt.enableCompression, diff --git a/pods/server/src/server_http.ts b/pods/server/src/server_http.ts index 298dca77f2..0705223178 100644 --- a/pods/server/src/server_http.ts +++ b/pods/server/src/server_http.ts @@ -89,7 +89,6 @@ const rpcHandler = new RPCHandler() const backpressureSize = 100 * 1024 /** * @public - * @param sessionFactory - * @param port - * @param host - */ @@ -403,27 +402,24 @@ export function startHttpServer ( try { const token = (req.query.token as string) ?? (req.headers.authorization ?? '').split(' ')[1] decodeToken(token) - const ws = sessions.workspaces.get(req.query.workspace as WorkspaceUuid) - if (ws !== undefined) { - // push the data to body - void retrieveJson(req) - .then((data) => { - if (Array.isArray(data)) { - sessions.broadcastAll(ws, data as Tx[]) - } else { - sessions.broadcastAll(ws, [data as unknown as Tx]) - } - res.end() - }) - .catch((err) => { - ctx.error('JSON parse error', { err }) - res.writeHead(400, {}) - res.end() - }) - } else { - res.writeHead(404, {}) - res.end() - } + + const ws = req.query.workspace as WorkspaceUuid + + // push the data to body + void retrieveJson(req) + .then((data) => { + if (Array.isArray(data)) { + sessions.broadcastAll(ws, data as Tx[]) + } else { + sessions.broadcastAll(ws, [data as unknown as Tx]) + } + res.end() + }) + .catch((err) => { + ctx.error('JSON parse error', { err }) + res.writeHead(400, {}) + res.end() + }) } catch (err: any) { Analytics.handleError(err) ctx.error('error', { err }) diff --git a/server/account/src/__tests__/utils.test.ts b/server/account/src/__tests__/utils.test.ts index ed6fe5b6c6..a1fb371c0b 100644 --- a/server/account/src/__tests__/utils.test.ts +++ b/server/account/src/__tests__/utils.test.ts @@ -347,10 +347,6 @@ describe('account utils', () => { }) describe('getEndpoint', () => { - const mockCtx = { - error: jest.fn() - } as unknown as MeasureContext - beforeEach(() => { jest.clearAllMocks() }) @@ -392,8 +388,7 @@ describe('account utils', () => { 'should handle workspace="%s" region="%s" kind=%s (%s)', (workspace, region, kind, transactors, expected, description) => { ;(getMetadata as jest.Mock).mockReturnValue(transactors) - expect(getEndpoint(mockCtx, workspace, region, kind)).toBe(expected) - expect(mockCtx.error).not.toHaveBeenCalled() + expect(getEndpoint(workspace as WorkspaceUuid, region, kind)).toBe(expected) } ) @@ -401,21 +396,17 @@ describe('account utils', () => { const transactors = 'http://internal:3000;http://external:3000;' ;(getMetadata as jest.Mock).mockReturnValue(transactors) - expect(getEndpoint(mockCtx, 'workspace1', 'nonexistent', EndpointKind.Internal)).toBe('http://internal:3000') - - expect(mockCtx.error).toHaveBeenCalledWith('No transactors for the target region, will use default region', { - group: 'nonexistent' - }) + expect(getEndpoint('workspace1' as WorkspaceUuid, 'nonexistent', EndpointKind.Internal)).toBe( + 'http://internal:3000' + ) }) test('should throw error when no transactors available', () => { ;(getMetadata as jest.Mock).mockReturnValue('http://internal:3000;http://external:3000;us') - expect(() => getEndpoint(mockCtx, 'workspace1', 'nonexistent', EndpointKind.Internal)).toThrow( + expect(() => getEndpoint('workspace1' as WorkspaceUuid, 'nonexistent', EndpointKind.Internal)).toThrow( 'Please provide transactor endpoint url' ) - - expect(mockCtx.error).toHaveBeenCalledWith('No transactors for the default region') }) }) diff --git a/server/account/src/collections/mongo.ts b/server/account/src/collections/mongo.ts index 99efaa14bc..49c2e2dd94 100644 --- a/server/account/src/collections/mongo.ts +++ b/server/account/src/collections/mongo.ts @@ -12,7 +12,17 @@ // See the License for the specific language governing permissions and // limitations under the License. // -import { UUID } from 'mongodb' +import { + type AccountRole, + type Data, + type Person, + type Version, + type WorkspaceMemberInfo, + type WorkspaceUuid, + AccountUuid, + buildSocialIdString, + SocialKey +} from '@hcengineering/core' import type { Collection, CreateIndexesOptions, @@ -22,38 +32,28 @@ import type { OptionalUnlessRequiredId, Sort as RawSort } from 'mongodb' -import { - type Person, - type WorkspaceMemberInfo, - buildSocialIdString, - SocialKey, - type AccountRole, - type Data, - type Version, - type WorkspaceUuid, - AccountUuid -} from '@hcengineering/core' +import { UUID } from 'mongodb' import type { - DbCollection, - Query, - Operations, - WorkspaceOperation, - AccountDB, Account, - SocialId, - WorkspaceInvite, - OTP, - WorkspaceStatus, + AccountDB, AccountEvent, - WorkspaceData, - WorkspaceInfoWithStatus, - WorkspaceStatusData, - Sort, + DbCollection, + Integration, + IntegrationSecret, Mailbox, MailboxSecret, - Integration, - IntegrationSecret + Operations, + OTP, + Query, + SocialId, + Sort, + WorkspaceData, + WorkspaceInfoWithStatus, + WorkspaceInvite, + WorkspaceOperation, + WorkspaceStatus, + WorkspaceStatusData } from '../types' import { isShallowEqual } from '../utils' @@ -716,6 +716,17 @@ export class MongoAccountDB implements AccountDB { return assignment?.role ?? null } + async getWorkspaceRoles (accountId: AccountUuid): Promise> { + const assignment = await this.workspaceMembers.find({ + accountUuid: accountId + }) + + return assignment.reduce>((acc, it) => { + acc.set(it.workspaceUuid, it.role) + return acc + }, new Map()) + } + async getWorkspaceMembers (workspaceId: WorkspaceUuid): Promise { return (await this.workspaceMembers.find({ workspaceUuid: workspaceId })).map((wmi) => ({ person: wmi.accountUuid, diff --git a/server/account/src/collections/postgres.ts b/server/account/src/collections/postgres.ts index bde9e0865c..68371d370a 100644 --- a/server/account/src/collections/postgres.ts +++ b/server/account/src/collections/postgres.ts @@ -490,8 +490,6 @@ export class PostgresAccountDB implements AccountDB { if (res.count === 1) { console.log(`Applying migration: ${name}`) await client.unsafe(ddl) - } else { - console.log(`Migration ${name} already applied`) } }) } @@ -535,12 +533,19 @@ export class PostgresAccountDB implements AccountDB { } async getWorkspaceRole (accountUuid: AccountUuid, workspaceUuid: WorkspaceUuid): Promise { - const res: any = await this + const res = await this .client`SELECT role FROM ${this.client(this.getWsMembersTableName())} WHERE workspace_uuid = ${workspaceUuid} AND account_uuid = ${accountUuid}` return res[0]?.role ?? null } + async getWorkspaceRoles (accountUuid: AccountUuid): Promise> { + const res = await this + .client`SELECT workspace_uuid, role FROM ${this.client(this.getWsMembersTableName())} WHERE account_uuid = ${accountUuid}` + + return new Map(res.map((it) => [it.workspace_uuid as WorkspaceUuid, it.role])) + } + async getWorkspaceMembers (workspaceUuid: WorkspaceUuid): Promise { const res: any = await this .client`SELECT account_uuid, role FROM ${this.client(this.getWsMembersTableName())} WHERE workspace_uuid = ${workspaceUuid}` diff --git a/server/account/src/operations.ts b/server/account/src/operations.ts index 84c9528bd5..d4e92ffb02 100644 --- a/server/account/src/operations.ts +++ b/server/account/src/operations.ts @@ -16,6 +16,7 @@ import { Analytics } from '@hcengineering/analytics' import { AccountInfo, AccountRole, + type AccountUuid, type Branding, buildSocialIdString, concatLink, @@ -28,18 +29,19 @@ import { SocialIdType, systemAccountUuid, type WorkspaceMemberInfo, - type WorkspaceUuid, - type AccountUuid + type WorkspaceUuid } from '@hcengineering/core' import platform, { getMetadata, PlatformError, Severity, Status, translate } from '@hcengineering/platform' import { decodeTokenVerbose, generateToken } from '@hcengineering/server-token' import { isAdminEmail } from './admin' import { accountPlugin } from './plugin' +import { type AccountServiceMethods, getServiceMethods } from './serviceOperations' import type { AccountDB, AccountMethodHandler, LoginInfo, + LoginInfoWithWorkspaces, Mailbox, MailboxOptions, Meta, @@ -55,6 +57,7 @@ import { checkInvite, cleanEmail, confirmEmail, + confirmHulyIds, createAccount, createWorkspaceRecord, doJoinByInvite, @@ -63,6 +66,7 @@ import { getAccount, getEmailSocialId, getEndpoint, + getEndpointInfo, getFrontUrl, getInviteEmail, getMailUrl, @@ -70,9 +74,12 @@ import { getRegions, getRolePower, getWorkspaceById, + getWorkspaceByUrl, + getWorkspaceEndpoint, getWorkspaceInfoWithStatusById, getWorkspaceInvite, getWorkspaceRole, + getWorkspaceRoles, GUEST_ACCOUNT, isEmail, isOtpValid, @@ -89,11 +96,8 @@ import { verifyAllowedRole, verifyAllowedServices, verifyPassword, - wrap, - getWorkspaceByUrl, - confirmHulyIds + wrap } from './utils' -import { type AccountServiceMethods, getServiceMethods } from './serviceOperations' // Note: it is IMPORTANT to always destructure params passed here to avoid sending extra params // to the database layer when searching/inserting as they may contain SQL injection @@ -383,7 +387,7 @@ export async function createWorkspace ( socialId: socialId._id, name: getPersonName(person), token: generateToken(account, workspaceUuid), - endpoint: getEndpoint(ctx, workspaceUuid, region, EndpointKind.External), + endpoint: getEndpoint(workspaceUuid, region, EndpointKind.External), workspace: workspaceUuid, workspaceUrl, role: AccountRole.Owner @@ -1326,7 +1330,7 @@ export async function getLoginInfoByToken ( return { ...loginInfo, workspace: workspaceUuid, - endpoint: getEndpoint(ctx, workspace.uuid, workspace.region, EndpointKind.External), + endpoint: getEndpoint(workspace.uuid, workspace.region, EndpointKind.External), role: AccountRole.DocGuest } } @@ -1342,7 +1346,7 @@ export async function getLoginInfoByToken ( ...loginInfo, workspace: workspace.uuid, workspaceDataId: workspace.dataId, - endpoint: getEndpoint(ctx, workspace.uuid, workspace.region, EndpointKind.External), + endpoint: getEndpoint(workspace.uuid, workspace.region, EndpointKind.External), role } } else { @@ -1350,6 +1354,101 @@ export async function getLoginInfoByToken ( } } +/** + * Validates the token and returns the decoded account information. + */ +export async function getLoginWithWorkspaceInfo ( + ctx: MeasureContext, + db: AccountDB, + branding: Branding | null, + token: string +): Promise { + let accountUuid: AccountUuid + let extra: any + try { + ;({ account: accountUuid, extra } = decodeTokenVerbose(ctx, token)) + } catch (err: any) { + Analytics.handleError(err) + ctx.error('Invalid token', { token }) + throw new PlatformError(new Status(Severity.ERROR, platform.status.Unauthorized, {})) + } + + if (accountUuid == null) { + throw new PlatformError(new Status(Severity.ERROR, platform.status.AccountNotFound, { account: accountUuid })) + } + + const isDocGuest = accountUuid === GUEST_ACCOUNT && extra?.guest === 'true' + const isSystem = accountUuid === systemAccountUuid + let socialIds: SocialId[] = [] + + if (!isDocGuest && !isSystem) { + // Any confirmed social ID will do + socialIds = await db.socialId.find({ personUuid: accountUuid, verifiedOn: { $gt: 0 } }) + if (socialIds.length === 0) { + return { + account: accountUuid, + workspaces: {}, + socialIds: [] + } + } + } + + let person: Person | null + if (isDocGuest) { + person = { + uuid: accountUuid, + firstName: 'Guest', + lastName: 'User' + } + } else if (isSystem) { + person = { + uuid: accountUuid, + firstName: 'System', + lastName: 'User' + } + } else { + person = await db.person.findOne({ uuid: accountUuid }) + } + + if (person == null) { + throw new PlatformError(new Status(Severity.ERROR, platform.status.InternalServerError, {})) + } + + const userWorkspaces = (await db.getAccountWorkspaces(accountUuid)).filter((it) => isActiveMode(it.status.mode)) + const roles: Map = await getWorkspaceRoles(db, accountUuid) + + const info = getEndpointInfo() + const loginInfo: LoginInfoWithWorkspaces = { + account: accountUuid, + name: getPersonName(person), + socialId: socialIds[0]?._id, + token, + workspaces: Object.fromEntries( + isSystem || isDocGuest + ? [] + : userWorkspaces.map((it, idx) => [ + it.uuid, + { + url: it.url, + dataId: it.dataId, + mode: it.status.mode, + endpoint: getWorkspaceEndpoint(info, it.uuid, it.region), + role: roles.get(it.uuid) ?? null, + version: { + versionMajor: it.status.versionMajor, + versionMinor: it.status.versionMinor, + versionPatch: it.status.versionPatch + }, + progress: it.status.processingProgress + } + ]) + ), + socialIds + } + + return loginInfo +} + export async function getSocialIds ( ctx: MeasureContext, db: AccountDB, @@ -1650,7 +1749,6 @@ async function exchangeGuestToken ( if (tokenObj.account == null) { // Check if it's old guest token const oldGuestEmail = '#guest@hc.engineering' - const guestAccount = 'b6996120-416f-49cd-841e-e4a5d2e49c9b' as PersonUuid const { linkId, guest, email, workspace: workspaceUrl } = tokenObj as any if (linkId == null || guest == null || email !== oldGuestEmail || workspaceUrl == null) { @@ -1663,7 +1761,7 @@ async function exchangeGuestToken ( throw new PlatformError(new Status(Severity.ERROR, platform.status.WorkspaceNotFound, { workspaceUrl })) } - return generateToken(guestAccount, workspace.uuid, { linkId, guest: 'true' }) + return generateToken(GUEST_ACCOUNT as PersonUuid, workspace.uuid, { linkId, guest: 'true' }) } return token @@ -1699,6 +1797,7 @@ export type AccountMethods = | 'getWorkspaceInfo' | 'getWorkspacesInfo' | 'getLoginInfoByToken' + | 'getLoginWithWorkspaceInfo' | 'getSocialIds' | 'getPerson' | 'getWorkspaceMembers' @@ -1759,6 +1858,7 @@ export function getMethods (hasSignUp: boolean = true): Partial Promise unassignWorkspace: (accountId: AccountUuid, workspaceId: WorkspaceUuid) => Promise getWorkspaceRole: (accountId: AccountUuid, workspaceId: WorkspaceUuid) => Promise + getWorkspaceRoles: (accountId: AccountUuid) => Promise> getWorkspaceMembers: (workspaceId: WorkspaceUuid) => Promise getAccountWorkspaces: (accountId: AccountUuid) => Promise getPendingWorkspace: ( @@ -276,6 +281,23 @@ export interface LoginInfo { token?: string } +export interface LoginInfoWorkspace { + url: string + dataId?: WorkspaceDataId + mode: WorkspaceMode + version: WorkspaceVersion + endpoint: EndpointInfo + role: AccountRole | null + + progress?: number +} + +export interface LoginInfoWithWorkspaces extends LoginInfo { + // Information necessary to handle user <--> transactor connectivity. + workspaces: Record + socialIds: SocialId[] +} + export interface WorkspaceLoginInfo extends LoginInfo { workspace: WorkspaceUuid workspaceUrl: string diff --git a/server/account/src/utils.ts b/server/account/src/utils.ts index 5f79395b96..1a2fede24d 100644 --- a/server/account/src/utils.ts +++ b/server/account/src/utils.ts @@ -13,54 +13,54 @@ // limitations under the License. // import { + AccountRole, + AccountUuid, Branding, concatLink, generateId, groupByArray, + isActiveMode, MeasureContext, - AccountRole, roleOrder, SocialIdType, - WorkspaceUuid, - WorkspaceMode, SocialKey, systemAccountUuid, - type WorkspaceInfoWithStatus as WorkspaceInfoWithStatusCore, - isActiveMode, - type PersonUuid, - type PersonId, + WorkspaceMode, + WorkspaceUuid, type Person, - AccountUuid + type PersonId, + type PersonUuid, + type WorkspaceInfoWithStatus as WorkspaceInfoWithStatusCore } from '@hcengineering/core' import { getMongoClient } from '@hcengineering/mongo' // TODO: get rid of this import later import platform, { getMetadata, PlatformError, Severity, Status, translate } from '@hcengineering/platform' import { getDBClient } from '@hcengineering/postgres' -import otpGenerator from 'otp-generator' import { pbkdf2Sync, randomBytes } from 'crypto' +import otpGenerator from 'otp-generator' +import { Analytics } from '@hcengineering/analytics' +import { sharedPipelineContextVars } from '@hcengineering/server-pipeline' +import { decodeTokenVerbose, generateToken, TokenError } from '@hcengineering/server-token' import { MongoAccountDB } from './collections/mongo' import { PostgresAccountDB } from './collections/postgres' import { accountPlugin } from './plugin' -import { sharedPipelineContextVars } from '@hcengineering/server-pipeline' import { + AccountEventType, AccountMethodHandler, + Integration, + LoginInfo, + Meta, OtpInfo, - WorkspaceInvite, WorkspaceInfoWithStatus, + WorkspaceInvite, + WorkspaceLoginInfo, + WorkspaceStatus, type Account, type AccountDB, type RegionInfo, type SocialId, - type Workspace, - LoginInfo, - WorkspaceLoginInfo, - WorkspaceStatus, - AccountEventType, - Meta, - Integration + type Workspace } from './types' -import { Analytics } from '@hcengineering/analytics' -import { TokenError, decodeTokenVerbose, generateToken } from '@hcengineering/server-token' export const GUEST_ACCOUNT = 'b6996120-416f-49cd-841e-e4a5d2e49c9b' @@ -172,7 +172,7 @@ export enum EndpointKind { External } -const toTransactor = (line: string): { internalUrl: string, region: string, externalUrl: string } => { +const toTransactor = (line: string): EndpointInfo => { const [internalUrl, externalUrl, region] = line .split(';') .map((it) => it.trim()) @@ -236,33 +236,46 @@ export const _getRegions = (): RegionInfo[] => { return _regionInfo } -export const getEndpoint = ( - ctx: MeasureContext, - workspace: string, - region: string | undefined, - kind: EndpointKind -): string => { - const byRegions = groupByArray(getEndpoints().map(toTransactor), (it) => it.region) - let transactors = (byRegions.get(region ?? '') ?? []) - .map((it) => (kind === EndpointKind.Internal ? it.internalUrl : it.externalUrl)) - .flat() +export interface EndpointInfo { + internalUrl: string + externalUrl: string + region: string +} + +export function getEndpointInfo (): Map { + return groupByArray(getEndpoints().map(toTransactor), (it) => it.region) +} + +export const selectKind = (kind: EndpointKind, it: EndpointInfo): string => { + return kind === EndpointKind.Internal ? it.internalUrl : it.externalUrl +} + +export const getEndpoint = (workspace: WorkspaceUuid, region: string | undefined, kind: EndpointKind): string => { + const hash = hashWorkspace(workspace) + const _endpointInfo = getEndpointInfo() + + let transactors = _endpointInfo.get(region ?? '') ?? [] - // This is really bad if (transactors.length === 0) { - ctx.error('No transactors for the target region, will use default region', { group: region }) - - transactors = (byRegions.get('') ?? []) - .map((it) => (kind === EndpointKind.Internal ? it.internalUrl : it.externalUrl)) - .flat() + console.warn('No transactors for the target region, will use default region', { group: region }) + transactors = _endpointInfo.get('') ?? [] } if (transactors.length === 0) { - ctx.error('No transactors for the default region') throw new Error('Please provide transactor endpoint url') } + return selectKind(kind, transactors[Math.abs(hash % transactors.length)]) +} + +export const getWorkspaceEndpoint = ( + info: Map, + workspace: WorkspaceUuid, + region: string | undefined +): EndpointInfo => { const hash = hashWorkspace(workspace) - return transactors[Math.abs(hash % transactors.length)] + const byRegion = info.get(region ?? '') ?? [] + return byRegion[Math.abs(hash % byRegion.length)] } export function getAllTransactors (kind: EndpointKind): string[] { @@ -541,7 +554,7 @@ export async function selectWorkspace ( // Guest mode select workspace return { account: accountUuid, - endpoint: getEndpoint(ctx, workspace.uuid, workspace.region, getKind(workspace.region)), + endpoint: getEndpoint(workspace.uuid, workspace.region, getKind(workspace.region)), token, workspace: workspace.uuid, workspaceUrl: workspace.url, @@ -572,7 +585,7 @@ export async function selectWorkspace ( return { account: accountUuid, token: generateToken(accountUuid, workspace.uuid, extra), - endpoint: getEndpoint(ctx, workspace.uuid, workspace.region, getKind(workspace.region)), + endpoint: getEndpoint(workspace.uuid, workspace.region, getKind(workspace.region)), workspace: workspace.uuid, workspaceUrl: workspace.url, role: AccountRole.Owner @@ -604,7 +617,7 @@ export async function selectWorkspace ( return { account: accountUuid, token: generateToken(accountUuid, workspace.uuid, extra), - endpoint: getEndpoint(ctx, workspace.uuid, workspace.region, getKind(workspace.region)), + endpoint: getEndpoint(workspace.uuid, workspace.region, getKind(workspace.region)), workspace: workspace.uuid, workspaceUrl: workspace.url, workspaceDataId: workspace.dataId, @@ -1413,6 +1426,13 @@ export async function getWorkspaceRole ( return await db.getWorkspaceRole(account, workspace) } +export async function getWorkspaceRoles ( + db: AccountDB, + account: AccountUuid +): Promise> { + return await db.getWorkspaceRoles(account) +} + export function generatePassword (len: number = 24): string { return randomBytes(len).toString('base64').slice(0, len) } diff --git a/server/backup/src/service.ts b/server/backup/src/service.ts index 4027f5b3be..eef12fcb80 100644 --- a/server/backup/src/service.ts +++ b/server/backup/src/service.ts @@ -332,7 +332,7 @@ class BackupWorker { }, getConnection: async () => { if (pipeline === undefined) { - pipeline = await this.pipelineFactory(ctx, wsIds, true, () => {}, null, null) + pipeline = await this.pipelineFactory(ctx, wsIds, () => {}, null, null) } return wrapPipeline(ctx, pipeline, wsIds) }, @@ -484,7 +484,7 @@ export async function doRestoreWorkspace ( cleanIndexState, getConnection: async () => { if (pipeline === undefined) { - pipeline = await pipelineFactory(ctx, wsIds, true, () => {}, null, null) + pipeline = await pipelineFactory(ctx, wsIds, () => {}, null, null) } return wrapPipeline(ctx, pipeline, wsIds) }, diff --git a/server/core/src/stats.ts b/server/core/src/stats.ts index 99a026f2ff..cdcfc072dd 100644 --- a/server/core/src/stats.ts +++ b/server/core/src/stats.ts @@ -81,7 +81,7 @@ export function initStatisticsContext ( logFile?: string logConsole?: boolean factory?: () => MeasureMetricsContext - getUsers?: () => WorkspaceStatistics[] + getStats?: () => WorkspaceStatistics[] statsUrl?: string serviceName?: () => string } @@ -144,7 +144,7 @@ export function initStatisticsContext ( cpu: getCPUInfo(), memory: getMemoryInfo(), stats: metricsContext.metrics, - workspaces: ops?.getUsers?.() + workspaces: ops?.getStats?.() } const statData = JSON.stringify(data) diff --git a/server/core/src/types.ts b/server/core/src/types.ts index 1cce904f89..f596f01e8c 100644 --- a/server/core/src/types.ts +++ b/server/core/src/types.ts @@ -46,12 +46,12 @@ import { type SearchQuery, type SearchResult, type SessionData, + type SocialId, type Space, type Timestamp, type Tx, type TxFactory, type TxResult, - type WorkspaceDataId, type WorkspaceIds, type WorkspaceUuid } from '@hcengineering/core' @@ -62,7 +62,7 @@ import type { Token } from '@hcengineering/server-token' import { type Readable } from 'stream' import type { DbAdapter, DomainHelper } from './adapter' -import type { StatisticsElement } from './stats' +import type { StatisticsElement, WorkspaceStatistics } from './stats' import { type StorageAdapter } from './storage' import { type PlatformQueue } from './queue' @@ -239,7 +239,6 @@ export interface Pipeline { export type PipelineFactory = ( ctx: MeasureContext, ws: WorkspaceIds, - upgrade: boolean, broadcast: BroadcastFunc, branding: Branding | null, communicationApi: CommunicationApi | null @@ -559,7 +558,7 @@ export interface ClientSessionCtx { * @public */ export interface Session { - workspace: Workspace + workspace: WorkspaceIds createTime: number // Session restore information @@ -587,8 +586,11 @@ export interface Session { // Client methods ping: (ctx: ClientSessionCtx) => Promise getUser: () => AccountUuid + getUserSocialIds: () => PersonId[] + getSocialIds: () => SocialId[] + loadModel: (ctx: ClientSessionCtx, lastModelTx: Timestamp, hash?: string) => Promise loadModelRaw: (ctx: ClientSessionCtx, lastModelTx: Timestamp, hash?: string) => Promise getRawAccount: () => Account @@ -663,61 +665,25 @@ export function disableLogging (): void { LOGGING_ENABLED = false } -interface TickHandler { - ticks: number - operation: () => void -} - -/** - * @public - */ -export interface Workspace { - context: MeasureContext - id: string - token: string // Account workspace update token. - pipeline: Promise | Pipeline - communicationApi: Promise | CommunicationApi - - tickHash: number - - tickHandlers: Map - - sessions: Map - upgrade: boolean - - closing?: Promise - softShutdown: number - workspaceInitCompleted: boolean - - workspaceName: string - workspaceUuid: WorkspaceUuid - workspaceUrl: string - workspaceDataId?: WorkspaceDataId - branding: Branding | null -} - export interface AddSessionActive { session: Session context: MeasureContext workspaceId: WorkspaceUuid } -export type AddSessionResponse = - | AddSessionActive +export type GetWorkspaceResponse = | { upgrade: true, progress?: number } | { error: any, terminate?: boolean, specialError?: 'archived' | 'migration' } -export type SessionFactory = (token: Token, workspace: Workspace, account: Account) => Session +export type AddSessionResponse = AddSessionActive | GetWorkspaceResponse /** * @public */ export interface SessionManager { - workspaces: Map + // workspaces: Map sessions: Map - createSession: SessionFactory - addSession: ( ctx: MeasureContext, ws: ConnectionSocket, @@ -726,18 +692,10 @@ export interface SessionManager { sessionId: string | undefined ) => Promise - broadcastAll: (workspace: Workspace, tx: Tx[], targets?: string[]) => void + broadcastAll: (workspace: WorkspaceUuid, tx: Tx[], targets?: string[]) => void close: (ctx: MeasureContext, ws: ConnectionSocket, workspaceId: WorkspaceUuid) => Promise - closeAll: ( - wsId: WorkspaceUuid, - workspace: Workspace, - code: number, - reason: 'upgrade' | 'shutdown', - ignoreSocket?: ConnectionSocket - ) => Promise - forceClose: (wsId: WorkspaceUuid, ignoreSocket?: ConnectionSocket) => Promise closeWorkspaces: (ctx: MeasureContext) => Promise @@ -773,6 +731,8 @@ export interface SessionManager { service: Session, ws: ConnectionSocket ) => ClientSessionCtx + + getStatistics: () => WorkspaceStatistics[] } export const pingConst = 'ping' diff --git a/server/core/src/utils.ts b/server/core/src/utils.ts index 0e8f36743b..696f973e23 100644 --- a/server/core/src/utils.ts +++ b/server/core/src/utils.ts @@ -7,7 +7,6 @@ import core, { type Account, type AccountUuid, type BackupClient, - type Branding, type BrandingMap, type BulkUpdateEvent, type Class, @@ -166,7 +165,6 @@ export class SessionDataImpl implements SessionData { readonly admin: boolean | undefined, _broadcast: SessionData['broadcast'] | undefined, readonly workspace: WorkspaceIds, - readonly branding: Branding | null, readonly isAsyncContext: boolean, _removedMap: Map, Doc> | undefined, _contextCache: Map | undefined, @@ -244,7 +242,6 @@ export function wrapPipeline ( true, { targets: {}, txes: [] }, wsIds, - null, true, undefined, undefined, diff --git a/server/indexer/src/indexer/indexer.ts b/server/indexer/src/indexer/indexer.ts index 5caff44c5e..a8bf946eae 100644 --- a/server/indexer/src/indexer/indexer.ts +++ b/server/indexer/src/indexer/indexer.ts @@ -559,7 +559,6 @@ export class FullTextIndexPipeline implements FullTextPipeline { true, undefined, this.workspace, - null, false, undefined, undefined, diff --git a/server/middleware/src/triggers.ts b/server/middleware/src/triggers.ts index 918a435eed..65815847ae 100644 --- a/server/middleware/src/triggers.ts +++ b/server/middleware/src/triggers.ts @@ -221,7 +221,6 @@ export class TriggersMiddleware extends BaseMiddleware implements Middleware { sctx.admin, { txes: [], targets: {} }, this.context.workspace, - this.context.branding, true, sctx.removedMap, sctx.contextCache, diff --git a/server/server-pipeline/src/pipeline.ts b/server/server-pipeline/src/pipeline.ts index a0f139c99d..050c80b075 100644 --- a/server/server-pipeline/src/pipeline.ts +++ b/server/server-pipeline/src/pipeline.ts @@ -110,7 +110,7 @@ export function createServerPipeline ( }, extensions?: Partial ): PipelineFactory { - return (ctx, workspace, upgrade, broadcast, branding, communicationApi) => { + return (ctx, workspace, broadcast, branding, communicationApi) => { const metricsCtx = opt.usePassedCtx === true ? ctx : metrics const wsMetrics = metricsCtx.newChild('๐Ÿงฒ session', {}) const conf = getConfig(metrics, dbUrl, wsMetrics, opt, extensions) @@ -176,7 +176,7 @@ export function createBackupPipeline ( externalStorage: StorageAdapter } ): PipelineFactory { - return (ctx, workspace, upgrade, broadcast, branding, communicationApi) => { + return (ctx, workspace, broadcast, branding, communicationApi) => { const metricsCtx = opt.usePassedCtx === true ? ctx : metrics const wsMetrics = metricsCtx.newChild('๐Ÿงฒ backup', {}) const conf = getConfig(metrics, dbUrl, wsMetrics, { @@ -229,7 +229,7 @@ export async function getServerPipeline ( }) // TODO: Communication API ?? - return await pipelineFactory(ctx, wsUrl, true, () => {}, null, null) + return await pipelineFactory(ctx, wsUrl, () => {}, null, null) } const txAdapterFactories: Record = {} diff --git a/server/server/package.json b/server/server/package.json index bb6b73f48e..6c04a197d8 100644 --- a/server/server/package.json +++ b/server/server/package.json @@ -47,6 +47,7 @@ "@hcengineering/platform": "^0.6.11", "@hcengineering/rpc": "^0.6.5", "@hcengineering/server-core": "^0.6.1", - "@hcengineering/server-token": "^0.6.11" + "@hcengineering/server-token": "^0.6.11", + "utf-8-validate": "^6.0.4" } } diff --git a/server/server/src/client.ts b/server/server/src/client.ts index ed9ee6460c..678e02760f 100644 --- a/server/server/src/client.ts +++ b/server/server/src/client.ts @@ -13,6 +13,21 @@ // limitations under the License. // +import type { LoginInfoWithWorkspaces } from '@hcengineering/account-client' +import { + RequestEvent as CommunicationEvent, + SessionData as CommunicationSession, + EventResult +} from '@hcengineering/communication-sdk-types' +import { + FindLabelsParams, + FindMessagesGroupsParams, + FindMessagesParams, + FindNotificationContextParams, + FindNotificationsParams, + Message, + MessagesGroup +} from '@hcengineering/communication-types' import { AccountUuid, generateId, @@ -32,11 +47,13 @@ import { type SearchQuery, type SearchResult, type SessionData, + type SocialId, type Timestamp, type Tx, type TxCUD, type TxResult, - type WorkspaceDataId + type WorkspaceDataId, + type WorkspaceIds } from '@hcengineering/core' import { PlatformError, unknownError } from '@hcengineering/platform' import { @@ -48,24 +65,9 @@ import { type Pipeline, type Session, type SessionRequest, - type StatisticsElement, - type Workspace + type StatisticsElement } from '@hcengineering/server-core' import { type Token } from '@hcengineering/server-token' -import { - FindMessagesGroupsParams, - FindMessagesParams, - Message, - MessagesGroup, - FindNotificationContextParams, - FindNotificationsParams, - FindLabelsParams -} from '@hcengineering/communication-types' -import { - RequestEvent as CommunicationEvent, - SessionData as CommunicationSession, - EventResult -} from '@hcengineering/communication-sdk-types' const useReserveContext = (process.env.USE_RESERVE_CTX ?? 'true') === 'true' @@ -93,8 +95,9 @@ export class ClientSession implements Session { constructor ( protected readonly token: Token, - readonly workspace: Workspace, + readonly workspace: WorkspaceIds, readonly account: Account, + readonly info: LoginInfoWithWorkspaces, readonly allowUpload: boolean ) { this.isAdmin = this.token.extra?.admin === 'true' @@ -108,6 +111,10 @@ export class ClientSession implements Session { return this.account.socialIds } + getSocialIds (): SocialId[] { + return this.info.socialIds + } + getRawAccount (): Account { return this.account } @@ -142,18 +149,16 @@ export class ClientSession implements Session { } includeSessionContext (ctx: ClientSessionCtx): void { - const dataId = this.workspace.workspaceDataId ?? (this.workspace.workspaceUuid as unknown as WorkspaceDataId) + const dataId = this.workspace.dataId ?? (this.workspace.uuid as unknown as WorkspaceDataId) const contextData = new SessionDataImpl( this.account, this.sessionId, this.isAdmin, undefined, { - uuid: this.workspace.workspaceUuid, - url: this.workspace.workspaceUrl, + ...this.workspace, dataId }, - this.workspace.branding, false, undefined, undefined, diff --git a/server/server/src/sessionManager.ts b/server/server/src/sessionManager.ts index a0608e3269..35e1f7c8a3 100644 --- a/server/server/src/sessionManager.ts +++ b/server/server/src/sessionManager.ts @@ -13,7 +13,11 @@ // limitations under the License. // -import { getClient as getAccountClient, isWorkspaceLoginInfo } from '@hcengineering/account-client' +import { + getClient as getAccountClient, + type LoginInfoWithWorkspaces, + type LoginInfoWorkspace +} from '@hcengineering/account-client' import { Analytics } from '@hcengineering/analytics' import { type ServerApi as CommunicationApi } from '@hcengineering/communication-sdk-types' import core, { @@ -28,13 +32,13 @@ import core, { pickPrimarySocialId, platformNow, platformNowDiff, + SocialIdType, systemAccountUuid, TxFactory, Version, versionToString, withContext, WorkspaceEvent, - type Account, type AccountUuid, type Branding, type BrandingMap, @@ -62,33 +66,25 @@ import { type AddSessionResponse, type ClientSessionCtx, type ConnectionSocket, + type GetWorkspaceResponse, type PlatformQueue, type PlatformQueueProducer, type QueueWorkspaceMessage, type Session, - type SessionFactory, - type Workspace + type UserStatistics, + type WorkspaceStatistics } from '@hcengineering/server-core' import { generateToken, type Token } from '@hcengineering/server-token' +import { Workspace, type PipelinePair } from './workspace' import { WorkspaceIds } from '@hcengineering/core' +import { ClientSession } from './client' import { sendResponse } from './utils' const ticksPerSecond = 20 const workspaceSoftShutdownTicks = 15 * ticksPerSecond -const guestAccount = 'b6996120-416f-49cd-841e-e4a5d2e49c9b' -function timeoutPromise (time: number): { promise: Promise, cancelHandle: () => void } { - let timer: any - return { - promise: new Promise((resolve) => { - timer = setTimeout(resolve, time) - }), - cancelHandle: () => { - clearTimeout(timer) - } - } -} +const guestAccount = 'b6996120-416f-49cd-841e-e4a5d2e49c9b' /** * @public @@ -104,7 +100,7 @@ export class TSessionManager implements SessionManager { readonly workspaces = new Map() checkInterval: any - sessions = new Map() + sessions = new Map() reconnectIds = new Set() maintenanceTimer: any @@ -113,14 +109,10 @@ export class TSessionManager implements SessionManager { modelVersion = process.env.MODEL_VERSION ?? '' serverVersion = process.env.VERSION ?? '' - oldClientErrors: number = 0 - clientErrors: number = 0 - lastClients: string[] = [] workspaceProducer: PlatformQueueProducer usersProducer: PlatformQueueProducer constructor ( readonly ctx: MeasureContext, - readonly sessionFactory: SessionFactory, readonly timeouts: Timeouts, readonly brandingMap: BrandingMap, readonly profiling: @@ -174,7 +166,7 @@ export class TSessionManager implements SessionManager { } const event: TxWorkspaceEvent = this.createMaintenanceWarning() for (const ws of this.workspaces.values()) { - this.broadcastAll(ws, [event]) + this.doBroadcast(ws, [event]) } } @@ -198,6 +190,13 @@ export class TSessionManager implements SessionManager { handleTick (): void { const now = Date.now() + this.handleWorkspaceTick() + + this.handleSessionTick(now) + this.ticks++ + } + + private handleWorkspaceTick (): void { for (const [wsId, workspace] of this.workspaces.entries()) { if (this.ticks % (60 * ticksPerSecond) === workspace.tickHash) { try { @@ -236,54 +235,6 @@ export class TSessionManager implements SessionManager { s[1].session.current = { find: 0, tx: 0 } } - const lastRequestDiff = now - s[1].session.lastRequest - - let timeout = 60000 - if (s[1].session.getUser() === systemAccountUuid) { - timeout = timeout * 10 - } - - const isCurrentUserTick = this.ticks % ticksPerSecond === s[1].tickHash - - if (isCurrentUserTick) { - if (lastRequestDiff > timeout) { - this.ctx.warn('session hang, closing...', { wsId, user: s[1].session.getUser() }) - - // Force close workspace if only one client and it hang. - void this.close(this.ctx, s[1].socket, wsId).catch((err) => { - this.ctx.error('failed to close', err) - }) - continue - } - if ( - lastRequestDiff + (1 / 10) * lastRequestDiff > this.timeouts.pingTimeout && - now - s[1].session.lastPing > this.timeouts.pingTimeout - ) { - // We need to check state and close socket if it broken - // And ping other wize - s[1].session.lastPing = now - if (s[1].socket.checkState()) { - void s[1].socket.send( - workspace.context, - { result: pingConst }, - s[1].session.binaryMode, - s[1].session.useCompression - ) - } - } - for (const r of s[1].session.requests.values()) { - const sec = Math.round((now - r.start) / 1000) - if (sec > 0 && sec % 30 === 0) { - this.ctx.warn('request hang found', { - sec, - wsId, - total: s[1].session.requests.size, - user: s[1].session.getUser(), - ...cutObjectArray(r.params) - }) - } - } - } } // Wait some time for new client to appear before closing workspace. @@ -291,30 +242,95 @@ export class TSessionManager implements SessionManager { workspace.softShutdown-- if (workspace.softShutdown <= 0) { this.ctx.warn('closing workspace, no users', { - workspace: workspace.workspaceUuid, + workspace: workspace.wsId.url, wsId, upgrade: workspace.upgrade }) - workspace.closing = this.performWorkspaceCloseCheck(workspace, wsId) + workspace.closing = this.performWorkspaceCloseCheck(workspace) } } else { workspace.softShutdown = workspaceSoftShutdownTicks } - - if (this.clientErrors !== this.oldClientErrors) { - this.ctx.warn('connection errors during interval', { - diff: this.clientErrors - this.oldClientErrors, - errors: this.clientErrors, - lastClients: this.lastClients - }) - this.oldClientErrors = this.clientErrors - } } - this.ticks++ } - createSession (token: Token, workspace: Workspace, account: Account): Session { - return this.sessionFactory(token, workspace, account) + private handleSessionTick (now: number): void { + for (const s of this.sessions.values()) { + const isCurrentUserTick = this.ticks % ticksPerSecond === s.tickHash + + if (isCurrentUserTick) { + const wsId = s.session.workspace.uuid + const lastRequestDiff = now - s.session.lastRequest + + let timeout = 60000 + if (s.session.getUser() === systemAccountUuid) { + timeout = timeout * 10 + } + if (lastRequestDiff > timeout) { + this.ctx.warn('session hang, closing...', { wsId, user: s.session.getUser() }) + + // Force close workspace if only one client and it hang. + void this.close(this.ctx, s.socket, wsId).catch((err) => { + this.ctx.error('failed to close', err) + }) + continue + } + if ( + lastRequestDiff + (1 / 10) * lastRequestDiff > this.timeouts.pingTimeout && + now - s.session.lastPing > this.timeouts.pingTimeout + ) { + // We need to check state and close socket if it broken + // And ping other wize + s.session.lastPing = now + if (s.socket.checkState()) { + void s.socket.send(this.ctx, { result: pingConst }, s.session.binaryMode, s.session.useCompression) + } + } + for (const r of s.session.requests.values()) { + const sec = Math.round((now - r.start) / 1000) + if (sec > 0 && sec % 30 === 0) { + this.ctx.warn('request hang found', { + sec, + wsId, + total: s.session.requests.size, + user: s.session.getUser(), + ...cutObjectArray(r.params) + }) + } + } + } + } + } + + createSession (token: Token, workspace: WorkspaceIds, info: LoginInfoWithWorkspaces): Session { + let primarySocialId: PersonId + let role: AccountRole = info.workspaces[workspace.uuid]?.role ?? AccountRole.User + switch (info.account) { + case systemAccountUuid: + primarySocialId = core.account.System + role = AccountRole.Owner + break + case guestAccount: + primarySocialId = '' as PersonId + role = AccountRole.DocGuest + break + default: + primarySocialId = pickPrimarySocialId(info.socialIds)._id + } + + return new ClientSession( + token, + workspace, + { + uuid: info.account, + socialIds: info.socialIds.map((it) => it._id), + primarySocialId, + fullSocialIds: [], + role + }, + info, + token.extra?.mode === 'backup' + ) } async getWorkspaceInfo (token: string, updateLastVisit = true): Promise { @@ -328,44 +344,10 @@ export class TSessionManager implements SessionManager { } } - async getAccount (token: string): Promise { + async getLoginWithWorkspaceInfo (token: string): Promise { try { const accountClient = getAccountClient(this.accountsUrl, token) - const loginInfo = await accountClient.getLoginInfoByToken() - - if (!isWorkspaceLoginInfo(loginInfo)) { - return - } - - if (loginInfo.account === guestAccount) { - return { - uuid: loginInfo.account, - role: loginInfo.role, - primarySocialId: '' as PersonId, - socialIds: [], - fullSocialIds: [] - } - } - - if (loginInfo.account === systemAccountUuid) { - return { - uuid: loginInfo.account, - role: loginInfo.role, - primarySocialId: core.account.System, - socialIds: [core.account.System], - fullSocialIds: [] - } - } - - const socialIds = await accountClient.getSocialIds() - - return { - uuid: loginInfo.account, - role: loginInfo.role, - primarySocialId: pickPrimarySocialId(socialIds)._id, - socialIds: socialIds.map((si) => si._id), - fullSocialIds: socialIds - } + return await accountClient.getLoginWithWorkspaceInfo() } catch (err: any) { if (err?.cause?.code === 'ECONNRESET' || err?.cause?.code === 'ECONNREFUSED') { return undefined @@ -382,67 +364,43 @@ export class TSessionManager implements SessionManager { tickCounter = 0 - @withContext('๐Ÿ“ฒ add-session') - async addSession ( + async getWorkspace ( ctx: MeasureContext, - ws: ConnectionSocket, + workspaceUuid: WorkspaceUuid, + workspaceInfo: LoginInfoWorkspace | undefined, token: Token, - rawToken: string, - sessionId: string | undefined - ): Promise { - const { workspace: workspaceUuid } = token - - let workspaceInfo: WorkspaceInfoWithStatus | undefined - try { - workspaceInfo = await this.getWorkspaceInfo(rawToken) - } catch (err: any) { - this.updateConnectErrorInfo(token) - return { error: err } - } - + ws: ConnectionSocket + ): Promise<{ workspace?: Workspace, resp?: GetWorkspaceResponse }> { if (workspaceInfo === undefined) { - return { error: new Error('Workspace not found or not available'), terminate: true } + return { resp: { error: new Error('Workspace not found or not available'), terminate: true } } } if (isArchivingMode(workspaceInfo.mode)) { // No access to disabled workspaces for regular users - return { error: new Error('Workspace is archived'), terminate: true, specialError: 'archived' } + return { resp: { error: new Error('Workspace is archived'), terminate: true, specialError: 'archived' } } } if (isMigrationMode(workspaceInfo.mode)) { // No access to disabled workspaces for regular users - return { error: new Error('Workspace is in region migration'), terminate: true, specialError: 'migration' } + return { + resp: { error: new Error('Workspace is in region migration'), terminate: true, specialError: 'migration' } + } } if (isRestoringMode(workspaceInfo.mode)) { // No access to disabled workspaces for regular users - return { error: new Error('Workspace is in backup restore'), terminate: true, specialError: 'migration' } + return { + resp: { error: new Error('Workspace is in backup restore'), terminate: true, specialError: 'migration' } + } } - if (workspaceInfo.isDisabled === true && token.account !== systemAccountUuid && token.extra?.admin !== 'true') { - // No access to disabled workspaces for regular users - return { error: new Error('Workspace not found or not available'), terminate: true } - } - - if (isWorkspaceCreating(workspaceInfo.mode) && token.account !== systemAccountUuid) { + if (isWorkspaceCreating(workspaceInfo.mode)) { // No access to workspace for token. - return { error: new Error(`Workspace during creation phase ${token.account} ${token.workspace}`) } - } - - let account: Account | undefined - try { - account = await this.getAccount(rawToken) - } catch (err: any) { - this.updateConnectErrorInfo(token) - return { error: err } - } - - if (account === undefined) { - return { error: new Error('Account not found or not available'), terminate: true } + return { resp: { error: new Error(`Workspace during creation phase...${workspaceUuid}`) } } } const wsVersion: Data = { - major: workspaceInfo.versionMajor, - minor: workspaceInfo.versionMinor, - patch: workspaceInfo.versionPatch + major: workspaceInfo.version.versionMajor, + minor: workspaceInfo.version.versionMinor, + patch: workspaceInfo.version.versionPatch } if ( @@ -454,13 +412,10 @@ export class TSessionManager implements SessionManager { ctx.warn('Model version mismatch', { version: this.modelVersion, workspaceVersion: versionToString(wsVersion), - workspace: workspaceInfo.uuid, - workspaceUrl: workspaceInfo.url, - account: token.account, - extra: JSON.stringify(token.extra ?? {}) + workspace: workspaceUuid }) // Version mismatch, return upgrading. - return { upgrade: true, progress: workspaceInfo.mode === 'upgrading' ? workspaceInfo.processingProgress ?? 0 : 0 } + return { resp: { upgrade: true, progress: workspaceInfo.mode === 'upgrading' ? workspaceInfo.progress ?? 0 : 0 } } } let workspace = this.workspaces.get(workspaceUuid) @@ -470,111 +425,131 @@ export class TSessionManager implements SessionManager { workspace = this.workspaces.get(workspaceUuid) - const oldSession = sessionId !== undefined ? workspace?.sessions?.get(sessionId) : undefined - if (oldSession !== undefined) { - // Just close old socket for old session id. - await this.close(ctx, oldSession.socket, workspaceUuid) - } - - const workspaceName = workspaceInfo.name ?? workspaceInfo.url ?? workspaceInfo.uuid - const branding = - (workspaceInfo.branding !== undefined - ? Object.values(this.brandingMap).find((b) => b.key === workspaceInfo?.branding) - : null) ?? null + const branding = null if (workspace === undefined) { ctx.warn('open workspace', { account: token.account, - workspace: workspaceInfo.uuid, - wsUrl: workspaceInfo.url, + workspace: workspaceUuid, ...token.extra }) - workspace = this.createWorkspace( - ctx.parent ?? ctx, - ctx, - token, - workspaceInfo.url ?? workspaceInfo.uuid, - workspaceName, - workspaceInfo.uuid, - workspaceInfo.dataId, - branding - ) - await this.workspaceProducer.send(workspaceInfo.uuid, [workspaceEvents.open()]) + workspace = this.createWorkspace(ctx.parent ?? ctx, ctx, token, workspaceInfo.url, workspaceInfo.dataId, branding) + await this.workspaceProducer.send(workspaceUuid, [workspaceEvents.open()]) } - let pipeline: Pipeline if (token.extra?.model === 'upgrade') { if (workspace.upgrade) { ctx.warn('reconnect workspace in upgrade', { account: token.account, - workspace: workspaceInfo.uuid, + workspace: workspaceUuid, wsUrl: workspaceInfo.url }) - pipeline = await ctx.with('๐Ÿ’ค wait-pipeline', {}, () => (workspace as Workspace).pipeline) } else { ctx.warn('reconnect workspace in upgrade switch', { email: token.account, - workspace: workspaceInfo.uuid, + workspace: workspaceUuid, wsUrl: workspaceInfo.url }) // We need to wait in case previous upgrade connection is already closing. - pipeline = await this.switchToUpgradeSession( - token, - sessionId, - ctx.parent ?? ctx, - workspaceInfo.uuid, - workspace, - ws, - workspaceInfo.url ?? workspaceInfo.uuid, - workspaceName - ) + await this.switchToUpgradeSession(token, ctx.parent ?? ctx, workspace, ws) } } else { if (workspace.upgrade) { ctx.warn('connect during upgrade', { account: token.account, - workspace: workspace.workspaceUuid, + workspace: workspace.wsId.url, sessionUsers: Array.from(workspace.sessions.values()).map((it) => it.session.getUser()), sessionData: Array.from(workspace.sessions.values()).map((it) => it.socket.data()) }) - return { upgrade: true } - } - - try { - if (workspace.pipeline instanceof Promise) { - pipeline = await ctx.with('๐Ÿ’ค wait-pipeline', {}, () => (workspace as Workspace).pipeline) - workspace.pipeline = pipeline - } else { - pipeline = workspace.pipeline - } - } catch (err: any) { - // Failed to create pipeline, etc - Analytics.handleError(err) - this.workspaces.delete(workspaceInfo.uuid) - throw err + return { resp: { upgrade: true } } } } + return { workspace } + } - const session = this.createSession(token, workspace, account) + @withContext('๐Ÿ“ฒ add-session') + async addSession ( + ctx: MeasureContext, + ws: ConnectionSocket, + token: Token, + rawToken: string, + sessionId: string | undefined + ): Promise { + let account: LoginInfoWithWorkspaces | undefined + + try { + account = await this.getLoginWithWorkspaceInfo(rawToken) + } catch (err: any) { + return { error: err } + } + + if (account === undefined) { + return { error: new Error('Account not found or not available'), terminate: true } + } + + let wsInfo = account.workspaces[token.workspace] + + if (wsInfo === undefined) { + // In case of guest or system account + // We need to get workspace info for system account. + const workspaceInfo = await this.getWorkspaceInfo(rawToken, false) + if (workspaceInfo === undefined) { + return { error: new Error('Workspace not found or not available'), terminate: true } + } + wsInfo = { + url: workspaceInfo.url, + mode: workspaceInfo.mode, + dataId: workspaceInfo.dataId, + version: { + versionMajor: workspaceInfo.versionMajor, + versionMinor: workspaceInfo.versionMinor, + versionPatch: workspaceInfo.versionPatch + }, + role: AccountRole.Owner, + endpoint: { externalUrl: '', internalUrl: '', region: workspaceInfo.region ?? '' }, + progress: workspaceInfo.processingProgress + } + } + const { workspace, resp } = await this.getWorkspace(ctx, token.workspace, wsInfo, token, ws) + if (resp !== undefined) { + return resp + } + + if (workspace === undefined || account === undefined) { + // Should not happen + return { error: new Error('Workspace not found or not available'), terminate: true } + } + + const oldSession = sessionId !== undefined ? workspace.sessions?.get(sessionId) : undefined + if (oldSession !== undefined) { + // Just close old socket for old session id. + await this.close(ctx, oldSession.socket, workspace.wsId.uuid) + } + + const session = this.createSession(token, workspace.wsId, account) session.sessionId = sessionId !== undefined && (sessionId ?? '').trim().length > 0 ? sessionId : generateId() session.sessionInstanceId = generateId() - this.sessions.set(ws.id, { session, socket: ws }) + const tickHash = this.tickCounter % ticksPerSecond + + this.sessions.set(ws.id, { session, socket: ws, tickHash }) // We need to delete previous session with Id if found. this.tickCounter++ - workspace.sessions.set(session.sessionId, { session, socket: ws, tickHash: this.tickCounter % ticksPerSecond }) + workspace.sessions.set(session.sessionId, { session, socket: ws, tickHash }) - const accountUuid = account.uuid - await this.usersProducer.send(workspaceInfo.uuid, [ - userEvents.login({ - user: accountUuid, - sessions: this.countUserSessions(workspace, accountUuid), - socialIds: account.socialIds - }) - ]) + const accountUuid = account.account + if (accountUuid !== systemAccountUuid && accountUuid !== guestAccount) { + await this.usersProducer.send(workspace.wsId.uuid, [ + userEvents.login({ + user: accountUuid, + sessions: this.countUserSessions(workspace, accountUuid), + socialIds: account.socialIds.map((it) => it._id) + }) + ]) + } // Mark workspace as init completed and we had at least one client. if (!workspace.workspaceInitCompleted) { @@ -588,76 +563,43 @@ export class TSessionManager implements SessionManager { ctx.error('failed to send maintenance warning', err) }) } - return { session, context: workspace.context, workspaceId: workspaceInfo.uuid } - } - - private updateConnectErrorInfo (token: Token): void { - this.clientErrors++ - if (!this.lastClients.includes(token.account)) { - this.lastClients = [token.account, ...this.lastClients.slice(0, 9)] - } + return { session, context: workspace.context, workspaceId: workspace.wsId.uuid } } private async switchToUpgradeSession ( token: Token, - sessionId: string | undefined, ctx: MeasureContext, - workspaceUuid: WorkspaceUuid, workspace: Workspace, - ws: ConnectionSocket, - workspaceUrl: string, - workspaceName: string - ): Promise { + ws: ConnectionSocket + ): Promise { if (LOGGING_ENABLED) { - ctx.info('reloading workspace', { workspaceName, token: JSON.stringify(token) }) + ctx.info('reloading workspace', { url: workspace.wsId.url, token: JSON.stringify(token) }) } // Mark as upgrade, to prevent any new clients to connect during close workspace.upgrade = true // If upgrade client is used. // Drop all existing clients - workspace.closing = this.closeAll(workspaceUuid, workspace, 0, 'upgrade') - await workspace.closing - workspace.closing = undefined - // Wipe workspace and update values. - workspace.workspaceName = workspaceName - if (!workspace.upgrade) { - // This is previous workspace, intended to be closed. - workspace.id = generateId() - workspace.sessions = new Map() - } - - const workspaceIds: WorkspaceIds = { - uuid: workspace.workspaceUuid, - url: workspace.workspaceUrl, - dataId: workspace.workspaceDataId - } - workspace.communicationApi = await this.communicationApiFactory(ctx, workspaceIds, (ctx, sessionIds, result) => { - this.broadcastSessions(ctx, workspace, sessionIds, result) - }) - // Re-create pipeline. - workspace.pipeline = this.pipelineFactory( - ctx, - workspaceIds, - true, - (ctx, tx, targets, exclude) => { - this.broadcastAll(workspace, tx, targets, exclude) - }, - workspace.branding, - workspace.communicationApi - ) - return await workspace.pipeline + await this.doCloseAll(workspace, 0, 'upgrade', ws) } - broadcastAll (workspace: Workspace, tx: Tx[], target?: string | string[], exclude?: string[]): void { - if (workspace.upgrade) { + broadcastAll (workspace: WorkspaceUuid, tx: Tx[], target?: string | string[], exclude?: string[]): void { + const ws = this.workspaces.get(workspace) + if (ws === undefined) { + return + } + this.doBroadcast(ws, tx, target, exclude) + } + + doBroadcast (ws: Workspace, tx: Tx[], target?: string | string[], exclude?: string[]): void { + if (ws.upgrade) { return } if (target !== undefined && !Array.isArray(target)) { target = [target] } const ctx = this.ctx.newChild('๐Ÿ“ฌ broadcast-all', {}) - const sessions = [...workspace.sessions.values()].filter((it) => { + const sessions = [...ws.sessions.values()].filter((it) => { if (it === undefined) { return false } @@ -760,8 +702,6 @@ export class TSessionManager implements SessionManager { pipelineCtx: MeasureContext, token: Token, workspaceUrl: string, - workspaceName: string, - workspaceUuid: WorkspaceUuid | undefined, workspaceDataId: WorkspaceDataId | undefined, branding: Branding | null ): Workspace { @@ -772,40 +712,36 @@ export class TSessionManager implements SessionManager { dataId: workspaceDataId, url: workspaceUrl } - const communicationApi = this.communicationApiFactory(pipelineCtx, workspaceIds, (ctx, sessionIds, result) => { - this.broadcastSessions(ctx, workspace, sessionIds, result) - }) - const factory = async (): Promise => { - return await this.pipelineFactory( + const factory = async (): Promise => { + const communicationApi = await this.communicationApiFactory( + pipelineCtx, + workspaceIds, + (ctx, sessionIds, result) => { + this.broadcastSessions(ctx, workspace, sessionIds, result) + } + ) + const pipeline = await this.pipelineFactory( pipelineCtx, workspaceIds, - upgrade, (ctx, tx, targets, exclude) => { - this.broadcastAll(workspace, tx, targets, exclude) + this.broadcastAll(workspaceIds.uuid, tx, targets, exclude) }, branding, - await communicationApi + communicationApi ) + return { pipeline, communicationApi } } - const workspace: Workspace = { + const workspace: Workspace = new Workspace( context, - id: generateId(), - pipeline: factory(), - communicationApi, - sessions: new Map(), - softShutdown: workspaceSoftShutdownTicks, - upgrade, - workspaceUuid: token.workspace, - workspaceName, - workspaceUrl, - workspaceDataId, - branding, - workspaceInitCompleted: false, - tickHash: this.tickCounter % ticksPerSecond, - tickHandlers: new Map(), - token: generateToken(systemAccountUuid, token.workspace) - } + generateToken(systemAccountUuid, token.workspace), + factory, + this.tickCounter % ticksPerSecond, + workspaceSoftShutdownTicks, + workspaceIds, + branding + ) + workspace.upgrade = upgrade this.workspaces.set(token.workspace, workspace) return workspace @@ -882,8 +818,9 @@ export class TSessionManager implements SessionManager { const sessionRef = this.sessions.get(ws.id) if (sessionRef !== undefined) { ctx.info('bye happen', { - workspace: workspace?.workspaceName, - user: sessionRef.session.getUser(), + workspace: workspace?.wsId.url, + userId: sessionRef.session.getUser(), + user: sessionRef.session.getSocialIds().find((it) => it.type !== SocialIdType.HULY)?.value, binary: sessionRef.session.binaryMode, compression: sessionRef.session.useCompression, totalTime: Date.now() - sessionRef.session.createTime, @@ -896,7 +833,7 @@ export class TSessionManager implements SessionManager { workspace.sessions.delete(sessionRef.session.sessionId) const userUuid = sessionRef.session.getUser() - await this.usersProducer.send(workspace.workspaceUuid, [ + await this.usersProducer.send(workspaceUuid, [ userEvents.logout({ user: userUuid, sessions: this.countUserSessions(workspace, userUuid), @@ -904,28 +841,28 @@ export class TSessionManager implements SessionManager { }) ]) - const pipeline = workspace.pipeline instanceof Promise ? await workspace.pipeline : workspace.pipeline - const communicationApi = - workspace.communicationApi instanceof Promise ? await workspace.communicationApi : workspace.communicationApi - if (this.doHandleTick) { workspace.tickHandlers.set(sessionRef.session.sessionId, { ticks: this.timeouts.reconnectTimeout * ticksPerSecond, operation: () => { this.reconnectIds.delete(sessionRef.session.sessionId) - void communicationApi.closeSession(sessionRef.session.sessionId) const user = sessionRef.session.getUser() if (workspace !== undefined) { const another = Array.from(workspace.sessions.values()).findIndex((p) => p.session.getUser() === user) if (another === -1 && !workspace.upgrade) { - void this.trySetStatus( - workspace.context, - pipeline, - communicationApi, - sessionRef.session, - false, - workspace.workspaceUuid - ).catch(() => {}) + void workspace.with(async (pipeline, communicationApi) => { + await communicationApi.closeSession(sessionRef.session.sessionId) + if (user !== guestAccount && user !== systemAccountUuid) { + await this.trySetStatus( + workspace.context, + pipeline, + communicationApi, + sessionRef.session, + false, + workspaceUuid + ).catch(() => {}) + } + }) } } } @@ -944,9 +881,9 @@ export class TSessionManager implements SessionManager { async forceClose (wsId: WorkspaceUuid, ignoreSocket?: ConnectionSocket): Promise { const ws = this.workspaces.get(wsId) if (ws !== undefined) { - this.ctx.warn('force-close', { name: ws.workspaceName }) + this.ctx.warn('force-close', { name: ws.wsId.url }) ws.upgrade = true // We need to similare upgrade to refresh all clients. - ws.closing = this.closeAll(wsId, ws, 99, 'force-close', ignoreSocket) + ws.closing = this.doCloseAll(ws, 99, 'force-close', ignoreSocket) this.workspaces.delete(wsId) await ws.closing ws.closing = undefined @@ -955,8 +892,7 @@ export class TSessionManager implements SessionManager { } } - async closeAll ( - wsId: WorkspaceUuid, + async doCloseAll ( workspace: Workspace, code: number, reason: 'upgrade' | 'shutdown' | 'force-close', @@ -964,16 +900,15 @@ export class TSessionManager implements SessionManager { ): Promise { if (LOGGING_ENABLED) { this.ctx.info('closing workspace', { - workspace: workspace.id, - wsName: workspace.workspaceName, + url: workspace.wsId.url, + uuid: workspace.wsId.uuid, code, - reason, - wsId + reason }) } const sessions = Array.from(workspace.sessions) - workspace.sessions = new Map() + workspace.sessions.clear() const closeS = (s: Session, webSocket: ConnectionSocket): void => { s.workspaceClosed = true @@ -987,9 +922,8 @@ export class TSessionManager implements SessionManager { if (LOGGING_ENABLED) { this.ctx.warn('Clients disconnected. Closing Workspace...', { - wsId, - workspace: workspace.id, - wsName: workspace.workspaceName + url: workspace.wsId.url, + uuid: workspace.wsId.uuid }) } @@ -999,9 +933,11 @@ export class TSessionManager implements SessionManager { closeS(s[1].session, s[1].socket) }) - await closeWorkspace(this.ctx, workspace) - if (LOGGING_ENABLED) { - this.ctx.warn('Workspace closed...', { workspace: workspace.id, wsId, wsName: workspace.workspaceName }) + if (reason !== 'upgrade') { + await workspace.close(this.ctx) + if (LOGGING_ENABLED) { + this.ctx.warn('Workspace closed...', { uuid: workspace.wsId.uuid, url: workspace.wsId.url }) + } } } @@ -1021,36 +957,34 @@ export class TSessionManager implements SessionManager { async closeWorkspaces (ctx: MeasureContext): Promise { clearInterval(this.checkInterval) for (const w of this.workspaces) { - await this.closeAll(w[0], w[1], 1, 'shutdown') + await this.doCloseAll(w[1], 1, 'shutdown') } await this.workspaceProducer.close() await this.usersProducer.close() } - private async performWorkspaceCloseCheck (workspace: Workspace, wsUuid: WorkspaceUuid): Promise { - const wsUID = workspace.id - const logParams = { wsUuid, workspace: workspace.id, wsName: workspace.workspaceName } + private async performWorkspaceCloseCheck (workspace: Workspace): Promise { + const uuid = workspace.wsId.uuid + const logParams = { uuid, url: workspace.wsId.url } if (workspace.sessions.size === 0) { if (LOGGING_ENABLED) { this.ctx.warn('no sessions for workspace', logParams) } try { if (workspace.sessions.size === 0) { - await closeWorkspace(this.ctx, workspace) + await workspace.close(this.ctx) - if (this.workspaces.get(wsUuid)?.id === wsUID) { - this.workspaces.delete(wsUuid) - } + this.workspaces.delete(uuid) workspace.context.end() if (LOGGING_ENABLED) { this.ctx.warn('Closed workspace', logParams) } - await this.workspaceProducer.send(workspace.workspaceUuid, [workspaceEvents.down()]) + await this.workspaceProducer.send(workspace.wsId.uuid, [workspaceEvents.down()]) } } catch (err: any) { Analytics.handleError(err) - this.workspaces.delete(wsUuid) + this.workspaces.delete(uuid) if (LOGGING_ENABLED) { this.ctx.error('failed', { ...logParams, error: err }) } @@ -1091,7 +1025,7 @@ export class TSessionManager implements SessionManager { sendPong: () => { ws.sendPong() }, - socialStringsToUsers: this.getActiveSocialStringsToUsersMap(service.workspace.workspaceUuid), + socialStringsToUsers: this.getActiveSocialStringsToUsersMap(service.workspace.uuid), sendError: (reqId, msg, error: Status) => sendResponse(sendCtx, service, ws, { id: reqId, @@ -1130,7 +1064,7 @@ export class TSessionManager implements SessionManager { service: S, ws: ConnectionSocket, request: Request, - workspace: WorkspaceUuid + workspaceId: WorkspaceUuid ): Promise { const userCtx = requestCtx.newChild('๐Ÿ“ž client', {}) @@ -1144,7 +1078,8 @@ export class TSessionManager implements SessionManager { const delta = platformNow() - request.time requestCtx.measure('msg-receive-delta', delta) } - if (service.workspace.closing !== undefined) { + const workspace = this.workspaces.get(workspaceId) + if (workspace === undefined || workspace.closing !== undefined) { await ws.send( ctx, { @@ -1163,12 +1098,12 @@ export class TSessionManager implements SessionManager { if (request.id === -2 && request.method === 'forceClose') { // TODO: we chould allow this only for admin or system accounts let done = false - const wsRef = this.workspaces.get(workspace) + const wsRef = this.workspaces.get(workspaceId) if (wsRef?.upgrade ?? false) { done = true - this.ctx.warn('FORCE CLOSE', { workspace }) + this.ctx.warn('FORCE CLOSE', { workspace: workspaceId }) // In case of upgrade, we need to force close workspace not in interval handler - await this.forceClose(workspace, ws) + await this.forceClose(workspaceId, ws) } const forceCloseResponse: Response = { id: request.id, @@ -1188,13 +1123,6 @@ export class TSessionManager implements SessionManager { return } - const pipeline = - service.workspace.pipeline instanceof Promise ? await service.workspace.pipeline : service.workspace.pipeline - const communicationApi = - service.workspace.communicationApi instanceof Promise - ? await service.workspace.communicationApi - : service.workspace.communicationApi - const f = (service as any)[request.method] try { const params = [...request.params] @@ -1203,12 +1131,14 @@ export class TSessionManager implements SessionManager { await ws.backpressure(ctx) } - await ctx.with('๐Ÿงจ process', {}, (callTx) => - f.apply(service, [ - this.createOpContext(callTx, userCtx, pipeline, communicationApi, request.id, service, ws), - ...params - ]) - ) + await workspace.with(async (pipeline, communicationApi) => { + await ctx.with('๐Ÿงจ process', {}, (callTx) => + f.apply(service, [ + this.createOpContext(callTx, userCtx, pipeline, communicationApi, request.id, service, ws), + ...params + ]) + ) + }) } catch (err: any) { Analytics.handleError(err) if (LOGGING_ENABLED) { @@ -1246,7 +1176,8 @@ export class TSessionManager implements SessionManager { const st = Date.now() return userCtx .with('๐Ÿงญ handleRPC', {}, async (ctx) => { - if (service.workspace.closing !== undefined) { + const workspace = this.workspaces.get(service.workspace.uuid) + if (workspace === undefined || workspace.closing !== undefined) { throw new Error('Workspace is closing') } @@ -1256,15 +1187,11 @@ export class TSessionManager implements SessionManager { start: st }) - const pipeline = - service.workspace.pipeline instanceof Promise ? await service.workspace.pipeline : service.workspace.pipeline - const communicationApi = - service.workspace.communicationApi instanceof Promise - ? await service.workspace.communicationApi - : service.workspace.communicationApi try { - const uctx = this.createOpContext(ctx, userCtx, pipeline, communicationApi, reqId, service, ws) - await operation(uctx) + await workspace.with(async (pipeline, communicationApi) => { + const uctx = this.createOpContext(ctx, userCtx, pipeline, communicationApi, reqId, service, ws) + await operation(uctx) + }) } catch (err: any) { Analytics.handleError(err) if (LOGGING_ENABLED) { @@ -1289,11 +1216,36 @@ export class TSessionManager implements SessionManager { }) } + entryToUserStats = (session: Session, socket: ConnectionSocket): UserStatistics => { + return { + current: session.current, + mins5: session.mins5, + userId: session.getUser(), + sessionId: socket.id, + total: session.total, + data: socket.data + } + } + + workspaceToWorkspaceStats = (ws: Workspace): WorkspaceStatistics => { + return { + clientsTotal: new Set(Array.from(ws.sessions.values()).map((it) => it.session.getUser())).size, + sessionsTotal: ws.sessions.size, + workspaceName: ws.wsId.url, + wsId: ws.wsId.uuid, + sessions: Array.from(ws.sessions.values()).map((it) => this.entryToUserStats(it.session, it.socket)) + } + } + + getStatistics (): WorkspaceStatistics[] { + return Array.from(this.workspaces.values()).map((it) => this.workspaceToWorkspaceStats(it)) + } + private async handleHello( request: Request, service: S, ctx: MeasureContext, - workspace: WorkspaceUuid, + workspace: Workspace, ws: ConnectionSocket, requestCtx: MeasureContext ): Promise { @@ -1304,12 +1256,14 @@ export class TSessionManager implements SessionManager { if (LOGGING_ENABLED) { ctx.info('hello happen', { - workspace, - user: service.getUser(), + workspace: workspace.wsId.url, + workspaceId: workspace.wsId.uuid, + userId: service.getUser(), + user: service.getSocialIds().find((it) => it.type !== SocialIdType.HULY)?.value, binary: service.binaryMode, compression: service.useCompression, timeToHello: Date.now() - service.createTime, - workspaceUsers: this.workspaces.get(workspace)?.sessions?.size, + workspaceUsers: workspace.sessions.size, totalUsers: this.sessions.size }) } @@ -1317,33 +1271,31 @@ export class TSessionManager implements SessionManager { if (reconnect) { this.reconnectIds.delete(service.sessionId) } - const pipeline = - service.workspace.pipeline instanceof Promise ? await service.workspace.pipeline : service.workspace.pipeline - const communicationApi = - service.workspace.communicationApi instanceof Promise - ? await service.workspace.communicationApi - : service.workspace.communicationApi - const helloResponse: HelloResponse = { - id: -1, - result: 'hello', - binary: service.binaryMode, - reconnect, - serverVersion: this.serverVersion, - lastTx: pipeline.context.lastTx, - lastHash: pipeline.context.lastHash, - account: service.getRawAccount(), - useCompression: service.useCompression - } - await ws.send(requestCtx, helloResponse, false, false) - // We do not need to wait for set-status, just return session to client - const _workspace = service.workspace - if (helloResponse.account.role !== AccountRole.DocGuest) { - void ctx - .with('set-status', {}, (ctx) => - this.trySetStatus(ctx, pipeline, communicationApi, service, true, _workspace.workspaceUuid) - ) - .catch(() => {}) + const account = service.getRawAccount() + await workspace.with(async (pipeline, communicationApi) => { + const helloResponse: HelloResponse = { + id: -1, + result: 'hello', + binary: service.binaryMode, + reconnect, + serverVersion: this.serverVersion, + lastTx: pipeline.context.lastTx, + lastHash: pipeline.context.lastHash, + account, + useCompression: service.useCompression + } + await ws.send(requestCtx, helloResponse, false, false) + }) + if (account.uuid !== guestAccount && account.uuid !== systemAccountUuid) { + void workspace.with(async (pipeline, communicationApi) => { + // We do not need to wait for set-status, just return session to client + await ctx + .with('set-status', {}, (ctx) => + this.trySetStatus(ctx, pipeline, communicationApi, service, true, service.workspace.uuid) + ) + .catch(() => {}) + }) } } catch (err: any) { ctx.error('error', { err }) @@ -1353,7 +1305,6 @@ export class TSessionManager implements SessionManager { export function createSessionManager ( ctx: MeasureContext, - sessionFactory: SessionFactory, brandingMap: BrandingMap, timeouts: Timeouts, profiling: @@ -1371,7 +1322,6 @@ export function createSessionManager ( ): SessionManager { return new TSessionManager( ctx, - sessionFactory, timeouts, brandingMap ?? null, profiling, @@ -1387,7 +1337,6 @@ export function createSessionManager ( export interface SessionManagerOptions extends Partial { pipelineFactory: PipelineFactory communicationApiFactory: CommunicationApiFactory - sessionFactory: SessionFactory brandingMap: BrandingMap enableCompression?: boolean accountsUrl: string @@ -1404,7 +1353,6 @@ export interface SessionManagerOptions extends Partial { export function startSessionManager (ctx: MeasureContext, opt: SessionManagerOptions): SessionManager { const sessions = createSessionManager( ctx, - opt.sessionFactory, opt.brandingMap, { pingTimeout: opt.pingTimeout ?? 10000, @@ -1420,33 +1368,3 @@ export function startSessionManager (ctx: MeasureContext, opt: SessionManagerOpt ) return sessions } - -async function closeWorkspace (ctx: MeasureContext, workspace: Workspace): Promise { - const closePipeline = async (): Promise => { - try { - await ctx.with('close-pipeline', {}, async () => { - await (await workspace.pipeline).close() - }) - } catch (err: any) { - Analytics.handleError(err) - ctx.error('close-pipeline-error', { error: err }) - } - } - - const closeCommunicationApi = async (): Promise => { - try { - await ctx.with('close-communication-api', {}, async () => { - await (await workspace.communicationApi).close() - }) - } catch (err: any) { - Analytics.handleError(err) - ctx.error('close-pipeline-error', { error: err }) - } - } - await ctx.with('closing', {}, async () => { - const to = timeoutPromise(120000) - const closePromises = [closePipeline(), closeCommunicationApi()] - await Promise.race([Promise.all(closePromises), to.promise]) - to.cancelHandle() - }) -} diff --git a/server/server/src/stats.ts b/server/server/src/stats.ts index 2254802d20..f3814ca4f3 100644 --- a/server/server/src/stats.ts +++ b/server/server/src/stats.ts @@ -20,22 +20,8 @@ export function getStatistics (ctx: MeasureContext, sessions: SessionManager, ad } data.statistics.totalClients = sessions.sessions.size if (admin) { - for (const [k, vv] of sessions.workspaces) { - data.statistics.activeSessions[k] = { - sessions: Array.from(vv.sessions.entries()).map(([k, v]) => ({ - userId: v.session.getUser(), - data: v.socket.data(), - mins5: v.session.mins5, - total: v.session.total, - current: v.session.current, - upgrade: v.session.isUpgradeClient() - })), - name: vv.workspaceName, - wsId: vv.workspaceUuid, - sessionsTotal: vv.sessions.size, - upgrading: vv.upgrade, - closing: vv.closing !== undefined - } + for (const wsStats of sessions.getStatistics()) { + data.statistics.activeSessions[wsStats.wsId] = wsStats } } diff --git a/server/server/src/workspace.ts b/server/server/src/workspace.ts new file mode 100644 index 0000000000..575dac0909 --- /dev/null +++ b/server/server/src/workspace.ts @@ -0,0 +1,130 @@ +// +// Copyright ยฉ 2022, 2023 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 { Analytics } from '@hcengineering/analytics' +import { type ServerApi as CommunicationApi } from '@hcengineering/communication-sdk-types' +import { type Branding, type MeasureContext, type WorkspaceIds } from '@hcengineering/core' +import type { ConnectionSocket, Pipeline, Session } from '@hcengineering/server-core' + +interface TickHandler { + ticks: number + operation: () => void +} + +export interface PipelinePair { + pipeline: Pipeline + communicationApi: CommunicationApi +} +export type WorkspacePipelineFactory = () => Promise + +/** + * @public + */ +export class Workspace { + pipeline?: PipelinePair | Promise + upgrade: boolean = false + closing?: Promise + + workspaceInitCompleted: boolean = false + + softShutdown: number + + sessions = new Map() + tickHandlers = new Map() + + operations: number = 0 + + constructor ( + readonly context: MeasureContext, + readonly token: string, // Account workspace update token. + readonly factory: WorkspacePipelineFactory, + + readonly tickHash: number, + + softShutdown: number, + + readonly wsId: WorkspaceIds, + readonly branding: Branding | null + ) { + this.softShutdown = softShutdown + } + + private getPipelinePair (): PipelinePair | Promise { + if (this.pipeline === undefined) { + this.pipeline = this.factory() + } + return this.pipeline + } + + async with(op: (pipeline: Pipeline, communicationApi: CommunicationApi) => Promise): Promise { + this.operations++ + let pair = this.getPipelinePair() + if (pair instanceof Promise) { + pair = await pair + this.pipeline = pair + } + try { + return await op(pair.pipeline, pair.communicationApi) + } finally { + this.operations-- + } + } + + async close (ctx: MeasureContext): Promise { + if (this.pipeline === undefined) { + return + } + const { pipeline, communicationApi } = await this.pipeline + const closePipeline = async (): Promise => { + try { + await ctx.with('close-pipeline', {}, async () => { + await pipeline.close() + }) + } catch (err: any) { + Analytics.handleError(err) + ctx.error('close-pipeline-error', { error: err }) + } + } + + const closeCommunicationApi = async (): Promise => { + try { + await ctx.with('close-communication-api', {}, async () => { + await communicationApi.close() + }) + } catch (err: any) { + Analytics.handleError(err) + ctx.error('close-pipeline-error', { error: err }) + } + } + await ctx.with('closing', {}, async () => { + const to = timeoutPromise(120000) + const closePromises = [closePipeline(), closeCommunicationApi()] + await Promise.race([Promise.all(closePromises), to.promise]) + to.cancelHandle() + }) + } +} + +function timeoutPromise (time: number): { promise: Promise, cancelHandle: () => void } { + let timer: any + return { + promise: new Promise((resolve) => { + timer = setTimeout(resolve, time) + }), + cancelHandle: () => { + clearTimeout(timer) + } + } +} diff --git a/server/workspace-service/src/ws-operations.ts b/server/workspace-service/src/ws-operations.ts index 9ca0fcd3bd..7b13b87c86 100644 --- a/server/workspace-service/src/ws-operations.ts +++ b/server/workspace-service/src/ws-operations.ts @@ -318,7 +318,6 @@ export async function upgradeWorkspaceWith ( true, undefined, wsIds, - null, true, undefined, undefined,