Extract postgres base package (#9397)

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
This commit is contained in:
Andrey Sobolev
2025-06-30 12:23:21 +07:00
committed by GitHub
parent 90ef532d64
commit f44b9c427d
39 changed files with 637 additions and 536 deletions
@@ -14,6 +14,7 @@ import { PostgresAdapter } from '../storage'
import { convertArrayParams, decodeArray, filterProjection } from '../utils'
import { genMinModel, test, type ComplexClass } from './minmodel'
import { createDummyClient, type TypedQuery } from './utils'
import { ConnectionMgr } from '@hcengineering/postgres-base'
describe('array conversion', () => {
it('should handle undefined parameters', () => {
@@ -149,6 +150,7 @@ function createTestContext (): { adapter: PostgresAdapter, ctx: MeasureMetricsCo
modelDb.addTxes(ctx, minModel, true)
const adapter = new PostgresAdapter(
c,
new ConnectionMgr(c),
{
url: () => 'test',
close: () => {}
@@ -25,8 +25,13 @@ import core, {
type WorkspaceUuid
} from '@hcengineering/core'
import { type DbAdapter, wrapAdapterToClient } from '@hcengineering/server-core'
import { createPostgresAdapter, createPostgresTxAdapter } from '..'
import { getDBClient, type PostgresClientReference, shutdownPostgres } from '../utils'
import {
createPostgresAdapter,
createPostgresTxAdapter,
getDBClient,
shutdownPostgres,
type PostgresClientReference
} from '..'
import { genMinModel } from './minmodel'
import { createTaskModel, type Task, type TaskComment, taskPlugin } from './tasks'
@@ -40,7 +45,7 @@ describe('postgres operations', () => {
const baseDbUri: string = process.env.DB_URL ?? 'postgresql://root@localhost:26257/defaultdb?sslmode=disable'
let dbUuid = crypto.randomUUID() as WorkspaceUuid
let dbUri: string = baseDbUri.replace('defaultdb', dbUuid)
const clientRef: PostgresClientReference = getDBClient(contextVars, baseDbUri)
const clientRef: PostgresClientReference = getDBClient(baseDbUri)
let hierarchy: Hierarchy
let model: ModelDb
let client: Client
@@ -49,7 +54,7 @@ describe('postgres operations', () => {
afterAll(async () => {
clientRef.close()
await shutdownPostgres(contextVars)
await shutdownPostgres()
})
beforeEach(async () => {
@@ -90,7 +95,6 @@ describe('postgres operations', () => {
const mctx = new MeasureMetricsContext('', {})
const txStorage = await createPostgresTxAdapter(
mctx,
contextVars,
hierarchy,
dbUri,
{
@@ -110,7 +114,6 @@ describe('postgres operations', () => {
const ctx = new MeasureMetricsContext('client', {})
const serverStorage = await createPostgresAdapter(
ctx,
contextVars,
hierarchy,
dbUri,
{
+1 -1
View File
@@ -1,4 +1,4 @@
import type { DBClient } from '../client'
import type { DBClient } from '@hcengineering/postgres-base'
export interface TypedQuery {
query: string
-29
View File
@@ -1,29 +0,0 @@
import type postgres from 'postgres'
import type { ParameterOrJSON } from 'postgres'
import { convertArrayParams, doFetchTypes, getPrepare } from './utils'
export type DBResult = any[] & { count: number }
export interface DBClient {
execute: (query: string, parameters?: ParameterOrJSON<any>[] | undefined) => Promise<DBResult>
release: () => void
reserve: () => Promise<DBClient>
raw: () => postgres.Sql
}
export function createDBClient (client: postgres.Sql, release: () => void = () => {}): DBClient {
return {
execute: (query, parameters) =>
client.unsafe(query, doFetchTypes ? parameters : convertArrayParams(parameters), getPrepare()),
release,
reserve: async () => {
const reserved = await client.reserve()
return createDBClient(reserved, () => {
reserved.release()
})
},
raw: () => client
}
}
+6 -5
View File
@@ -13,19 +13,20 @@
// limitations under the License.
//
import { getDBClient, retryTxn } from '@hcengineering/postgres-base'
import type { WorkspaceDestroyAdapter } from '@hcengineering/server-core'
import { domainSchemas } from './schemas'
import { getDBClient, retryTxn } from './utils'
export { createDBClient } from './client'
export { getDocFieldsByDomains, translateDomain } from './schemas'
export * from './storage'
export { convertDoc, createTables, getDBClient, retryTxn, setDBExtraOptions, shutdownPostgres } from './utils'
export { convertDoc, createTables } from './utils'
export * from '@hcengineering/postgres-base'
export function createPostgreeDestroyAdapter (url: string): WorkspaceDestroyAdapter {
return {
deleteWorkspace: async (ctx, contextVars, workspaceUuid): Promise<void> => {
const client = getDBClient(contextVars, url)
deleteWorkspace: async (ctx, workspaceUuid): Promise<void> => {
const client = getDBClient(url)
try {
if (workspaceUuid == null) {
throw new Error('Workspace uuid is not defined')
+31 -230
View File
@@ -30,7 +30,6 @@ import core, {
DOMAIN_TX,
type FindOptions,
type FindResult,
generateId,
groupByArray,
type Hierarchy,
isOperator,
@@ -62,6 +61,13 @@ import core, {
type WorkspaceIds,
type WorkspaceUuid
} from '@hcengineering/core'
import {
type ConnectionMgr,
createDBClient,
type DBClient,
doFetchTypes,
getDBClient
} from '@hcengineering/postgres-base'
import {
calcHashHash,
type DbAdapter,
@@ -72,7 +78,6 @@ import {
type TxAdapter
} from '@hcengineering/server-core'
import type postgres from 'postgres'
import { createDBClient, type DBClient } from './client'
import {
getDocFieldsByDomains,
getSchema,
@@ -88,10 +93,8 @@ import {
createTables,
DBCollectionHelper,
type DBDoc,
doFetchTypes,
escape,
filterProjection,
getDBClient,
inferType,
isDataField,
isOwner,
@@ -128,212 +131,6 @@ async function * createCursorGenerator (
}
}
class ConnectionInfo {
// It should preserve at least one available connection in pool, other connection should be closed
available: DBClient[] = []
released: boolean = false
constructor (
readonly mgrId: string,
readonly connectionId: string,
protected readonly client: DBClient,
readonly managed: boolean
) {}
async withReserve (action: (reservedClient: DBClient) => Promise<any>, forced: boolean = false): Promise<any> {
let reserved: DBClient | undefined
// Check if we have at least one available connection and reserve one more if required.
if (this.available.length === 0) {
if (this.managed || forced) {
reserved = await this.client.reserve()
}
} else {
reserved = this.available.shift() as DBClient
}
try {
// Use reserved or pool
return await action(reserved ?? this.client)
} catch (err: any) {
console.error(err)
throw err
} finally {
if (this.released) {
try {
reserved?.release()
} catch (err: any) {
console.error('failed to release', err)
}
} else if (reserved !== undefined) {
if (this.available.length > 0) {
reserved?.release()
} else {
this.available.push(reserved)
}
}
}
}
release (): void {
for (const c of [...this.available]) {
c.release()
}
this.available = []
}
}
class ConnectionMgr {
constructor (
protected readonly client: DBClient,
protected readonly connections: () => Map<string, ConnectionInfo>,
readonly mgrId: string
) {}
async write (id: string | undefined, fn: (client: DBClient) => Promise<any>): Promise<void> {
const backoffInterval = 25 // millis
const maxTries = 5
let tries = 0
const realId = id ?? generateId()
const connection = this.getConnection(realId, false)
try {
while (true) {
const retry: boolean | Error = await connection.withReserve(async (client) => {
tries++
try {
await client.execute('BEGIN;')
await fn(client)
await client.execute('COMMIT;')
return true
} catch (err: any) {
await client.execute('ROLLBACK;')
console.error({ message: 'failed to process tx', error: err.message, cause: err })
if (!this.isRetryableError(err) || tries === maxTries) {
return err
} else {
console.log('Transaction failed. Retrying.')
console.log(err.message)
return false
}
}
}, true)
if (retry === true) {
break
}
if (retry instanceof Error) {
// Pass it to exit
throw retry
}
// Retry for a timeout
await new Promise((resolve) => setTimeout(resolve, backoffInterval))
}
} finally {
if (!connection.managed) {
// We need to relase in case it temporaty connection was used
connection.release()
}
}
}
async retry (id: string | undefined, fn: (client: DBClient) => Promise<any>): Promise<any> {
const backoffInterval = 25 // millis
const maxTries = 5
let tries = 0
const realId = id ?? generateId()
// Will reuse reserved if had and use new one if not
const connection = this.getConnection(realId, false)
try {
while (true) {
const retry: false | { result: any } | Error = await connection.withReserve(async (client) => {
tries++
try {
return { result: await fn(client) }
} catch (err: any) {
console.error({ message: 'failed to process sql', error: err.message, cause: err })
if (!this.isRetryableError(err) || tries === maxTries) {
return err
} else {
console.log('Read Transaction failed. Retrying.')
console.log(err.message)
return false
}
}
})
if (retry instanceof Error) {
// Pass it to exit
throw retry
}
if (retry === false) {
// Retry for a timeout
await new Promise((resolve) => setTimeout(resolve, backoffInterval))
continue
}
return retry.result
}
} finally {
if (!connection.managed) {
// We need to relase in case it temporaty connection was used
connection.release()
}
}
}
release (id: string): void {
const conn = this.connections().get(id)
if (conn !== undefined) {
conn.released = true
this.connections().delete(id) // We need to delete first
conn.release()
} else {
console.log('wrne')
}
}
close (): void {
const cnts = this.connections()
for (const [k, conn] of Array.from(cnts.entries()).filter(
([, it]: [string, ConnectionInfo]) => it.mgrId === this.mgrId
)) {
cnts.delete(k)
try {
conn.release()
} catch (err: any) {
console.error('failed to release connection')
}
}
}
getConnection (id: string, managed: boolean = true): ConnectionInfo {
let conn = this.connections().get(id)
if (conn === undefined) {
conn = new ConnectionInfo(this.mgrId, id, this.client, managed)
}
if (managed) {
this.connections().set(id, conn)
}
return conn
}
private isRetryableError (err: any): boolean {
const msg: string = err?.message ?? ''
return (
err.code === '40001' || // Retry transaction
err.code === '55P03' || // Lock not available
err.code === 'CONNECTION_CLOSED' || // This error is thrown if the connection was closed without an error.
err.code === 'CONNECTION_DESTROYED' || // This error is thrown for any queries that were pending when the timeout to sql.end({ timeout: X }) was reached. If the DB client is being closed completely retry will result in CONNECTION_ENDED which is not retried so should be fine.
msg.includes('RETRY_SERIALIZABLE')
)
}
}
class ValuesVariables {
index: number = 1
values: any[] = []
@@ -409,12 +206,10 @@ abstract class PostgresAdapterBase implements DbAdapter {
protected readonly _helper: DBCollectionHelper
protected readonly tableFields = new Map<string, string[]>()
protected connections = new Map<string, ConnectionInfo>()
mgr: ConnectionMgr
constructor (
protected readonly client: DBClient,
protected readonly mgr: ConnectionMgr,
protected readonly refClient: {
url: () => string
close: () => void
@@ -425,15 +220,12 @@ abstract class PostgresAdapterBase implements DbAdapter {
readonly mgrId: string
) {
this._helper = new DBCollectionHelper(this.client, this.workspaceId)
this.mgr = new ConnectionMgr(client, () => this.connections, mgrId)
}
reserveContext (id: string): () => void {
const conn = this.mgr.getConnection(id, true)
this.mgr.getConnection(id, true)
return () => {
conn.released = true
conn.release()
this.connections.delete(id) // We need to delete first
this.mgr.release(id) // We need to release first
}
}
@@ -1863,8 +1655,6 @@ export class PostgresAdapter extends PostgresAdapterBase {
domains?: string[],
excludeDomains?: string[]
): Promise<void> {
this.connections = contextVars.cntInfoPG ?? new Map<string, ConnectionInfo>()
contextVars.cntInfoPG = this.connections
let resultDomains = [...(domains ?? this.hierarchy.domains()), 'kanban']
if (excludeDomains !== undefined) {
resultDomains = resultDomains.filter((it) => !excludeDomains.includes(it))
@@ -2186,9 +1976,6 @@ class PostgresTxAdapter extends PostgresAdapterBase implements TxAdapter {
domains?: string[],
excludeDomains?: string[]
): Promise<void> {
this.connections = contextVars.cntInfoPG ?? new Map<string, ConnectionInfo>()
contextVars.cntInfoPG = this.connections
const resultDomains = domains ?? [DOMAIN_TX, DOMAIN_MODEL_TX]
await initRateLimit.exec(async () => {
const url = this.refClient.url()
@@ -2263,31 +2050,45 @@ function prepareJsonValue (tkey: string, valType: string): { tlkey: string, arro
*/
export async function createPostgresAdapter (
ctx: MeasureContext,
contextVars: Record<string, any>,
hierarchy: Hierarchy,
url: string,
wsIds: WorkspaceIds,
modelDb: ModelDb
): Promise<DbAdapter> {
const client = getDBClient(contextVars, url)
const client = getDBClient(url)
const connection = await client.getClient()
return new PostgresAdapter(createDBClient(connection), client, wsIds.uuid, hierarchy, modelDb, 'default-' + wsIds.url)
return new PostgresAdapter(
createDBClient(connection),
client.mgr,
client,
wsIds.uuid,
hierarchy,
modelDb,
'default-' + wsIds.url
)
}
/**
* @public
*/
export async function createPostgresTxAdapter (
ctx: MeasureContext,
contextVars: Record<string, any>,
hierarchy: Hierarchy,
url: string,
wsIds: WorkspaceIds,
modelDb: ModelDb
): Promise<TxAdapter> {
const client = getDBClient(contextVars, url)
const client = getDBClient(url)
const connection = await client.getClient()
return new PostgresTxAdapter(createDBClient(connection), client, wsIds.uuid, hierarchy, modelDb, 'tx' + wsIds.url)
return new PostgresTxAdapter(
createDBClient(connection),
client.mgr,
client,
wsIds.uuid,
hierarchy,
modelDb,
'tx' + wsIds.url
)
}
function isPersonAccount (tx: Tx): boolean {
+3 -167
View File
@@ -21,7 +21,6 @@ import core, {
type DocumentUpdate,
type Domain,
type FieldIndexConfig,
generateId,
type MeasureContext,
type MixinUpdate,
platformNow,
@@ -31,10 +30,9 @@ import core, {
systemAccountUuid,
type WorkspaceUuid
} from '@hcengineering/core'
import { PlatformError, unknownStatus } from '@hcengineering/platform'
import { type DomainHelperOperations } from '@hcengineering/server-core'
import postgres, { type Options, type ParameterOrJSON } from 'postgres'
import type { DBClient } from './client'
import type postgres from 'postgres'
import { type ParameterOrJSON } from 'postgres'
import {
addSchema,
type DataType,
@@ -46,22 +44,12 @@ import {
type SchemaAndFields,
translateDomain
} from './schemas'
import { retryTxn, type DBClient } from '@hcengineering/postgres-base'
const clientRefs = new Map<string, ClientRef>()
const loadedDomains = new Set<string>()
let loadedTables = new Set<string>()
export async function retryTxn (
pool: postgres.Sql,
operation: (client: postgres.TransactionSql) => Promise<any>
): Promise<any> {
await pool.begin(async (client) => {
const result = await operation(client)
return result
})
}
export const NumericTypes = [
core.class.TypeNumber,
core.class.TypeTimestamp,
@@ -190,158 +178,6 @@ async function createTable (client: postgres.Sql, domain: string): Promise<void>
}
}
/**
* @public
*/
export async function shutdownPostgres (contextVars: Record<string, any>): Promise<void> {
const connections: Map<string, PostgresClientReferenceImpl> | undefined =
contextVars.pgConnections ?? new Map<string, PostgresClientReferenceImpl>()
if (connections === undefined) {
return
}
for (const c of connections.values()) {
c.close(true)
}
connections.clear()
}
export interface PostgresClientReference {
getClient: () => Promise<postgres.Sql>
close: () => void
url: () => string
}
class PostgresClientReferenceImpl {
count: number
client: postgres.Sql | Promise<postgres.Sql>
constructor (
readonly connectionString: string,
client: postgres.Sql | Promise<postgres.Sql>,
readonly onclose: () => void
) {
this.count = 0
this.client = client
}
url (): string {
return this.connectionString
}
async getClient (): Promise<postgres.Sql> {
if (this.client instanceof Promise) {
this.client = await this.client
}
return this.client
}
close (force: boolean = false): void {
this.count--
if (this.count === 0 || force) {
if (force) {
this.count = 0
}
void (async () => {
this.onclose()
const cl = await this.client
await cl.end({ timeout: 1 })
})()
}
}
addRef (): void {
this.count++
}
}
export class ClientRef implements PostgresClientReference {
id = generateId()
constructor (readonly client: PostgresClientReferenceImpl) {
clientRefs.set(this.id, this)
}
url (): string {
return this.client.url()
}
closed = false
async getClient (): Promise<postgres.Sql> {
if (!this.closed) {
return await this.client.getClient()
} else {
throw new PlatformError(unknownStatus('DB client is already closed'))
}
}
close (): void {
// Do not allow double close of mongo connection client
if (!this.closed) {
clientRefs.delete(this.id)
this.closed = true
this.client.close()
}
}
}
export let dbExtraOptions: Partial<Options<any>> = {}
export function setDBExtraOptions (options: Partial<Options<any>>): void {
dbExtraOptions = options
}
export function getPrepare (): { prepare: boolean } {
return { prepare: dbExtraOptions.prepare ?? false }
}
export const doFetchTypes = true
/**
* Initialize a workspace connection to DB
* @public
*/
export function getDBClient (
contextVars: Record<string, any>,
connectionString: string,
database?: string
): PostgresClientReference {
const extraOptions = JSON.parse(process.env.POSTGRES_OPTIONS ?? '{}')
const key = `${connectionString}${extraOptions}`
const connections = contextVars.pgConnections ?? new Map<string, PostgresClientReferenceImpl>()
contextVars.pgConnections = connections
let existing = connections.get(key)
if (existing === undefined) {
const sql = postgres(connectionString, {
connection: {
application_name: 'transactor'
},
database,
max: 10,
min: 2,
connect_timeout: 30,
idle_timeout: 0,
transform: {
undefined: null
},
debug: false,
notice: false,
onnotice (notice) {},
onparameter (key, value) {},
...dbExtraOptions,
...extraOptions,
fetch_types: doFetchTypes
})
existing = new PostgresClientReferenceImpl(connectionString, sql, () => {
connections.delete(key)
})
connections.set(key, existing)
}
// Add reference and return once closable
existing.addRef()
return new ClientRef(existing)
}
export function convertDoc<T extends Doc> (
domain: string,
doc: T,