mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-30 21:45:01 +02:00
UBERF-4319: Performance changes (#4474)
Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
This commit is contained in:
@@ -47,14 +47,23 @@ export async function createPipeline (
|
||||
let broadcastHook: HandledBroadcastFunc = (): Tx[] => {
|
||||
return []
|
||||
}
|
||||
const storage = await createServerStorage(conf, {
|
||||
upgrade,
|
||||
broadcast: (tx: Tx[], targets?: string[]) => {
|
||||
const sendTx = broadcastHook?.(tx, targets) ?? tx
|
||||
broadcast(sendTx, targets)
|
||||
}
|
||||
})
|
||||
const pipeline = PipelineImpl.create(ctx, storage, constructors, broadcast)
|
||||
const storage = await ctx.with(
|
||||
'create-server-storage',
|
||||
{},
|
||||
async (ctx) =>
|
||||
await createServerStorage(ctx, conf, {
|
||||
upgrade,
|
||||
broadcast: (tx: Tx[], targets?: string[]) => {
|
||||
const sendTx = broadcastHook?.(tx, targets) ?? tx
|
||||
broadcast(sendTx, targets)
|
||||
}
|
||||
})
|
||||
)
|
||||
const pipeline = ctx.with(
|
||||
'create pipeline',
|
||||
{},
|
||||
async (ctx) => await PipelineImpl.create(ctx, storage, constructors, broadcast)
|
||||
)
|
||||
const pipelineResult = await pipeline
|
||||
broadcastHook = (tx, targets) => {
|
||||
return pipelineResult.handleBroadcast(tx, targets)
|
||||
@@ -92,7 +101,7 @@ class PipelineImpl implements Pipeline {
|
||||
let current: Middleware | undefined
|
||||
for (let index = constructors.length - 1; index >= 0; index--) {
|
||||
const element = constructors[index]
|
||||
current = await element(ctx, broadcast, this.storage, current)
|
||||
current = await ctx.with('build chain', {}, async (ctx) => await element(ctx, broadcast, this.storage, current))
|
||||
}
|
||||
return current
|
||||
}
|
||||
|
||||
+90
-65
@@ -20,16 +20,15 @@ import core, {
|
||||
Class,
|
||||
ClassifierKind,
|
||||
Collection,
|
||||
DOMAIN_DOC_INDEX_STATE,
|
||||
DOMAIN_MODEL,
|
||||
DOMAIN_TX,
|
||||
Doc,
|
||||
DocumentQuery,
|
||||
DocumentUpdate,
|
||||
Domain,
|
||||
DOMAIN_DOC_INDEX_STATE,
|
||||
DOMAIN_MODEL,
|
||||
DOMAIN_TX,
|
||||
FindOptions,
|
||||
FindResult,
|
||||
generateId,
|
||||
Hierarchy,
|
||||
IndexingUpdateEvent,
|
||||
LoadModelResponse,
|
||||
@@ -37,13 +36,16 @@ import core, {
|
||||
Mixin,
|
||||
ModelDb,
|
||||
Ref,
|
||||
SearchOptions,
|
||||
SearchQuery,
|
||||
SearchResult,
|
||||
ServerStorage,
|
||||
StorageIterator,
|
||||
Timestamp,
|
||||
Tx,
|
||||
TxApplyIf,
|
||||
TxCollectionCUD,
|
||||
TxCUD,
|
||||
TxCollectionCUD,
|
||||
TxFactory,
|
||||
TxProcessor,
|
||||
TxRemoveDoc,
|
||||
@@ -52,9 +54,7 @@ import core, {
|
||||
TxWorkspaceEvent,
|
||||
WorkspaceEvent,
|
||||
WorkspaceId,
|
||||
SearchQuery,
|
||||
SearchOptions,
|
||||
SearchResult
|
||||
generateId
|
||||
} from '@hcengineering/core'
|
||||
import { MinioService } from '@hcengineering/minio'
|
||||
import { getResource } from '@hcengineering/platform'
|
||||
@@ -74,7 +74,6 @@ import type {
|
||||
ObjectDDParticipant,
|
||||
TriggerControl
|
||||
} from './types'
|
||||
import { createFindAll } from './utils'
|
||||
|
||||
/**
|
||||
* @public
|
||||
@@ -165,12 +164,16 @@ class TServerStorage implements ServerStorage {
|
||||
const adapter = this.getAdapter(lastDomain as Domain)
|
||||
const toDelete = part.filter((it) => it._class === core.class.TxRemoveDoc).map((it) => it.objectId)
|
||||
|
||||
const toDeleteDocs = await adapter.load(lastDomain as Domain, toDelete)
|
||||
const toDeleteDocs = await ctx.with(
|
||||
'adapter-load',
|
||||
{ domain: lastDomain },
|
||||
async () => await adapter.load(lastDomain as Domain, toDelete)
|
||||
)
|
||||
for (const ddoc of toDeleteDocs) {
|
||||
removedDocs.set(ddoc._id, ddoc)
|
||||
}
|
||||
|
||||
const r = await adapter.tx(...part)
|
||||
const r = await ctx.with('adapter-tx', {}, async () => await adapter.tx(...part))
|
||||
if (Array.isArray(r)) {
|
||||
result.push(...r)
|
||||
} else {
|
||||
@@ -370,13 +373,13 @@ class TServerStorage implements ServerStorage {
|
||||
query: DocumentQuery<T>,
|
||||
options?: FindOptions<T>
|
||||
): Promise<FindResult<T>> {
|
||||
return await ctx.with('find-all', {}, (ctx) => {
|
||||
const domain = this.hierarchy.getDomain(clazz)
|
||||
if (query?.$search !== undefined) {
|
||||
return ctx.with('full-text-find-all', {}, (ctx) => this.fulltext.findAll(ctx, clazz, query, options))
|
||||
}
|
||||
return ctx.with('db-find-all', { d: domain }, () => this.getAdapter(domain).findAll(clazz, query, options))
|
||||
})
|
||||
const domain = this.hierarchy.getDomain(clazz)
|
||||
if (query?.$search !== undefined) {
|
||||
return await ctx.with('client-fulltext-find-all', {}, (ctx) => this.fulltext.findAll(ctx, clazz, query, options))
|
||||
}
|
||||
return await ctx.with('client-find-all', { _class: clazz }, () =>
|
||||
this.getAdapter(domain).findAll(clazz, query, options)
|
||||
)
|
||||
}
|
||||
|
||||
async searchFulltext (ctx: MeasureContext, query: SearchQuery, options: SearchOptions): Promise<SearchResult> {
|
||||
@@ -550,13 +553,13 @@ class TServerStorage implements ServerStorage {
|
||||
): Promise<FindResult<T>> =>
|
||||
findAll(mctx, clazz, query, options)
|
||||
|
||||
const removed = await ctx.with('process-remove', {}, () => this.processRemove(ctx, txes, findAll, removedMap))
|
||||
const collections = await ctx.with('process-collection', {}, () =>
|
||||
const removed = await ctx.with('process-remove', {}, (ctx) => this.processRemove(ctx, txes, findAll, removedMap))
|
||||
const collections = await ctx.with('process-collection', {}, (ctx) =>
|
||||
this.processCollection(ctx, txes, findAll, removedMap)
|
||||
)
|
||||
const moves = await ctx.with('process-move', {}, () => this.processMove(ctx, txes, findAll))
|
||||
const moves = await ctx.with('process-move', {}, (ctx) => this.processMove(ctx, txes, findAll))
|
||||
|
||||
const triggerControl: Omit<TriggerControl, 'txFactory'> = {
|
||||
const triggerControl: Omit<TriggerControl, 'txFactory' | 'ctx'> = {
|
||||
removedMap,
|
||||
workspace: this.workspace,
|
||||
fx: triggerFx.fx,
|
||||
@@ -572,15 +575,25 @@ class TServerStorage implements ServerStorage {
|
||||
triggerFx.fx(() => f(adapter, this.workspace))
|
||||
},
|
||||
findAll: fAll(ctx),
|
||||
findAllCtx: findAll,
|
||||
modelDb: this.modelDb,
|
||||
hierarchy: this.hierarchy,
|
||||
apply: async (tx, broadcast) => {
|
||||
return await this.apply(ctx, tx, broadcast)
|
||||
},
|
||||
applyCtx: async (ctx, tx, broadcast) => {
|
||||
return await this.apply(ctx, tx, broadcast)
|
||||
}
|
||||
}
|
||||
const triggers = await ctx.with('process-triggers', {}, async (ctx) => {
|
||||
const result: Tx[] = []
|
||||
result.push(...(await this.triggers.apply(ctx, txes, triggerControl)))
|
||||
result.push(
|
||||
...(await this.triggers.apply(ctx, txes, {
|
||||
...triggerControl,
|
||||
ctx,
|
||||
findAll: fAll(ctx)
|
||||
}))
|
||||
)
|
||||
return result
|
||||
})
|
||||
|
||||
@@ -685,7 +698,18 @@ class TServerStorage implements ServerStorage {
|
||||
|
||||
async processTxes (ctx: MeasureContext, txes: Tx[]): Promise<[TxResult, Tx[]]> {
|
||||
// store tx
|
||||
const _findAll = createFindAll(this)
|
||||
const _findAll: ServerStorage['findAll'] = async <T extends Doc>(
|
||||
ctx: MeasureContext,
|
||||
clazz: Ref<Class<T>>,
|
||||
query: DocumentQuery<T>,
|
||||
options?: FindOptions<T>
|
||||
): Promise<FindResult<T>> => {
|
||||
const domain = this.hierarchy.getDomain(clazz)
|
||||
if (query?.$search !== undefined) {
|
||||
return await ctx.with('full-text-find-all', {}, (ctx) => this.fulltext.findAll(ctx, clazz, query, options))
|
||||
}
|
||||
return await ctx.with('find-all', { _class: clazz }, () => this.getAdapter(domain).findAll(clazz, query, options))
|
||||
}
|
||||
const txToStore: Tx[] = []
|
||||
const modelTx: Tx[] = []
|
||||
const applyTxes: Tx[] = []
|
||||
@@ -748,7 +772,7 @@ class TServerStorage implements ServerStorage {
|
||||
}
|
||||
|
||||
async tx (ctx: MeasureContext, tx: Tx): Promise<[TxResult, Tx[]]> {
|
||||
return await this.processTxes(ctx, [tx])
|
||||
return await ctx.with('client-tx', { _class: tx._class }, async (ctx) => await this.processTxes(ctx, [tx]))
|
||||
}
|
||||
|
||||
find (domain: Domain): StorageIterator {
|
||||
@@ -801,6 +825,7 @@ export interface ServerStorageOptions {
|
||||
* @public
|
||||
*/
|
||||
export async function createServerStorage (
|
||||
ctx: MeasureContext,
|
||||
conf: DbConfiguration,
|
||||
options: ServerStorageOptions
|
||||
): Promise<ServerStorage> {
|
||||
@@ -809,63 +834,65 @@ export async function createServerStorage (
|
||||
const adapters = new Map<string, DbAdapter>()
|
||||
const modelDb = new ModelDb(hierarchy)
|
||||
|
||||
console.timeLog(conf.workspace.name, 'create server storage')
|
||||
const storageAdapter = conf.storageFactory?.()
|
||||
|
||||
for (const key in conf.adapters) {
|
||||
const adapterConf = conf.adapters[key]
|
||||
adapters.set(key, await adapterConf.factory(hierarchy, adapterConf.url, conf.workspace, modelDb, storageAdapter))
|
||||
console.timeLog(conf.workspace.name, 'adapter', key)
|
||||
}
|
||||
|
||||
const txAdapter = adapters.get(conf.domains[DOMAIN_TX]) as TxAdapter
|
||||
if (txAdapter === undefined) {
|
||||
console.log('no txadapter found')
|
||||
}
|
||||
|
||||
console.timeLog(conf.workspace.name, 'begin get model')
|
||||
const model = await txAdapter.getModel()
|
||||
console.timeLog(conf.workspace.name, 'get model')
|
||||
for (const tx of model) {
|
||||
try {
|
||||
hierarchy.tx(tx)
|
||||
await triggers.tx(tx)
|
||||
} catch (err: any) {
|
||||
console.error('failed to apply model transaction, skipping', JSON.stringify(tx), err)
|
||||
const model = await ctx.with('get model', {}, async (ctx) => {
|
||||
const model = await txAdapter.getModel()
|
||||
for (const tx of model) {
|
||||
try {
|
||||
hierarchy.tx(tx)
|
||||
await triggers.tx(tx)
|
||||
} catch (err: any) {
|
||||
console.error('failed to apply model transaction, skipping', JSON.stringify(tx), err)
|
||||
}
|
||||
}
|
||||
}
|
||||
console.timeLog(conf.workspace.name, 'finish hierarchy')
|
||||
|
||||
for (const tx of model) {
|
||||
try {
|
||||
await modelDb.tx(tx)
|
||||
} catch (err: any) {
|
||||
console.error('failed to apply model transaction, skipping', JSON.stringify(tx), err)
|
||||
for (const tx of model) {
|
||||
try {
|
||||
await modelDb.tx(tx)
|
||||
} catch (err: any) {
|
||||
console.error('failed to apply model transaction, skipping', JSON.stringify(tx), err)
|
||||
}
|
||||
}
|
||||
}
|
||||
console.timeLog(conf.workspace.name, 'finish local model')
|
||||
return model
|
||||
})
|
||||
|
||||
for (const [adn, adapter] of adapters) {
|
||||
await adapter.init(model)
|
||||
console.timeLog(conf.workspace.name, 'finish init adapter', adn)
|
||||
await ctx.with('init-adapter', { name: adn }, async (ctx) => {
|
||||
await adapter.init(model)
|
||||
})
|
||||
}
|
||||
|
||||
const fulltextAdapter = await conf.fulltextAdapter.factory(
|
||||
conf.fulltextAdapter.url,
|
||||
conf.workspace,
|
||||
conf.metrics.newChild('fulltext', {})
|
||||
const fulltextAdapter = await ctx.with(
|
||||
'create full text adapter',
|
||||
{},
|
||||
async (ctx) =>
|
||||
await conf.fulltextAdapter.factory(
|
||||
conf.fulltextAdapter.url,
|
||||
conf.workspace,
|
||||
conf.metrics.newChild('🗒️ fulltext', {})
|
||||
)
|
||||
)
|
||||
console.timeLog(conf.workspace.name, 'finish fulltext adapter')
|
||||
|
||||
const metrics = conf.metrics.newChild('server-storage', {})
|
||||
const metrics = conf.metrics.newChild('📔 server-storage', {})
|
||||
|
||||
const contentAdapter = await createContentAdapter(
|
||||
conf.contentAdapters,
|
||||
conf.defaultContentAdapter,
|
||||
conf.workspace,
|
||||
metrics.newChild('content', {})
|
||||
const contentAdapter = await ctx.with(
|
||||
'create content adapter',
|
||||
{},
|
||||
async (ctx) =>
|
||||
await createContentAdapter(
|
||||
conf.contentAdapters,
|
||||
conf.defaultContentAdapter,
|
||||
conf.workspace,
|
||||
metrics.newChild('content', {})
|
||||
)
|
||||
)
|
||||
console.timeLog(conf.workspace.name, 'finish content adapter')
|
||||
|
||||
const defaultAdapter = adapters.get(conf.defaultAdapter)
|
||||
if (defaultAdapter === undefined) {
|
||||
@@ -877,7 +904,6 @@ export async function createServerStorage (
|
||||
throw new Error('No storage adapter')
|
||||
}
|
||||
const stages = conf.fulltextAdapter.stages(fulltextAdapter, storage, storageAdapter, contentAdapter)
|
||||
console.timeLog(conf.workspace.name, 'finish index pipeline stages')
|
||||
|
||||
const indexer = new FullTextIndexPipeline(
|
||||
defaultAdapter,
|
||||
@@ -903,7 +929,6 @@ export async function createServerStorage (
|
||||
options.broadcast?.([tx])
|
||||
}
|
||||
)
|
||||
console.timeLog(conf.workspace.name, 'finish create indexer')
|
||||
return new FullTextIndex(
|
||||
hierarchy,
|
||||
fulltextAdapter,
|
||||
|
||||
@@ -70,7 +70,17 @@ export class Triggers {
|
||||
if (matches.length > 0) {
|
||||
await ctx.with(resource, {}, async (ctx) => {
|
||||
for (const tx of matches) {
|
||||
result.push(...(await trigger(tx, { ...ctrl, txFactory: new TxFactory(tx.modifiedBy, true) })))
|
||||
result.push(
|
||||
...(await trigger(tx, {
|
||||
...ctrl,
|
||||
ctx,
|
||||
txFactory: new TxFactory(tx.modifiedBy, true),
|
||||
findAll: async (clazz, query, options) => await ctrl.findAllCtx(ctx, clazz, query, options),
|
||||
apply: async (tx, broadcast) => {
|
||||
return await ctrl.applyCtx(ctx, tx, broadcast)
|
||||
}
|
||||
}))
|
||||
)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -113,13 +113,23 @@ export interface Pipeline extends LowLevelStorage {
|
||||
* @public
|
||||
*/
|
||||
export interface TriggerControl {
|
||||
ctx: MeasureContext
|
||||
workspace: WorkspaceId
|
||||
txFactory: TxFactory
|
||||
findAll: Storage['findAll']
|
||||
findAllCtx: <T extends Doc>(
|
||||
ctx: MeasureContext,
|
||||
_class: Ref<Class<T>>,
|
||||
query: DocumentQuery<T>,
|
||||
options?: FindOptions<T>
|
||||
) => Promise<FindResult<T>>
|
||||
hierarchy: Hierarchy
|
||||
modelDb: ModelDb
|
||||
removedMap: Map<Ref<Doc>, Doc>
|
||||
|
||||
// // An object cache,
|
||||
// getCachedObject: <T extends Doc>(_class: Ref<Class<T>>, _id: Ref<T>) => Promise<T | undefined>
|
||||
|
||||
fulltextFx: (f: (adapter: FullTextAdapter) => Promise<void>) => void
|
||||
// Since we don't have other storages let's consider adapter is MinioClient
|
||||
// Later can be replaced with generic one with bucket encapsulated inside.
|
||||
@@ -128,6 +138,7 @@ export interface TriggerControl {
|
||||
|
||||
// Bulk operations in case trigger require some
|
||||
apply: (tx: Tx[], broadcast: boolean) => Promise<TxResult>
|
||||
applyCtx: (ctx: MeasureContext, tx: Tx[], broadcast: boolean) => Promise<TxResult>
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -1,24 +0,0 @@
|
||||
import {
|
||||
Class,
|
||||
Doc,
|
||||
DocumentQuery,
|
||||
FindOptions,
|
||||
FindResult,
|
||||
MeasureContext,
|
||||
Ref,
|
||||
ServerStorage
|
||||
} from '@hcengineering/core'
|
||||
|
||||
/**
|
||||
* @public
|
||||
*/
|
||||
export function createFindAll (storage: ServerStorage): ServerStorage['findAll'] {
|
||||
return async <T extends Doc>(
|
||||
ctx: MeasureContext,
|
||||
clazz: Ref<Class<T>>,
|
||||
query: DocumentQuery<T>,
|
||||
options?: FindOptions<T>
|
||||
): Promise<FindResult<T>> => {
|
||||
return await storage.findAll(ctx, clazz, query, options)
|
||||
}
|
||||
}
|
||||
@@ -78,7 +78,9 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
next?: Middleware
|
||||
): Promise<SpaceSecurityMiddleware> {
|
||||
const res = new SpaceSecurityMiddleware(broadcast, storage, next)
|
||||
await res.init(ctx)
|
||||
await ctx.with('space chain', {}, async (ctx) => {
|
||||
await res.init(ctx)
|
||||
})
|
||||
return res
|
||||
}
|
||||
|
||||
|
||||
@@ -160,8 +160,8 @@ describe('mongo operations', () => {
|
||||
workspace: getWorkspaceId(dbId, ''),
|
||||
storageFactory: () => createNullStorageFactory()
|
||||
}
|
||||
const serverStorage = await createServerStorage(conf, { upgrade: false })
|
||||
const ctx = new MeasureMetricsContext('client', {})
|
||||
const serverStorage = await createServerStorage(ctx, conf, { upgrade: false })
|
||||
client = await createClient(async (handler) => {
|
||||
const st: ClientConnection = {
|
||||
findAll: async (_class, query, options) => await serverStorage.findAll(ctx, _class, query, options),
|
||||
@@ -174,7 +174,8 @@ describe('mongo operations', () => {
|
||||
upload: async (domain: Domain, docs: Doc[]) => {},
|
||||
clean: async (domain: Domain, docs: Ref<Doc>[]) => {},
|
||||
loadModel: async () => txes,
|
||||
getAccount: async () => ({}) as any
|
||||
getAccount: async () => ({}) as any,
|
||||
measure: async () => async () => ({ time: 0, serverTime: 0 })
|
||||
}
|
||||
return st
|
||||
})
|
||||
|
||||
@@ -80,14 +80,20 @@ export class APMMeasureContext implements MeasureContext {
|
||||
}
|
||||
}
|
||||
|
||||
async error (err: any): Promise<void> {
|
||||
async error (message: string, ...args: any[]): Promise<void> {
|
||||
this.logger.error(message, args)
|
||||
|
||||
await new Promise<void>((resolve) => {
|
||||
this.agent.captureError(err, () => {
|
||||
this.agent.captureError({ message, params: args }, () => {
|
||||
resolve()
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
async info (message: string, ...args: any[]): Promise<void> {
|
||||
this.logger.info(message, args)
|
||||
}
|
||||
|
||||
end (): void {
|
||||
this.transaction?.end()
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import { writeFile } from 'fs/promises'
|
||||
|
||||
const apmUrl = process.env.APM_SERVER_URL
|
||||
const metricsFile = process.env.METRICS_FILE
|
||||
// const logsRoot = process.env.LOGS_ROOT
|
||||
const metricsConsole = (process.env.METRICS_CONSOLE ?? 'false') === 'true'
|
||||
|
||||
const METRICS_UPDATE_INTERVAL = !metricsConsole ? 1000 : 60000
|
||||
|
||||
@@ -135,7 +135,7 @@ export async function initModel (
|
||||
const result = await db.collection(DOMAIN_TX).insertMany(model as Document[])
|
||||
logger.log(`${result.insertedCount} model transactions inserted.`)
|
||||
|
||||
logger.log('creating data...')
|
||||
logger.log('creating data...', transactorUrl)
|
||||
const connection = (await connect(transactorUrl, workspaceId, undefined, {
|
||||
model: 'upgrade'
|
||||
})) as unknown as CoreClient & BackupClient
|
||||
|
||||
@@ -56,6 +56,7 @@ export class ClientSession implements Session {
|
||||
total: StatisticsElement = { find: 0, tx: 0 }
|
||||
current: StatisticsElement = { find: 0, tx: 0 }
|
||||
mins5: StatisticsElement = { find: 0, tx: 0 }
|
||||
measures: { id: string, message: string, time: 0 }[] = []
|
||||
|
||||
constructor (
|
||||
protected readonly broadcast: BroadcastCall,
|
||||
|
||||
+94
-46
@@ -26,7 +26,7 @@ import core, {
|
||||
type WorkspaceId
|
||||
} from '@hcengineering/core'
|
||||
import { unknownError } from '@hcengineering/platform'
|
||||
import { readRequest, type HelloRequest, type HelloResponse, type Response } from '@hcengineering/rpc'
|
||||
import { readRequest, type HelloRequest, type HelloResponse, type Request, type Response } from '@hcengineering/rpc'
|
||||
import type { Pipeline, SessionContext } from '@hcengineering/server-core'
|
||||
import { type Token } from '@hcengineering/server-token'
|
||||
// import WebSocket, { RawData } from 'ws'
|
||||
@@ -155,14 +155,14 @@ class TSessionManager implements SessionManager {
|
||||
}
|
||||
|
||||
async addSession (
|
||||
ctx: MeasureContext,
|
||||
baseCtx: MeasureContext,
|
||||
ws: ConnectionSocket,
|
||||
token: Token,
|
||||
pipelineFactory: PipelineFactory,
|
||||
productId: string,
|
||||
sessionId?: string
|
||||
): Promise<Session> {
|
||||
return await ctx.with('add-session', {}, async (ctx) => {
|
||||
): Promise<{ session: Session, context: MeasureContext } | { upgrade: true }> {
|
||||
return await baseCtx.with('📲 add-session', {}, async (ctx) => {
|
||||
const wsString = toWorkspaceString(token.workspace, '@')
|
||||
|
||||
let workspace = this.workspaces.get(wsString)
|
||||
@@ -170,22 +170,29 @@ class TSessionManager implements SessionManager {
|
||||
workspace = this.workspaces.get(wsString)
|
||||
|
||||
if (workspace === undefined) {
|
||||
workspace = this.createWorkspace(ctx, pipelineFactory, token)
|
||||
workspace = this.createWorkspace(baseCtx, pipelineFactory, token)
|
||||
}
|
||||
|
||||
let pipeline: Pipeline
|
||||
if (token.extra?.model === 'upgrade') {
|
||||
if (workspace.upgrade) {
|
||||
pipeline = await ctx.with('pipeline', {}, async () => await (workspace as Workspace).pipeline)
|
||||
pipeline = await ctx.with(
|
||||
'💤 wait ' + token.workspace.name,
|
||||
{},
|
||||
async () => await (workspace as Workspace).pipeline
|
||||
)
|
||||
} else {
|
||||
pipeline = await this.createUpgradeSession(token, sessionId, ctx, wsString, workspace, pipelineFactory, ws)
|
||||
}
|
||||
} else {
|
||||
if (workspace.upgrade) {
|
||||
ws.close()
|
||||
throw new Error('Upgrade in progress....')
|
||||
return { upgrade: true }
|
||||
}
|
||||
pipeline = await ctx.with('pipeline', {}, async () => await (workspace as Workspace).pipeline)
|
||||
pipeline = await ctx.with(
|
||||
'💤 wait ' + token.workspace.name,
|
||||
{},
|
||||
async () => await (workspace as Workspace).pipeline
|
||||
)
|
||||
}
|
||||
|
||||
const session = this.createSession(token, pipeline)
|
||||
@@ -204,7 +211,7 @@ class TSessionManager implements SessionManager {
|
||||
session.useCompression
|
||||
)
|
||||
}
|
||||
return session
|
||||
return { session, context: workspace.context }
|
||||
})
|
||||
}
|
||||
|
||||
@@ -222,7 +229,7 @@ class TSessionManager implements SessionManager {
|
||||
}
|
||||
// If upgrade client is used.
|
||||
// Drop all existing clients
|
||||
await this.closeAll(ctx, wsString, workspace, 0, 'upgrade')
|
||||
await this.closeAll(wsString, workspace, 0, 'upgrade')
|
||||
// Wipe workspace and update values.
|
||||
if (!workspace.upgrade) {
|
||||
// This is previous workspace, intended to be closed.
|
||||
@@ -238,10 +245,10 @@ class TSessionManager implements SessionManager {
|
||||
}
|
||||
|
||||
broadcastAll (workspace: Workspace, tx: Tx[], targets?: string[]): void {
|
||||
if (workspace?.upgrade ?? false) {
|
||||
if (workspace.upgrade) {
|
||||
return
|
||||
}
|
||||
const ctx = this.ctx.newChild('broadcast-all', {})
|
||||
const ctx = this.ctx.newChild('📬 broadcast-all', {})
|
||||
const sessions = [...workspace.sessions.values()]
|
||||
function send (): void {
|
||||
for (const session of sessions.splice(0, 1)) {
|
||||
@@ -266,9 +273,11 @@ class TSessionManager implements SessionManager {
|
||||
|
||||
private createWorkspace (ctx: MeasureContext, pipelineFactory: PipelineFactory, token: Token): Workspace {
|
||||
const upgrade = token.extra?.model === 'upgrade'
|
||||
const context = ctx.newChild('🧲 ' + token.workspace.name, {})
|
||||
const workspace: Workspace = {
|
||||
context,
|
||||
id: generateId(),
|
||||
pipeline: pipelineFactory(ctx, token.workspace, upgrade, (tx, targets) => {
|
||||
pipeline: pipelineFactory(context, token.workspace, upgrade, (tx, targets) => {
|
||||
this.broadcastAll(workspace, tx, targets)
|
||||
}),
|
||||
sessions: new Map(),
|
||||
@@ -309,13 +318,7 @@ class TSessionManager implements SessionManager {
|
||||
} catch {}
|
||||
}
|
||||
|
||||
async close (
|
||||
ctx: MeasureContext,
|
||||
ws: ConnectionSocket,
|
||||
workspaceId: WorkspaceId,
|
||||
code: number,
|
||||
reason: string
|
||||
): Promise<void> {
|
||||
async close (ws: ConnectionSocket, workspaceId: WorkspaceId, code: number, reason: string): Promise<void> {
|
||||
// if (LOGGING_ENABLED) console.log(workspaceId.name, `closing websocket, code: ${code}, reason: ${reason}`)
|
||||
const wsid = toWorkspaceString(workspaceId)
|
||||
const workspace = this.workspaces.get(wsid)
|
||||
@@ -340,7 +343,7 @@ class TSessionManager implements SessionManager {
|
||||
const user = sessionRef.session.getUser()
|
||||
const another = Array.from(workspace.sessions.values()).findIndex((p) => p.session.getUser() === user)
|
||||
if (another === -1) {
|
||||
await this.setStatus(ctx, sessionRef.session, false)
|
||||
await this.setStatus(workspace.context, sessionRef.session, false)
|
||||
}
|
||||
if (!workspace.upgrade) {
|
||||
// Wait some time for new client to appear before closing workspace.
|
||||
@@ -355,13 +358,7 @@ class TSessionManager implements SessionManager {
|
||||
}
|
||||
}
|
||||
|
||||
async closeAll (
|
||||
ctx: MeasureContext,
|
||||
wsId: string,
|
||||
workspace: Workspace,
|
||||
code: number,
|
||||
reason: 'upgrade' | 'shutdown'
|
||||
): Promise<void> {
|
||||
async closeAll (wsId: string, workspace: Workspace, code: number, reason: 'upgrade' | 'shutdown'): Promise<void> {
|
||||
if (LOGGING_ENABLED) console.timeLog(wsId, `closing workspace ${workspace.id}, code: ${code}, reason: ${reason}`)
|
||||
|
||||
const sessions = Array.from(workspace.sessions)
|
||||
@@ -371,19 +368,10 @@ class TSessionManager implements SessionManager {
|
||||
s.workspaceClosed = true
|
||||
if (reason === 'upgrade') {
|
||||
// Override message handler, to wait for upgrading response from clients.
|
||||
await webSocket.send(
|
||||
ctx,
|
||||
{
|
||||
result: {
|
||||
_class: core.class.TxModelUpgrade
|
||||
}
|
||||
},
|
||||
s.binaryResponseMode,
|
||||
false
|
||||
)
|
||||
await this.sendUpgrade(workspace.context, webSocket, s.binaryResponseMode)
|
||||
}
|
||||
webSocket.close()
|
||||
await this.setStatus(ctx, s, false)
|
||||
await this.setStatus(workspace.context, s, false)
|
||||
}
|
||||
|
||||
if (LOGGING_ENABLED) console.timeLog(wsId, workspace.id, 'Clients disconnected. Closing Workspace...')
|
||||
@@ -403,12 +391,25 @@ class TSessionManager implements SessionManager {
|
||||
console.timeEnd(wsId)
|
||||
}
|
||||
|
||||
private async sendUpgrade (ctx: MeasureContext, webSocket: ConnectionSocket, binary: boolean): Promise<void> {
|
||||
await webSocket.send(
|
||||
ctx,
|
||||
{
|
||||
result: {
|
||||
_class: core.class.TxModelUpgrade
|
||||
}
|
||||
},
|
||||
binary,
|
||||
false
|
||||
)
|
||||
}
|
||||
|
||||
async closeWorkspaces (ctx: MeasureContext): Promise<void> {
|
||||
if (this.checkInterval !== undefined) {
|
||||
clearInterval(this.checkInterval)
|
||||
}
|
||||
for (const w of this.workspaces) {
|
||||
await this.closeAll(ctx, w[0], w[1], 1, 'shutdown')
|
||||
await this.closeAll(w[0], w[1], 1, 'shutdown')
|
||||
}
|
||||
}
|
||||
|
||||
@@ -432,6 +433,7 @@ class TSessionManager implements SessionManager {
|
||||
if (this.workspaces.get(wsid)?.id === wsUID) {
|
||||
this.workspaces.delete(wsid)
|
||||
}
|
||||
workspace.context.end()
|
||||
if (LOGGING_ENABLED) {
|
||||
console.timeLog(workspaceId.name, 'Closed workspace', wsUID)
|
||||
}
|
||||
@@ -459,7 +461,7 @@ class TSessionManager implements SessionManager {
|
||||
if (LOGGING_ENABLED) console.log(workspaceId.name, `server broadcasting to ${workspace.sessions.size} clients...`)
|
||||
|
||||
const sessions = [...workspace.sessions.values()]
|
||||
const ctx = this.ctx.newChild('broadcast', {})
|
||||
const ctx = this.ctx.newChild('📭 broadcast', {})
|
||||
function send (): void {
|
||||
for (const sessionRef of sessions.splice(0, 1)) {
|
||||
if (sessionRef.session.sessionId !== from?.sessionId) {
|
||||
@@ -496,7 +498,7 @@ class TSessionManager implements SessionManager {
|
||||
msg: any,
|
||||
workspace: string
|
||||
): Promise<void> {
|
||||
const userCtx = requestCtx.newChild('client', { workspace }) as SessionContext
|
||||
const userCtx = requestCtx.newChild('📞 client', {}) as SessionContext
|
||||
userCtx.sessionId = service.sessionInstanceId ?? ''
|
||||
|
||||
// Calculate total number of clients
|
||||
@@ -504,8 +506,8 @@ class TSessionManager implements SessionManager {
|
||||
|
||||
const st = Date.now()
|
||||
try {
|
||||
await userCtx.with('handleRequest', {}, async (ctx) => {
|
||||
const request = await ctx.with('read', {}, async () => readRequest(msg, false))
|
||||
await userCtx.with('🧭 handleRequest', {}, async (ctx) => {
|
||||
const request = await ctx.with('📥 read', {}, async () => readRequest(msg, false))
|
||||
if (request.id === -1 && request.method === 'hello') {
|
||||
const hello = request as HelloRequest
|
||||
service.binaryResponseMode = hello.binary ?? false
|
||||
@@ -536,6 +538,10 @@ class TSessionManager implements SessionManager {
|
||||
await ws.send(ctx, helloResponse, false, false)
|
||||
return
|
||||
}
|
||||
if (request.method === 'measure' || request.method === 'measure-done') {
|
||||
await this.handleMeasure<S>(service, request, ctx, ws)
|
||||
return
|
||||
}
|
||||
service.requests.set(reqId, {
|
||||
id: reqId,
|
||||
params: request,
|
||||
@@ -545,10 +551,15 @@ class TSessionManager implements SessionManager {
|
||||
ws.close()
|
||||
return
|
||||
}
|
||||
|
||||
const f = (service as any)[request.method]
|
||||
try {
|
||||
const params = [...request.params]
|
||||
const result = await ctx.with('call', {}, async (callTx) => f.apply(service, [callTx, ...params]))
|
||||
|
||||
const result =
|
||||
service.measureCtx?.ctx !== undefined
|
||||
? await f.apply(service, [service.measureCtx?.ctx, ...params])
|
||||
: await ctx.with('🧨 process', {}, async (callTx) => f.apply(service, [callTx, ...params]))
|
||||
|
||||
const resp: Response<any> = { id: request.id, result }
|
||||
|
||||
@@ -575,6 +586,43 @@ class TSessionManager implements SessionManager {
|
||||
service.requests.delete(reqId)
|
||||
}
|
||||
}
|
||||
|
||||
private async handleMeasure<S extends Session>(
|
||||
service: S,
|
||||
request: Request<any[]>,
|
||||
ctx: MeasureContext,
|
||||
ws: ConnectionSocket
|
||||
): Promise<void> {
|
||||
let serverTime = 0
|
||||
if (request.method === 'measure') {
|
||||
service.measureCtx = { ctx: ctx.newChild('📶 ' + request.params[0], {}), time: Date.now() }
|
||||
} else {
|
||||
if (service.measureCtx !== undefined) {
|
||||
serverTime = Date.now() - service.measureCtx.time
|
||||
service.measureCtx.ctx.end(serverTime)
|
||||
}
|
||||
}
|
||||
try {
|
||||
const resp: Response<any> = { id: request.id, result: request.method === 'measure' ? 'started' : serverTime }
|
||||
|
||||
await handleSend(
|
||||
ctx,
|
||||
ws,
|
||||
resp,
|
||||
this.sessions.size < 100 ? 10000 : 1001,
|
||||
service.binaryResponseMode,
|
||||
service.useCompression
|
||||
)
|
||||
} catch (err: any) {
|
||||
if (LOGGING_ENABLED) console.error(err)
|
||||
const resp: Response<any> = {
|
||||
id: request.id,
|
||||
error: unknownError(err),
|
||||
result: JSON.parse(JSON.stringify(err?.stack))
|
||||
}
|
||||
await ws.send(ctx, resp, service.binaryResponseMode, service.useCompression)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function handleSend (
|
||||
|
||||
@@ -169,11 +169,11 @@ export function startHttpServer (
|
||||
if (ws.readyState !== ws.OPEN) {
|
||||
return
|
||||
}
|
||||
const smsg = await ctx.with('serialize', {}, async () => serialize(msg, binary))
|
||||
const smsg = await ctx.with('📦 serialize', {}, async () => serialize(msg, binary))
|
||||
|
||||
ctx.measure('send-data', smsg.length)
|
||||
|
||||
await ctx.with('socket-send', {}, async (ctx) => {
|
||||
await ctx.with('📤 socket-send', {}, async (ctx) => {
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
ws.send(smsg, { binary, compress: compression }, (err) => {
|
||||
if (err != null) {
|
||||
@@ -191,6 +191,10 @@ export function startHttpServer (
|
||||
buffer?.push(msg)
|
||||
})
|
||||
const session = await sessions.addSession(ctx, cs, token, pipelineFactory, productId, sessionId)
|
||||
if ('upgrade' in session) {
|
||||
cs.close()
|
||||
return
|
||||
}
|
||||
// eslint-disable-next-line @typescript-eslint/no-misused-promises
|
||||
ws.on('message', (msg: RawData) => {
|
||||
let buff: any | undefined
|
||||
@@ -200,22 +204,21 @@ export function startHttpServer (
|
||||
buff = Buffer.concat(msg).toString()
|
||||
}
|
||||
if (buff !== undefined) {
|
||||
void handleRequest(ctx, session, cs, buff, token.workspace.name)
|
||||
void handleRequest(session.context, session.session, cs, buff, token.workspace.name)
|
||||
}
|
||||
})
|
||||
// eslint-disable-next-line @typescript-eslint/no-misused-promises
|
||||
ws.on('close', (code: number, reason: Buffer) => {
|
||||
if (session.workspaceClosed ?? false) {
|
||||
if (session.session.workspaceClosed ?? false) {
|
||||
return
|
||||
}
|
||||
// remove session after 1seconds, give a time to reconnect.
|
||||
// if (LOGGING_ENABLED) console.log(token.workspace.name, `client "${token.email}" closed ${code === 1000 ? 'normally' : 'abnormally'}`)
|
||||
void sessions.close(ctx, cs, token.workspace, code, reason.toString())
|
||||
void sessions.close(cs, token.workspace, code, reason.toString())
|
||||
})
|
||||
const b = buffer
|
||||
buffer = undefined
|
||||
for (const msg of b) {
|
||||
await handleRequest(ctx, session, cs, msg, token.workspace.name)
|
||||
await handleRequest(session.context, session.session, cs, msg, token.workspace.name)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -226,7 +229,6 @@ export function startHttpServer (
|
||||
try {
|
||||
const payload = decodeToken(token ?? '')
|
||||
const sessionId = url.searchParams.get('sessionId')
|
||||
// if (LOGGING_ENABLED) console.log(payload.workspace.name, 'client connected with payload', payload, sessionId)
|
||||
|
||||
if (payload.workspace.productId !== productId) {
|
||||
throw new Error('Invalid workspace product')
|
||||
|
||||
+6
-15
@@ -59,6 +59,8 @@ export interface Session {
|
||||
total: StatisticsElement
|
||||
current: StatisticsElement
|
||||
mins5: StatisticsElement
|
||||
|
||||
measureCtx?: { ctx: MeasureContext, time: number }
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -107,6 +109,7 @@ export function disableLogging (): void {
|
||||
* @public
|
||||
*/
|
||||
export interface Workspace {
|
||||
context: MeasureContext
|
||||
id: string
|
||||
pipeline: Promise<Pipeline>
|
||||
sessions: Map<string, { session: Session, socket: ConnectionSocket }>
|
||||
@@ -130,25 +133,13 @@ export interface SessionManager {
|
||||
pipelineFactory: PipelineFactory,
|
||||
productId: string,
|
||||
sessionId?: string
|
||||
) => Promise<Session>
|
||||
) => Promise<{ session: Session, context: MeasureContext } | { upgrade: true }>
|
||||
|
||||
broadcastAll: (workspace: Workspace, tx: Tx[], targets?: string[]) => void
|
||||
|
||||
close: (
|
||||
ctx: MeasureContext,
|
||||
ws: ConnectionSocket,
|
||||
workspaceId: WorkspaceId,
|
||||
code: number,
|
||||
reason: string
|
||||
) => Promise<void>
|
||||
close: (ws: ConnectionSocket, workspaceId: WorkspaceId, code: number, reason: string) => Promise<void>
|
||||
|
||||
closeAll: (
|
||||
ctx: MeasureContext,
|
||||
wsId: string,
|
||||
workspace: Workspace,
|
||||
code: number,
|
||||
reason: 'upgrade' | 'shutdown'
|
||||
) => Promise<void>
|
||||
closeAll: (wsId: string, workspace: Workspace, code: number, reason: 'upgrade' | 'shutdown') => Promise<void>
|
||||
|
||||
closeWorkspaces: (ctx: MeasureContext) => Promise<void>
|
||||
|
||||
|
||||
Reference in New Issue
Block a user