Merge remote-tracking branch 'origin/develop' into staging

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
This commit is contained in:
Andrey Sobolev
2024-10-28 16:13:49 +07:00
40 changed files with 876 additions and 776 deletions
+6 -8
View File
@@ -1316,9 +1316,6 @@ dependencies:
'@types/pdfjs-dist':
specifier: 2.10.378
version: 2.10.378
'@types/pg':
specifier: ^8.11.6
version: 8.11.6
'@types/png-chunks-extract':
specifier: ^1.0.2
version: 1.0.2
@@ -25518,7 +25515,7 @@ packages:
dev: false
file:projects/account.tgz(@types/node@20.11.19)(esbuild@0.20.1)(ts-node@10.9.2):
resolution: {integrity: sha512-duOFkzCRppNR6RuL9NSfuJoOTnsRRmC8QO8FmT80PlX4+I9X22KpWvgWHM/1I2RKzjDt1AL+W8zH4toqwL/D5g==, tarball: file:projects/account.tgz}
resolution: {integrity: sha512-m2pyxgv41Godk3uDJXacbY/eH+YjKchrlqN268CXDaZNIu0xwQoctNjlxfdxn5Ut/+oXjNP50bFDA+SDgpTOTw==, tarball: file:projects/account.tgz}
id: file:projects/account.tgz
name: '@rush-temp/account'
version: 0.0.0
@@ -25540,6 +25537,7 @@ packages:
node-fetch: 2.7.0
otp-generator: 4.0.1
pg: 8.12.0
postgres: 3.4.4
prettier: 3.2.5
ts-jest: 29.1.2(esbuild@0.20.1)(jest@29.7.0)(typescript@5.3.3)
typescript: 5.3.3
@@ -25556,7 +25554,6 @@ packages:
- kerberos
- mongodb-client-encryption
- node-notifier
- pg-native
- snappy
- socks
- supports-color
@@ -31622,7 +31619,7 @@ packages:
dev: false
file:projects/postgres.tgz(esbuild@0.20.1)(ts-node@10.9.2):
resolution: {integrity: sha512-qZVG4Pk9RAvQfkKRB1iQPPOWsKFvFuFyNBtD/ksT/eW0/ByyGABvYMeQPNenrh3ZH2n/hZ5lH5DON042MXebPg==, tarball: file:projects/postgres.tgz}
resolution: {integrity: sha512-/6jhoJPjD7X4j5u87Epy0kCDI5gedKzlzCM/anvCuB57Igtpi6Ji1cX6LrzyJ33gbswI8tlygVs4miiYWToa5g==, tarball: file:projects/postgres.tgz}
id: file:projects/postgres.tgz
name: '@rush-temp/postgres'
version: 0.0.0
@@ -31639,6 +31636,7 @@ packages:
eslint-plugin-promise: 6.1.1(eslint@8.56.0)
jest: 29.7.0(@types/node@20.11.19)(ts-node@10.9.2)
pg: 8.12.0
postgres: 3.4.4
prettier: 3.2.5
ts-jest: 29.1.2(esbuild@0.20.1)(jest@29.7.0)(typescript@5.3.3)
typescript: 5.3.3
@@ -36075,7 +36073,7 @@ packages:
dev: false
file:projects/tool.tgz(bufferutil@4.0.8)(utf-8-validate@6.0.4):
resolution: {integrity: sha512-dgOs2ysRWXhSAqAdgds7tTLssdROYsTidpPcozKPLpcQnZsmn0Gs42t/y577rd0diONe8AWyXmUcfJMpj0b0mg==, tarball: file:projects/tool.tgz}
resolution: {integrity: sha512-Xiug9crD6dZ0BSoLCOQVEQE05/LitZ2QEPOeqZP0k9K31mffLqyYbCZfbRMvWMzBoFolNR8yUC38wIPvkRal1g==, tarball: file:projects/tool.tgz}
id: file:projects/tool.tgz
name: '@rush-temp/tool'
version: 0.0.0
@@ -36107,6 +36105,7 @@ packages:
mime-types: 2.1.35
mongodb: 6.9.0-dev.20241016.sha.3d5bd513
pg: 8.12.0
postgres: 3.4.4
prettier: 3.2.5
ts-jest: 29.1.2(esbuild@0.20.1)(jest@29.7.0)(typescript@5.3.3)
ts-node: 10.9.2(@types/node@20.11.19)(typescript@5.3.3)
@@ -36126,7 +36125,6 @@ packages:
- kerberos
- mongodb-client-encryption
- node-notifier
- pg-native
- snappy
- socks
- supports-color
+2 -3
View File
@@ -52,8 +52,7 @@
"@types/request": "~2.48.8",
"jest": "^29.7.0",
"ts-jest": "^29.1.1",
"@types/jest": "^29.5.5",
"@types/pg": "^8.11.6"
"@types/jest": "^29.5.5"
},
"dependencies": {
"@elastic/elasticsearch": "^7.14.0",
@@ -157,7 +156,7 @@
"libphonenumber-js": "^1.9.46",
"mime-types": "~2.1.34",
"mongodb": "6.9.0-dev.20241016.sha.3d5bd513",
"pg": "8.12.0",
"postgres": "^3.4.4",
"ws": "^8.18.0"
}
}
+12 -32
View File
@@ -26,11 +26,12 @@ import {
retryTxn,
translateDomain
} from '@hcengineering/postgres'
import { type DBDoc } from '@hcengineering/postgres/types/utils'
import { getTransactorEndpoint } from '@hcengineering/server-client'
import { generateToken } from '@hcengineering/server-token'
import { connect } from '@hcengineering/server-tool'
import { type MongoClient } from 'mongodb'
import { type Pool } from 'pg'
import type postgres from 'postgres'
export async function moveFromMongoToPG (
accountDb: AccountDB,
@@ -64,7 +65,7 @@ export async function moveFromMongoToPG (
async function moveWorkspace (
accountDb: AccountDB,
mongo: MongoClient,
pgClient: Pool,
pgClient: postgres.Sql,
ws: Workspace,
region: string
): Promise<void> {
@@ -84,17 +85,16 @@ async function moveWorkspace (
for (const collection of collections) {
const cursor = collection.find()
const domain = translateDomain(collection.collectionName)
const current = await pgClient.query(`SELECT _id FROM ${domain} WHERE "workspaceId" = $1`, [ws.workspace])
const currentIds = new Set(current.rows.map((r) => r._id))
const current = await pgClient`SELECT _id FROM ${pgClient(domain)} WHERE "workspaceId" = ${ws.workspace}`
const currentIds = new Set(current.map((r) => r._id))
console.log('move domain', domain)
const docs: Doc[] = []
const fields = getDocFieldsByDomains(domain)
const filedsWithData = [...fields, 'data']
const insertFields: string[] = []
const insertFields: string[] = ['workspaceId']
for (const field of filedsWithData) {
insertFields.push(`"${field}"`)
insertFields.push(field)
}
const insertStr = insertFields.join(', ')
while (true) {
while (docs.length < 50000) {
const doc = (await cursor.next()) as Doc | null
@@ -105,35 +105,15 @@ async function moveWorkspace (
if (docs.length === 0) break
while (docs.length > 0) {
const part = docs.splice(0, 500)
const values: any[] = []
const vars: string[] = []
let index = 1
const values: DBDoc[] = []
for (let i = 0; i < part.length; i++) {
const doc = part[i]
const variables: string[] = []
const d = convertDoc(domain, doc, ws.workspace)
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(',')
try {
await retryTxn(pgClient, async (client) => {
await client.query(
`INSERT INTO ${translateDomain(domain)} ("workspaceId", ${insertStr}) VALUES ${vals}`,
values
)
})
} catch (err) {
console.log('error when move doc to', domain, err)
continue
values.push(d)
}
await retryTxn(pgClient, async (client) => {
await client`INSERT INTO ${client(translateDomain(domain))} ${client(values, insertFields)}`
})
}
}
}
+1 -1
View File
@@ -294,7 +294,7 @@ async function migrateDocSections (client: MigrationClient): Promise<void> {
try {
const ydoc = await loadCollaborativeDoc(ctx, storage, client.workspaceId, document.content)
if (ydoc === undefined) {
ctx.error('collaborative document content not found', { document: document.title })
// no content, ignore
continue
}
+43 -6
View File
@@ -29,7 +29,7 @@ import {
SearchResult
} from '../storage'
import { Tx } from '../tx'
import { genMinModel, test, TestMixin } from './minmodel'
import { createDoc, deleteDoc, genMinModel, test, TestMixin, updateDoc } from './minmodel'
const txes = genMinModel()
@@ -59,17 +59,17 @@ class ClientModel extends ModelDb implements Client {
async close (): Promise<void> {}
}
async function createModel (): Promise<{ model: ClientModel, hierarchy: Hierarchy, txDb: TxDb }> {
async function createModel (modelTxes: Tx[] = txes): Promise<{ model: ClientModel, hierarchy: Hierarchy, txDb: TxDb }> {
const hierarchy = new Hierarchy()
for (const tx of txes) {
for (const tx of modelTxes) {
hierarchy.tx(tx)
}
const model = new ClientModel(hierarchy)
for (const tx of txes) {
for (const tx of modelTxes) {
await model.tx(tx)
}
const txDb = new TxDb(hierarchy)
for (const tx of txes) await txDb.tx(tx)
for (const tx of modelTxes) await txDb.tx(tx)
return { model, hierarchy, txDb }
}
@@ -78,7 +78,7 @@ describe('memdb', () => {
const { txDb } = await createModel()
const result = await txDb.findAll(core.class.Tx, {})
expect(result.length).toBe(txes.filter((tx) => tx._class === core.class.TxCreateDoc).length)
expect(result.length).toBe(txes.length)
})
it('should create space', async () => {
@@ -396,4 +396,41 @@ describe('memdb', () => {
expect(e).toEqual(new Error('createDoc cannot be used for objects inherited from AttachedDoc'))
}
})
it('has correct accounts', async () => {
const modTxes = [...txes]
modTxes.push(
createDoc(core.class.Account, {
email: 'system_admin',
role: AccountRole.Owner
})
)
const system1Account = createDoc(core.class.Account, {
email: 'system1',
role: AccountRole.Maintainer
})
modTxes.push(system1Account)
const user1Account = createDoc(core.class.Account, {
email: 'user1',
role: AccountRole.User
})
modTxes.push(user1Account)
modTxes.push(updateDoc(core.class.Account, core.space.Model, system1Account.objectId, { email: 'user1' }))
modTxes.push(deleteDoc(core.class.Account, core.space.Model, user1Account.objectId))
const { model } = await createModel(modTxes)
expect(model.getAccountByEmail('system_admin')).not.toBeUndefined()
expect(model.getAccountByEmail('system_admin')?.role).toBe(AccountRole.Owner)
expect(model.getAccountByEmail('system1')).toBeUndefined()
expect(model.getAccountByEmail('user1')).not.toBeUndefined()
expect(model.getAccountByEmail('user1')?.role).toBe(AccountRole.Maintainer)
})
})
+16 -3
View File
@@ -15,10 +15,10 @@
import type { IntlString, Plugin } from '@hcengineering/platform'
import { plugin } from '@hcengineering/platform'
import type { Arr, Class, Data, Doc, Interface, Mixin, Obj, Ref } from '../classes'
import type { Arr, Class, Data, Doc, Interface, Mixin, Obj, Ref, Space } from '../classes'
import { AttachedDoc, ClassifierKind, DOMAIN_MODEL } from '../classes'
import core from '../component'
import type { TxCUD, TxCreateDoc } from '../tx'
import type { DocumentUpdate, TxCUD, TxCreateDoc, TxRemoveDoc, TxUpdateDoc } from '../tx'
import { DOMAIN_TX, TxFactory } from '../tx'
const txFactory = new TxFactory(core.account.System)
@@ -31,10 +31,23 @@ function createInterface (_interface: Ref<Interface<Doc>>, attributes: Data<Inte
return txFactory.createTxCreateDoc(core.class.Interface, core.space.Model, attributes, _interface)
}
export function createDoc<T extends Doc> (_class: Ref<Class<T>>, attributes: Data<T>): TxCreateDoc<Doc> {
export function createDoc<T extends Doc> (_class: Ref<Class<T>>, attributes: Data<T>): TxCreateDoc<T> {
return txFactory.createTxCreateDoc(_class, core.space.Model, attributes)
}
export function updateDoc<T extends Doc> (
_class: Ref<Class<T>>,
space: Ref<Space>,
objectId: Ref<T>,
operations: DocumentUpdate<T>
): TxUpdateDoc<Doc> {
return txFactory.createTxUpdateDoc(_class, space, objectId, operations)
}
export function deleteDoc<T extends Doc> (_class: Ref<Class<T>>, space: Ref<Space>, objectId: Ref<T>): TxRemoveDoc<Doc> {
return txFactory.createTxRemoveDoc(_class, space, objectId)
}
export interface TestMixin extends Doc {
arr: Arr<string>
}
+7
View File
@@ -441,6 +441,8 @@ async function buildModel (
)
})
userTx.sort(compareTxes)
let txes = systemTx.concat(userTx)
if (modelFilter !== undefined) {
txes = await modelFilter(txes)
@@ -473,3 +475,8 @@ function getLastTxTime (txes: Tx[]): number {
}
return lastTxTime
}
function compareTxes (a: Tx, b: Tx): number {
const result = a._id.localeCompare(b._id)
return result !== 0 ? result : a.modifiedOn - b.modifiedOn
}
+37 -6
View File
@@ -32,7 +32,7 @@ export abstract class MemDb extends TxProcessor implements Storage {
private readonly objectById = new Map<Ref<Doc>, Doc>()
private readonly accountByPersonId = new Map<Ref<Doc>, Account[]>()
private readonly accountByEmail = new Map<string, Account>()
private readonly accountByEmail = new Map<string, [string, Account][]>()
constructor (protected readonly hierarchy: Hierarchy) {
super()
@@ -83,7 +83,14 @@ export abstract class MemDb extends TxProcessor implements Storage {
}
getAccountByEmail (email: Account['email']): Account | undefined {
return this.accountByEmail.get(email)
const accounts = this.accountByEmail.get(email)
if (accounts === undefined || accounts.length === 0) {
return undefined
}
if (accounts.length > 0) {
return accounts[accounts.length - 1][1]
}
}
findObject<T extends Doc>(_id: Ref<T>): T | undefined {
@@ -225,6 +232,14 @@ export abstract class MemDb extends TxProcessor implements Storage {
)
}
addAccount (account: Account): void {
if (!this.accountByEmail.has(account.email)) {
this.accountByEmail.set(account.email, [])
}
this.accountByEmail.get(account.email)?.push([account._id, account])
}
addDoc (doc: Doc): void {
this.hierarchy.getAncestors(doc._class).forEach((_class) => {
const arr = this.getObjectsByClass(_class)
@@ -232,7 +247,9 @@ export abstract class MemDb extends TxProcessor implements Storage {
})
if (this.hierarchy.isDerived(doc._class, core.class.Account)) {
const account = doc as Account
this.accountByEmail.set(account.email, account)
this.addAccount(account)
if (account.person !== undefined) {
this.accountByPersonId.set(account.person, [...(this.accountByPersonId.get(account.person) ?? []), account])
}
@@ -240,6 +257,19 @@ export abstract class MemDb extends TxProcessor implements Storage {
this.objectById.set(doc._id, doc)
}
delAccount (account: Account): void {
const accounts = this.accountByEmail.get(account.email)
if (accounts !== undefined) {
const newAccounts = accounts.filter((it) => it[0] !== account._id)
if (newAccounts.length === 0) {
this.accountByEmail.delete(account.email)
} else {
this.accountByEmail.set(account.email, newAccounts)
}
}
}
delDoc (_id: Ref<Doc>): void {
const doc = this.objectById.get(_id)
if (doc === undefined) {
@@ -251,7 +281,8 @@ export abstract class MemDb extends TxProcessor implements Storage {
})
if (this.hierarchy.isDerived(doc._class, core.class.Account)) {
const account = doc as Account
this.accountByEmail.delete(account.email)
this.delAccount(account)
if (account.person !== undefined) {
const acc = this.accountByPersonId.get(account.person) ?? []
this.accountByPersonId.set(
@@ -280,8 +311,8 @@ export abstract class MemDb extends TxProcessor implements Storage {
}
} else if (newEmail !== undefined) {
const account = doc as Account
this.accountByEmail.delete(account.email)
this.accountByEmail.set(newEmail, account)
this.delAccount(account)
this.addAccount({ ...account, email: newEmail })
}
}
}
-3
View File
@@ -23,9 +23,6 @@
min-width: 0;
border: 1px solid var(--theme-divider-color); // var(--global-surface-02-BorderColor);
border-radius: var(--small-focus-BorderRadius);
border-bottom-right-radius: 0;
border-top-right-radius: 0;
border-right: 0;
&:not(.modal) {
background-color: var(--theme-panel-color); // var(--global-surface-02-BackgroundColor);
@@ -64,6 +64,7 @@
props = $panelstore.panel
})
}
$panelstore.panel.refit = fitPopupInstance
} else {
props = undefined
}
@@ -132,6 +133,7 @@
if (!keepSize && props?.element === 'content') {
keepSize = true
resizeObserver(contentPanel, checkResize)
if (!contentPanel.hasAttribute('data-id')) contentPanel.setAttribute('data-id', 'contentPanel')
}
}
+15 -13
View File
@@ -174,7 +174,7 @@
}
}
const handleScroll = (event: MouseEvent): void => {
const handleScroll = (event: PointerEvent): void => {
scrolling = false
if (
(divBar == null && isScrolling === 'vertical') ||
@@ -185,7 +185,7 @@
}
const rectScroll = divScroll.getBoundingClientRect()
if (isScrolling === 'vertical') {
let Y = event.clientY - dXY
let Y = Math.round(event.clientY) - dXY
if (Y < rectScroll.top + shiftTop + 2) Y = rectScroll.top + shiftTop + 2
if (Y > rectScroll.bottom - divBar.clientHeight - shiftBottom - 2) {
Y = rectScroll.bottom - divBar.clientHeight - shiftBottom - 2
@@ -201,7 +201,7 @@
divScroll.scrollTop = (divScroll.scrollHeight - divScroll.clientHeight) * procBar
}
} else if (isScrolling === 'horizontal') {
let X = event.clientX - dXY
let X = Math.round(event.clientX) - dXY
if (X < rectScroll.left + 2 + shiftLeft) X = rectScroll.left + 2 + shiftLeft
if (X > rectScroll.right - divBarH.clientWidth - (mask !== 'none' ? 12 : 2) - shiftRight) {
X = rectScroll.right - divBarH.clientWidth - (mask !== 'none' ? 12 : 2) - shiftRight
@@ -215,18 +215,18 @@
}
}
const onScrollEnd = (): void => {
document.removeEventListener('mousemove', handleScroll)
document.removeEventListener('pointermove', handleScroll)
document.body.style.userSelect = 'auto'
document.body.style.webkitUserSelect = 'auto'
document.removeEventListener('mouseup', onScrollEnd)
document.removeEventListener('pointerup', onScrollEnd)
isScrolling = false
}
const onScrollStart = (event: MouseEvent, direction: 'vertical' | 'horizontal'): void => {
const onScrollStart = (event: PointerEvent, direction: 'vertical' | 'horizontal'): void => {
if (divScroll == null) return
scrolling = false
dXY = direction === 'vertical' ? event.offsetY : event.offsetX
document.addEventListener('mouseup', onScrollEnd)
document.addEventListener('mousemove', handleScroll)
dXY = Math.round(direction === 'vertical' ? event.offsetY : event.offsetX)
document.addEventListener('pointerup', onScrollEnd)
document.addEventListener('pointermove', handleScroll)
document.body.style.userSelect = 'none'
document.body.style.webkitUserSelect = 'none'
isScrolling = direction
@@ -666,10 +666,10 @@
class:hovered={isScrolling === 'vertical'}
class:reverse={scrollDirection === 'vertical-reverse'}
bind:this={divBar}
on:mousedown|stopPropagation={(ev) => {
on:pointerdown|stopPropagation={(ev) => {
onScrollStart(ev, 'vertical')
}}
on:mouseleave={checkFade}
on:pointerleave={checkFade}
/>
{/if}
{#if horizontal && maskH !== 'none'}
@@ -687,10 +687,10 @@
class="bar-horizontal"
class:hovered={isScrolling === 'horizontal'}
bind:this={divBarH}
on:mousedown|stopPropagation={(ev) => {
on:pointerdown|stopPropagation={(ev) => {
onScrollStart(ev, 'horizontal')
}}
on:mouseleave={checkFade}
on:pointerleave={checkFade}
/>
{/if}
</div>
@@ -864,9 +864,11 @@
overscroll-behavior: none;
}
&::-webkit-scrollbar:vertical {
display: none;
width: 0;
}
&::-webkit-scrollbar:horizontal {
display: none;
height: 0;
}
+117 -88
View File
@@ -24,6 +24,7 @@
separatorsStore,
SeparatorState
} from '..'
import { panelstore } from '../panelup'
export let prevElementSize: SeparatedItem | undefined = undefined
export let nextElementSize: SeparatedItem | undefined = undefined
@@ -52,6 +53,7 @@
let isSeparate: boolean = false
let excludedIndexes: number[] = []
let correctedIndex: number = index
let realIndex: number = index
let offset: number = 0
let separatorsSizes: number[] | null = null
const separatorsWide: { before: number, after: number, total: number } = { before: 0, after: 0, total: 0 }
@@ -80,7 +82,7 @@
if (prevElementSize !== undefined) prevElSize = prevElementSize
if (nextElementSize !== undefined) nextElSize = nextElementSize
setTimeout(() => {
if (!parentElement && separator) parentElement = separator.parentElement
if (parentElement === null && separator != null) parentElement = separator.parentElement
checkSibling(true)
calculateSeparators()
})
@@ -109,21 +111,34 @@
const sizePx = direction === 'horizontal' ? rect.width : rect.height
element.setAttribute('data-size', `${sizePx}`)
if (sState === SeparatorState.NORMAL) {
if (separators) separators[index + (next ? 1 : 0)].size = pxToRem(sizePx)
if (separators != null) separators[index + (next ? 1 : 0)].size = pxToRem(sizePx)
if (next) nextElSize.size = typeof size === 'number' ? pxToRem(sizePx) : size
else prevElSize.size = typeof size === 'number' ? pxToRem(sizePx) : size
}
}
const getStyles = (
element: Element | null,
dropStyles: string[] = ['min-width', 'max-width', 'width']
): Map<string, string> => {
const result = new Map<string, string>()
const style = element != null ? element.getAttribute('style') : null
if (style !== null) {
style
.replace(/ /g, '')
.split(';')
.filter((f) => f !== '')
.forEach((st) => result.set(st.split(':')[0], st.split(':')[1]))
dropStyles.forEach((key) => result.delete(key))
}
return result
}
const generateMap = (): void => {
if (parentElement === null) return
if (parentElement == null || separators === null || separatorsSizes === null) return
const children: Element[] = Array.from(parentElement.children)
if (children.length > 1 && separators !== null && separatorsSizes !== null) {
const elements = children.filter(
(el) =>
!el.classList.contains('antiSeparator') && (el.hasAttribute('data-size') || el.hasAttribute('data-auto'))
)
const hasSep = elements.filter((el) => el.hasAttribute('data-float')).map((el) => el.getAttribute('data-float'))
if (children.length > 1) {
const hasSep = children.filter((el) => el.hasAttribute('data-float')).map((el) => el.getAttribute('data-float'))
const excluded = separators
.filter((separ) => separ.float !== undefined && !hasSep.includes(separ.float))
.map((separ) => separ.float)
@@ -132,43 +147,49 @@
if (excluded.includes(separ.float)) excludedIndexes.push(i)
})
correctedIndex = index - excludedIndexes.filter((i) => i < index).length
realIndex = correctedIndex
const sm: SeparatedElement[] = []
let ind: number = 0
elements.forEach((element, i) => {
if (separators && excluded.includes(separators[i].float)) ind++
const styles = new Map<string, string>()
const dropStyles = ['min-width', 'max-width', 'width']
const style = elements[i] ? elements[i].getAttribute('style') : null
if (style !== null) {
style
.replace(/ /g, '')
.split(';')
.filter((f) => f !== '')
.forEach((st) => styles.set(st.split(':')[0], st.split(':')[1]))
dropStyles.forEach((key) => styles.delete(key))
let drop: number = 0
children.forEach((element, i) => {
if (separators != null) {
if (separators[ind]?.float !== undefined && excluded.includes(separators[ind].float)) {
ind++
drop++
}
const styles: Map<string, string> = getStyles(element)
const rect = element.getBoundingClientRect()
const size = direction === 'horizontal' ? rect.width : rect.height
const sep = element.classList.contains('antiSeparator')
const extra = !(sep || element.hasAttribute('data-size') || element.hasAttribute('data-auto'))
if (extra) realIndex++
if (!sep) {
sm.push({
id: extra ? -1 : ind,
element,
styles,
minSize: extra
? size
: typeof separators[ind].minSize === 'number'
? remToPx(separators[ind].minSize as number)
: remToPx(20),
maxSize: extra
? size
: typeof separators[ind].maxSize === 'number'
? remToPx(separators[ind].maxSize as number)
: -1,
size,
begin: ind - drop <= correctedIndex,
resize: false,
float: extra ? undefined : separators[ind].float
})
if (!extra) ind++
}
}
const rect = element.getBoundingClientRect()
const size = direction === 'horizontal' ? rect.width : rect.height
if (separators) {
sm.push({
id: ind,
element,
styles,
minSize:
typeof separators[ind].minSize === 'number' ? remToPx(separators[ind].minSize as number) : remToPx(20),
maxSize: typeof separators[ind].maxSize === 'number' ? remToPx(separators[ind].maxSize as number) : -1,
size,
begin: i <= correctedIndex,
resize: false,
float: separators[ind].float
})
}
ind++
})
separatorMap = sm
const cropIndex = correctedIndex - excludedIndexes.filter((ex) => ex < correctedIndex).length
const startBoxes = separatorMap.filter((_, i) => i < cropIndex + 1)
const endBoxes = separatorMap.slice(cropIndex + 1, sm.length)
const startBoxes = separatorMap.filter((sm) => sm.begin)
const endBoxes = separatorMap.filter((sm) => !sm.begin)
containers.minStart = startBoxes.map((box) => box.minSize).reduce((prev, a) => prev + a, 0)
containers.minEnd = endBoxes.map((box) => box.minSize).reduce((prev, a) => prev + a, 0)
containers.maxStart =
@@ -232,7 +253,7 @@
}
if (isSeparate) style += 'pointer-events:none;'
item.element.setAttribute('style', style)
if (final) {
if (final && item.id !== -1) {
const rect = item.element.getBoundingClientRect()
item.element.setAttribute(
item.maxSize === -1 ? 'data-auto' : 'data-size',
@@ -246,7 +267,7 @@
const resizeContainer = (id: number, min: number, max: number, count: number, stretch: boolean = false): number => {
const diff = max - min
if (diff) {
if (diff !== 0) {
const size = min + (count >= diff ? (stretch ? diff : 0) : stretch ? count : diff - count)
separatorMap[id].size = size
separatorMap[id].resize = true
@@ -260,13 +281,13 @@
return 0
}
function mouseMove (event: MouseEvent) {
function pointerMove (event: PointerEvent): void {
if (sState === SeparatorState.NORMAL) normalMouseMove(event)
else if (sState === SeparatorState.FLOAT) floatMouseMove(event)
}
const preparePanel = (): void => {
if (!parentElement || parentSize === null) return
if (parentElement === null || parentSize === null) return
setSize(parentElement, panel.size === 'auto' ? 'auto' : remToPx(panel.size))
const s = separator.getBoundingClientRect()
if (s) {
@@ -282,9 +303,9 @@
parentElement.style.pointerEvents = 'none'
}
function floatMouseMove (event: MouseEvent) {
function floatMouseMove (event: PointerEvent): void {
if (!isSeparate || parentSize === null || parentElement === null) return
const coord: number = direction === 'horizontal' ? event.x - offset : event.y - offset
const coord: number = Math.round(direction === 'horizontal' ? event.x - offset : event.y - offset)
const parentCoord: number = coord - parentSize.start
const min = remToPx(panel.minSize === 'auto' ? 10 : panel.minSize)
const max = remToPx(panel.maxSize === 'auto' ? 30 : panel.maxSize)
@@ -304,9 +325,9 @@
setSize(parentElement, newCoord)
}
function normalMouseMove (event: MouseEvent) {
if (!isSeparate || separatorMap === null || parentSize === null || separatorsSizes === null) return
const coord: number = direction === 'horizontal' ? event.x - offset : event.y - offset
function normalMouseMove (event: PointerEvent): void {
if (!isSeparate || separatorMap === undefined || parentSize === null || separatorsSizes === null) return
const coord: number = Math.round(direction === 'horizontal' ? event.x - offset : event.y - offset)
let parentCoord: number = coord - parentSize.start
let prevCoord: number = separatorMap
.filter((f) => f.begin)
@@ -330,22 +351,22 @@
if (remains !== 0) {
const reverse = remains < 0
if (reverse) remains = Math.abs(remains)
const minusId = correctedIndex + (reverse ? 1 : 0)
const plusId = correctedIndex + (reverse ? 0 : 1)
const minusId = realIndex + (reverse ? 1 : 0)
const plusId = realIndex + (reverse ? 0 : 1)
const minusAutoBoxes = separatorMap.filter(
(s, i) => s.maxSize === -1 && ((!reverse && i < correctedIndex) || (reverse && i > correctedIndex + 1))
(s, i) => s.maxSize === -1 && ((!reverse && i < realIndex) || (reverse && i > realIndex + 1))
)
const minusBoxes = separatorMap.filter(
(s, i) => s.maxSize !== -1 && ((!reverse && i < correctedIndex) || (reverse && i > correctedIndex + 1))
(s, i) => s.maxSize !== -1 && ((!reverse && i < realIndex) || (reverse && i > realIndex + 1))
)
const minusBox = separatorMap[minusId]
const startMinus = separatorMap[minusId].maxSize === -1
const plusAutoBoxes = separatorMap.filter(
(s, i) => s.maxSize === -1 && ((!reverse && i > correctedIndex + 1) || (reverse && i < correctedIndex))
(s, i) => s.maxSize === -1 && ((!reverse && i > realIndex + 1) || (reverse && i < realIndex))
)
const plusBoxes = separatorMap.filter(
(s, i) => s.maxSize !== -1 && ((!reverse && i > correctedIndex + 1) || (reverse && i < correctedIndex))
(s, i) => s.maxSize !== -1 && ((!reverse && i > realIndex + 1) || (reverse && i < realIndex))
)
const plusBox = separatorMap[plusId]
const startPlus = separatorMap[plusId].maxSize === -1
@@ -354,46 +375,53 @@
if (startMinus && minusBox.size - minusBox.minSize > 0) {
remains = resizeContainer(minusId, minusBox.minSize, minusBox.size, remains)
}
if (remains && minusAutoBoxes.length > 0) {
if (remains > 0 && minusAutoBoxes.length > 0) {
minusAutoBoxes.forEach((box) => {
if (remains) remains = resizeContainer(box.id, box.minSize, box.size, remains)
if (remains > 0) remains = resizeContainer(box.id, box.minSize, box.size, remains)
})
}
if (remains && !startMinus && minusBox.size - minusBox.minSize > 0) {
if (remains > 0 && !startMinus && minusBox.size - minusBox.minSize > 0) {
remains = resizeContainer(minusId, minusBox.minSize, minusBox.size, remains)
}
if (remains && minusBoxes.length > 0) {
if (remains > 0 && minusBoxes.length > 0) {
minusBoxes.forEach((box) => {
if (remains) remains = resizeContainer(box.id, box.minSize, box.size, remains)
if (remains > 0) remains = resizeContainer(box.id, box.minSize, box.size, remains)
})
}
let needAdd: number = Math.abs(diff) - remains
// Find for stretch
if (needAdd && startPlus) needAdd = stretchContainer(plusId, plusBox.size + needAdd)
if (needAdd && plusAutoBoxes.length > 0) {
if (needAdd > 0 && startPlus) needAdd = stretchContainer(plusId, plusBox.size + needAdd)
if (needAdd > 0 && plusAutoBoxes.length > 0) {
const div = needAdd / plusAutoBoxes.length
plusAutoBoxes.forEach((box) => (needAdd = stretchContainer(box.id, box.size + div)))
}
if (needAdd && plusBox.maxSize - plusBox.size > 0) {
if (needAdd > 0 && plusBox.maxSize - plusBox.size > 0) {
needAdd = resizeContainer(plusId, plusBox.size, plusBox.maxSize, needAdd, true)
}
if (needAdd && plusBoxes.length > 0) {
if (needAdd > 0 && plusBoxes.length > 0) {
plusBoxes.forEach((box) => {
if (needAdd) needAdd = resizeContainer(box.id, box.size, box.maxSize, needAdd, true)
if (needAdd > 0) needAdd = resizeContainer(box.id, box.size, box.maxSize, needAdd, true)
})
}
separatorMap = separatorMap
}
applyStyles()
if ($panelstore.panel?.refit !== undefined) $panelstore.panel.refit()
}
function mouseUp () {
function pointerUp (): void {
finalSeparation()
document.removeEventListener('pointermove', pointerMove)
document.removeEventListener('pointerup', pointerUp)
}
function finalSeparation (): void {
isSeparate = false
if (sState === SeparatorState.NORMAL) {
applyStyles(true)
if (index !== -1 && separators && separatorMap) {
if (index !== -1 && separators != null && separatorMap != null) {
let ind: number = 0
const sep: SeparatedItem[] = []
separatorMap = separatorMap.filter((sm) => sm.id !== -1)
separators.forEach((sm, i) => {
let save = false
if (excludedIndexes.includes(i)) {
@@ -412,17 +440,20 @@
})
saveSeparator(name, false, sep)
}
} else if (sState === SeparatorState.FLOAT && parentElement) {
} else if (sState === SeparatorState.FLOAT && parentElement != null) {
parentElement.style.pointerEvents = 'all'
saveSeparator(name, float, panel)
}
document.body.style.cursor = ''
document.removeEventListener('mousemove', mouseMove)
document.removeEventListener('mouseup', mouseUp)
}
function mouseDown (event: MouseEvent) {
if (!parentElement) return
function pointerDown (event: PointerEvent): void {
prepareSeparation(event)
document.addEventListener('pointermove', pointerMove)
document.addEventListener('pointerup', pointerUp)
}
function prepareSeparation (event: PointerEvent): void {
if (parentElement == null) return
if (sState === SeparatorState.FLOAT && parentElement === null) {
checkParent()
return
@@ -430,7 +461,7 @@
checkSibling()
return
}
offset = direction === 'horizontal' ? event.offsetX : event.offsetY
offset = Math.round(direction === 'horizontal' ? event.offsetX : event.offsetY)
const p = parentElement.getBoundingClientRect()
parentSize =
direction === 'horizontal'
@@ -441,33 +472,31 @@
generateMap()
applyStyles(true)
} else if (sState === SeparatorState.FLOAT) preparePanel()
document.addEventListener('mousemove', mouseMove)
document.addEventListener('mouseup', mouseUp)
document.body.style.cursor = direction === 'horizontal' ? 'col-resize' : 'row-resize'
}
const checkSibling = (start: boolean = false): void => {
if (separator === null) return
if ((prevElement === null || start) && separator) {
if ((prevElement === null || start) && separator != null) {
prevElement = separator.previousElementSibling as HTMLElement
}
if ((nextElement === null || start) && separator) {
if ((nextElement === null || start) && separator != null) {
nextElement = separator.nextElementSibling as HTMLElement
}
if (separators && prevElement && separators[index].float !== undefined) {
if (separators != null && prevElement != null && separators[index].float !== undefined) {
prevElement.setAttribute('data-float', separators[index].float ?? '')
}
if (separators && nextElement && separators[index + 1].float !== undefined) {
if (separators != null && nextElement != null && separators[index + 1].float !== undefined) {
nextElement.setAttribute('data-float', separators[index + 1].float ?? '')
}
}
const checkParent = (): void => {
if (parentElement === null && separator) parentElement = separator.parentElement as HTMLElement
if (parentElement && typeof float === 'string') parentElement.setAttribute('data-float', float)
if (parentElement === null && separator != null) parentElement = separator.parentElement as HTMLElement
if (parentElement != null && typeof float === 'string') parentElement.setAttribute('data-float', float)
}
const calculateSeparators = (): void => {
if (parentElement) {
if (parentElement != null) {
const elements: Element[] = Array.from(parentElement.children)
separatorsSizes = elements
.filter((el) => el.classList.contains('antiSeparator'))
@@ -483,7 +512,7 @@
if (parentElement == null || checkElements || sState !== SeparatorState.NORMAL) return
checkElements = true
setTimeout(() => {
if (parentElement != null && separators) {
if (parentElement != null && separators != null) {
const children: Element[] = Array.from(parentElement.children)
let totalSize: number = 0
let ind: number = 0
@@ -513,9 +542,9 @@
let ind: number = 0
reverseSep.forEach((separ, i) => {
const pass = excluded.includes(separ.float)
if (diff > 0 && !pass && separators) {
if (diff > 0 && !pass && separators != null) {
const box = rects.get(reverseSep.length - ind - 1)
if (box) {
if (box != null) {
const minSize: number = remToPx(separ.minSize === 'auto' ? 20 : separ.minSize)
const forCrop = box.size - minSize
if (forCrop > 0) {
@@ -545,7 +574,7 @@
}
onMount(() => {
if (separator) {
if (separator != null) {
parentElement = separator.parentElement as HTMLElement
if (sState === SeparatorState.FLOAT) checkParent()
else if (sState === SeparatorState.NORMAL) {
@@ -586,7 +615,7 @@
class:short
class:hovered={isSeparate}
data-size={separatorSize}
on:mousedown|stopPropagation={mouseDown}
on:pointerdown|stopPropagation={pointerDown}
/>
{/if}
+1
View File
@@ -8,6 +8,7 @@ export interface PanelProps {
_class: string
element?: PopupAlignment
rightSection?: AnyComponent
refit?: () => void
}
export const panelstore = writable<{ panel?: PanelProps | undefined }>({ panel: undefined })
@@ -152,9 +152,9 @@
}}
>
{#if personInfo}
<Avatar name={person?.name ?? personInfo.name} {person} size={'full'} />
<Avatar name={person?.name ?? personInfo.name} {person} size={'full'} showStatus={false} />
{:else if hoveredRoomX === x && hoveredRoomY === y}
<Avatar name={meName} person={meAvatar} size={'full'} />
<Avatar name={meName} person={meAvatar} size={'full'} showStatus={false} />
{/if}
</div>
{/each}
@@ -90,7 +90,7 @@
<style lang="scss">
.container {
padding: var(--spacing-2) var(--spacing-2) var(--spacing-2_5);
min-height: calc(4.75rem + 0.5px);
min-height: calc(4.875rem - 0.5px);
border-bottom: 1px solid var(--theme-divider-color);
}
</style>
@@ -2,6 +2,4 @@
import PlanView from './PlanView.svelte'
</script>
<div class="hulyPanels-container">
<PlanView on:change />
</div>
<PlanView on:change />
@@ -86,20 +86,18 @@
{#if $deviceInfo.navigator.visible}
<ToDosNavigator bind:mode bind:tag bind:currentDate />
<Separator
name={'time'}
float={$deviceInfo.navigator.float}
index={0}
disabledWhen={['panel-aside']}
color={'var(--theme-divider-color)'}
/>
<Separator name={'time'} float={$deviceInfo.navigator.float} index={0} color={'var(--theme-divider-color)'} />
{/if}
<div class="flex-col w-full clear-mins" class:left-divider={!$deviceInfo.navigator.visible} bind:this={mainPanel}>
<ToDos {mode} {tag} bind:currentDate />
</div>
{#if visibleCalendar}
<Separator name={'time'} index={1} color={'transparent'} separatorSize={0} short />
<div class="flex-col clear-mins" bind:this={replacedPanel}>
<PlanningCalendar {dragItem} bind:currentDate displayedDaysCount={5} on:dragDrop={drop} />
</div>
<PlanningCalendar
{dragItem}
bind:element={replacedPanel}
bind:currentDate
displayedDaysCount={5}
on:dragDrop={drop}
/>
{/if}
@@ -27,6 +27,7 @@
export let dragItem: ToDo | null = null
export let currentDate: Date = new Date()
export let displayedDaysCount = 1
export let element: HTMLElement | undefined = undefined
export let createComponent: AnyComponent | undefined = calendar.component.CreateEvent
const q = createQuery()
@@ -175,6 +176,7 @@
<div
class="hulyComponent modal"
bind:this={element}
use:resizeObserver={(element) => {
showLabel = showLabel ? element.clientWidth > rem(3.5) + 399 : element.clientWidth > rem(3.5) + 400
}}
+1 -1
View File
@@ -28,7 +28,7 @@ export function getNearest (events: WorkSlot[]): WorkSlot | undefined {
export const timeSeparators: DefSeparators = [
{ minSize: 18, size: 18, maxSize: 22.5, float: 'navigator' },
null,
{ minSize: 20, size: 41.25, maxSize: 90 }
{ minSize: 25, size: 41.25, maxSize: 90 }
]
/**
@@ -872,7 +872,7 @@
application: currentApplication?._id
}}
/>
<div class="workbench-container inner">
<div class="workbench-container inner" class:rounded={$sidebarStore.variant === SidebarVariant.EXPANDED}>
{#if mainNavigator}
<!-- svelte-ignore a11y-click-events-have-key-events -->
<!-- svelte-ignore a11y-no-static-element-interactions -->
@@ -932,6 +932,7 @@
<div
bind:this={contentPanel}
class={navigatorModel === undefined ? 'hulyPanels-container' : 'hulyComponent overflow-hidden'}
data-id={'contentPanel'}
>
{#if currentApplication && currentApplication.component}
<Component
@@ -972,11 +973,11 @@
</div>
{/if}
</div>
{#if $sidebarStore.variant === SidebarVariant.EXPANDED}
<Separator name={'main'} index={0} color={'transparent'} separatorSize={0} short />
{/if}
<WidgetsBar />
</div>
{#if $sidebarStore.variant === SidebarVariant.EXPANDED}
<Separator name={'main'} index={0} color={'transparent'} separatorSize={0} short />
{/if}
<WidgetsBar />
<Dock />
<div bind:this={cover} class="cover" />
<TooltipInstance />
@@ -1002,12 +1003,15 @@
min-height: 0;
width: 100%;
height: 100%;
background-color: var(--theme-statusbar-color);
background-color: var(--theme-panel-color);
touch-action: none;
&.inner {
background-color: var(--theme-navpanel-color);
border-radius: 0 var(--medium-BorderRadius) var(--medium-BorderRadius) 0;
&.rounded {
border-radius: 0 var(--medium-BorderRadius) var(--medium-BorderRadius) 0;
}
}
&:not(.inner)::after {
position: absolute;
@@ -1015,8 +1019,6 @@
inset: 0;
border: 1px solid var(--theme-divider-color);
border-radius: var(--medium-BorderRadius);
border-bottom-right-radius: 0;
border-top-right-radius: 0;
pointer-events: none;
}
.antiPanel-application {
@@ -17,6 +17,7 @@
import { WidgetPreference, SidebarEvent, TxSidebarEvent, OpenSidebarWidgetParams } from '@hcengineering/workbench'
import { Tx } from '@hcengineering/core'
import { onMount } from 'svelte'
import { panelstore } from '@hcengineering/ui'
import workbench from '../../plugin'
import { createWidgetTab, openWidget, sidebarStore, SidebarVariant } from '../../sidebar'
@@ -33,7 +34,8 @@
preferences = res
})
$: size = $sidebarStore.variant === SidebarVariant.MINI ? 'mini' : undefined
$: mini = $sidebarStore.variant === SidebarVariant.MINI
$: if ((!mini || mini) && $panelstore.panel?.refit !== undefined) $panelstore.panel.refit()
function txListener (tx: Tx): void {
if (tx._class === workbench.class.TxSidebarEvent) {
@@ -59,8 +61,8 @@
})
</script>
<div class="antiPanel-component antiComponent root size-{size}" id="sidebar">
{#if $sidebarStore.variant === SidebarVariant.MINI}
<div class="antiPanel-application vertical root" class:mini id="sidebar">
{#if mini}
<SidebarMini {widgets} {preferences} />
{:else if $sidebarStore.variant === SidebarVariant.EXPANDED}
<SidebarExpanded {widgets} {preferences} />
@@ -69,10 +71,11 @@
<style lang="scss">
.root {
position: relative;
background-color: var(--theme-panel-color);
flex-direction: row;
min-width: 25rem;
border-radius: 0 var(--medium-BorderRadius) var(--medium-BorderRadius) 0;
&.size-mini {
&.mini {
width: 3.5rem !important;
min-width: 3.5rem !important;
max-width: 3.5rem !important;
@@ -96,84 +96,78 @@
}
</script>
<div class="root">
<div class="content">
{#if widget?.component}
<div class="component" use:resizeObserver={resize}>
{#if widget.headerLabel}
<Header
allowFullsize={false}
type="type-aside"
hideBefore={true}
hideActions={false}
hideDescription={true}
adaptive="disabled"
closeOnEscape={false}
on:close={() => {
if (widget !== undefined) {
closeWidget(widget._id)
}
}}
>
<Breadcrumbs items={[{ label: widget.headerLabel }]} currentOnly />
</Header>
{/if}
<Component
is={widget?.component}
props={{ tab, widgetState, height: componentHeight, width: componentWidth, widget }}
<div class="content">
{#if widget?.component}
<div class="component" use:resizeObserver={resize}>
{#if widget.headerLabel}
<Header
allowFullsize={false}
type="type-aside"
hideBefore={true}
hideActions={false}
hideDescription={true}
adaptive="disabled"
closeOnEscape={false}
on:close={() => {
if (widget !== undefined) {
closeWidget(widget._id)
}
}}
/>
</div>
{/if}
</div>
{#if widget !== undefined && tabs.length > 0}
<SidebarTabs
{tabs}
selected={tab?.id}
{widget}
on:close={(e) => {
void handleTabClose(e.detail, widget)
}}
on:open={(e) => {
handleTabOpen(e.detail, widget)
}}
/>
>
<Breadcrumbs items={[{ label: widget.headerLabel }]} currentOnly />
</Header>
{/if}
<Component
is={widget?.component}
props={{ tab, widgetState, height: componentHeight, width: componentWidth, widget }}
on:close={() => {
if (widget !== undefined) {
closeWidget(widget._id)
}
}}
/>
</div>
{/if}
<WidgetsBar {widgets} {preferences} selected={widgetId} />
</div>
{#if widget !== undefined && tabs.length > 0}
<SidebarTabs
{tabs}
selected={tab?.id}
{widget}
on:close={(e) => {
void handleTabClose(e.detail, widget)
}}
on:open={(e) => {
handleTabOpen(e.detail, widget)
}}
/>
{/if}
<WidgetsBar {widgets} {preferences} selected={widgetId} />
<style lang="scss">
.root {
display: flex;
flex: 1;
height: 100%;
overflow: hidden;
}
.content {
display: flex;
flex-direction: column;
border-top: 1px solid var(--theme-divider-color);
border-top: 1px solid transparent; // var(--theme-divider-color);
border-right: 1px solid var(--theme-divider-color);
overflow: auto;
flex: 1;
width: calc(100% - 3.5rem);
height: 100%;
background: var(--theme-panel-color);
border-right: 1px solid var(--global-ui-BorderColor);
min-width: 0;
min-height: 0;
background-color: var(--theme-panel-color);
border-left: none;
}
.component {
position: relative;
display: flex;
flex-direction: column;
flex: 1;
flex-grow: 1;
overflow: hidden;
border-bottom: 1px solid transparent;
background-color: var(--theme-panel-color);
}
</style>
@@ -89,9 +89,8 @@
width: 30px;
min-width: 30px;
max-width: 30px;
flex: 1;
height: 100%;
border-right: 1px solid var(--theme-divider-color);
border-top: 1px solid var(--theme-divider-color);
gap: 0.25rem;
align-items: center;
padding: 0.25rem 0;
@@ -95,13 +95,12 @@
.root {
display: flex;
flex-direction: column;
justify-content: space-between;
padding-block: var(--spacing-2);
height: 100%;
padding: 0.5rem 0;
width: 3.5rem;
min-width: 3.5rem;
max-width: 3.5rem;
border-top: 1px solid var(--theme-divider-color);
border-radius: 0 var(--medium-BorderRadius) var(--medium-BorderRadius) 0;
overflow-y: auto;
}
@@ -336,7 +336,7 @@ export function getTextPresenter (_class: Ref<Class<Doc>>, hierarchy: Hierarchy)
}
async function getSenderName (control: TriggerControl, sender: SenderInfo): Promise<string> {
if (sender._id === core.account.System) {
if (sender._id === core.account.System || sender._id === core.account.ConfigUser) {
return await translate(core.string.System, {})
}
+2 -3
View File
@@ -35,14 +35,13 @@
"jest": "^29.7.0",
"ts-jest": "^29.1.1",
"@types/jest": "^29.5.5",
"@types/otp-generator": "^4.0.2",
"@types/pg": "^8.11.6"
"@types/otp-generator": "^4.0.2"
},
"dependencies": {
"@hcengineering/mongo": "^0.6.1",
"@hcengineering/postgres": "^0.6.0",
"mongodb": "6.9.0-dev.20241016.sha.3d5bd513",
"pg": "8.12.0",
"postgres": "^3.4.4",
"@hcengineering/platform": "^0.6.11",
"@hcengineering/core": "^0.6.32",
"@hcengineering/contact": "^0.6.24",
+53 -41
View File
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
//
import { QueryResultRow, type Pool } from 'pg'
import postgres from 'postgres'
import { generateId, type Data, type Version } from '@hcengineering/core'
import type {
@@ -34,7 +34,7 @@ import type {
export class PostgresDbCollection<T extends Record<string, any>> implements DbCollection<T> {
constructor (
readonly name: string,
readonly client: Pool
readonly client: postgres.Sql
) {}
protected buildSelectClause (): string {
@@ -110,7 +110,7 @@ export class PostgresDbCollection<T extends Record<string, any>> implements DbCo
return `ORDER BY ${sortChunks.join(', ')}`
}
protected convertToObj (row: QueryResultRow): T {
protected convertToObj (row: any): T {
return row as T
}
@@ -131,9 +131,9 @@ export class PostgresDbCollection<T extends Record<string, any>> implements DbCo
}
const finalSql: string = sqlChunks.join(' ')
const result = await this.client.query(finalSql, whereValues)
const result = await this.client.unsafe(finalSql, whereValues)
return result.rows.map((row) => this.convertToObj(row))
return result.map((row) => this.convertToObj(row))
}
async findOne (query: Query<T>): Promise<T | null> {
@@ -150,7 +150,9 @@ export class PostgresDbCollection<T extends Record<string, any>> implements DbCo
const sql = `INSERT INTO ${this.name} (${keys.map((k) => `"${k}"`).join(', ')}) VALUES (${keys.map((_, idx) => `$${idx + 1}`).join(', ')})`
await this.client.query(sql, values)
await this.client.begin(async (client) => {
await client.unsafe(sql, values)
})
return id
}
@@ -194,7 +196,9 @@ export class PostgresDbCollection<T extends Record<string, any>> implements DbCo
}
const finalSql = sqlChunks.join(' ')
await this.client.query(finalSql, [...updateValues, ...whereValues])
await this.client.begin(async (client) => {
await client.unsafe(finalSql, [...updateValues, ...whereValues])
})
}
async deleteMany (query: Query<T>): Promise<void> {
@@ -206,12 +210,14 @@ export class PostgresDbCollection<T extends Record<string, any>> implements DbCo
}
const finalSql = sqlChunks.join(' ')
await this.client.query(finalSql, whereValues)
await this.client.begin(async (client) => {
await client.unsafe(finalSql, whereValues)
})
}
}
export class AccountPostgresDbCollection extends PostgresDbCollection<Account> implements DbCollection<Account> {
constructor (readonly client: Pool) {
constructor (readonly client: postgres.Sql) {
super('account', client)
}
@@ -251,7 +257,7 @@ export class AccountPostgresDbCollection extends PostgresDbCollection<Account> i
}
export class WorkspacePostgresDbCollection extends PostgresDbCollection<Workspace> implements WorkspaceDbCollection {
constructor (readonly client: Pool) {
constructor (readonly client: postgres.Sql) {
super('workspace', client)
}
@@ -349,9 +355,14 @@ export class WorkspacePostgresDbCollection extends PostgresDbCollection<Workspac
sqlChunks.push(`WHERE ${whereChunks.join(' AND ')}`)
const res = await this.client.query(sqlChunks.join(' '), values)
let count = 0
return res.rows[0].count
await this.client.begin(async (client) => {
const res = await client.unsafe(sqlChunks.join(' '), values)
count = res[0].count
})
return count
}
async getPendingWorkspace (
@@ -406,24 +417,25 @@ export class WorkspacePostgresDbCollection extends PostgresDbCollection<Workspac
// We must have all the conditions in the DB query and we cannot filter anything in the code
// because of possible concurrency between account services.
await this.client.query('BEGIN')
const res = await this.client.query(sqlChunks.join(' '), values)
if ((res.rowCount ?? 0) > 0) {
await this.client.query(
`UPDATE ${this.name} SET attempts = attempts + 1, "lastProcessingTime" = $1 WHERE _id = $2`,
[Date.now(), res.rows[0]._id]
)
}
let res: any | undefined
await this.client.begin(async (client) => {
res = await client.unsafe(sqlChunks.join(' '), values)
await this.client.query('COMMIT')
if ((res.length ?? 0) > 0) {
await client.unsafe(
`UPDATE ${this.name} SET attempts = attempts + 1, "lastProcessingTime" = $1 WHERE _id = $2`,
[Date.now(), res[0]._id]
)
}
})
return res.rows[0] as WorkspaceInfo
return res[0] as WorkspaceInfo
}
}
export class InvitePostgresDbCollection extends PostgresDbCollection<Invite> implements DbCollection<Invite> {
constructor (readonly client: Pool) {
constructor (readonly client: postgres.Sql) {
super('invite', client)
}
@@ -471,7 +483,7 @@ export class PostgresAccountDB implements AccountDB {
invite: PostgresDbCollection<Invite>
upgrade: PostgresDbCollection<UpgradeStatistic>
constructor (readonly client: Pool) {
constructor (readonly client: postgres.Sql) {
this.workspace = new WorkspacePostgresDbCollection(client)
this.account = new AccountPostgresDbCollection(client)
this.otp = new PostgresDbCollection<OtpRecord>('otp', client)
@@ -489,21 +501,21 @@ export class PostgresAccountDB implements AccountDB {
}
async migrate (name: string, ddl: string): Promise<void> {
const res = await this.client.query(
'INSERT INTO _account_applied_migrations (identifier, ddl) VALUES ($1, $2) ON CONFLICT DO NOTHING',
[name, ddl]
)
await this.client.begin(async (client) => {
const res =
await client`INSERT INTO _account_applied_migrations (identifier, ddl) VALUES (${name}, ${ddl}) ON CONFLICT DO NOTHING`
if (res.rowCount === 1) {
console.log(`Applying migration: ${name}`)
await this.client.query(ddl)
} else {
console.log(`Migration ${name} already applied`)
}
if (res.count === 1) {
console.log(`Applying migration: ${name}`)
await client.unsafe(ddl)
} else {
console.log(`Migration ${name} already applied`)
}
})
}
async _init (): Promise<void> {
await this.client.query(
await this.client.unsafe(
`
CREATE TABLE IF NOT EXISTS _account_applied_migrations (
identifier VARCHAR(255) NOT NULL PRIMARY KEY
@@ -515,15 +527,15 @@ export class PostgresAccountDB implements AccountDB {
}
async assignWorkspace (accountId: ObjectId, workspaceId: ObjectId): Promise<void> {
const sql = `INSERT INTO ${this.wsAssignmentName} (workspace, account) VALUES ($1, $2)`
await this.client.query(sql, [workspaceId, accountId])
await this.client.begin(async (client) => {
await client`INSERT INTO ${client(this.wsAssignmentName)} (workspace, account) VALUES (${workspaceId}, ${accountId})`
})
}
async unassignWorkspace (accountId: ObjectId, workspaceId: ObjectId): Promise<void> {
const sql = `DELETE FROM ${this.wsAssignmentName} WHERE workspace = $1 AND account = $2`
await this.client.query(sql, [workspaceId, accountId])
await this.client.begin(async (client) => {
await client`DELETE FROM ${client(this.wsAssignmentName)} WHERE workspace = ${workspaceId} AND account = ${accountId}`
})
}
getObjectId (id: string): ObjectId {
+1 -1
View File
@@ -1613,7 +1613,7 @@ class MongoTxAdapter extends MongoAdapterBase implements TxAdapter {
@withContext('get-model')
async getModel (ctx: MeasureContext): Promise<Tx[]> {
const txCollection = this.db.collection<Tx>(DOMAIN_TX)
const cursor = txCollection.find({ objectSpace: core.space.Model })
const cursor = txCollection.find({ objectSpace: core.space.Model }, { sort: { _id: 1, modifiedOn: 1 } })
const model = await toArray<Tx>(cursor)
// We need to put all core.account.System transactions first
const systemTx: Tx[] = []
+2 -2
View File
@@ -31,11 +31,11 @@
"jest": "^29.7.0",
"ts-jest": "^29.1.1",
"@types/jest": "^29.5.5",
"@types/node": "~20.11.16",
"@types/pg": "^8.11.6"
"@types/node": "~20.11.16"
},
"dependencies": {
"pg": "8.12.0",
"postgres": "^3.4.4",
"@hcengineering/core": "^0.6.32",
"@hcengineering/platform": "^0.6.11",
"@hcengineering/server-core": "^0.6.1"
@@ -60,7 +60,7 @@ describe('postgres operations', () => {
dbId = 'pg_testdb_' + generateId()
dbUri = baseDbUri + '/' + dbId
const client = await clientRef.getClient()
await client.query(`CREATE DATABASE ${dbId}`)
await client`CREATE DATABASE ${client(dbId)}`
} catch (err) {
console.error(err)
}
+281 -292
View File
@@ -65,15 +65,17 @@ import {
updateHashForDoc
} from '@hcengineering/server-core'
import { createHash } from 'crypto'
import { Pool, type PoolClient } from 'pg'
import type postgres from 'postgres'
import { type ValueType } from './types'
import {
convertDoc,
createTable,
DBCollectionHelper,
type DBDoc,
escapeBackticks,
getDBClient,
getDocFieldsByDomains,
inferType,
isDataField,
isOwner,
type JoinProps,
@@ -88,13 +90,12 @@ import {
abstract class PostgresAdapterBase implements DbAdapter {
protected readonly _helper: DBCollectionHelper
protected readonly tableFields = new Map<string, string[]>()
protected readonly queue: ((client: PoolClient) => Promise<any>)[] = []
protected readonly mutex = new Mutex()
protected readonly connections = new Map<string, PoolClient>()
protected readonly connections = new Map<string, postgres.ReservedSql>()
protected readonly retryTxn = async (
connection: Pool | PoolClient,
fn: (client: Pool | PoolClient) => Promise<any>
connection: postgres.ReservedSql,
fn: (client: postgres.ReservedSql) => Promise<any>
): Promise<void> => {
await this.mutex.runExclusive(async () => {
await this.processOps(connection, fn)
@@ -102,7 +103,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
}
constructor (
protected readonly client: Pool,
protected readonly client: postgres.Sql,
protected readonly refClient: PostgresClientReference,
protected readonly workspaceId: WorkspaceId,
protected readonly hierarchy: Hierarchy,
@@ -111,6 +112,23 @@ abstract class PostgresAdapterBase implements DbAdapter {
this._helper = new DBCollectionHelper(this.client, this.workspaceId)
}
protected async withConnection<T>(
ctx: MeasureContext,
operation: (client: postgres.ReservedSql) => Promise<T>
): Promise<T> {
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<void> {
if (ctx.id === undefined) return
const conn = this.connections.get(ctx.id)
@@ -120,47 +138,42 @@ abstract class PostgresAdapterBase implements DbAdapter {
}
}
protected async getConnection (ctx: MeasureContext): Promise<PoolClient | Pool> {
if (ctx.id === undefined) return this.client
protected async getConnection (ctx: MeasureContext): Promise<postgres.ReservedSql | undefined> {
if (ctx.id === undefined) return
const conn = this.connections.get(ctx.id)
if (conn !== undefined) return conn
const client = await this.client.connect()
const client = await this.client.reserve()
this.connections.set(ctx.id, client)
return client
}
private async processOps (client: Pool | PoolClient, operation: (client: PoolClient) => Promise<any>): Promise<void> {
private async processOps (
client: postgres.ReservedSql,
operation: (client: postgres.ReservedSql) => Promise<any>
): Promise<void> {
const backoffInterval = 100 // millis
const maxTries = 5
let tries = 0
const _client = client instanceof Pool ? await client.connect() : client
while (true) {
await client.unsafe('BEGIN;')
tries++
try {
while (true) {
await _client.query('BEGIN;')
tries++
try {
const result = await operation(client)
await client.unsafe('COMMIT;')
return result
} catch (err: any) {
await client.unsafe('ROLLBACK;')
try {
const result = await operation(_client)
await _client.query('COMMIT;')
return result
} catch (err: any) {
await _client.query('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))
}
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))
}
}
} finally {
if (client instanceof Pool) {
_client.release()
}
}
}
@@ -169,15 +182,15 @@ abstract class PostgresAdapterBase implements DbAdapter {
query: DocumentQuery<T>,
options?: Pick<FindOptions<T>, 'sort' | 'limit' | 'projection'>
): Promise<Iterator<T>> {
const client = await this.client.connect()
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<void> => {
if (closed) return
try {
await client.query(`CLOSE ${cursorName}`)
await client.query('COMMIT;')
await client.unsafe(`CLOSE ${cursorName}`)
await client.unsafe('COMMIT;')
} finally {
client.release()
closed = true
@@ -195,17 +208,17 @@ abstract class PostgresAdapterBase implements DbAdapter {
sqlChunks.push(`LIMIT ${options.limit}`)
}
const finalSql: string = sqlChunks.join(' ')
await client.query('BEGIN;')
await client.query(`DECLARE ${cursorName} ${finalSql}`)
await client.unsafe('BEGIN;')
await client.unsafe(`DECLARE ${cursorName} ${finalSql}`)
}
const next = async (count: number): Promise<T[] | null> => {
const result = await client.query(`FETCH ${count} FROM ${cursorName}`)
if (result.rows.length === 0) {
const result = await client.unsafe(`FETCH ${count} FROM ${cursorName}`)
if (result.length === 0) {
await close(cursorName)
return null
}
return result.rows.map(parseDoc<T>)
return result.map((p) => parseDoc(p as any))
}
await init()
@@ -244,8 +257,8 @@ abstract class PostgresAdapterBase implements DbAdapter {
sqlChunks.push(`LIMIT ${options.limit}`)
}
const finalSql: string = [select, ...sqlChunks].join(' ')
const result = await this.client.query(finalSql)
return result.rows.map((p) => parseDocWithProjection(p, options?.projection))
const result = await this.client.unsafe(finalSql)
return result.map((p) => parseDocWithProjection(p as any, options?.projection))
}
buildRawOrder<T extends Doc>(domain: string, sort: SortingQuery<T>): string {
@@ -290,48 +303,58 @@ abstract class PostgresAdapterBase implements DbAdapter {
if ((operations as any)['%hash%'] === undefined) {
;(operations as any)['%hash%'] = null
}
await this.retryTxn(this.client, async (client) => {
const res = await client.query(`SELECT * FROM ${translateDomain(domain)} WHERE ${translatedQuery} FOR UPDATE`)
const docs = res.rows.map(parseDoc)
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)
let paramsIndex = 3
const params: any[] = [doc._id, this.workspaceId.name]
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)
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))
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)
let paramsIndex = 3
const params: any[] = [doc._id, this.workspaceId.name]
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)
}
} 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)
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
)
}
await client.query(
`UPDATE ${translateDomain(domain)} SET ${updates.join(', ')} WHERE _id = $1 AND "workspaceId" = $2`,
params
)
}
})
})
} finally {
conn.release()
}
}
async rawDeleteMany<T extends Doc>(domain: Domain, query: DocumentQuery<T>): Promise<void> {
const translatedQuery = this.buildRawQuery(domain, query)
await this.retryTxn(this.client, async (client) => {
await client.query(`DELETE FROM ${translateDomain(domain)} WHERE ${translatedQuery}`)
})
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()
}
}
async findAll<T extends Doc>(
@@ -354,14 +377,14 @@ abstract class PostgresAdapterBase implements DbAdapter {
sqlChunks.push(this.buildJoinString(joins))
}
sqlChunks.push(`WHERE ${this.buildQuery(_class, domain, query, joins, options)}`)
const connection = await this.getConnection(ctx)
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.query(totalSql)
const parsed = Number.parseInt(totalResult.rows[0]?.count ?? '')
const totalResult = await connection.unsafe(totalSql)
const parsed = Number.parseInt(totalResult[0].count)
total = Number.isNaN(parsed) ? 0 : parsed
}
if (options?.sort !== undefined) {
@@ -372,14 +395,14 @@ abstract class PostgresAdapterBase implements DbAdapter {
}
const finalSql: string = [select, ...sqlChunks].join(' ')
const result = await connection.query(finalSql)
const result = await connection.unsafe(finalSql)
if (options?.lookup === undefined) {
return toFindResult(
result.rows.map((p) => parseDocWithProjection(p, options?.projection)),
result.map((p) => parseDocWithProjection(p as any, options?.projection)),
total
)
} else {
const res = this.parseLookup<T>(result.rows, joins, options?.projection)
const res = this.parseLookup<T>(result, joins, options?.projection)
return toFindResult(res, total)
}
} catch (err) {
@@ -1015,15 +1038,15 @@ abstract class PostgresAdapterBase implements DbAdapter {
}
let initialized: boolean = false
let client: PoolClient
let client: postgres.ReservedSql
let mode: 'hashed' | 'non_hashed' = 'hashed'
let cursorName = getCursorName()
const bulkUpdate = new Map<Ref<Doc>, string>()
const close = async (cursorName: string): Promise<void> => {
try {
await client.query(`CLOSE ${cursorName}`)
await client.query('COMMIT')
await client.unsafe(`CLOSE ${cursorName}`)
await client.unsafe('COMMIT')
} catch (err) {
ctx.error('Error while closing cursor', { cursorName, err })
} finally {
@@ -1033,20 +1056,20 @@ abstract class PostgresAdapterBase implements DbAdapter {
const init = async (projection: string, query: string): Promise<void> => {
cursorName = getCursorName()
client = await this.client.connect()
await client.query('BEGIN')
await client.query(
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<Doc[]> => {
const result = await client.query(`FETCH ${limit} FROM ${cursorName}`)
if (result.rows.length === 0) {
const result = await client.unsafe(`FETCH ${limit} FROM ${cursorName}`)
if (result.length === 0) {
return []
}
return result.rows.filter((it) => it != null).map((it) => parseDoc(it))
return result.filter((it) => it != null).map((it) => parseDoc(it as any))
}
const flush = async (flush = false): Promise<void> => {
@@ -1066,10 +1089,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
if (!initialized) {
if (recheck === true) {
await this.retryTxn(client, async (client) => {
await client.query(
`UPDATE ${translateDomain(domain)} SET jsonb_set(data, '{%hash%}', 'NULL', true) WHERE "workspaceId" = $1 AND data ->> '%hash%' IS NOT NULL`,
[this.workspaceId.name]
)
await client`UPDATE ${client(translateDomain(domain))} SET jsonb_set(data, '{%hash%}', 'NULL', true) WHERE "workspaceId" = ${this.workspaceId.name} AND data ->> '%hash%' IS NOT NULL`
})
}
await init('_id, data', "data ->> '%hash%' IS NOT NULL AND data ->> '%hash%' <> ''")
@@ -1136,12 +1156,10 @@ abstract class PostgresAdapterBase implements DbAdapter {
if (docs.length === 0) {
return []
}
const connection = await this.getConnection(ctx)
const res = await connection.query(
`SELECT * FROM ${translateDomain(domain)} WHERE _id = ANY($1) AND "workspaceId" = $2`,
[docs, this.workspaceId.name]
)
return res.rows as Doc[]
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 as any as Doc[]
})
}
@@ -1158,57 +1176,55 @@ abstract class PostgresAdapterBase implements DbAdapter {
}
const insertStr = insertFields.join(', ')
const onConflictStr = onConflict.join(', ')
const connection = await this.getConnection(ctx)
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])
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(', ')})`)
}
values.push(d.data)
variables.push(`$${index++}`)
vars.push(`(${variables.join(', ')})`)
}
const vals = vars.join(',')
await this.retryTxn(connection, async (client) => {
await client.query(
`INSERT INTO ${translateDomain(domain)} ("workspaceId", ${insertStr}) VALUES ${vals}
ON CONFLICT ("workspaceId", _id) DO UPDATE SET ${onConflictStr};`,
values
)
})
}
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<Doc>[]): Promise<void> {
const connection = await this.getConnection(ctx)
await connection.query(`DELETE FROM ${translateDomain(domain)} WHERE _id = ANY($1) AND "workspaceId" = $2`, [
docs,
this.workspaceId.name
])
const connection = (await this.getConnection(ctx)) ?? this.client
await connection`DELETE FROM ${connection(translateDomain(domain))} WHERE _id = ANY(${docs}) AND "workspaceId" = ${this.workspaceId.name}`
}
async groupBy<T>(ctx: MeasureContext, domain: Domain, field: string): Promise<Set<T>> {
const connection = await this.getConnection(ctx)
const connection = (await this.getConnection(ctx)) ?? this.client
const key = isDataField(domain, field) ? `data ->> '${field}'` : `"${field}"`
const result = await ctx.with('groupBy', { domain }, async (ctx) => {
try {
const result = await connection.query(
const result = await connection.unsafe(
`SELECT DISTINCT ${key} as ${field} FROM ${translateDomain(domain)} WHERE "workspaceId" = $1`,
[this.workspaceId.name]
)
return new Set(result.rows.map((r) => r[field]))
return new Set(result.map((r) => r[field]))
} catch (err) {
ctx.error('Error while grouping by', { domain, field })
throw err
@@ -1219,85 +1235,67 @@ abstract class PostgresAdapterBase implements DbAdapter {
async update (ctx: MeasureContext, domain: Domain, operations: Map<Ref<Doc>, DocumentUpdate<Doc>>): Promise<void> {
const ids = Array.from(operations.keys())
const connection = await this.getConnection(ctx)
await this.retryTxn(connection, async (client) => {
try {
const res = await client.query(
`SELECT * FROM ${translateDomain(domain)} WHERE _id = ANY($1) AND "workspaceId" = $2 FOR UPDATE`,
[ids, this.workspaceId.name]
)
const docs = res.rows.map(parseDoc)
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)
await this.withConnection(ctx, async (client) => {
await 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))
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 updates: string[] = []
let paramsIndex = 3
const { extractedFields, remainingData } = parseUpdate(domain, op)
const params: any[] = [doc._id, this.workspaceId.name]
for (const key in extractedFields) {
updates.push(`"${key}" = $${paramsIndex++}`)
params.push((extractedFields as any)[key])
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}`
}
if (Object.keys(remainingData).length > 0) {
updates.push(`data = $${paramsIndex++}`)
params.push(converted.data)
}
await client.query(
`UPDATE ${translateDomain(domain)} SET ${updates.join(', ')} WHERE _id = $1 AND "workspaceId" = $2`,
params
)
} catch (err) {
ctx.error('Error while updating', { domain, operations, err })
throw err
}
} catch (err) {
ctx.error('Error while updating', { domain, operations, err })
throw err
}
})
})
}
async insert (ctx: MeasureContext, domain: string, docs: Doc[]): Promise<TxResult> {
const fields = getDocFieldsByDomains(domain)
const filedsWithData = [...fields, 'data']
const insertFields: string[] = []
const columns: string[] = ['workspaceId']
for (const field of filedsWithData) {
insertFields.push(`"${field}"`)
columns.push(field)
}
const insertStr = insertFields.join(', ')
while (docs.length > 0) {
const part = docs.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++}`)
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)
}
values.push(d.data)
variables.push(`$${index++}`)
vars.push(`(${variables.join(', ')})`)
await this.retryTxn(connection, async (client) => {
await client`INSERT INTO ${client(translateDomain(domain))} ${client(values, columns)}`
})
}
const vals = vars.join(',')
const connection = await this.getConnection(ctx)
await this.retryTxn(connection, async (client) => {
await client.query(
`INSERT INTO ${translateDomain(domain)} ("workspaceId", ${insertStr}) VALUES ${vals}`,
values
)
})
}
})
return {}
}
}
@@ -1357,27 +1355,24 @@ class PostgresAdapter extends PostgresAdapterBase {
private async txMixin (ctx: MeasureContext, tx: TxMixin<Doc, Doc>): Promise<TxResult> {
await ctx.with('tx-mixin', { _class: tx.objectClass, mixin: tx.mixin }, async (ctx) => {
const connection = await this.getConnection(ctx)
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)
const domain = this.hierarchy.getDomain(tx.objectClass)
const converted = convertDoc(domain, doc, this.workspaceId.name)
const updates: string[] = ['"modifiedBy" = $1', '"modifiedOn" = $2']
let paramsIndex = 5
const { extractedFields } = parseUpdate(domain, tx.attributes as Partial<Doc>)
const params: any[] = [tx.modifiedBy, tx.modifiedOn, tx.objectId, this.workspaceId.name]
for (const key in extractedFields) {
updates.push(`"${key}" = $${paramsIndex++}`)
params.push(converted[key])
}
updates.push(`data = $${paramsIndex++}`)
params.push(converted.data)
await client.query(
`UPDATE ${translateDomain(domain)} SET ${updates.join(', ')} WHERE _id = $3 AND "workspaceId" = $4`,
params
)
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<Doc>)
const columns = new Set<string>()
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 {}
@@ -1411,7 +1406,7 @@ class PostgresAdapter extends PostgresAdapterBase {
protected async txCreateDoc (ctx: MeasureContext, tx: TxCreateDoc<Doc>): Promise<TxResult> {
const doc = TxProcessor.createDoc2Doc(tx)
return await ctx.with('create-doc', { _class: doc._class }, async (_ctx) => {
return await this.insert(_ctx, translateDomain(this.hierarchy.getDomain(doc._class)), [doc])
return await this.insert(_ctx, this.hierarchy.getDomain(doc._class), [doc])
})
}
@@ -1419,39 +1414,35 @@ class PostgresAdapter extends PostgresAdapterBase {
return await ctx.with('tx-update-doc', { _class: tx.objectClass }, async (_ctx) => {
if (isOperator(tx.operations)) {
let doc: Doc | undefined
const ops = { '%hash%': null, ...tx.operations }
const ops: any = { '%hash%': null, ...tx.operations }
return await _ctx.with(
'update with operations',
{ operations: JSON.stringify(Object.keys(tx.operations)) },
async (ctx) => {
const connection = await this.getConnection(ctx)
await this.retryTxn(connection, async (client) => {
doc = await this.findDoc(ctx, client, tx.objectClass, tx.objectId, true)
if (doc === undefined) return {}
TxProcessor.applyUpdate(doc, ops)
const domain = this.hierarchy.getDomain(tx.objectClass)
const converted = convertDoc(domain, doc, this.workspaceId.name)
const updates: string[] = ['"modifiedBy" = $1', '"modifiedOn" = $2']
let paramsIndex = 5
const { extractedFields, remainingData } = parseUpdate(domain, ops)
const params: any[] = [tx.modifiedBy, tx.modifiedOn, tx.objectId, this.workspaceId.name]
for (const key in extractedFields) {
updates.push(`"${key}" = $${paramsIndex++}`)
params.push(converted[key])
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 }
}
if (Object.keys(remainingData).length > 0) {
updates.push(`data = $${paramsIndex++}`)
params.push(converted.data)
}
await client.query(
`UPDATE ${translateDomain(domain)} SET ${updates.join(', ')} WHERE _id = $3 AND "workspaceId" = $4`,
params
)
return {}
})
if (tx.retrieve === true && doc !== undefined) {
return { object: doc }
}
return {}
}
)
} else {
@@ -1483,60 +1474,58 @@ class PostgresAdapter extends PostgresAdapterBase {
let dataUpdated = false
for (const key in remainingData) {
if (ops[key] === undefined) continue
from = `jsonb_set(${from}, '{${key}}', $${paramsIndex++}::jsonb, true)`
params.push(JSON.stringify((remainingData as any)[key]))
const val = (remainingData as any)[key]
from = `jsonb_set(${from}, '{${key}}', to_jsonb($${paramsIndex++}${inferType(val)}) , true)`
params.push(val)
dataUpdated = true
}
if (dataUpdated) {
updates.push(`data = ${from}`)
}
try {
const connection = await this.getConnection(_ctx)
await this.retryTxn(connection, async (client) => {
await client.query(
`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 }
await this.withConnection(ctx, async (connection) => {
try {
await this.retryTxn(connection, async (client) => {
await 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)
}
} catch (err) {
console.error(err)
}
})
return {}
})
}
private async findDoc (
ctx: MeasureContext,
client: Pool | PoolClient,
client: postgres.Sql | postgres.ReservedSql,
_class: Ref<Class<Doc>>,
_id: Ref<Doc>,
forUpdate: boolean = false
): Promise<Doc | undefined> {
return await ctx.with('find-doc', { _class }, async () => {
let query = `SELECT * FROM ${translateDomain(this.hierarchy.getDomain(_class))} WHERE _id = $1 AND "workspaceId" = $2`
if (forUpdate) {
query += ' FOR UPDATE'
}
const res = await client.query(query, [_id, this.workspaceId.name])
const dbDoc = res.rows[0]
return dbDoc !== undefined ? parseDoc(dbDoc) : undefined
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]
return dbDoc !== undefined ? parseDoc(dbDoc as any) : undefined
})
}
protected async txRemoveDoc (ctx: MeasureContext, tx: TxRemoveDoc<Doc>): Promise<TxResult> {
await ctx.with('tx-remove-doc', { _class: tx.objectClass }, async (_ctx) => {
const domain = translateDomain(this.hierarchy.getDomain(tx.objectClass))
const connection = await this.getConnection(_ctx)
await this.retryTxn(connection, async (client) => {
await client.query(`DELETE FROM ${domain} WHERE _id = $1 AND "workspaceId" = $2`, [
tx.objectId,
this.workspaceId.name
])
await this.withConnection(_ctx, async (connection) => {
await this.retryTxn(connection, async (client) => {
await client`DELETE FROM ${client(domain)} WHERE _id = ${tx.objectId} AND "workspaceId" = ${this.workspaceId.name}`
})
})
})
return {}
@@ -1563,10 +1552,10 @@ class PostgresTxAdapter extends PostgresAdapterBase implements TxAdapter {
}
async getModel (ctx: MeasureContext): Promise<Tx[]> {
const res = await this.client.query(
`SELECT * FROM ${translateDomain(DOMAIN_TX)} WHERE "workspaceId" = '${this.workspaceId.name}' AND data->>'objectSpace' = '${core.space.Model}' ORDER BY _id ASC, "modifiedOn" ASC`
)
const model = res.rows.map((p) => parseDoc<Tx>(p))
const res = await this
.client`SELECT * FROM ${this.client(translateDomain(DOMAIN_TX))} WHERE "workspaceId" = ${this.workspaceId.name} AND data->>'objectSpace' = ${core.space.Model} ORDER BY _id ASC, "modifiedOn" ASC`
const model = res.map((p) => parseDoc<Tx>(p as any))
// We need to put all core.account.System transactions first
const systemTx: Tx[] = []
const userTx: Tx[] = []
@@ -1586,8 +1575,8 @@ export async function createPostgresAdapter (
modelDb: ModelDb
): Promise<DbAdapter> {
const client = getDBClient(url)
const pool = await client.getClient()
const adapter = new PostgresAdapter(pool, client, workspaceId, hierarchy, modelDb)
const connection = await client.getClient()
const adapter = new PostgresAdapter(connection, client, workspaceId, hierarchy, modelDb)
return adapter
}
@@ -1602,8 +1591,8 @@ export async function createPostgresTxAdapter (
modelDb: ModelDb
): Promise<TxAdapter> {
const client = getDBClient(url)
const pool = await client.getClient()
const adapter = new PostgresTxAdapter(pool, client, workspaceId, hierarchy, modelDb)
const connection = await client.getClient()
const adapter = new PostgresTxAdapter(connection, client, workspaceId, hierarchy, modelDb)
await adapter.init()
return adapter
}
+61 -75
View File
@@ -29,7 +29,7 @@ import core, {
} from '@hcengineering/core'
import { PlatformError, unknownStatus } from '@hcengineering/platform'
import { type DomainHelperOperations } from '@hcengineering/server-core'
import { Pool, type PoolClient } from 'pg'
import postgres from 'postgres'
import { defaultSchema, domainSchemas, getSchema } from './schemas'
const connections = new Map<string, PostgresClientReferenceImpl>()
@@ -43,53 +43,23 @@ process.on('exit', () => {
const clientRefs = new Map<string, ClientRef>()
export async function retryTxn (pool: Pool, operation: (client: PoolClient) => Promise<any>): Promise<any> {
const backoffInterval = 100 // millis
const maxTries = 5
let tries = 0
const client = await pool.connect()
try {
while (true) {
await client.query('BEGIN;')
tries++
try {
const result = await operation(client)
await client.query('COMMIT;')
return result
} catch (err: any) {
await client.query('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))
}
}
}
} finally {
client.release()
}
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 async function createTable (client: Pool, domains: string[]): Promise<void> {
export async function createTable (client: postgres.Sql, domains: string[]): Promise<void> {
if (domains.length === 0) {
return
}
const mapped = domains.map((p) => translateDomain(p))
const inArr = mapped.map((it) => `'${it}'`).join(', ')
const exists = await client.query(`
SELECT table_name
FROM information_schema.tables
WHERE table_name IN (${inArr})
`)
const toCreate = mapped.filter((it) => !exists.rows.map((it) => it.table_name).includes(it))
await retryTxn(client, async (client) => {
for (const domain of toCreate) {
for (const domain of mapped) {
const schema = getSchema(domain)
const fields: string[] = []
for (const key in schema) {
@@ -97,28 +67,28 @@ export async function createTable (client: Pool, domains: string[]): Promise<voi
fields.push(`"${key}" ${val[0]} ${val[1] ? 'NOT NULL' : ''}`)
}
const colums = fields.join(', ')
await client.query(
`CREATE TABLE ${domain} (
const res = await client.unsafe(`CREATE TABLE IF NOT EXISTS ${domain} (
"workspaceId" text NOT NULL,
${colums},
data JSONB NOT NULL,
PRIMARY KEY("workspaceId", _id)
)`
)
if (schema.attachedTo !== undefined) {
await client.query(`
CREATE INDEX ${domain}_attachedTo ON ${domain} ("attachedTo")
)`)
if (res.count > 0) {
if (schema.attachedTo !== undefined) {
await client.unsafe(`
CREATE INDEX ${domain}_attachedTo ON ${domain} ("attachedTo")
`)
}
await client.unsafe(`
CREATE INDEX ${domain}_class ON ${domain} (_class)
`)
await client.unsafe(`
CREATE INDEX ${domain}_space ON ${domain} (space)
`)
await client.unsafe(`
CREATE INDEX ${domain}_idxgin ON ${domain} USING GIN (data)
`)
}
await client.query(`
CREATE INDEX ${domain}_class ON ${domain} (_class)
`)
await client.query(`
CREATE INDEX ${domain}_space ON ${domain} (space)
`)
await client.query(`
CREATE INDEX ${domain}_idxgin ON ${domain} USING GIN (data)
`)
}
})
}
@@ -134,23 +104,23 @@ export async function shutdown (): Promise<void> {
}
export interface PostgresClientReference {
getClient: () => Promise<Pool>
getClient: () => Promise<postgres.Sql>
close: () => void
}
class PostgresClientReferenceImpl {
count: number
client: Pool | Promise<Pool>
client: postgres.Sql | Promise<postgres.Sql>
constructor (
client: Pool | Promise<Pool>,
client: postgres.Sql | Promise<postgres.Sql>,
readonly onclose: () => void
) {
this.count = 0
this.client = client
}
async getClient (): Promise<Pool> {
async getClient (): Promise<postgres.Sql> {
if (this.client instanceof Promise) {
this.client = await this.client
}
@@ -183,7 +153,7 @@ export class ClientRef implements PostgresClientReference {
}
closed = false
async getClient (): Promise<Pool> {
async getClient (): Promise<postgres.Sql> {
if (!this.closed) {
return await this.client.getClient()
} else {
@@ -211,15 +181,18 @@ export function getDBClient (connectionString: string, database?: string): Postg
let existing = connections.get(key)
if (existing === undefined) {
const pool = new Pool({
const sql = postgres(connectionString, {
connectionString,
application_name: 'transactor',
database,
max: 10,
transform: {
undefined: null
},
...extraOptions
})
existing = new PostgresClientReferenceImpl(pool, () => {
existing = new PostgresClientReferenceImpl(sql, () => {
connections.delete(key)
})
connections.set(key, existing)
@@ -258,6 +231,19 @@ export function convertDoc<T extends Doc> (domain: string, doc: T, workspaceId:
return res
}
export function inferType (val: any): string {
if (typeof val === 'string') {
return '::text'
}
if (typeof val === 'number') {
return '::numeric'
}
if (typeof val === 'boolean') {
return '::boolean'
}
return ''
}
export function parseUpdate<T extends Doc> (
domain: string,
ops: DocumentUpdate<T> | MixinUpdate<Doc, T>
@@ -303,20 +289,22 @@ export function isOwner (account: Account): boolean {
export class DBCollectionHelper implements DomainHelperOperations {
constructor (
protected readonly client: Pool,
protected readonly client: postgres.Sql,
protected readonly workspaceId: WorkspaceId
) {}
async dropIndex (domain: Domain, name: string): Promise<void> {}
domains = new Set<Domain>()
async create (domain: Domain): Promise<void> {}
async exists (domain: Domain): Promise<boolean> {
const exists = await this.client.query(`
const exists = await this.client`
SELECT table_name
FROM information_schema.tables
WHERE table_name = '${translateDomain(domain)}'
`)
return exists.rows.length > 0
WHERE table_name = '${this.client(translateDomain(domain))}'
`
return exists.length > 0
}
async listDomains (): Promise<Set<Domain>> {
@@ -325,17 +313,15 @@ export class DBCollectionHelper implements DomainHelperOperations {
async createIndex (domain: Domain, value: string | FieldIndexConfig<Doc>, options?: { name: string }): Promise<void> {}
async dropIndex (domain: Domain, name: string): Promise<void> {}
async listIndexes (domain: Domain): Promise<{ name: string }[]> {
return []
}
async estimatedCount (domain: Domain): Promise<number> {
const res = await this.client.query(`SELECT COUNT(_id) FROM ${translateDomain(domain)} WHERE "workspaceId" = $1`, [
this.workspaceId.name
])
return res.rows[0].count
const res = await this
.client`SELECT COUNT(_id) FROM ${this.client(translateDomain(domain))} WHERE "workspaceId" = ${this.workspaceId.name}`
return res.count
}
}
+3 -1
View File
@@ -340,7 +340,9 @@ export async function upgradeModel (
let i = 0
for (const op of migrateOperations) {
const t = Date.now()
await ctx.with(op[0], {}, () => op[1].upgrade(migrateState, async () => connection, logger))
await ctx.with(op[0], {}, async () => {
await op[1].upgrade(migrateState, async () => connection, logger)
})
const tdelta = Date.now() - t
if (tdelta > 0) {
logger.log('upgrade:', { operation: op[0], time: tdelta, workspaceId: workspaceId.name })
@@ -34,7 +34,7 @@ function wrapPipeline (ctx: MeasureContext, pipeline: Pipeline, wsUrl: Workspace
{ targets: {}, txes: [] },
wsUrl,
null,
false,
true,
new Map(),
new Map(),
pipeline.context.modelDb
@@ -304,7 +304,7 @@ export async function upgradeWorkspaceWith (
{ targets: {}, txes: [] },
wsUrl,
null,
false,
true,
new Map(),
new Map(),
pipeline.context.modelDb
+13 -17
View File
@@ -18,7 +18,7 @@ import postgres from 'postgres'
import * as db from './db'
import { toUUID } from './encodings'
import { selectStorage } from './storage'
import { type UUID } from './types'
import { type BlobRequest, type WorkspaceRequest, type UUID } from './types'
import { copyVideo, deleteVideo } from './video'
const expires = 86400
@@ -39,13 +39,9 @@ export function getBlobURL (request: Request, workspace: string, name: string):
return new URL(path, request.url).toString()
}
export async function handleBlobGet (
request: Request,
env: Env,
ctx: ExecutionContext,
workspace: string,
name: string
): Promise<Response> {
export async function handleBlobGet (request: BlobRequest, env: Env, ctx: ExecutionContext): Promise<Response> {
const { workspace, name } = request
const sql = postgres(env.HYPERDRIVE.connectionString)
const { bucket } = selectStorage(env, workspace)
@@ -82,13 +78,9 @@ export async function handleBlobGet (
return response
}
export async function handleBlobHead (
request: Request,
env: Env,
ctx: ExecutionContext,
workspace: string,
name: string
): Promise<Response> {
export async function handleBlobHead (request: BlobRequest, env: Env, ctx: ExecutionContext): Promise<Response> {
const { workspace, name } = request
const sql = postgres(env.HYPERDRIVE.connectionString)
const { bucket } = selectStorage(env, workspace)
@@ -106,7 +98,9 @@ export async function handleBlobHead (
return new Response(null, { headers, status: 200 })
}
export async function deleteBlob (env: Env, workspace: string, name: string): Promise<Response> {
export async function handleBlobDelete (request: BlobRequest, env: Env, ctx: ExecutionContext): Promise<Response> {
const { workspace, name } = request
const sql = postgres(env.HYPERDRIVE.connectionString)
try {
@@ -120,13 +114,15 @@ export async function deleteBlob (env: Env, workspace: string, name: string): Pr
}
}
export async function postBlobFormData (request: Request, env: Env, workspace: string): Promise<Response> {
export async function handleUploadFormData (request: WorkspaceRequest, env: Env): Promise<Response> {
const contentType = request.headers.get('Content-Type')
if (contentType === null || !contentType.includes('multipart/form-data')) {
console.error({ error: 'expected multipart/form-data' })
return error(400, 'expected multipart/form-data')
}
const { workspace } = request
const sql = postgres(env.HYPERDRIVE.connectionString)
let formData: FormData
+8 -6
View File
@@ -14,15 +14,17 @@
//
import { getBlobURL } from './blob'
import { type BlobRequest } from './types'
const prefferedImageFormats = ['webp', 'avif', 'jpeg', 'png']
export async function getImage (
request: Request,
workspace: string,
name: string,
transform: string
): Promise<Response> {
export async function handleImageGet (request: BlobRequest): Promise<Response> {
const {
workspace,
name,
params: { transform }
} = request
const Accept = request.headers.get('Accept') ?? 'image/*'
const image: Record<string, string> = {}
+78 -49
View File
@@ -13,61 +13,90 @@
// limitations under the License.
//
import { type IRequest, Router, error, html } from 'itty-router'
import {
deleteBlob as handleBlobDelete,
handleBlobGet,
handleBlobHead,
postBlobFormData as handleUploadFormData
} from './blob'
import { WorkerEntrypoint } from 'cloudflare:workers'
import { type IRequestStrict, type RequestHandler, Router, error, html } from 'itty-router'
import { handleBlobDelete, handleBlobGet, handleBlobHead, handleUploadFormData } from './blob'
import { cors } from './cors'
import { getImage as handleImageGet } from './image'
import { getVideoMeta as handleVideoMetaGet } from './video'
import { handleImageGet } from './image'
import { handleVideoMetaGet } from './video'
import { handleSignAbort, handleSignComplete, handleSignCreate } from './sign'
import { type BlobRequest, type WorkspaceRequest } from './types'
const { preflight, corsify } = cors({
maxAge: 86400
})
export default {
async fetch (request, env, ctx): Promise<Response> {
const router = Router<IRequest>({
before: [preflight],
finally: [corsify]
})
const router = Router<IRequestStrict, [Env, ExecutionContext], Response>({
before: [preflight],
finally: [corsify]
})
router
.get('/blob/:workspace/:name', ({ params }) => handleBlobGet(request, env, ctx, params.workspace, params.name))
.head('/blob/:workspace/:name', ({ params }) => handleBlobHead(request, env, ctx, params.workspace, params.name))
.delete('/blob/:workspace/:name', ({ params }) => handleBlobDelete(env, params.workspace, params.name))
// Image
.get('/image/:transform/:workspace/:name', ({ params }) =>
handleImageGet(request, params.workspace, params.name, params.transform)
)
// Video
.get('/video/:workspace/:name/meta', ({ params }) =>
handleVideoMetaGet(request, env, ctx, params.workspace, params.name)
)
// Form Data
.post('/upload/form-data/:workspace', ({ params }) => handleUploadFormData(request, env, params.workspace))
// Signed URL
.post('/upload/signed-url/:workspace/:name', ({ params }) =>
handleSignCreate(request, env, ctx, params.workspace, params.name)
)
.put('/upload/signed-url/:workspace/:name', ({ params }) =>
handleSignComplete(request, env, ctx, params.workspace, params.name)
)
.delete('/upload/signed-url/:workspace/:name', ({ params }) =>
handleSignAbort(request, env, ctx, params.workspace, params.name)
)
.all('/', () =>
html(
`Huly&reg; Datalake&trade; <a href="https://huly.io">https://huly.io</a>
&copy; 2024 <a href="https://hulylabs.com">Huly Labs</a>`
)
)
.all('*', () => error(404))
return await router.fetch(request).catch(error)
const withWorkspace: RequestHandler<WorkspaceRequest> = (request: WorkspaceRequest) => {
if (request.params.workspace === undefined || request.params.workspace === '') {
return error(400, 'Missing workspace')
}
} satisfies ExportedHandler<Env>
request.workspace = decodeURIComponent(request.params.workspace)
}
const withBlob: RequestHandler<BlobRequest> = (request: BlobRequest) => {
if (request.params.name === undefined || request.params.name === '') {
return error(400, 'Missing blob name')
}
request.workspace = decodeURIComponent(request.params.name)
}
router
.get('/blob/:workspace/:name', withBlob, handleBlobGet)
.head('/blob/:workspace/:name', withBlob, handleBlobHead)
.delete('/blob/:workspace/:name', withBlob, handleBlobDelete)
// Image
.get('/image/:transform/:workspace/:name', withBlob, handleImageGet)
// Video
.get('/video/:workspace/:name/meta', withBlob, handleVideoMetaGet)
// Form Data
.post('/upload/form-data/:workspace', withWorkspace, handleUploadFormData)
// Signed URL
.post('/upload/signed-url/:workspace/:name', withBlob, handleSignCreate)
.put('/upload/signed-url/:workspace/:name', withBlob, handleSignComplete)
.delete('/upload/signed-url/:workspace/:name', withBlob, handleSignAbort)
.all('/', () =>
html(
`Huly&reg; Datalake&trade; <a href="https://huly.io">https://huly.io</a>
&copy; 2024 <a href="https://hulylabs.com">Huly Labs</a>`
)
)
.all('*', () => error(404))
export default class DatalakeWorker extends WorkerEntrypoint<Env> {
async fetch (request: Request): Promise<Response> {
return await router.fetch(request, this.env, this.ctx).catch(error)
}
async getBlob (workspace: string, name: string): Promise<ArrayBuffer> {
const request = new Request(`https://datalake/blob/${workspace}/${name}`)
const response = await router.fetch(request)
if (!response.ok) {
console.error({ error: 'datalake error: ' + response.statusText, workspace, name })
throw new Error(`Failed to fetch blob: ${response.statusText}`)
}
return await response.arrayBuffer()
}
async putBlob (workspace: string, name: string, data: ArrayBuffer | Blob | string, type: string): Promise<void> {
const request = new Request(`https://datalake/upload/form-data/${workspace}`)
const body = new FormData()
const blob = new Blob([data], { type })
body.set('file', blob, name)
const response = await router.fetch(request, { method: 'POST', body })
if (!response.ok) {
console.error({ error: 'datalake error: ' + response.statusText, workspace, name })
throw new Error(`Failed to fetch blob: ${response.statusText}`)
}
}
}
+9 -22
View File
@@ -17,7 +17,7 @@ import { AwsClient } from 'aws4fetch'
import { error } from 'itty-router'
import { handleBlobUploaded } from './blob'
import { type UUID } from './types'
import { type BlobRequest, type UUID } from './types'
import { selectStorage, type Storage } from './storage'
const S3_SIGNED_LINK_TTL = 3600
@@ -39,13 +39,8 @@ function getS3Client (storage: Storage): AwsClient {
})
}
export async function handleSignCreate (
request: Request,
env: Env,
ctx: ExecutionContext,
workspace: string,
name: string
): Promise<Response> {
export async function handleSignCreate (request: BlobRequest, env: Env, ctx: ExecutionContext): Promise<Response> {
const { workspace, name } = request
const storage = selectStorage(env, workspace)
const accountId = env.R2_ACCOUNT_ID
@@ -78,13 +73,9 @@ export async function handleSignCreate (
return new Response(signed.url, { status: 200, headers })
}
export async function handleSignComplete (
request: Request,
env: Env,
ctx: ExecutionContext,
workspace: string,
name: string
): Promise<Response> {
export async function handleSignComplete (request: BlobRequest, env: Env, ctx: ExecutionContext): Promise<Response> {
const { workspace, name } = request
const { bucket } = selectStorage(env, workspace)
const key = signBlobKey(workspace, name)
@@ -117,13 +108,9 @@ export async function handleSignComplete (
return new Response(null, { status: 201 })
}
export async function handleSignAbort (
request: Request,
env: Env,
ctx: ExecutionContext,
workspace: string,
name: string
): Promise<Response> {
export async function handleSignAbort (request: BlobRequest, env: Env, ctx: ExecutionContext): Promise<Response> {
const { workspace, name } = request
const key = signBlobKey(workspace, name)
// Check if the blob has been uploaded
+11
View File
@@ -13,10 +13,21 @@
// limitations under the License.
//
import { type IRequestStrict } from 'itty-router'
export type Location = 'weur' | 'eeur' | 'wnam' | 'enam' | 'apac'
export type UUID = string & { __uuid: true }
export type WorkspaceRequest = {
workspace: string
} & IRequestStrict
export type BlobRequest = {
workspace: string
name: string
} & IRequestStrict
export interface CloudflareResponse {
success: boolean
errors: any
+4 -8
View File
@@ -15,7 +15,7 @@
import { error, json } from 'itty-router'
import { type CloudflareResponse, type StreamUploadResponse } from './types'
import { type BlobRequest, type CloudflareResponse, type StreamUploadResponse } from './types'
export type StreamUploadState = 'ready' | 'error' | 'inprogress' | 'queued' | 'downloading' | 'pendingupload'
@@ -42,13 +42,9 @@ function streamBlobKey (workspace: string, name: string): string {
return `v/${workspace}/${name}`
}
export async function getVideoMeta (
request: Request,
env: Env,
ctx: ExecutionContext,
workspace: string,
name: string
): Promise<Response> {
export async function handleVideoMetaGet (request: BlobRequest, env: Env, ctx: ExecutionContext): Promise<Response> {
const { workspace, name } = request
const key = streamBlobKey(workspace, name)
const streamInfo = await env.datalake_blobs.get<StreamBlobInfo>(key, { type: 'json' })