// // Copyright © 2024 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 core, { AccountRole, type Class, type Doc, type DocInfo, type DocumentQuery, type DocumentUpdate, type Domain, DOMAIN_MODEL, DOMAIN_MODEL_TX, DOMAIN_SPACE, DOMAIN_TX, type FindOptions, type FindResult, generateId, groupByArray, type Hierarchy, isOperator, type Iterator, type Lookup, type MeasureContext, type ModelDb, type ObjQueryType, type Projection, type Ref, type ReverseLookups, type SessionData, type SortingQuery, type StorageIterator, toFindResult, type Tx, type TxCreateDoc, type TxCUD, type TxMixin, TxProcessor, type TxRemoveDoc, type TxResult, type TxUpdateDoc, type WithLookup, type WorkspaceId } from '@hcengineering/core' import { type DbAdapter, type DbAdapterHandler, type DomainHelperOperations, estimateDocSize, type ServerFindOptions, type TxAdapter, updateHashForDoc } from '@hcengineering/server-core' import { createHash } from 'crypto' import type postgres from 'postgres' import { getDocFieldsByDomains, translateDomain } from './schemas' import { type ValueType } from './types' import { convertDoc, createTables, DBCollectionHelper, type DBDoc, escapeBackticks, getDBClient, inferType, isDataField, isOwner, type JoinProps, parseDoc, parseDocWithProjection, parseUpdate, type PostgresClientReference } from './utils' abstract class PostgresAdapterBase implements DbAdapter { protected readonly _helper: DBCollectionHelper protected readonly tableFields = new Map() protected readonly connections = new Map>() protected readonly retryTxn = async ( client: postgres.ReservedSql, fn: (client: postgres.ReservedSql) => Promise ): Promise => { const backoffInterval = 100 // millis const maxTries = 5 let tries = 0 while (true) { await client.unsafe('BEGIN;') tries++ try { const result = await fn(client) await client.unsafe('COMMIT;') return result } catch (err: any) { await client.unsafe('ROLLBACK;') if (err.code !== '40001' || tries === maxTries) { throw err } else { console.log('Transaction failed. Retrying.') console.log(err.message) await new Promise((resolve) => setTimeout(resolve, tries * backoffInterval)) } } } } constructor ( protected readonly client: postgres.Sql, protected readonly refClient: PostgresClientReference, protected readonly workspaceId: WorkspaceId, protected readonly hierarchy: Hierarchy, protected readonly modelDb: ModelDb ) { this._helper = new DBCollectionHelper(this.client, this.workspaceId) } protected async withConnection( ctx: MeasureContext, operation: (client: postgres.ReservedSql) => Promise ): Promise { const connection = await this.getConnection(ctx) if (connection !== undefined) { return await operation(connection) } else { const client = await this.client.reserve() try { return await operation(client) } finally { client.release() } } } async closeContext (ctx: MeasureContext): Promise { if (ctx.id === undefined) return const conn = this.connections.get(ctx.id) if (conn !== undefined) { if (conn instanceof Promise) { ;(await conn).release() } else { conn.release() } this.connections.delete(ctx.id) } } protected async getConnection (ctx: MeasureContext): Promise { if (ctx.id === undefined) return const conn = this.connections.get(ctx.id) if (conn !== undefined) return await conn const client = this.client.reserve() this.connections.set(ctx.id, client) return await client } async traverse( _domain: Domain, query: DocumentQuery, options?: Pick, 'sort' | 'limit' | 'projection'> ): Promise> { const client = await this.client.reserve() let closed = false const cursorName = `cursor_${translateDomain(this.workspaceId.name)}_${translateDomain(_domain)}_${generateId()}` const close = async (cursorName: string): Promise => { if (closed) return try { await client.unsafe(`CLOSE ${cursorName}`) await client.unsafe('COMMIT;') } finally { client.release() closed = true } } const init = async (): Promise => { const domain = translateDomain(_domain) const sqlChunks: string[] = [`CURSOR FOR SELECT * FROM ${domain}`] sqlChunks.push(`WHERE ${this.buildRawQuery(domain, query, options)}`) if (options?.sort !== undefined) { sqlChunks.push(this.buildRawOrder(domain, options.sort)) } if (options?.limit !== undefined) { sqlChunks.push(`LIMIT ${options.limit}`) } const finalSql: string = sqlChunks.join(' ') await client.unsafe('BEGIN;') await client.unsafe(`DECLARE ${cursorName} ${finalSql}`) } const next = async (count: number): Promise => { const result = await client.unsafe(`FETCH ${count} FROM ${cursorName}`) if (result.length === 0) { await close(cursorName) return null } return result.map((p) => parseDoc(p as any, _domain)) } await init() return { next, close: async () => { await close(cursorName) } } } helper (): DomainHelperOperations { return this._helper } on?: ((handler: DbAdapterHandler) => void) | undefined abstract init (): Promise async close (): Promise { for (const c of this.connections.values()) { if (c instanceof Promise) { ;(await c).release() } else { c.release() } } this.refClient.close() } async rawFindAll(_domain: Domain, query: DocumentQuery, options?: FindOptions): Promise { const domain = translateDomain(_domain) const select = `SELECT * FROM ${domain}` const sqlChunks: string[] = [] sqlChunks.push(`WHERE ${this.buildRawQuery(domain, query, options)}`) if (options?.sort !== undefined) { sqlChunks.push(this.buildRawOrder(domain, options.sort)) } if (options?.limit !== undefined) { sqlChunks.push(`LIMIT ${options.limit}`) } const finalSql: string = [select, ...sqlChunks].join(' ') const result = await this.client.unsafe(finalSql) return result.map((p) => parseDocWithProjection(p as any, domain, options?.projection)) } buildRawOrder(domain: string, sort: SortingQuery): string { const res: string[] = [] for (const key in sort) { const val = sort[key] if (val === undefined) { continue } if (typeof val === 'number') { res.push(`${this.transformKey(domain, core.class.Doc, key, false)} ${val === 1 ? 'ASC' : 'DESC'}`) } else { // todo handle custom sorting } } return `ORDER BY ${res.join(', ')}` } buildRawQuery(domain: string, query: DocumentQuery, options?: FindOptions): string { const res: string[] = [] res.push(`"workspaceId" = '${this.workspaceId.name}'`) for (const key in query) { const value = query[key] const tkey = this.transformKey(domain, core.class.Doc, key, false) const translated = this.translateQueryValue(tkey, value, 'common') if (translated !== undefined) { res.push(translated) } } return res.join(' AND ') } async rawUpdate( domain: Domain, query: DocumentQuery, operations: DocumentUpdate ): Promise { const translatedQuery = this.buildRawQuery(domain, query) if ((operations as any).$set !== undefined) { ;(operations as any) = { ...(operations as any).$set } } const isOps = isOperator(operations) if ((operations as any)['%hash%'] === undefined) { ;(operations as any)['%hash%'] = null } if (isOps) { const conn = await this.client.reserve() try { await this.retryTxn(conn, async (client) => { const res = await client.unsafe( `SELECT * FROM ${translateDomain(domain)} WHERE ${translatedQuery} FOR UPDATE` ) const docs = res.map((p) => parseDoc(p as any, domain)) for (const doc of docs) { if (doc === undefined) continue const prevAttachedTo = (doc as any).attachedTo TxProcessor.applyUpdate(doc, operations) const converted = convertDoc(domain, doc, this.workspaceId.name) const params: any[] = [doc._id, this.workspaceId.name] let paramsIndex = params.length + 1 const updates: string[] = [] const { extractedFields, remainingData } = parseUpdate(domain, operations) const newAttachedTo = (doc as any).attachedTo if (Object.keys(extractedFields).length > 0) { for (const key in extractedFields) { const val = (extractedFields as any)[key] if (key === 'attachedTo' && val === prevAttachedTo) continue updates.push(`"${key}" = $${paramsIndex++}`) params.push(val) } } else if (prevAttachedTo !== undefined && prevAttachedTo !== newAttachedTo) { updates.push(`"attachedTo" = $${paramsIndex++}`) params.push(newAttachedTo) } if (Object.keys(remainingData).length > 0) { updates.push(`data = $${paramsIndex++}`) params.push(converted.data) } await client.unsafe( `UPDATE ${translateDomain(domain)} SET ${updates.join(', ')} WHERE _id = $1 AND "workspaceId" = $2`, params ) } }) } finally { conn.release() } } else { await this.rawUpdateDoc(domain, query, operations) } } private async rawUpdateDoc( domain: Domain, query: DocumentQuery, operations: DocumentUpdate ): Promise { const translatedQuery = this.buildRawQuery(domain, query) const updates: string[] = [] const params: any[] = [] let paramsIndex = params.length + 1 const { extractedFields, remainingData } = parseUpdate(domain, operations) const { space, attachedTo, ...ops } = operations as any for (const key in extractedFields) { updates.push(`"${key}" = $${paramsIndex++}`) params.push((extractedFields as any)[key]) } let from = 'data' let dataUpdated = false for (const key in remainingData) { if (ops[key] === undefined) continue const val = (remainingData as any)[key] from = `jsonb_set(${from}, '{${key}}', coalesce(to_jsonb($${paramsIndex++}${inferType(val)}), 'null') , true)` params.push(val) dataUpdated = true } if (dataUpdated) { updates.push(`data = ${from}`) } const conn = await this.client.reserve() try { await this.retryTxn(conn, async (client) => { await client.unsafe( `UPDATE ${translateDomain(domain)} SET ${updates.join(', ')} WHERE ${translatedQuery}`, params ) }) } catch (err) { console.error(err, { domain, params, updates }) } finally { conn.release() } } async rawDeleteMany(domain: Domain, query: DocumentQuery): Promise { const translatedQuery = this.buildRawQuery(domain, query) const conn = await this.client.reserve() try { await this.retryTxn(conn, async (client) => { await client.unsafe(`DELETE FROM ${translateDomain(domain)} WHERE ${translatedQuery}`) }) } finally { conn.release() } } findAll( ctx: MeasureContext, _class: Ref>, query: DocumentQuery, options?: ServerFindOptions ): Promise> { return ctx.with('findAll', { _class }, async () => { try { const domain = translateDomain(options?.domain ?? this.hierarchy.getDomain(_class)) const sqlChunks: string[] = [] const joins = this.buildJoin(_class, options?.lookup) if (options?.domainLookup !== undefined) { const baseDomain = translateDomain(this.hierarchy.getDomain(_class)) const domain = translateDomain(options.domainLookup.domain) const key = options.domainLookup.field const as = `dl_lookup_${domain}_${key}` joins.push({ isReverse: false, table: domain, path: options.domainLookup.field, toAlias: as, toField: '_id', fromField: key, fromAlias: baseDomain, toClass: undefined }) } const select = `SELECT ${this.getProjection(_class, domain, options?.projection, joins)} FROM ${domain}` const secJoin = this.addSecurity(query, domain, ctx.contextData) if (secJoin !== undefined) { sqlChunks.push(secJoin) } if (joins.length > 0) { sqlChunks.push(this.buildJoinString(joins)) } sqlChunks.push(`WHERE ${this.buildQuery(_class, domain, query, joins, options)}`) const connection = (await this.getConnection(ctx)) ?? this.client let total = options?.total === true ? 0 : -1 if (options?.total === true) { const totalReq = `SELECT COUNT(${domain}._id) as count FROM ${domain}` const totalSql = [totalReq, ...sqlChunks].join(' ') const totalResult = await connection.unsafe(totalSql) const parsed = Number.parseInt(totalResult[0].count) total = Number.isNaN(parsed) ? 0 : parsed } if (options?.sort !== undefined) { sqlChunks.push(this.buildOrder(_class, domain, options.sort, joins)) } if (options?.limit !== undefined) { sqlChunks.push(`LIMIT ${options.limit}`) } const finalSql: string = [select, ...sqlChunks].join(' ') const result = await connection.unsafe(finalSql) if (options?.lookup === undefined && options?.domainLookup === undefined) { return toFindResult( result.map((p) => parseDocWithProjection(p as any, domain, options?.projection)), total ) } else { const res = this.parseLookup(result, joins, options?.projection, domain) return toFindResult(res, total) } } catch (err) { ctx.error('Error in findAll', { err }) throw err } }) } addSecurity(query: DocumentQuery, domain: string, sessionContext: SessionData): string | undefined { if (sessionContext !== undefined && sessionContext.isTriggerCtx !== true) { if (sessionContext.admin !== true && sessionContext.account !== undefined) { const acc = sessionContext.account if (isOwner(acc) || acc.role === AccountRole.DocGuest) { return } if (query.space === acc._id) return const key = domain === DOMAIN_SPACE ? '_id' : domain === DOMAIN_TX ? "data ->> 'objectSpace'" : 'space' const privateCheck = domain === DOMAIN_SPACE ? ' OR sec.private = false' : '' const q = `(sec.members @> '{"${acc._id}"}' OR sec."_class" = '${core.class.SystemSpace}'${privateCheck})` return `INNER JOIN ${translateDomain(DOMAIN_SPACE)} AS sec ON sec._id = ${domain}.${key} AND sec."workspaceId" = '${this.workspaceId.name}' AND ${q}` } } } private parseLookup( rows: any[], joins: JoinProps[], projection: Projection | undefined, domain: string ): WithLookup[] { const map = new Map, WithLookup>() const modelJoins: JoinProps[] = [] const reverseJoins: JoinProps[] = [] const simpleJoins: JoinProps[] = [] for (const join of joins) { if (join.table === DOMAIN_MODEL) { modelJoins.push(join) } else if (join.isReverse) { reverseJoins.push(join) } else { simpleJoins.push(join) } } for (const row of rows) { /* eslint-disable @typescript-eslint/consistent-type-assertions */ let doc: WithLookup = map.get(row._id) ?? ({ _id: row._id, $lookup: {} } as WithLookup) const lookup: Record = doc.$lookup as Record let joinIndex: number | undefined let skip = false try { for (const column in row) { if (column.startsWith('reverse_lookup_')) { if (row[column] != null) { const join = reverseJoins.find((j) => j.toAlias === column) if (join === undefined) { continue } const res = this.getLookupValue(join.path, lookup, false) if (res === undefined) continue const { obj, key } = res const parsed = row[column].map((p: any) => parseDoc(p, domain)) obj[key] = parsed } } else if (column.startsWith('lookup_')) { const keys = column.split('_') let key = keys[keys.length - 1] if (keys[keys.length - 2] === '') { key = '_' + key } if (key === 'workspaceId') { continue } if (key === '_id') { if (row[column] === null) { skip = true continue } joinIndex = joinIndex === undefined ? 0 : ++joinIndex skip = false } if (skip) { continue } const join = simpleJoins[joinIndex ?? 0] const res = this.getLookupValue(join.path, lookup) if (res === undefined) continue const { obj, key: p } = res if (key === 'data') { obj[p] = { ...obj[p], ...row[column] } } else { if (key === 'attachedTo' && row[column] === 'NULL') { continue } else { obj[p][key] = row[column] === 'NULL' ? null : row[column] } } } else { joinIndex = undefined if (!map.has(row._id)) { if (column === 'workspaceId') { continue } if (column === 'data') { const data = row[column] if (projection !== undefined) { if (projection !== undefined) { for (const key in data) { if (!Object.prototype.hasOwnProperty.call(projection, key) || (projection as any)[key] === 0) { // eslint-disable-next-line @typescript-eslint/no-dynamic-delete delete data[key] } } } } doc = { ...doc, ...data } } else { if (column === 'createdOn' || column === 'modifiedOn') { const val = Number.parseInt(row[column]) ;(doc as any)[column] = Number.isNaN(val) ? null : val } else { ;(doc as any)[column] = row[column] === 'NULL' ? null : row[column] } } } } } } catch (err) { console.log(err) throw err } for (const modelJoin of modelJoins) { const res = this.getLookupValue(modelJoin.path, lookup) if (res === undefined) continue const { obj, key } = res const val = this.getModelLookupValue(doc, modelJoin, simpleJoins) if (val !== undefined && modelJoin.toClass !== undefined) { const res = this.modelDb.findAllSync(modelJoin.toClass, { [modelJoin.toField]: (doc as any)[modelJoin.fromField] }) obj[key] = modelJoin.isReverse ? res : res[0] } } map.set(row._id, doc) } return Array.from(map.values()) } private getLookupValue ( fullPath: string, obj: Record, shouldCreate: boolean = true ): | { obj: any key: string } | undefined { const path = fullPath.split('.') for (let i = 0; i < path.length; i++) { const p = path[i] if (i > 0) { if (obj.$lookup === undefined) { obj.$lookup = {} } obj = obj.$lookup } if (obj[p] === undefined) { if (!shouldCreate && i < path.length - 1) { return } else { obj[p] = {} } } if (i === path.length - 1) { return { obj, key: p } } obj = obj[p] } } private getModelLookupValue(doc: WithLookup, join: JoinProps, simpleJoins: JoinProps[]): any { if (join.fromAlias.startsWith('lookup_')) { const simple = simpleJoins.find((j) => j.toAlias === join.fromAlias) if (simple !== undefined) { const val = this.getLookupValue(simple.path, doc.$lookup ?? {}) if (val !== undefined) { const data = val.obj[val.key] return data[join.fromField] } } } else { return (doc as any)[join.fromField] } } private buildJoinString (value: JoinProps[]): string { const res: string[] = [] for (const val of value) { if (val.isReverse) continue if (val.table === DOMAIN_MODEL) continue res.push( `LEFT JOIN ${val.table} AS ${val.toAlias} ON ${val.fromAlias}.${val.fromField} = ${val.toAlias}."${val.toField}" AND ${val.toAlias}."workspaceId" = '${this.workspaceId.name}'` ) if (val.classes !== undefined) { if (val.classes.length === 1) { res.push(`AND ${val.toAlias}._class = '${val.classes[0]}'`) } else { res.push(`AND ${val.toAlias}._class IN (${val.classes.map((c) => `'${c}'`).join(', ')})`) } } } return res.join(' ') } private buildJoin(clazz: Ref>, lookup: Lookup | undefined): JoinProps[] { const res: JoinProps[] = [] if (lookup !== undefined) { this.buildJoinValue(clazz, lookup, res) } return res } private buildJoinValue( clazz: Ref>, lookup: Lookup, res: JoinProps[], parentKey?: string, parentAlias?: string ): void { const baseDomain = parentAlias ?? translateDomain(this.hierarchy.getDomain(clazz)) for (const key in lookup) { if (key === '_id') { this.getReverseLookupValue(baseDomain, lookup, res, parentKey) continue } const value = (lookup as any)[key] const _class = Array.isArray(value) ? value[0] : value const nested = Array.isArray(value) ? value[1] : undefined const domain = translateDomain(this.hierarchy.getDomain(_class)) const tkey = domain === DOMAIN_MODEL ? key : this.transformKey(baseDomain, clazz, key) const as = `lookup_${domain}_${parentKey !== undefined ? parentKey + '_lookup_' + key : key}` res.push({ isReverse: false, table: domain, path: parentKey !== undefined ? `${parentKey}.${key}` : key, toAlias: as, toField: '_id', fromField: tkey, fromAlias: baseDomain, toClass: _class }) if (nested !== undefined) { this.buildJoinValue(_class, nested, res, key, as) } } } private getReverseLookupValue ( parentDomain: string, lookup: ReverseLookups, result: JoinProps[], parent?: string ): void { const lid = lookup?._id ?? {} for (const key in lid) { const value = lid[key] let _class: Ref> let attr = 'attachedTo' if (Array.isArray(value)) { _class = value[0] attr = value[1] } else { _class = value } const domain = translateDomain(this.hierarchy.getDomain(_class)) const desc = this.hierarchy .getDescendants(this.hierarchy.getBaseClass(_class)) .filter((it) => !this.hierarchy.isMixin(it)) const as = `reverse_lookup_${domain}_${parent !== undefined ? parent + '_lookup_' + key : key}` result.push({ isReverse: true, table: domain, toAlias: as, toField: attr, classes: desc, path: parent !== undefined ? `${parent}.${key}` : key, fromAlias: parentDomain, toClass: _class, fromField: '_id' }) } } private buildOrder( _class: Ref>, baseDomain: string, sort: SortingQuery, joins: JoinProps[] ): string { const res: string[] = [] for (const key in sort) { const val = sort[key] if (val === undefined) { continue } if (typeof val === 'number') { res.push(`${this.getKey(_class, baseDomain, key, joins)} ${val === 1 ? 'ASC' : 'DESC'}`) } else { // todo handle custom sorting } } return `ORDER BY ${res.join(', ')}` } private buildQuery( _class: Ref>, baseDomain: string, query: DocumentQuery, joins: JoinProps[], options?: ServerFindOptions ): string { const res: string[] = [] res.push(`${baseDomain}."workspaceId" = '${this.workspaceId.name}'`) if (options?.skipClass !== true) { query._class = this.fillClass(_class, query) as any } for (const key in query) { if (options?.skipSpace === true && key === 'space') { continue } if (options?.skipClass === true && key === '_class') { continue } const value = query[key] if (value === undefined) continue const valueType = this.getValueType(_class, key) const tkey = this.getKey(_class, baseDomain, key, joins, valueType === 'dataArray') const translated = this.translateQueryValue(tkey, value, valueType) if (translated !== undefined) { res.push(translated) } } return res.join(' AND ') } private getValueType(_class: Ref>, key: string): ValueType { const splitted = key.split('.') const mixinOrKey = splitted[0] const domain = this.hierarchy.getDomain(_class) if (this.hierarchy.isMixin(mixinOrKey as Ref>)) { key = splitted.slice(1).join('.') const attr = this.hierarchy.findAttribute(mixinOrKey as Ref>, key) if (attr !== undefined && attr.type._class === core.class.ArrOf) { return isDataField(domain, key) ? 'dataArray' : 'array' } return 'common' } else { const attr = this.hierarchy.findAttribute(_class, key) if (attr !== undefined && attr.type._class === core.class.ArrOf) { return isDataField(domain, key) ? 'dataArray' : 'array' } return 'common' } } private fillClass( _class: Ref>, query: DocumentQuery ): ObjQueryType | undefined { let value: any = query._class const baseClass = this.hierarchy.getBaseClass(_class) if (baseClass !== core.class.Doc) { const classes = this.hierarchy.getDescendants(baseClass).filter((it) => !this.hierarchy.isMixin(it)) // Only replace if not specified if (value === undefined) { value = { $in: classes } } else if (typeof value === 'string') { if (!classes.includes(value as Ref>)) { value = classes.length === 1 ? classes[0] : { $in: classes } } } else if (typeof value === 'object' && value !== null) { let descendants: Ref>[] = classes if (Array.isArray(value.$in)) { const classesIds = new Set(classes) descendants = value.$in.filter((c: Ref>) => classesIds.has(c)) } if (value != null && Array.isArray(value.$nin)) { const excludedClassesIds = new Set>>(value.$nin) descendants = descendants.filter((c) => !excludedClassesIds.has(c)) } if (value.$ne != null) { descendants = descendants.filter((c) => c !== value?.$ne) } const desc = descendants.filter((it: any) => !this.hierarchy.isMixin(it as Ref>)) value = desc.length === 1 ? desc[0] : { $in: desc } } if (baseClass !== _class) { // Add an mixin to be exists flag ;(query as any)[_class] = { $exists: true } } } else { // No need to pass _class in case of fixed domain search. return undefined } if (value?.$in?.length === 1 && value?.$nin === undefined) { value = value.$in[0] } return value } private getKey( _class: Ref>, baseDomain: string, key: string, joins: JoinProps[], isDataArray: boolean = false ): string { if (key.startsWith('$lookup')) { return this.transformLookupKey(baseDomain, key, joins, isDataArray) } return `${baseDomain}.${this.transformKey(baseDomain, _class, key, isDataArray)}` } private transformLookupKey (domain: string, key: string, joins: JoinProps[], isDataArray: boolean = false): string { const arr = key.split('.').filter((p) => p !== '$lookup') const tKey = arr.pop() ?? '' const path = arr.join('.') const join = joins.find((p) => p.path === path) if (join === undefined) { throw new Error(`Can't fined join for path: ${path}`) } if (join.isReverse) { return `${join.toAlias}->'${tKey}'` } const res = isDataField(domain, tKey) ? (isDataArray ? `data->'${tKey}'` : `data#>>'{${tKey}}'`) : key return `${join.toAlias}.${res}` } private transformKey( domain: string, _class: Ref>, key: string, isDataArray: boolean = false ): string { if (!isDataField(domain, key)) return `"${key}"` const arr = key.split('.').filter((p) => p) let tKey = '' let isNestedField = false for (let i = 0; i < arr.length; i++) { const element = arr[i] if (element === '$lookup') { tKey += arr[++i] + '_lookup' } else if (this.hierarchy.isMixin(element as Ref>)) { isNestedField = true tKey += `${element}` if (i !== arr.length - 1) { tKey += "'->'" } } else { tKey += arr[i] if (i !== arr.length - 1) { tKey += ',' } } // Check if key is belong to mixin class, we need to add prefix. tKey = this.checkMixinKey(tKey, _class, isDataArray) } return isDataArray || isNestedField ? `data->'${tKey}'` : `data#>>'{${tKey}}'` } private checkMixinKey(key: string, _class: Ref>, isDataArray: boolean): string { if (!key.includes('.')) { try { const attr = this.hierarchy.findAttribute(_class, key) if (attr !== undefined && this.hierarchy.isMixin(attr.attributeOf)) { // It is mixin if (isDataArray) { key = `${attr.attributeOf}'->'${key}` } else { key = `${attr.attributeOf},${key}` } } } catch (err: any) { // ignore, if } } return key } private translateQueryValue (tkey: string, value: any, type: ValueType): string | undefined { if (value === null) { return `${tkey} IS NULL` } else if (typeof value === 'object' && !Array.isArray(value)) { // we can have multiple criteria for one field const res: string[] = [] for (const operator in value) { const val = value[operator] switch (operator) { case '$ne': if (val === null) { res.push(`${tkey} IS NOT NULL`) } else { res.push(`${tkey} != '${val}'`) } break case '$gt': res.push(`${tkey} > '${val}'`) break case '$gte': res.push(`${tkey} >= '${val}'`) break case '$lt': res.push(`${tkey} < '${val}'`) break case '$lte': res.push(`${tkey} <= '${val}'`) break case '$in': res.push( type !== 'common' ? `${tkey} ?| array[${val.length > 0 ? val.map((v: any) => `'${v}'`).join(', ') : 'NULL'}]` : `${tkey} IN (${val.length > 0 ? val.map((v: any) => `'${v}'`).join(', ') : 'NULL'})` ) break case '$nin': if (val.length > 0) { res.push(`${tkey} NOT IN (${val.map((v: any) => `'${v}'`).join(', ')})`) } break case '$like': res.push(`${tkey} ILIKE '${escapeBackticks(val)}'`) break case '$exists': res.push(`${tkey} IS ${val === true ? 'NOT NULL' : 'NULL'}`) break case '$regex': res.push(`${tkey} SIMILAR TO '${val}'`) break case '$options': break case '$all': res.push(`${tkey} @> ARRAY[${value}]`) break default: res.push(`${tkey} @> '[${JSON.stringify(value)}]'`) break } } return res.length === 0 ? undefined : res.join(' AND ') } return type === 'common' ? `${tkey} = '${value}'` : type === 'array' ? `${tkey} @> '${typeof value === 'string' ? '{"' + value + '"}' : value}'` : `${tkey} @> '${typeof value === 'string' ? '"' + value + '"' : value}'` } private getProjectionsAliases (join: JoinProps): string[] { if (join.table === DOMAIN_MODEL) return [] if (join.isReverse) { let classsesQuery = '' if (join.classes !== undefined) { if (join.classes.length === 1) { classsesQuery = ` AND ${join.toAlias}._class = '${join.classes[0]}'` } else { classsesQuery = ` AND ${join.toAlias}._class IN (${join.classes.map((c) => `'${c}'`).join(', ')})` } } return [ `(SELECT jsonb_agg(${join.toAlias}.*) FROM ${join.table} AS ${join.toAlias} WHERE ${join.fromAlias}.${join.fromField} = ${join.toAlias}."${join.toField}" ${classsesQuery}) AS ${join.toAlias}` ] } const fields = getDocFieldsByDomains(join.table) const res: string[] = [] for (const key of [...fields, 'data']) { res.push(`${join.toAlias}."${key}" as "lookup_${join.path.replaceAll('.', '_')}_${key}"`) } return res } private getProjection( _class: Ref>, baseDomain: string, projection: Projection | undefined, joins: JoinProps[] ): string | '*' { if (projection === undefined && joins.length === 0) return `${baseDomain}.*` const res: string[] = [] let dataAdded = false if (projection === undefined) { res.push(`${baseDomain}.*`) } else { for (const key in projection) { if (isDataField(baseDomain, key)) { if (!dataAdded) { res.push(`${baseDomain}.data as data`) dataAdded = true } } else { res.push(`${baseDomain}."${key}" AS "${key}"`) } } } for (const join of joins) { res.push(...this.getProjectionsAliases(join)) } return res.join(', ') } async tx (ctx: MeasureContext, ...tx: Tx[]): Promise { return [] } find (_ctx: MeasureContext, domain: Domain, recheck?: boolean): StorageIterator { const ctx = _ctx.newChild('find', { domain }) const getCursorName = (): string => { return `cursor_${translateDomain(this.workspaceId.name)}_${translateDomain(domain)}_${mode}` } let initialized: boolean = false let client: postgres.ReservedSql let mode: 'hashed' | 'non_hashed' = 'hashed' let cursorName = getCursorName() const bulkUpdate = new Map, string>() const close = async (cursorName: string): Promise => { try { await client.unsafe(`CLOSE ${cursorName}`) await client.unsafe('COMMIT') } catch (err) { ctx.error('Error while closing cursor', { cursorName, err }) } finally { client.release() } } const init = async (projection: string, query: string): Promise => { cursorName = getCursorName() client = await this.client.reserve() await client.unsafe('BEGIN') await client.unsafe( `DECLARE ${cursorName} CURSOR FOR SELECT ${projection} FROM ${translateDomain(domain)} WHERE "workspaceId" = $1 AND ${query}`, [this.workspaceId.name] ) } const next = async (limit: number): Promise => { const result = await client.unsafe(`FETCH ${limit} FROM ${cursorName}`) if (result.length === 0) { return [] } return result.filter((it) => it != null).map((it) => parseDoc(it as any, domain)) } const flush = async (flush = false): Promise => { if (bulkUpdate.size > 1000 || flush) { if (bulkUpdate.size > 0) { await ctx.with('bulk-write-find', {}, () => { const updates = new Map(Array.from(bulkUpdate.entries()).map((it) => [it[0], { '%hash%': it[1] }])) return this.update(ctx, domain, updates) }) } bulkUpdate.clear() } } return { next: async () => { if (!initialized) { if (recheck === true) { await this.retryTxn(client, async (client) => { await client`UPDATE ${client(translateDomain(domain))} SET '%hash%' = NULL WHERE "workspaceId" = ${this.workspaceId.name} AND '%hash%' IS NOT NULL` }) } await init('_id, data', "'%hash%' IS NOT NULL AND '%hash%' <> ''") initialized = true } let docs = await ctx.with('next', { mode }, () => next(50)) if (docs.length === 0 && mode === 'hashed') { await close(cursorName) mode = 'non_hashed' await init('*', "'%hash%' IS NULL OR '%hash%' = ''") docs = await ctx.with('next', { mode }, () => next(50)) } if (docs.length === 0) { return [] } const result: DocInfo[] = [] for (const d of docs) { let digest: string | null = (d as any)['%hash%'] if ('%hash%' in d) { delete d['%hash%'] } const pos = (digest ?? '').indexOf('|') if (digest == null || digest === '') { const cs = ctx.newChild('calc-size', {}) const size = estimateDocSize(d) cs.end() const hash = createHash('sha256') updateHashForDoc(hash, d) digest = hash.digest('base64') bulkUpdate.set(d._id, `${digest}|${size.toString(16)}`) await ctx.with('flush', {}, () => flush()) result.push({ id: d._id, hash: digest, size }) } else { result.push({ id: d._id, hash: digest.slice(0, pos), size: parseInt(digest.slice(pos + 1), 16) }) } } return result }, close: async () => { await ctx.with('flush', {}, () => flush(true)) await close(cursorName) ctx.end() } } } load (ctx: MeasureContext, domain: Domain, docs: Ref[]): Promise { return ctx.with('load', { domain }, async () => { if (docs.length === 0) { return [] } const connection = (await this.getConnection(ctx)) ?? this.client const res = await connection`SELECT * FROM ${connection(translateDomain(domain))} WHERE _id = ANY(${docs}) AND "workspaceId" = ${this.workspaceId.name}` return res.map((p) => parseDocWithProjection(p as any, domain)) }) } upload (ctx: MeasureContext, domain: Domain, docs: Doc[]): Promise { return ctx.with('upload', { domain }, async (ctx) => { const arr = docs.concat() const fields = getDocFieldsByDomains(domain) const filedsWithData = [...fields, 'data'] const insertFields: string[] = [] const onConflict: string[] = [] for (const field of filedsWithData) { insertFields.push(`"${field}"`) onConflict.push(`"${field}" = EXCLUDED."${field}"`) } const insertStr = insertFields.join(', ') const onConflictStr = onConflict.join(', ') await this.withConnection(ctx, async (connection) => { while (arr.length > 0) { const part = arr.splice(0, 500) const values: any[] = [] const vars: string[] = [] let index = 1 for (let i = 0; i < part.length; i++) { const doc = part[i] const variables: string[] = [] const d = convertDoc(domain, doc, this.workspaceId.name) values.push(d.workspaceId) variables.push(`$${index++}`) for (const field of fields) { values.push(d[field]) variables.push(`$${index++}`) } values.push(d.data) variables.push(`$${index++}`) vars.push(`(${variables.join(', ')})`) } const vals = vars.join(',') await this.retryTxn(connection, async (client) => { await client.unsafe( `INSERT INTO ${translateDomain(domain)} ("workspaceId", ${insertStr}) VALUES ${vals} ON CONFLICT ("workspaceId", _id) DO UPDATE SET ${onConflictStr};`, values ) }) } }) }) } async clean (ctx: MeasureContext, domain: Domain, docs: Ref[]): Promise { const connection = (await this.getConnection(ctx)) ?? this.client await connection`DELETE FROM ${connection(translateDomain(domain))} WHERE _id = ANY(${docs}) AND "workspaceId" = ${this.workspaceId.name}` } groupBy( ctx: MeasureContext, domain: Domain, field: string, query?: DocumentQuery

): Promise> { const key = isDataField(domain, field) ? `data ->> '${field}'` : `"${field}"` return ctx.with('groupBy', { domain }, async (ctx) => { const connection = (await this.getConnection(ctx)) ?? this.client try { const result = await connection.unsafe( `SELECT DISTINCT ${key} as ${field}, Count(*) AS count FROM ${translateDomain(domain)} WHERE ${this.buildRawQuery(domain, query ?? {})} GROUP BY ${key}` ) return new Map(result.map((r) => [r[field.toLocaleLowerCase()], parseInt(r.count)])) } catch (err) { ctx.error('Error while grouping by', { domain, field }) throw err } }) } update (ctx: MeasureContext, domain: Domain, operations: Map, DocumentUpdate>): Promise { const ids = Array.from(operations.keys()) return this.withConnection(ctx, (client) => { return this.retryTxn(client, async (client) => { try { const res = await client`SELECT * FROM ${client(translateDomain(domain))} WHERE _id = ANY(${ids}) AND "workspaceId" = ${this.workspaceId.name} FOR UPDATE` const docs = res.map((p) => parseDoc(p as any, domain)) const map = new Map(docs.map((d) => [d._id, d])) for (const [_id, ops] of operations) { const doc = map.get(_id) if (doc === undefined) continue const op = { ...ops } if ((op as any)['%hash%'] === undefined) { ;(op as any)['%hash%'] = null } TxProcessor.applyUpdate(doc, op) const converted = convertDoc(domain, doc, this.workspaceId.name) const columns: string[] = [] const { extractedFields, remainingData } = parseUpdate(domain, op) for (const key in extractedFields) { columns.push(key) } if (Object.keys(remainingData).length > 0) { columns.push('data') } columns.push('modifiedBy') columns.push('modifiedOn') await client`UPDATE ${client(translateDomain(domain))} SET ${client( converted, columns )} WHERE _id = ${doc._id} AND "workspaceId" = ${this.workspaceId.name}` } } catch (err) { ctx.error('Error while updating', { domain, operations, err }) throw err } }) }) } async insert (ctx: MeasureContext, domain: string, docs: Doc[]): Promise { const fields = getDocFieldsByDomains(domain) const filedsWithData = [...fields, 'data'] const columns: string[] = ['workspaceId'] for (const field of filedsWithData) { columns.push(field) } await this.withConnection(ctx, async (connection) => { while (docs.length > 0) { const part = docs.splice(0, 500) const values: DBDoc[] = [] for (let i = 0; i < part.length; i++) { const doc = part[i] const d = convertDoc(domain, doc, this.workspaceId.name) values.push(d) } await this.retryTxn(connection, async (client) => { await client`INSERT INTO ${client(translateDomain(domain))} ${client(values, columns)}` }) } }) return {} } } class PostgresAdapter extends PostgresAdapterBase { async init (domains?: string[], excludeDomains?: string[]): Promise { let resultDomains = domains ?? this.hierarchy.domains() if (excludeDomains !== undefined) { resultDomains = resultDomains.filter((it) => !excludeDomains.includes(it)) } await createTables(this.client, resultDomains) this._helper.domains = new Set(resultDomains as Domain[]) } private async process (ctx: MeasureContext, tx: Tx): Promise { switch (tx._class) { case core.class.TxCreateDoc: return await this.txCreateDoc(ctx, tx as TxCreateDoc) case core.class.TxUpdateDoc: return await this.txUpdateDoc(ctx, tx as TxUpdateDoc) case core.class.TxRemoveDoc: await this.txRemoveDoc(ctx, tx as TxRemoveDoc) break case core.class.TxMixin: return await this.txMixin(ctx, tx as TxMixin) case core.class.TxApplyIf: return undefined default: console.error('Unknown/Unsupported operation:', tx._class, tx) break } } private async txMixin (ctx: MeasureContext, tx: TxMixin): Promise { await ctx.with('tx-mixin', { _class: tx.objectClass, mixin: tx.mixin }, async (ctx) => { await this.withConnection(ctx, async (connection) => { await this.retryTxn(connection, async (client) => { const doc = await this.findDoc(ctx, client, tx.objectClass, tx.objectId, true) if (doc === undefined) return TxProcessor.updateMixin4Doc(doc, tx) ;(doc as any)['%hash%'] = null const domain = this.hierarchy.getDomain(tx.objectClass) const converted = convertDoc(domain, doc, this.workspaceId.name) const { extractedFields } = parseUpdate(domain, tx.attributes as Partial) const columns = new Set() for (const key in extractedFields) { columns.add(key) } columns.add('modifiedBy') columns.add('modifiedOn') columns.add('data') await client`UPDATE ${client(translateDomain(domain))} SET ${client(converted, Array.from(columns))} WHERE _id = ${tx.objectId} AND "workspaceId" = ${this.workspaceId.name}` }) }) }) return {} } async tx (ctx: MeasureContext, ...txes: Tx[]): Promise { const result: TxResult[] = [] const h = this.hierarchy const byDomain = groupByArray(txes, (it) => { if (TxProcessor.isExtendsCUD(it._class)) { return h.findDomain((it as TxCUD).objectClass) } return undefined }) for (const [domain, txs] of byDomain) { if (domain === undefined) { continue } for (const tx of txs) { const res = await this.process(ctx, tx) if (res !== undefined) { result.push(res) } } } return result } protected async txCreateDoc (ctx: MeasureContext, tx: TxCreateDoc): Promise { const doc = TxProcessor.createDoc2Doc(tx) return await ctx.with('create-doc', { _class: doc._class }, (_ctx) => { return this.insert(_ctx, this.hierarchy.getDomain(doc._class), [doc]) }) } protected txUpdateDoc (ctx: MeasureContext, tx: TxUpdateDoc): Promise { return ctx.with('tx-update-doc', { _class: tx.objectClass }, (_ctx) => { if (isOperator(tx.operations)) { let doc: Doc | undefined const ops: any = { '%hash%': null, ...tx.operations } return _ctx.with( 'update with operations', { operations: JSON.stringify(Object.keys(tx.operations)) }, async (ctx) => { return await this.withConnection(ctx, async (connection) => { await this.retryTxn(connection, async (client) => { doc = await this.findDoc(ctx, client, tx.objectClass, tx.objectId, true) if (doc === undefined) return {} ops.modifiedBy = tx.modifiedBy ops.modifiedOn = tx.modifiedOn TxProcessor.applyUpdate(doc, ops) const domain = this.hierarchy.getDomain(tx.objectClass) const converted = convertDoc(domain, doc, this.workspaceId.name) const columns: string[] = [] const { extractedFields, remainingData } = parseUpdate(domain, ops) for (const key in extractedFields) { columns.push(key) } if (Object.keys(remainingData).length > 0) { columns.push('data') } await client`UPDATE ${client(translateDomain(domain))} SET ${client(converted, columns)} WHERE _id = ${tx.objectId} AND "workspaceId" = ${this.workspaceId.name}` }) if (tx.retrieve === true && doc !== undefined) { return { object: doc } } return {} }) } ) } else { return this.updateDoc(_ctx, tx, tx.retrieve ?? false) } }) } private updateDoc(ctx: MeasureContext, tx: TxUpdateDoc, retrieve: boolean): Promise { return ctx.with('update jsonb_set', {}, async (_ctx) => { const updates: string[] = ['"modifiedBy" = $1', '"modifiedOn" = $2'] const params: any[] = [tx.modifiedBy, tx.modifiedOn, tx.objectId, this.workspaceId.name] let paramsIndex = params.length + 1 const domain = this.hierarchy.getDomain(tx.objectClass) const { extractedFields, remainingData } = parseUpdate(domain, tx.operations) const { space, attachedTo, ...ops } = tx.operations as any if (ops['%hash%'] === undefined) { ops['%hash%'] = null } for (const key in extractedFields) { updates.push(`"${key}" = $${paramsIndex++}`) params.push((extractedFields as any)[key]) } let from = 'data' let dataUpdated = false for (const key in remainingData) { if (ops[key] === undefined) continue const val = (remainingData as any)[key] from = `jsonb_set(${from}, '{${key}}', coalesce(to_jsonb($${paramsIndex++}${inferType(val)}), 'null') , true)` params.push(val) dataUpdated = true } if (dataUpdated) { updates.push(`data = ${from}`) } await this.withConnection(ctx, async (connection) => { try { await this.retryTxn(connection, (client) => client.unsafe( `UPDATE ${translateDomain(this.hierarchy.getDomain(tx.objectClass))} SET ${updates.join(', ')} WHERE _id = $3 AND "workspaceId" = $4`, params ) ) if (retrieve) { const object = await this.findDoc(_ctx, connection, tx.objectClass, tx.objectId) return { object } } } catch (err) { console.error(err, { tx, params, updates }) } }) return {} }) } private findDoc ( ctx: MeasureContext, client: postgres.Sql | postgres.ReservedSql, _class: Ref>, _id: Ref, forUpdate: boolean = false ): Promise { return ctx.with('find-doc', { _class }, async () => { const res = await client`SELECT * FROM ${this.client(translateDomain(this.hierarchy.getDomain(_class)))} WHERE _id = ${_id} AND "workspaceId" = ${this.workspaceId.name} ${ forUpdate ? client` FOR UPDATE` : client`` }` const dbDoc = res[0] const domain = this.hierarchy.getDomain(_class) return dbDoc !== undefined ? parseDoc(dbDoc as any, domain) : undefined }) } protected async txRemoveDoc (ctx: MeasureContext, tx: TxRemoveDoc): Promise { await ctx.with('tx-remove-doc', { _class: tx.objectClass }, async (_ctx) => { const domain = translateDomain(this.hierarchy.getDomain(tx.objectClass)) await this.withConnection(_ctx, async (connection) => { await this.retryTxn( connection, (client) => client`DELETE FROM ${client(domain)} WHERE _id = ${tx.objectId} AND "workspaceId" = ${this.workspaceId.name}` ) }) }) return {} } } class PostgresTxAdapter extends PostgresAdapterBase implements TxAdapter { async init (domains?: string[], excludeDomains?: string[]): Promise { const resultDomains = domains ?? [DOMAIN_TX, DOMAIN_MODEL_TX] await createTables(this.client, resultDomains) this._helper.domains = new Set(resultDomains as Domain[]) } override async tx (ctx: MeasureContext, ...tx: Tx[]): Promise { if (tx.length === 0) { return [] } try { const modelTxes: Tx[] = [] const baseTxes: Tx[] = [] for (const _tx of tx) { if (_tx.objectSpace === core.space.Model) { modelTxes.push(_tx) } else { baseTxes.push(_tx) } } if (modelTxes.length > 0) { await this.insert(ctx, DOMAIN_MODEL_TX, modelTxes) } if (baseTxes.length > 0) { await this.insert(ctx, DOMAIN_TX, baseTxes) } } catch (err) { console.error(err) } return [] } async getModel (ctx: MeasureContext): Promise { const res = await this .client`SELECT * FROM ${this.client(translateDomain(DOMAIN_MODEL_TX))} WHERE "workspaceId" = ${this.workspaceId.name} ORDER BY _id ASC, "modifiedOn" ASC` const model = res.map((p) => parseDoc(p as any, DOMAIN_MODEL_TX)) // We need to put all core.account.System transactions first const systemTx: Tx[] = [] const userTx: Tx[] = [] model.forEach((tx) => (tx.modifiedBy === core.account.System && !isPersonAccount(tx) ? systemTx : userTx).push(tx)) return systemTx.concat(userTx) } } /** * @public */ export async function createPostgresAdapter ( ctx: MeasureContext, hierarchy: Hierarchy, url: string, workspaceId: WorkspaceId, modelDb: ModelDb ): Promise { const client = getDBClient(url) const connection = await client.getClient() const adapter = new PostgresAdapter(connection, client, workspaceId, hierarchy, modelDb) return adapter } /** * @public */ export async function createPostgresTxAdapter ( ctx: MeasureContext, hierarchy: Hierarchy, url: string, workspaceId: WorkspaceId, modelDb: ModelDb ): Promise { const client = getDBClient(url) const connection = await client.getClient() const adapter = new PostgresTxAdapter(connection, client, workspaceId, hierarchy, modelDb) await adapter.init() return adapter } function isPersonAccount (tx: Tx): boolean { return ( (tx._class === core.class.TxCreateDoc || tx._class === core.class.TxUpdateDoc || tx._class === core.class.TxRemoveDoc) && ((tx as TxCUD).objectClass === 'contact:class:PersonAccount' || (tx as TxCUD).objectClass === 'contact:class:EmployeeAccount') ) }