mirror of
https://github.com/hcengineering/platform.git
synced 2026-09-28 04:25:03 +02:00
UBERF-6374: Improve server logging and improve startup performance (#5210)
This commit is contained in:
@@ -50,11 +50,6 @@ export interface RawDBAdapter {
|
||||
* @public
|
||||
*/
|
||||
export interface DbAdapter {
|
||||
/**
|
||||
* Method called after hierarchy is ready to use.
|
||||
*/
|
||||
init: (model: Tx[]) => Promise<void>
|
||||
|
||||
createIndexes: (domain: Domain, config: Pick<IndexingConfiguration<Doc>, 'indexes'>) => Promise<void>
|
||||
removeOldIndex: (domain: Domain, deletePattern: RegExp, keepPattern: RegExp) => Promise<void>
|
||||
|
||||
@@ -83,7 +78,7 @@ export interface DbAdapter {
|
||||
* @public
|
||||
*/
|
||||
export interface TxAdapter extends DbAdapter {
|
||||
getModel: () => Promise<Tx[]>
|
||||
getModel: (ctx: MeasureContext) => Promise<Tx[]>
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -63,7 +63,7 @@ export class FullTextIndex implements WithFind {
|
||||
readonly indexer: FullTextIndexPipeline,
|
||||
private readonly upgrade: boolean
|
||||
) {
|
||||
if (!upgrade) {
|
||||
if (!this.upgrade) {
|
||||
// Schedule indexing after consistency check
|
||||
void this.indexer.startIndexing()
|
||||
}
|
||||
|
||||
@@ -91,13 +91,12 @@ export class FullTextIndexPipeline implements FullTextPipeline {
|
||||
}
|
||||
|
||||
async cancel (): Promise<void> {
|
||||
console.log(this.workspace.name, 'Cancel indexing', this.indexId)
|
||||
this.cancelling = true
|
||||
clearTimeout(this.skippedReiterationTimeout)
|
||||
this.triggerIndexing()
|
||||
await this.indexing
|
||||
await this.flush(true)
|
||||
console.log(this.workspace.name, 'Indexing canceled', this.indexId)
|
||||
await this.metrics.info('Cancel indexing', { workspace: this.workspace.name, indexId: this.indexId })
|
||||
}
|
||||
|
||||
async markRemove (doc: DocIndexState): Promise<void> {
|
||||
@@ -336,7 +335,10 @@ export class FullTextIndexPipeline implements FullTextPipeline {
|
||||
try {
|
||||
this.hierarchy.getClass(core.class.DocIndexState)
|
||||
} catch (err: any) {
|
||||
console.log(this.workspace.name, 'Models is not upgraded to support indexer', this.indexId)
|
||||
await this.metrics.info('Models is not upgraded to support indexer', {
|
||||
indexId: this.indexId,
|
||||
workspace: this.workspace.name
|
||||
})
|
||||
return
|
||||
}
|
||||
await this.metrics.with('init-states', {}, async () => {
|
||||
@@ -367,12 +369,12 @@ export class FullTextIndexPipeline implements FullTextPipeline {
|
||||
_classes.forEach((it) => this.broadcastClasses.add(it))
|
||||
|
||||
if (this.triggerCounts > 0) {
|
||||
console.log('No wait, trigger counts', this.triggerCounts)
|
||||
await this.metrics.info('No wait, trigger counts', { triggerCount: this.triggerCounts })
|
||||
}
|
||||
|
||||
if (this.toIndex.size === 0 && this.stageChanged === 0 && this.triggerCounts === 0) {
|
||||
if (this.toIndex.size === 0) {
|
||||
console.log(this.workspace.name, 'Indexing complete', this.indexId)
|
||||
await this.metrics.info('Indexing complete', { indexId: this.indexId, workspace: this.workspace.name })
|
||||
}
|
||||
if (!this.cancelling) {
|
||||
// We need to send index update event
|
||||
@@ -398,7 +400,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
|
||||
}
|
||||
}
|
||||
}
|
||||
console.log(this.workspace.name, 'Exit indexer', this.indexId)
|
||||
await this.metrics.info('Exit indexer', { indexId: this.indexId, workspace: this.workspace.name })
|
||||
}
|
||||
|
||||
private async processIndex (ctx: MeasureContext): Promise<Ref<Class<Doc>>[]> {
|
||||
@@ -470,13 +472,12 @@ export class FullTextIndexPipeline implements FullTextPipeline {
|
||||
}
|
||||
|
||||
if (result.length > 0) {
|
||||
console.log(
|
||||
this.workspace.name,
|
||||
`Full text: Indexing ${this.indexId} ${st.stageId}`,
|
||||
Object.entries(this.currentStages)
|
||||
.map((it) => `${it[0]}:${it[1]}`)
|
||||
.join(' ')
|
||||
)
|
||||
await this.metrics.info('Full text: Indexing', {
|
||||
indexId: this.indexId,
|
||||
stageId: st.stageId,
|
||||
workspace: this.workspace.name,
|
||||
...this.currentStages
|
||||
})
|
||||
} else {
|
||||
// Nothing to index, check on next cycle.
|
||||
break
|
||||
@@ -528,7 +529,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
|
||||
}
|
||||
}
|
||||
} catch (err: any) {
|
||||
console.error(err)
|
||||
await this.metrics.error('error during index', { error: err })
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
@@ -38,7 +38,6 @@ import { type DbAdapter } from './adapter'
|
||||
* @public
|
||||
*/
|
||||
export class DummyDbAdapter implements DbAdapter {
|
||||
async init (model: Tx[]): Promise<void> {}
|
||||
async findAll<T extends Doc>(
|
||||
ctx: MeasureContext,
|
||||
_class: Ref<Class<T>>,
|
||||
@@ -99,16 +98,6 @@ class InMemoryAdapter extends DummyDbAdapter implements DbAdapter {
|
||||
async tx (ctx: MeasureContext, ...tx: Tx[]): Promise<TxResult[]> {
|
||||
return await this.modeldb.tx(...tx)
|
||||
}
|
||||
|
||||
async init (model: Tx[]): Promise<void> {
|
||||
for (const tx of model) {
|
||||
try {
|
||||
await this.modeldb.tx(tx)
|
||||
} catch (err: any) {
|
||||
console.error('skip broken TX', err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -31,8 +31,8 @@ import {
|
||||
type Tx,
|
||||
type TxResult
|
||||
} from '@hcengineering/core'
|
||||
import { createServerStorage } from './server'
|
||||
import { type DbConfiguration } from './configuration'
|
||||
import { createServerStorage } from './server'
|
||||
import {
|
||||
type BroadcastFunc,
|
||||
type HandledBroadcastFunc,
|
||||
@@ -67,12 +67,12 @@ export async function createPipeline (
|
||||
}
|
||||
})
|
||||
)
|
||||
const pipeline = ctx.with(
|
||||
'create pipeline',
|
||||
{},
|
||||
async (ctx) => await PipelineImpl.create(ctx, storage, constructors, broadcast)
|
||||
const pipelineResult = await PipelineImpl.create(
|
||||
ctx.newChild('pipeline-operations', {}),
|
||||
storage,
|
||||
constructors,
|
||||
broadcast
|
||||
)
|
||||
const pipelineResult = await pipeline
|
||||
broadcastHook = (tx, targets) => {
|
||||
return pipelineResult.handleBroadcast(tx, targets)
|
||||
}
|
||||
|
||||
@@ -56,18 +56,20 @@ export async function createServerStorage (
|
||||
|
||||
const storageAdapter = conf.storageFactory?.()
|
||||
|
||||
for (const key in conf.adapters) {
|
||||
const adapterConf = conf.adapters[key]
|
||||
adapters.set(
|
||||
key,
|
||||
await adapterConf.factory(ctx, hierarchy, adapterConf.url, conf.workspace, modelDb, storageAdapter)
|
||||
)
|
||||
}
|
||||
await ctx.with('create-adapters', {}, async (ctx) => {
|
||||
for (const key in conf.adapters) {
|
||||
const adapterConf = conf.adapters[key]
|
||||
adapters.set(
|
||||
key,
|
||||
await adapterConf.factory(ctx, hierarchy, adapterConf.url, conf.workspace, modelDb, storageAdapter)
|
||||
)
|
||||
}
|
||||
})
|
||||
|
||||
const txAdapter = adapters.get(conf.domains[DOMAIN_TX]) as TxAdapter
|
||||
|
||||
const model = await ctx.with('get model', {}, async (ctx) => {
|
||||
const model = await txAdapter.getModel()
|
||||
const model = await ctx.with('fetch-model', {}, async (ctx) => await txAdapter.getModel(ctx))
|
||||
for (const tx of model) {
|
||||
try {
|
||||
hierarchy.tx(tx)
|
||||
@@ -76,22 +78,10 @@ export async function createServerStorage (
|
||||
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)
|
||||
}
|
||||
}
|
||||
modelDb.addTxes(ctx, model, false)
|
||||
return model
|
||||
})
|
||||
|
||||
for (const [adn, adapter] of adapters) {
|
||||
await ctx.with('init-adapter', { name: adn }, async (ctx) => {
|
||||
await adapter.init(model)
|
||||
})
|
||||
}
|
||||
|
||||
const fulltextAdapter = await ctx.with(
|
||||
'create full text adapter',
|
||||
{},
|
||||
|
||||
@@ -60,10 +60,10 @@ import crypto from 'node:crypto'
|
||||
import { type DbAdapter } from '../adapter'
|
||||
import { type FullTextIndex } from '../fulltext'
|
||||
import serverCore from '../plugin'
|
||||
import { type ServiceAdaptersManager } from '../service'
|
||||
import { type StorageAdapter } from '../storage'
|
||||
import { type Triggers } from '../triggers'
|
||||
import type { FullTextAdapter, ObjectDDParticipant, ServerStorageOptions, TriggerControl } from '../types'
|
||||
import { type StorageAdapter } from '../storage'
|
||||
import { type ServiceAdaptersManager } from '../service'
|
||||
|
||||
export class TServerStorage implements ServerStorage {
|
||||
private readonly fulltext: FullTextIndex
|
||||
@@ -137,15 +137,11 @@ export class TServerStorage implements ServerStorage {
|
||||
}
|
||||
|
||||
async close (): Promise<void> {
|
||||
console.timeLog(this.workspace.name, 'closing')
|
||||
await this.fulltext.close()
|
||||
console.timeLog(this.workspace.name, 'closing adapters')
|
||||
for (const o of this.adapters.values()) {
|
||||
await o.close()
|
||||
}
|
||||
console.timeLog(this.workspace.name, 'closing fulltext')
|
||||
await this.fulltextAdapter.close()
|
||||
console.timeLog(this.workspace.name, 'closing service adapters')
|
||||
await this.serviceAdaptersManager.close()
|
||||
}
|
||||
|
||||
@@ -199,7 +195,7 @@ export class TServerStorage implements ServerStorage {
|
||||
const txCUD = TxProcessor.extractTx(tx) as TxCUD<Doc>
|
||||
if (!this.hierarchy.isDerived(txCUD._class, core.class.TxCUD)) {
|
||||
// Skip unsupported tx
|
||||
console.error('Unsupported transaction', tx)
|
||||
await ctx.error('Unsupported transaction', tx)
|
||||
continue
|
||||
}
|
||||
const domain = this.hierarchy.getDomain(txCUD.objectClass)
|
||||
@@ -397,7 +393,7 @@ export class TServerStorage implements ServerStorage {
|
||||
{ clazz, query, options }
|
||||
)
|
||||
if (Date.now() - st > 1000) {
|
||||
console.error('FindAll', Date.now() - st, clazz, query, options)
|
||||
await ctx.error('FindAll', { time: Date.now() - st, clazz, query, options })
|
||||
}
|
||||
return result
|
||||
}
|
||||
@@ -794,7 +790,7 @@ export class TServerStorage implements ServerStorage {
|
||||
await fx()
|
||||
}
|
||||
} catch (err: any) {
|
||||
console.log(err)
|
||||
await ctx.error('error process tx', { error: err })
|
||||
throw err
|
||||
} finally {
|
||||
onEnds.forEach((p) => {
|
||||
|
||||
@@ -59,8 +59,6 @@ class ElasticDataAdapter implements DbAdapter {
|
||||
return []
|
||||
}
|
||||
|
||||
async init (model: Tx[]): Promise<void> {}
|
||||
|
||||
async createIndexes (domain: Domain, config: Pick<IndexingConfiguration<Doc>, 'indexes'>): Promise<void> {}
|
||||
async removeOldIndex (domain: Domain, deletePattern: RegExp, keepPattern: RegExp): Promise<void> {}
|
||||
|
||||
|
||||
@@ -292,7 +292,9 @@ export function start (
|
||||
const token = req.query.token as string
|
||||
const payload = decodeToken(token)
|
||||
const admin = payload.extra?.admin === 'true'
|
||||
res.writeHead(200, { 'Content-Type': 'application/json' })
|
||||
res.status(200)
|
||||
res.setHeader('Content-Type', 'application/json')
|
||||
res.setHeader('Cache-Control', cacheControlNoCache)
|
||||
|
||||
const json = JSON.stringify({
|
||||
metrics: metricsAggregate((ctx as any).metrics),
|
||||
@@ -301,7 +303,6 @@ export function start (
|
||||
},
|
||||
admin
|
||||
})
|
||||
res.set('Cache-Control', 'private, no-cache')
|
||||
res.end(json)
|
||||
} catch (err) {
|
||||
console.error(err)
|
||||
|
||||
@@ -62,6 +62,8 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
|
||||
private spaceMeasureCtx!: MeasureContext
|
||||
|
||||
private spaceSecurityInit: Promise<void> | undefined
|
||||
|
||||
private readonly systemSpaces = [
|
||||
core.space.Configuration,
|
||||
core.space.DerivedTx,
|
||||
@@ -86,7 +88,7 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
): Promise<SpaceSecurityMiddleware> {
|
||||
const res = new SpaceSecurityMiddleware(broadcast, storage, next)
|
||||
res.spaceMeasureCtx = ctx.newChild('space chain', {})
|
||||
await res.init(res.spaceMeasureCtx)
|
||||
res.spaceSecurityInit = res.init(res.spaceMeasureCtx)
|
||||
return res
|
||||
}
|
||||
|
||||
@@ -124,6 +126,13 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
this.publicSpaces = spaces.filter((it) => !it.private).map((p) => p._id)
|
||||
}
|
||||
|
||||
async waitInit (): Promise<void> {
|
||||
if (this.spaceSecurityInit !== undefined) {
|
||||
await this.spaceSecurityInit
|
||||
this.spaceSecurityInit = undefined
|
||||
}
|
||||
}
|
||||
|
||||
private removeMemberSpace (member: Ref<Account>, space: Ref<Space>): void {
|
||||
const arr = this.allowedSpaces[member]
|
||||
if (arr !== undefined) {
|
||||
@@ -240,6 +249,8 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
}
|
||||
|
||||
private async handleUpdate (ctx: SessionContext, tx: TxCUD<Space>): Promise<void> {
|
||||
await this.waitInit()
|
||||
|
||||
const updateDoc = tx as TxUpdateDoc<Space>
|
||||
if (!this.storage.hierarchy.isDerived(updateDoc.objectClass, core.class.Space)) return
|
||||
|
||||
@@ -285,6 +296,7 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
}
|
||||
|
||||
private async handleTx (ctx: SessionContext, tx: TxCUD<Space>): Promise<void> {
|
||||
await this.waitInit()
|
||||
if (tx._class === core.class.TxCreateDoc) {
|
||||
this.handleCreate(tx)
|
||||
} else if (tx._class === core.class.TxUpdateDoc) {
|
||||
@@ -370,6 +382,7 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
}
|
||||
|
||||
async tx (ctx: SessionContext, tx: Tx): Promise<TxMiddlewareResult> {
|
||||
await this.waitInit()
|
||||
const account = await getUser(this.storage, ctx)
|
||||
if (account.role === AccountRole.Guest) {
|
||||
throw new PlatformError(new Status(Severity.ERROR, platform.status.Forbidden, {}))
|
||||
@@ -385,6 +398,7 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
|
||||
handleBroadcast (tx: Tx[], targets?: string[]): Tx[] {
|
||||
const process = async (): Promise<void> => {
|
||||
await this.waitInit()
|
||||
for (const t of tx) {
|
||||
if (this.storage.hierarchy.isDerived(t._class, core.class.TxCUD)) {
|
||||
await this.processTxSpaceDomain(t as TxCUD<Doc>)
|
||||
@@ -476,6 +490,8 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
query: DocumentQuery<T>,
|
||||
options?: FindOptions<T>
|
||||
): Promise<FindResult<T>> {
|
||||
await this.waitInit()
|
||||
|
||||
const domain = this.storage.hierarchy.getDomain(_class)
|
||||
const newQuery = query
|
||||
const account = await getUser(this.storage, ctx)
|
||||
@@ -509,6 +525,7 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
|
||||
query: SearchQuery,
|
||||
options: SearchOptions
|
||||
): Promise<SearchResult> {
|
||||
await this.waitInit()
|
||||
const newQuery = { ...query }
|
||||
const account = await getUser(this.storage, ctx)
|
||||
if (!isSystem(account)) {
|
||||
|
||||
@@ -1264,12 +1264,21 @@ class MongoTxAdapter extends MongoAdapterBase implements TxAdapter {
|
||||
return this.txColl
|
||||
}
|
||||
|
||||
async getModel (): Promise<Tx[]> {
|
||||
const cursor = this.db
|
||||
.collection(DOMAIN_TX)
|
||||
.find<Tx>({ objectSpace: core.space.Model })
|
||||
.sort({ _id: 1, modifiedOn: 1 })
|
||||
const model = await toArray(cursor)
|
||||
async getModel (ctx: MeasureContext): Promise<Tx[]> {
|
||||
const modelProjection = {
|
||||
'%hash%': 0,
|
||||
objectSpace: 0,
|
||||
createdBy: 0,
|
||||
space: 0
|
||||
}
|
||||
const cursor = await ctx.with('find', {}, async () =>
|
||||
this.db
|
||||
.collection<Tx>(DOMAIN_TX)
|
||||
.find({ objectSpace: core.space.Model })
|
||||
.sort({ _id: 1, modifiedOn: 1 })
|
||||
.project<Tx>(modelProjection)
|
||||
)
|
||||
const model = await ctx.with('to-array', {}, async () => await toArray<Tx>(cursor))
|
||||
// We need to put all core.account.System transactions first
|
||||
const systemTx: Tx[] = []
|
||||
const userTx: Tx[] = []
|
||||
@@ -1284,7 +1293,6 @@ class MongoTxAdapter extends MongoAdapterBase implements TxAdapter {
|
||||
(tx as TxCUD<Doc>).objectClass === 'contact:class:EmployeeAccount')
|
||||
)
|
||||
}
|
||||
|
||||
model.forEach((tx) => (tx.modifiedBy === core.account.System && !isPersonAccount(tx) ? systemTx : userTx).push(tx))
|
||||
return systemTx.concat(userTx)
|
||||
}
|
||||
|
||||
@@ -17,7 +17,7 @@ let metricsContext: MeasureContext | undefined
|
||||
/**
|
||||
* @public
|
||||
*/
|
||||
export function getMetricsContext (): MeasureContext {
|
||||
export function getMetricsContext (factory?: () => MeasureMetricsContext): MeasureContext {
|
||||
if (metricsContext !== undefined) {
|
||||
return metricsContext
|
||||
}
|
||||
@@ -25,7 +25,11 @@ export function getMetricsContext (): MeasureContext {
|
||||
console.info('please provide apm server url for monitoring')
|
||||
|
||||
const metrics = newMetrics()
|
||||
metricsContext = new MeasureMetricsContext('System', {}, {}, metrics)
|
||||
if (factory !== undefined) {
|
||||
metricsContext = factory()
|
||||
} else {
|
||||
metricsContext = new MeasureMetricsContext('System', {}, {}, metrics)
|
||||
}
|
||||
|
||||
if (metricsFile !== undefined || metricsConsole) {
|
||||
console.info('storing measurements into local file', metricsFile)
|
||||
|
||||
@@ -55,8 +55,6 @@ class StorageBlobAdapter implements DbAdapter {
|
||||
return []
|
||||
}
|
||||
|
||||
async init (model: Tx[]): Promise<void> {}
|
||||
|
||||
async createIndexes (domain: Domain, config: Pick<IndexingConfiguration<Doc>, 'indexes'>): Promise<void> {}
|
||||
async removeOldIndex (domain: Domain, deletePattern: RegExp, keepPattern: RegExp): Promise<void> {}
|
||||
|
||||
|
||||
@@ -268,7 +268,7 @@ async function fetchModelFromMongo (
|
||||
|
||||
const txAdapter = await createMongoTxAdapter(ctx, hierarchy, mongodbUri, workspaceId, modelDb)
|
||||
|
||||
const model = await ctx.with('get-model', {}, async () => await txAdapter.getModel())
|
||||
const model = await ctx.with('get-model', {}, async (ctx) => await txAdapter.getModel(ctx))
|
||||
|
||||
await ctx.with('build local model', {}, async () => {
|
||||
for (const tx of model) {
|
||||
|
||||
@@ -42,12 +42,9 @@ import {
|
||||
import { type SessionContext } from '@hcengineering/server-core'
|
||||
import { ClientSession } from '../client'
|
||||
import { startHttpServer } from '../server_http'
|
||||
import { disableLogging } from '../types'
|
||||
import { genMinModel } from './minmodel'
|
||||
|
||||
describe('server', () => {
|
||||
disableLogging()
|
||||
|
||||
async function getModelDb (): Promise<ModelDb> {
|
||||
const txes = genMinModel()
|
||||
const hierarchy = new Hierarchy()
|
||||
|
||||
@@ -48,6 +48,7 @@ import { type BroadcastCall, type Session, type SessionRequest, type StatisticsE
|
||||
* @public
|
||||
*/
|
||||
export class ClientSession implements Session {
|
||||
createTime = Date.now()
|
||||
requests = new Map<string, SessionRequest>()
|
||||
binaryResponseMode: boolean = false
|
||||
useCompression: boolean = true
|
||||
@@ -81,7 +82,7 @@ export class ClientSession implements Session {
|
||||
}
|
||||
|
||||
async loadModel (ctx: MeasureContext, lastModelTx: Timestamp, hash?: string): Promise<Tx[] | LoadModelResponse> {
|
||||
return await this._pipeline.storage.loadModel(lastModelTx, hash)
|
||||
return await ctx.with('load-model', {}, async () => await this._pipeline.storage.loadModel(lastModelTx, hash))
|
||||
}
|
||||
|
||||
async getAccount (ctx: MeasureContext): Promise<Account> {
|
||||
|
||||
+107
-60
@@ -145,7 +145,7 @@ class TSessionManager implements SessionManager {
|
||||
const now = Date.now()
|
||||
const diff = now - s[1].session.lastRequest
|
||||
if (diff > 60000 && this.ticks % 10 === 0) {
|
||||
console.log('session hang, closing...', h[0], s[1].session.getUser())
|
||||
void this.ctx.error('session hang, closing...', { sessionId: h[0], user: s[1].session.getUser() })
|
||||
void this.close(s[1].socket, h[1].workspaceId, 1001, 'CLIENT_HANGOUT')
|
||||
continue
|
||||
}
|
||||
@@ -160,7 +160,11 @@ class TSessionManager implements SessionManager {
|
||||
|
||||
for (const r of s[1].session.requests.values()) {
|
||||
if (now - r.start > 30000) {
|
||||
console.log(h[0], 'request hang found, 30sec', h[0], s[1].session.getUser(), r.params)
|
||||
void this.ctx.info('request hang found, 30sec', {
|
||||
sessionId: h[0],
|
||||
user: s[1].session.getUser(),
|
||||
...r.params
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -212,8 +216,9 @@ class TSessionManager implements SessionManager {
|
||||
return await baseCtx.with('📲 add-session', {}, async (ctx) => {
|
||||
const wsString = toWorkspaceString(token.workspace, '@')
|
||||
|
||||
let workspaceInfo =
|
||||
let workspaceInfo = await ctx.with('check-token', {}, async (ctx) =>
|
||||
accountsUrl !== '' ? await this.getWorkspaceInfo(accountsUrl, rawToken) : this.wsFromToken(token)
|
||||
)
|
||||
if (workspaceInfo === undefined && token.extra?.admin !== 'true') {
|
||||
// No access to workspace for token.
|
||||
return { error: new Error(`No access to workspace for token ${token.email} ${token.workspace.name}`) }
|
||||
@@ -222,6 +227,10 @@ class TSessionManager implements SessionManager {
|
||||
}
|
||||
|
||||
let workspace = this.workspaces.get(wsString)
|
||||
if (workspace?.closeTimeout !== undefined) {
|
||||
await ctx.info('Cancel workspace warm close', { wsString })
|
||||
clearTimeout(workspace?.closeTimeout)
|
||||
}
|
||||
await workspace?.closing
|
||||
workspace = this.workspaces.get(wsString)
|
||||
if (sessionId !== undefined && workspace?.sessions?.has(sessionId) === true) {
|
||||
@@ -278,7 +287,9 @@ class TSessionManager implements SessionManager {
|
||||
this.sessions.set(ws.id, { session, socket: ws })
|
||||
// We need to delete previous session with Id if found.
|
||||
workspace.sessions.set(session.sessionId, { session, socket: ws })
|
||||
await ctx.with('set-status', {}, () => this.setStatus(ctx, session, true))
|
||||
|
||||
// We do not need to wait for set-status, just return session to client
|
||||
void ctx.with('set-status', {}, (ctx) => this.setStatus(ctx, session, true))
|
||||
|
||||
if (this.timeMinutes > 0) {
|
||||
void ws.send(
|
||||
@@ -316,7 +327,7 @@ class TSessionManager implements SessionManager {
|
||||
workspaceName: string
|
||||
): Promise<Pipeline> {
|
||||
if (LOGGING_ENABLED) {
|
||||
console.log(workspaceName, 'reloading workspace', JSON.stringify(token))
|
||||
await ctx.info('reloading workspace', { workspaceName, token: JSON.stringify(token) })
|
||||
}
|
||||
// If upgrade client is used.
|
||||
// Drop all existing clients
|
||||
@@ -351,12 +362,16 @@ class TSessionManager implements SessionManager {
|
||||
for (const session of sessions.splice(0, 1)) {
|
||||
if (targets !== undefined && !targets.includes(session.session.getUser())) continue
|
||||
for (const _tx of tx) {
|
||||
void session.socket.send(
|
||||
ctx,
|
||||
{ result: _tx },
|
||||
session.session.binaryResponseMode,
|
||||
session.session.useCompression
|
||||
)
|
||||
try {
|
||||
void session.socket.send(
|
||||
ctx,
|
||||
{ result: _tx },
|
||||
session.session.binaryResponseMode,
|
||||
session.session.useCompression
|
||||
)
|
||||
} catch (err: any) {
|
||||
void ctx.error('error during send', { error: err })
|
||||
}
|
||||
}
|
||||
}
|
||||
if (sessions.length > 0) {
|
||||
@@ -377,11 +392,12 @@ class TSessionManager implements SessionManager {
|
||||
): Workspace {
|
||||
const upgrade = token.extra?.model === 'upgrade'
|
||||
const context = ctx.newChild('🧲 session', {})
|
||||
const pipelineCtx = context.newChild('🧲 pipeline-factory', {})
|
||||
const workspace: Workspace = {
|
||||
context,
|
||||
id: generateId(),
|
||||
pipeline: pipelineFactory(
|
||||
context,
|
||||
pipelineCtx,
|
||||
{ ...token.workspace, workspaceUrl, workspaceName },
|
||||
upgrade,
|
||||
(tx, targets) => {
|
||||
@@ -393,8 +409,6 @@ class TSessionManager implements SessionManager {
|
||||
workspaceId: token.workspace,
|
||||
workspaceName
|
||||
}
|
||||
if (LOGGING_ENABLED) console.time(workspaceName)
|
||||
if (LOGGING_ENABLED) console.timeLog(workspaceName, 'Creating Workspace:', workspace.id)
|
||||
this.workspaces.set(toWorkspaceString(token.workspace), workspace)
|
||||
return workspace
|
||||
}
|
||||
@@ -429,11 +443,12 @@ class TSessionManager implements SessionManager {
|
||||
}
|
||||
|
||||
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)
|
||||
if (workspace === undefined) {
|
||||
if (LOGGING_ENABLED) console.error(new Error('internal: cannot find sessions'))
|
||||
if (LOGGING_ENABLED) {
|
||||
await this.ctx.error('internal: cannot find sessions', { id: ws.id, workspace: workspaceId.name, code, reason })
|
||||
}
|
||||
return
|
||||
}
|
||||
const sessionRef = this.sessions.get(ws.id)
|
||||
@@ -458,7 +473,9 @@ class TSessionManager implements SessionManager {
|
||||
if (!workspace.upgrade) {
|
||||
// Wait some time for new client to appear before closing workspace.
|
||||
if (workspace.sessions.size === 0) {
|
||||
setTimeout(() => {
|
||||
clearTimeout(workspace.closeTimeout)
|
||||
void this.ctx.info('schedule warm closing', { workspace: workspace.workspaceName, wsid })
|
||||
workspace.closeTimeout = setTimeout(() => {
|
||||
void this.performWorkspaceCloseCheck(workspace, workspaceId, wsid)
|
||||
}, this.timeouts.shutdownWarmTimeout)
|
||||
}
|
||||
@@ -469,7 +486,15 @@ class TSessionManager implements SessionManager {
|
||||
}
|
||||
|
||||
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}`)
|
||||
if (LOGGING_ENABLED) {
|
||||
await this.ctx.info('closing workspace', {
|
||||
workspace: workspace.id,
|
||||
wsName: workspace.workspaceName,
|
||||
code,
|
||||
reason,
|
||||
wsId
|
||||
})
|
||||
}
|
||||
|
||||
const sessions = Array.from(workspace.sessions)
|
||||
workspace.sessions = new Map()
|
||||
@@ -484,21 +509,30 @@ class TSessionManager implements SessionManager {
|
||||
await this.setStatus(workspace.context, s, false)
|
||||
}
|
||||
|
||||
if (LOGGING_ENABLED) console.timeLog(wsId, workspace.id, 'Clients disconnected. Closing Workspace...')
|
||||
if (LOGGING_ENABLED) {
|
||||
await this.ctx.info('Clients disconnected. Closing Workspace...', {
|
||||
wsId,
|
||||
workspace: workspace.id,
|
||||
wsName: workspace.workspaceName
|
||||
})
|
||||
}
|
||||
await Promise.all(sessions.map((s) => closeS(s[1].session, s[1].socket)))
|
||||
|
||||
const closePipeline = async (): Promise<void> => {
|
||||
try {
|
||||
if (LOGGING_ENABLED) console.timeLog(wsId, 'closing pipeline')
|
||||
await (await workspace.pipeline).close()
|
||||
if (LOGGING_ENABLED) console.timeLog(wsId, 'closing pipeline done')
|
||||
await this.ctx.with('close-pipeline', {}, async () => {
|
||||
await (await workspace.pipeline).close()
|
||||
})
|
||||
} catch (err: any) {
|
||||
console.error(err)
|
||||
await this.ctx.error('close-pipeline-error', { error: err })
|
||||
}
|
||||
}
|
||||
await Promise.race([closePipeline(), timeoutPromise(15000)])
|
||||
if (LOGGING_ENABLED) console.timeLog(wsId, 'Workspace closed...')
|
||||
console.timeEnd(wsId)
|
||||
await this.ctx.with('closing', {}, async () => {
|
||||
await Promise.race([closePipeline(), timeoutPromise(15000)])
|
||||
})
|
||||
if (LOGGING_ENABLED) {
|
||||
await this.ctx.info('Workspace closed...', { workspace: workspace.id, wsId, wsName: workspace.workspaceName })
|
||||
}
|
||||
}
|
||||
|
||||
private async sendUpgrade (ctx: MeasureContext, webSocket: ConnectionSocket, binary: boolean): Promise<void> {
|
||||
@@ -530,31 +564,36 @@ class TSessionManager implements SessionManager {
|
||||
): Promise<void> {
|
||||
if (workspace.sessions.size === 0) {
|
||||
const wsUID = workspace.id
|
||||
const logParams = { wsid, workspace: workspace.id, wsName: workspaceId.name }
|
||||
if (LOGGING_ENABLED) {
|
||||
console.log(workspaceId.name, 'no sessions for workspace', wsid, wsUID)
|
||||
await this.ctx.info('no sessions for workspace', logParams)
|
||||
}
|
||||
|
||||
const waitAndClose = async (workspace: Workspace): Promise<void> => {
|
||||
try {
|
||||
const pl = await workspace.pipeline
|
||||
await Promise.race([pl, timeoutPromise(60000)])
|
||||
await Promise.race([pl.close(), timeoutPromise(60000)])
|
||||
if (workspace.closing === undefined) {
|
||||
const waitAndClose = async (workspace: Workspace): Promise<void> => {
|
||||
try {
|
||||
if (workspace.sessions.size === 0) {
|
||||
const pl = await workspace.pipeline
|
||||
await Promise.race([pl, timeoutPromise(60000)])
|
||||
await Promise.race([pl.close(), timeoutPromise(60000)])
|
||||
|
||||
if (this.workspaces.get(wsid)?.id === wsUID) {
|
||||
if (this.workspaces.get(wsid)?.id === wsUID) {
|
||||
this.workspaces.delete(wsid)
|
||||
}
|
||||
workspace.context.end()
|
||||
if (LOGGING_ENABLED) {
|
||||
await this.ctx.info('Closed workspace', logParams)
|
||||
}
|
||||
}
|
||||
} catch (err: any) {
|
||||
this.workspaces.delete(wsid)
|
||||
}
|
||||
workspace.context.end()
|
||||
if (LOGGING_ENABLED) {
|
||||
console.timeLog(workspaceId.name, 'Closed workspace', wsUID)
|
||||
}
|
||||
} catch (err: any) {
|
||||
this.workspaces.delete(wsid)
|
||||
if (LOGGING_ENABLED) {
|
||||
console.error(workspaceId.name, err)
|
||||
if (LOGGING_ENABLED) {
|
||||
await this.ctx.error('failed', { ...logParams, error: err })
|
||||
}
|
||||
}
|
||||
}
|
||||
workspace.closing = waitAndClose(workspace)
|
||||
}
|
||||
workspace.closing = waitAndClose(workspace)
|
||||
await workspace.closing
|
||||
}
|
||||
}
|
||||
@@ -562,13 +601,22 @@ class TSessionManager implements SessionManager {
|
||||
broadcast (from: Session | null, workspaceId: WorkspaceId, resp: Response<any>, target?: string[]): void {
|
||||
const workspace = this.workspaces.get(toWorkspaceString(workspaceId))
|
||||
if (workspace === undefined) {
|
||||
console.error(new Error('internal: cannot find sessions'))
|
||||
void this.ctx.error('internal: cannot find sessions', {
|
||||
workspaceId: workspaceId.name,
|
||||
target,
|
||||
userId: from?.getUser() ?? '$unknown'
|
||||
})
|
||||
return
|
||||
}
|
||||
if (workspace?.upgrade ?? false) {
|
||||
return
|
||||
}
|
||||
if (LOGGING_ENABLED) console.log(workspaceId.name, `server broadcasting to ${workspace.sessions.size} clients...`)
|
||||
if (LOGGING_ENABLED) {
|
||||
void this.ctx.info('server broadcasting to clients...', {
|
||||
workspace: workspaceId.name,
|
||||
count: workspace.sessions.size
|
||||
})
|
||||
}
|
||||
|
||||
const sessions = [...workspace.sessions.values()]
|
||||
const ctx = this.ctx.newChild('📭 broadcast', {})
|
||||
@@ -627,19 +675,14 @@ class TSessionManager implements SessionManager {
|
||||
service.useBroadcast = hello.broadcast ?? false
|
||||
|
||||
if (LOGGING_ENABLED) {
|
||||
console.timeLog(
|
||||
workspace,
|
||||
'hello happen',
|
||||
service.getUser(),
|
||||
'binary:',
|
||||
service.binaryResponseMode,
|
||||
'compression:',
|
||||
service.useCompression,
|
||||
'workspace users:',
|
||||
this.workspaces.get(workspace)?.sessions?.size,
|
||||
'total users:',
|
||||
this.sessions.size
|
||||
)
|
||||
await ctx.info('hello happen', {
|
||||
user: service.getUser(),
|
||||
binary: service.binaryResponseMode,
|
||||
compression: service.useCompression,
|
||||
timeToHello: Date.now() - service.createTime,
|
||||
workspaceUsers: this.workspaces.get(workspace)?.sessions?.size,
|
||||
totalUsers: this.sessions.size
|
||||
})
|
||||
}
|
||||
const helloResponse: HelloResponse = {
|
||||
id: -1,
|
||||
@@ -684,7 +727,9 @@ class TSessionManager implements SessionManager {
|
||||
service.useCompression
|
||||
)
|
||||
} catch (err: any) {
|
||||
if (LOGGING_ENABLED) console.error(err)
|
||||
if (LOGGING_ENABLED) {
|
||||
await this.ctx.error('error handle request', { error: err, request })
|
||||
}
|
||||
const resp: Response<any> = {
|
||||
id: request.id,
|
||||
error: unknownError(err),
|
||||
@@ -726,7 +771,9 @@ class TSessionManager implements SessionManager {
|
||||
service.useCompression
|
||||
)
|
||||
} catch (err: any) {
|
||||
if (LOGGING_ENABLED) console.error(err)
|
||||
if (LOGGING_ENABLED) {
|
||||
await ctx.error('error handle measure', { error: err, request })
|
||||
}
|
||||
const resp: Response<any> = {
|
||||
id: request.id,
|
||||
error: unknownError(err),
|
||||
|
||||
@@ -47,7 +47,9 @@ export function startHttpServer (
|
||||
enableCompression: boolean,
|
||||
accountsUrl: string
|
||||
): () => Promise<void> {
|
||||
if (LOGGING_ENABLED) console.log(`starting server on port ${port} ...`)
|
||||
if (LOGGING_ENABLED) {
|
||||
void ctx.info('starting server on', { port, productId, enableCompression, accountsUrl })
|
||||
}
|
||||
|
||||
const app = express()
|
||||
app.use(cors())
|
||||
@@ -209,21 +211,27 @@ export function startHttpServer (
|
||||
)
|
||||
if ('upgrade' in session || 'error' in session) {
|
||||
if ('error' in session) {
|
||||
console.error(session.error)
|
||||
void ctx.error('error', { error: session.error })
|
||||
}
|
||||
cs.close()
|
||||
return
|
||||
}
|
||||
// eslint-disable-next-line @typescript-eslint/no-misused-promises
|
||||
ws.on('message', (msg: RawData) => {
|
||||
let buff: any | undefined
|
||||
if (msg instanceof Buffer) {
|
||||
buff = msg?.toString()
|
||||
} else if (Array.isArray(msg)) {
|
||||
buff = Buffer.concat(msg).toString()
|
||||
}
|
||||
if (buff !== undefined) {
|
||||
void handleRequest(session.context, session.session, cs, buff, session.workspaceName)
|
||||
try {
|
||||
let buff: any | undefined
|
||||
if (msg instanceof Buffer) {
|
||||
buff = msg?.toString()
|
||||
} else if (Array.isArray(msg)) {
|
||||
buff = Buffer.concat(msg).toString()
|
||||
}
|
||||
if (buff !== undefined) {
|
||||
void handleRequest(session.context, session.session, cs, buff, session.workspaceName)
|
||||
}
|
||||
} catch (err: any) {
|
||||
if (LOGGING_ENABLED) {
|
||||
void ctx.error('message error', err)
|
||||
}
|
||||
}
|
||||
})
|
||||
// eslint-disable-next-line @typescript-eslint/no-misused-promises
|
||||
@@ -251,12 +259,17 @@ export function startHttpServer (
|
||||
const sessionId = url.searchParams.get('sessionId')
|
||||
|
||||
if (payload.workspace.productId !== productId) {
|
||||
if (LOGGING_ENABLED) {
|
||||
void ctx.error('invalid product', { required: payload.workspace.productId, productId })
|
||||
}
|
||||
throw new Error('Invalid workspace product')
|
||||
}
|
||||
|
||||
wss.handleUpgrade(request, socket, head, (ws) => wss.emit('connection', ws, request, payload, token, sessionId))
|
||||
} catch (err) {
|
||||
if (LOGGING_ENABLED) console.error('invalid token', err)
|
||||
} catch (err: any) {
|
||||
if (LOGGING_ENABLED) {
|
||||
void ctx.error('invalid token', err)
|
||||
}
|
||||
wss.handleUpgrade(request, socket, head, (ws) => {
|
||||
const resp: Response<any> = {
|
||||
id: -1,
|
||||
@@ -274,7 +287,9 @@ export function startHttpServer (
|
||||
}
|
||||
})
|
||||
httpServer.on('error', (err) => {
|
||||
if (LOGGING_ENABLED) console.error('server error', err)
|
||||
if (LOGGING_ENABLED) {
|
||||
void ctx.error('server error', err)
|
||||
}
|
||||
})
|
||||
|
||||
httpServer.listen(port)
|
||||
|
||||
@@ -35,6 +35,7 @@ export interface StatisticsElement {
|
||||
* @public
|
||||
*/
|
||||
export interface Session {
|
||||
createTime: number
|
||||
getUser: () => string
|
||||
pipeline: () => Pipeline
|
||||
ping: () => Promise<string>
|
||||
@@ -117,7 +118,9 @@ export interface Workspace {
|
||||
pipeline: Promise<Pipeline>
|
||||
sessions: Map<string, { session: Session, socket: ConnectionSocket }>
|
||||
upgrade: boolean
|
||||
|
||||
closing?: Promise<void>
|
||||
closeTimeout?: any
|
||||
|
||||
workspaceId: WorkspaceId
|
||||
workspaceName: string
|
||||
|
||||
Reference in New Issue
Block a user