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

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
This commit is contained in:
Andrey Sobolev
2024-12-06 18:22:02 +07:00
44 changed files with 1082 additions and 574 deletions
+11 -9
View File
@@ -300,22 +300,24 @@
"cwd": "${workspaceRoot}/dev/tool"
},
{
"name": "Debug tool upgrade PG(tests)",
"name": "Debug tool upgrade PG(Cockroach)",
"type": "node",
"request": "launch",
"args": ["src/__start.ts", "upgrade", "--force"],
"args": ["src/__start.ts", "upgrade-workspace", "w-haiodo-alex-staff-c-673ee7ab-87df5406ea-2b8b4d" ],
"env": {
"SERVER_SECRET": "secret",
"MINIO_ACCESS_KEY": "minioadmin",
"MINIO_SECRET_KEY": "minioadmin",
"MINIO_ENDPOINT": "localhost:9002",
"TRANSACTOR_URL": "ws://localhost:3334",
"ACCOUNT_DB_URL": "postgresql://postgres:example@localhost:5433",
"DB_URL": "postgresql://postgres:example@localhost:5433",
"MONGO_URL": "mongodb://localhost:27018",
"ACCOUNTS_URL": "http://localhost:3003",
"MINIO_ENDPOINT": "localhost:9000",
"TRANSACTOR_URL": "ws://localhost:3332",
"ACCOUNTS_URL": "http://localhost:3000",
"ACCOUNT_DB_URL": "mongodb://localhost:27017",
// "ACCOUNT_DB_URL": "postgresql://postgres:example@localhost:5433",
// "DB_URL": "postgresql://postgres:example@localhost:5433",
"DB_URL": "postgresql://root@host.docker.internal:26257/defaultdb?sslmode=disable",
"MONGO_URL": "mongodb://localhost:27017",
"TELEGRAM_DATABASE": "telegram-service",
"ELASTIC_URL": "http://localhost:9201",
"ELASTIC_URL": "http://localhost:9200",
"REKONI_URL": "http://localhost:4004",
"MODEL_VERSION": "0.6.287"
},
+4 -3
View File
@@ -22208,7 +22208,7 @@ packages:
dev: false
file:projects/collaboration.tgz(esbuild@0.20.1)(ts-node@10.9.2):
resolution: {integrity: sha512-krhgq1XiDnWKIP/HUM8VQgEzXdxLNfDf68lZgDl/Yl2tEFUu8yLYpzd1qWVMkl8N0dXyGts+DEFC7Ntns48lgA==, tarball: file:projects/collaboration.tgz}
resolution: {integrity: sha512-aJ4uMSpM7IB3wgrjVKYm4jR3IeBYSaFvYZZFjxkriMD1fAxvjr/WpKUWxQy0q2x3gZb4SoGLoiX2d2qj2/hdhA==, tarball: file:projects/collaboration.tgz}
id: file:projects/collaboration.tgz
name: '@rush-temp/collaboration'
version: 0.0.0
@@ -24713,7 +24713,7 @@ packages:
dev: false
file:projects/model-document.tgz:
resolution: {integrity: sha512-tSr57oIXY1fECAB/axaDBJLSh/RVC4BXacjVHQ4wx3y+buoNngZoX9kpJsbNxEjCpW8yyhWwO1+sseyBi9RJdg==, tarball: file:projects/model-document.tgz}
resolution: {integrity: sha512-5JcKBBX19mvQXZAg2p1z/qMSYqiR7py8mtNiHLaWKQpjHhMynutkkwGsnGq36Hb3Af8VtQiVxBZEq4rtcoph1Q==, tarball: file:projects/model-document.tgz}
name: '@rush-temp/model-document'
version: 0.0.0
dependencies:
@@ -28589,7 +28589,7 @@ packages:
dev: false
file:projects/server-backup.tgz(esbuild@0.20.1)(ts-node@10.9.2):
resolution: {integrity: sha512-e0MNgA1LeSOikxK3b6X/HMuzpyQ4MjAoJtHi84X28nxD+49sEy1NctCTptyqz2qWWicAZAghKV46Qi9N9+RnEw==, tarball: file:projects/server-backup.tgz}
resolution: {integrity: sha512-uuTU9Pa0R+AlAtxZo0mStW6hAnk/Ca3B++AoFAE5WDapBBxT0mVNEoDcAd0ATDKea0fo8mbRB4blnuoQa06NbA==, tarball: file:projects/server-backup.tgz}
id: file:projects/server-backup.tgz
name: '@rush-temp/server-backup'
version: 0.0.0
@@ -28604,6 +28604,7 @@ packages:
eslint-plugin-import: 2.29.1(eslint@8.56.0)
eslint-plugin-n: 15.7.0(eslint@8.56.0)
eslint-plugin-promise: 6.1.1(eslint@8.56.0)
fast-equals: 5.0.1
jest: 29.7.0(@types/node@20.11.19)(ts-node@10.9.2)
prettier: 3.2.5
prettier-plugin-svelte: 3.2.1(prettier@3.2.5)(svelte@4.2.11)
+2 -1
View File
@@ -42,7 +42,8 @@ async function doBackup (dirName: string, token: string, endpoint: string, works
ctx.info('do backup', { workspace, endpoint })
await backup(ctx, endpoint, wsid, storage, {
force: true,
recheck: false,
freshBackup: false,
clean: false,
skipDomains: [],
timeout: 0,
connectTimeout: 60 * 1000,
+9 -4
View File
@@ -639,6 +639,7 @@ export function devTool (
},
cmd.region,
true,
true,
5000, // 5 gigabytes per blob
async (storage, workspaceStorage) => {
if (cmd.remove) {
@@ -710,7 +711,8 @@ export function devTool (
})
},
cmd.region,
true,
false,
false,
100
)
) {
@@ -930,7 +932,8 @@ export function devTool (
)
.option('-bl, --blobLimit <blobLimit>', 'A blob size limit in megabytes (default 15mb)', '15')
.option('-f, --force', 'Force backup', false)
.option('-c, --recheck', 'Force hash recheck on server', false)
.option('-f, --fresh', 'Force fresh backup', false)
.option('-c, --clean', 'Force clean of old backup files, only with fresh backup option', false)
.option('-t, --timeout <timeout>', 'Connect timeout in seconds', '30')
.action(
async (
@@ -939,7 +942,8 @@ export function devTool (
cmd: {
skip: string
force: boolean
recheck: boolean
fresh: boolean
clean: boolean
timeout: string
include: string
blobLimit: string
@@ -951,7 +955,8 @@ export function devTool (
const endpoint = await getTransactorEndpoint(generateToken(systemAccountEmail, wsid), 'external')
await backup(toolCtx, endpoint, wsid, storage, {
force: cmd.force,
recheck: cmd.recheck,
freshBackup: cmd.fresh,
clean: cmd.clean,
include: cmd.include === '*' ? undefined : new Set(cmd.include.split(';').map((it) => it.trim())),
skipDomains: (cmd.skip ?? '').split(';').map((it) => it.trim()),
timeout: 0,
+2 -2
View File
@@ -416,8 +416,8 @@ export default function buildModel (enabled: string[] = ['*'], disabled: string[
{
label: testManagement.string.ConfigLabel,
description: testManagement.string.ConfigDescription,
enabled: false,
beta: false,
enabled: true,
beta: true,
classFilter: defaultFilter
}
],
+36 -29
View File
@@ -21,9 +21,9 @@ import core, {
DOMAIN_STATUS,
DOMAIN_TX,
generateId,
makeDocCollabId,
makeCollabJsonId,
makeCollabYdocId,
makeDocCollabId,
MeasureMetricsContext,
RateLimiter,
type AnyAttribute,
@@ -208,7 +208,7 @@ async function processMigrateContentFor (
const operations: { filter: MigrationDocumentQuery<Doc>, update: MigrateUpdate<Doc> }[] = []
for (const doc of docs) {
await rateLimiter.exec(async () => {
await rateLimiter.add(async () => {
const update: MigrateUpdate<Doc> = {}
for (const attribute of attributes) {
@@ -305,7 +305,7 @@ async function processMigrateJsonForDomain (
const operations: { filter: MigrationDocumentQuery<Doc>, update: MigrateUpdate<Doc> }[] = []
for (const doc of docs) {
await rateLimiter.exec(async () => {
await rateLimiter.add(async () => {
const update = await processMigrateJsonForDoc(ctx, doc, attributes, client, storageAdapter)
if (Object.keys(update).length > 0) {
operations.push({ filter: { _id: doc._id }, update })
@@ -385,8 +385,8 @@ async function processMigrateJsonForDoc (
const unset = update.$unset ?? {}
update.$unset = { ...unset, [attribute.name]: 1 }
} catch (err) {
ctx.warn('failed to process collaborative doc', { workspaceId, collabId, currentYdocId, err })
} catch (err: any) {
ctx.warn('failed to process collaborative doc', { workspaceId, collabId, currentYdocId, err: err.message })
}
}
@@ -427,36 +427,43 @@ export const coreOperation: MigrateOperation = {
state: 'remove-collection-txes',
func: async (client) => {
let processed = 0
let last = 0
const iterator = await client.traverse<TxCUD<Doc>>(DOMAIN_TX, {
_class: 'core:class:TxCollectionCUD' as Ref<Class<Doc>>
})
while (true) {
const txes = await iterator.next(200)
if (txes === null || txes.length === 0) break
processed += txes.length
try {
await client.create(
DOMAIN_TX,
txes.map((tx) => {
const { collection, objectId, objectClass } = tx
return {
collection,
attachedTo: objectId,
attachedToClass: objectClass,
...(tx as any).tx,
objectSpace: (tx as any).tx.objectSpace ?? tx.objectClass
}
try {
while (true) {
const txes = await iterator.next(1000)
if (txes === null || txes.length === 0) break
processed += txes.length
try {
await client.create(
DOMAIN_TX,
txes.map((tx) => {
const { collection, objectId, objectClass } = tx
return {
collection,
attachedTo: objectId,
attachedToClass: objectClass,
...(tx as any).tx,
objectSpace: (tx as any).tx.objectSpace ?? tx.objectClass
}
})
)
await client.deleteMany(DOMAIN_TX, {
_id: { $in: txes.map((it) => it._id) }
})
)
await client.deleteMany(DOMAIN_TX, {
_id: { $in: txes.map((it) => it._id) }
})
} catch (err: any) {
console.error(err)
} catch (err: any) {
console.error(err)
}
if (last !== Math.round(processed / 1000)) {
last = Math.round(processed / 1000)
console.log('processed', processed)
}
}
console.log('processed', processed)
} finally {
await iterator.close()
}
await iterator.close()
}
},
{
+1 -1
View File
@@ -120,7 +120,7 @@ describe('client', () => {
},
close: async () => {},
loadChunk: async (domain: Domain, idx?: number, recheck?: boolean) => ({
loadChunk: async (domain: Domain, idx?: number) => ({
idx: -1,
index: -1,
docs: [],
+1 -1
View File
@@ -60,7 +60,7 @@ export async function connect (handler: (tx: Tx) => void): Promise<ClientConnect
},
close: async () => {},
loadChunk: async (domain: Domain, idx?: number, recheck?: boolean) => ({
loadChunk: async (domain: Domain, idx?: number) => ({
idx: -1,
index: -1,
docs: [],
+1 -1
View File
@@ -17,7 +17,7 @@ export interface DocChunk {
* @public
*/
export interface BackupClient {
loadChunk: (domain: Domain, idx?: number, recheck?: boolean) => Promise<DocChunk>
loadChunk: (domain: Domain, idx?: number) => Promise<DocChunk>
closeChunk: (idx: number) => Promise<void>
loadDocs: (domain: Domain, docs: Ref<Doc>[]) => Promise<Doc[]>
+2 -2
View File
@@ -178,8 +178,8 @@ class ClientImpl implements AccountClient, BackupClient {
await this.conn.close()
}
async loadChunk (domain: Domain, idx?: number, recheck?: boolean): Promise<DocChunk> {
return await this.conn.loadChunk(domain, idx, recheck)
async loadChunk (domain: Domain, idx?: number): Promise<DocChunk> {
return await this.conn.loadChunk(domain, idx)
}
async closeChunk (idx: number): Promise<void> {
+1 -3
View File
@@ -25,7 +25,6 @@ import type { WorkspaceIdWithUrl } from './utils'
export interface DocInfo {
id: string
hash: string
size: number // Aprox size
}
/**
* @public
@@ -68,8 +67,7 @@ export interface SessionData {
*/
export interface LowLevelStorage {
// Low level streaming API to retrieve information
// If recheck is passed, all %hash% for documents, will be re-calculated.
find: (ctx: MeasureContext, domain: Domain, recheck?: boolean) => StorageIterator
find: (ctx: MeasureContext, domain: Domain) => StorageIterator
// Load passed documents from domain
load: (ctx: MeasureContext, domain: Domain, docs: Ref<Doc>[]) => Promise<Doc[]>
+1
View File
@@ -159,6 +159,7 @@ export async function tryMigrate (client: MigrationClient, plugin: string, migra
for (const migration of migrations) {
if (states.has(migration.state)) continue
try {
console.log('running migration', plugin, migration.state)
await migration.func(client)
} catch (err: any) {
console.error(err)
@@ -16,7 +16,7 @@
import { Analytics } from '@hcengineering/analytics'
import { resizeObserver } from '@hcengineering/ui'
import { onDestroy } from 'svelte'
import { drawing, type DrawingCmd, type DrawingData, type DrawingTool } from '../drawing'
import { drawing, type DrawingCmd, type DrawingData, type DrawingTool, type DrawTextCmd } from '../drawing'
import DrawingBoardToolbar from './DrawingBoardToolbar.svelte'
export let active = false
@@ -30,6 +30,7 @@
let penColor: string
let penWidth: number
let eraserWidth: number
let fontSize: number
let commands: DrawingCmd[] | undefined
let board: HTMLDivElement
let toolbar: HTMLDivElement
@@ -37,6 +38,7 @@
let oldReadonly: boolean
let oldDrawings: DrawingData[]
let modified = false
let changingCmdIndex: number | undefined
$: updateToolbarPosition(readonly, board, toolbar)
$: updateEditableState(drawings, readonly)
@@ -68,6 +70,7 @@
} else {
commands = undefined
}
changingCmdIndex = undefined
oldDrawings = drawings
oldReadonly = readonly
}
@@ -99,6 +102,40 @@
}
}
function addCommand (cmd: DrawingCmd): void {
if (commands !== undefined) {
commands = [...commands, cmd]
changingCmdIndex = undefined
modified = true
}
}
function showCommandProps (index: number): void {
changingCmdIndex = index
const anyCmd = commands?.[index]
if (anyCmd?.type === 'text') {
const cmd = anyCmd as DrawTextCmd
penColor = cmd.color
fontSize = cmd.fontSize
}
}
function changeCommand (index: number, cmd: DrawingCmd): void {
if (commands !== undefined) {
commands = commands.map((c, i) => (i === index ? cmd : c))
changingCmdIndex = undefined
modified = true
}
}
function deleteCommand (index: number): void {
if (commands !== undefined) {
commands = commands.filter((_, i) => i !== index)
changingCmdIndex = undefined
modified = true
}
}
onDestroy(() => {
saveDrawing()
})
@@ -121,9 +158,15 @@
penColor,
penWidth,
eraserWidth,
cmdAdded: () => {
modified = true
}
fontSize,
changingCmdIndex,
cmdAdded: addCommand,
cmdChanging: showCommandProps,
cmdChanged: changeCommand,
cmdUnchanged: () => {
changingCmdIndex = undefined
},
cmdDeleted: deleteCommand
}}
>
{#if !readonly}
@@ -134,6 +177,7 @@
bind:penColor
bind:penWidth
bind:eraserWidth
bind:fontSize
on:clear={() => {
commands = []
modified = true
@@ -28,6 +28,7 @@
import { createEventDispatcher, onMount } from 'svelte'
import IconEraser from './icons/Eraser.svelte'
import IconMove from './icons/Move.svelte'
import IconText from './icons/Text.svelte'
import { DrawingTool } from '../drawing'
import presentation from '../plugin'
@@ -40,13 +41,15 @@
color: 'drawingBoard.color',
colors: 'drawingBoard.colors',
penWidth: 'drawingBoard.penWidth',
eraserWidth: 'drawingBoard.eraserWidth'
eraserWidth: 'drawingBoard.eraserWidth',
fontSize: 'drawingBoard.fontSize'
}
export let tool: DrawingTool = 'pen'
export let penColor: string
export let penWidth: number
export let eraserWidth: number
export let fontSize: number
export let placeInside = false
export let showPanTool = false
export let toolbar: HTMLDivElement | undefined
@@ -130,8 +133,9 @@
if (!penColors.includes(penColor)) {
penColor = penColors[0] ?? defaultColor
}
penWidth = parseInt(localStorage.getItem(storageKey.penWidth) ?? '6')
penWidth = parseInt(localStorage.getItem(storageKey.penWidth) ?? '4')
eraserWidth = parseInt(localStorage.getItem(storageKey.eraserWidth) ?? '50')
fontSize = parseInt(localStorage.getItem(storageKey.fontSize) ?? '20')
})
function updatePenWidth (): void {
@@ -141,12 +145,17 @@
function updateEraserWidth (): void {
localStorage.setItem(storageKey.eraserWidth, eraserWidth.toString())
}
function updateFontSize (): void {
localStorage.setItem(storageKey.fontSize, fontSize.toString())
}
</script>
<div class="toolbar" class:inside={placeInside} bind:this={toolbar}>
<Button
icon={IconDelete}
kind="icon"
noFocus
on:click={() => {
tool = 'pen'
dispatch('clear')
@@ -156,6 +165,7 @@
<Button
icon={IconEdit}
kind="icon"
noFocus
selected={tool === 'pen'}
on:click={() => {
tool = 'pen'
@@ -164,6 +174,7 @@
<Button
icon={IconEraser}
kind="icon"
noFocus
selected={tool === 'erase'}
on:click={() => {
tool = 'erase'
@@ -173,24 +184,34 @@
<Button
icon={IconMove}
kind="icon"
noFocus
selected={tool === 'pan'}
on:click={() => {
tool = 'pan'
}}
/>
{/if}
<Button
icon={IconText}
kind="icon"
noFocus
selected={tool === 'text'}
on:click={() => {
tool = 'text'
}}
/>
<div class="divider buttons-divider" />
{#if tool === 'pen'}
<input
class="widthSelector"
type="range"
min={2}
max={18}
step={4}
max={20}
step={2}
bind:value={penWidth}
on:change={updatePenWidth}
/>
{:else}
{:else if tool === 'erase'}
<input
class="widthSelector"
type="range"
@@ -200,14 +221,27 @@
bind:value={eraserWidth}
on:change={updateEraserWidth}
/>
{:else if tool === 'text'}
<input
class="widthSelector"
type="range"
min={15}
max={35}
step={5}
bind:value={fontSize}
on:change={updateFontSize}
/>
{/if}
<div class="divider buttons-divider" />
{#each penColors as color}
<Button
kind="icon"
noFocus
selected={penColor === color}
on:click={() => {
tool = 'pen'
if (tool === 'erase') {
tool = 'pen'
}
selectColor(color)
}}
>
@@ -222,7 +256,7 @@
bind:value={penColor}
on:change={addColorPreset}
/>
<Button kind="icon" icon={IconMoreH} on:click={showMenu} />
<Button kind="icon" icon={IconMoreH} noFocus on:click={showMenu} />
</div>
</div>
@@ -242,7 +276,7 @@
border-radius: var(--small-BorderRadius);
border: 1px solid var(--theme-popup-divider);
box-shadow: 0.05rem 0.05rem 0.25rem rgba(0, 0, 0, 0.2);
z-index: 1;
z-index: 10;
}
}
@@ -0,0 +1,8 @@
<script lang="ts">
export let size: 'small' | 'medium' | 'large'
const fill: string = 'currentColor'
</script>
<svg class="svg-{size}" {fill} viewBox="0 0 16 16" xmlns="http://www.w3.org/2000/svg">
<path d="M5.032 13l.9-3h4.137l.9 3h1.775l-3-10H6.256l-3 10h1.776zm2.4-8h1.137l.9 3H6.532l.9-3z" />
</svg>
+423 -51
View File
@@ -22,26 +22,43 @@ export interface DrawingProps {
autoSize?: boolean
imageWidth?: number
imageHeight?: number
commandCount?: number
commands: DrawingCmd[]
offset?: Point
tool?: DrawingTool
penColor?: string
penWidth?: number
eraserWidth?: number
fontSize?: number
defaultCursor?: string
changingCmdIndex?: number
cmdAdded?: (cmd: DrawingCmd) => void
cmdChanging?: (index: number) => void
cmdUnchanged?: (index: number) => void
cmdChanged?: (index: number, cmd: DrawingCmd) => void
cmdDeleted?: (index: number) => void
panned?: (offset: Point) => void
}
export interface DrawingCmd {
type: 'line' | 'text'
}
export interface DrawTextCmd extends DrawingCmd {
text: string
pos: Point
fontSize: number
fontFace: string
color: string
}
export interface DrawLineCmd extends DrawingCmd {
lineWidth: number
erasing: boolean
penColor: string
points: Point[]
}
export type DrawingTool = 'pen' | 'erase' | 'pan'
export type DrawingTool = 'pen' | 'erase' | 'pan' | 'text'
interface Point {
x: number
@@ -52,6 +69,12 @@ function avgPoint (p1: Point, p2: Point): Point {
return { x: (p1.x + p2.x) / 2, y: (p1.y + p2.y) / 2 }
}
const maxTextLength = 500
const crossSvg = `<svg height="8" width="8" viewBox="0 0 16 16" fill="currentColor" xmlns="http://www.w3.org/2000/svg">
<path d="m1.29 2.71 5.3 5.29-5.3 5.29c-.92.92.49 2.34 1.41 1.41l5.3-5.29 5.29 5.3c.92.92 2.34-.49 1.41-1.41l-5.29-5.3 5.3-5.29c.92-.93-.49-2.34-1.42-1.42l-5.29 5.3-5.29-5.3c-.93-.92-2.34.49-1.42 1.42z"/>
</svg>`
class DrawState {
on = false
tool: DrawingTool = 'pen'
@@ -59,6 +82,8 @@ class DrawState {
penWidth = 4
eraserWidth = 30
minLineLength = 6
fontSize = 20
fontFace = '"IBM Plex Sans"'
center: Point = { x: 0, y: 0 }
offset: Point = { x: 0, y: 0 }
points: Point[] = []
@@ -78,16 +103,31 @@ class DrawState {
}
addPoint = (mouseX: number, mouseY: number): void => {
this.points.push({
x: mouseX * this.scale.x - this.offset.x - this.center.x,
y: mouseY * this.scale.y - this.offset.y - this.center.y
})
this.points.push(this.mouseToCanvasPoint({ x: mouseX, y: mouseY }))
}
mouseToCanvasPoint = (mouse: Point): Point => {
return {
x: mouse.x * this.scale.x - this.offset.x - this.center.x,
y: mouse.y * this.scale.y - this.offset.y - this.center.y
}
}
canvasToMousePoint = (canvas: Point): Point => {
return {
x: canvas.x / this.scale.x + this.offset.x + this.center.x,
y: canvas.y / this.scale.y + this.offset.y + this.center.y
}
}
isDrawingTool = (): boolean => {
return this.tool === 'pen' || this.tool === 'erase'
}
translateCtx = (): void => {
this.ctx.translate(this.offset.x + this.center.x, this.offset.y + this.center.y)
}
drawLive = (x: number, y: number, lastPoint = false): void => {
window.requestAnimationFrame(() => {
if (!lastPoint || this.points.length > 1) {
@@ -95,11 +135,11 @@ class DrawState {
}
const erasing = this.tool === 'erase'
this.ctx.save()
this.ctx.translate(this.offset.x + this.center.x, this.offset.y + this.center.y)
this.translateCtx()
this.ctx.beginPath()
this.ctx.lineCap = 'round'
this.ctx.strokeStyle = this.penColor
this.ctx.lineWidth = (erasing ? this.eraserWidth : this.penWidth) * this.lineScale()
this.ctx.lineWidth = erasing ? this.eraserWidth : this.penWidth
this.ctx.globalCompositeOperation = erasing ? 'destination-out' : 'source-over'
if (this.points.length === 1) {
this.drawPoint(this.points[0], erasing)
@@ -112,8 +152,16 @@ class DrawState {
}
drawCommand = (cmd: DrawingCmd): void => {
if (cmd.type === 'text') {
this.drawTextCommand(cmd as DrawTextCmd)
} else {
this.drawLineCommand(cmd as DrawLineCmd)
}
}
drawLineCommand = (cmd: DrawLineCmd): void => {
this.ctx.save()
this.ctx.translate(this.offset.x + this.center.x, this.offset.y + this.center.y)
this.translateCtx()
this.ctx.beginPath()
this.ctx.lineCap = 'round'
this.ctx.strokeStyle = cmd.penColor
@@ -130,6 +178,39 @@ class DrawState {
this.ctx.restore()
}
drawTextCommand = (cmd: DrawTextCmd): void => {
const p = { ...cmd.pos }
this.ctx.save()
this.translateCtx()
this.ctx.font = `${cmd.fontSize}px ${cmd.fontFace}`
this.ctx.fillStyle = cmd.color
this.ctx.textBaseline = 'top'
const lines = cmd.text.split('\n').map((l) => l.trim())
for (let i = 0; i < lines.length; i++) {
const line = lines[i]
this.ctx.fillText(line, p.x, p.y)
p.y += cmd.fontSize
}
this.ctx.restore()
}
isPointInText = (p: Point, cmd: DrawTextCmd): boolean => {
this.ctx.font = `${cmd.fontSize}px ${cmd.fontFace}`
const lines = cmd.text.split('\n').map((l) => l.trim())
for (let i = 0; i < lines.length; i++) {
if (p.y < cmd.pos.y + i * cmd.fontSize || p.y > cmd.pos.y + (i + 1) * cmd.fontSize) {
continue
}
const line = lines[i]
const metrics = this.ctx.measureText(line)
if (p.x < cmd.pos.x || p.x > cmd.pos.x + metrics.width) {
continue
}
return true
}
return false
}
drawPoint = (p: Point, erasing: boolean): void => {
let r = this.ctx.lineWidth / 2
if (!erasing) {
@@ -211,16 +292,20 @@ export function drawing (
draw.penColor = props.penColor ?? draw.penColor
draw.penWidth = props.penWidth ?? draw.penWidth
draw.eraserWidth = props.eraserWidth ?? draw.eraserWidth
draw.fontSize = props.fontSize ?? draw.fontSize
draw.offset = props.offset ?? draw.offset
updateCanvasCursor()
interface LiveTextBox {
pos: Point
box: HTMLDivElement
editor: HTMLDivElement
cmdIndex: number
}
let liveTextBox: LiveTextBox | undefined
let commands = props.commands
let commandCount = props.commandCount ?? 0
replayCommands({
offset: {
x: props.offset?.x ?? 0,
y: props.offset?.y ?? 0
}
})
replayCommands()
const resizeObserver = new ResizeObserver((entries) => {
for (const entry of entries) {
@@ -231,7 +316,7 @@ export function drawing (
canvas.height = Math.floor(entry.contentRect.height)
draw.center.x = canvas.width / 2
draw.center.y = canvas.height / 2
replayCommands({ offset: draw.offset })
replayCommands()
} else {
draw.scale = {
x: canvas.width / entry.contentRect.width,
@@ -288,15 +373,16 @@ export function drawing (
if (draw.on && draw.tool === 'pan') {
requestAnimationFrame(() => {
replayCommands({
offset: {
x: draw.offset.x + x - prevPos.x,
y: draw.offset.y + y - prevPos.y
}
})
draw.offset.x += x - prevPos.x
draw.offset.y += y - prevPos.y
replayCommands()
prevPos = { x, y }
})
}
if (draw.on && draw.tool === 'text') {
prevPos = { x, y }
}
}
canvas.onpointerup = (e) => {
@@ -308,9 +394,16 @@ export function drawing (
if (draw.on) {
if (draw.isDrawingTool()) {
draw.drawLive(e.offsetX, e.offsetY, true)
storeCommand()
storeLineCommand()
} else if (draw.tool === 'pan') {
props.panned?.(draw.offset)
} else if (draw.tool === 'text') {
if (liveTextBox !== undefined) {
storeTextCommand()
} else {
const cmdIndex = findTextCommand(prevPos)
props.cmdChanging?.(cmdIndex)
}
}
draw.on = false
}
@@ -328,16 +421,280 @@ export function drawing (
}
}
function storeCommand (): void {
function findTextCommand (mousePos: Point): number {
const pos = draw.mouseToCanvasPoint(mousePos)
for (let i = commands.length - 1; i >= 0; i--) {
const anyCmd = commands[i]
if (anyCmd.type === 'text') {
const cmd = anyCmd as DrawTextCmd
if (draw.isPointInText(pos, cmd)) {
return i
}
}
}
return -1
}
function makeLiveTextBox (cmdIndex: number): void {
let pos = prevPos
let existingCmd: DrawTextCmd | undefined
if (cmdIndex >= 0 && commands[cmdIndex]?.type === 'text') {
existingCmd = commands[cmdIndex] as DrawTextCmd
pos = draw.canvasToMousePoint(existingCmd.pos)
}
const padding = 6
const handleSize = 14
const box = document.createElement('div')
box.style.zIndex = '1'
box.style.position = 'absolute'
box.style.left = `calc(${pos.x}px - ${padding}px)`
box.style.top = `calc(${pos.y}px - ${padding}px)`
box.style.border = '1px solid var(--theme-editbox-focus-border)'
box.style.borderRadius = 'var(--small-BorderRadius)'
box.style.padding = `${padding}px`
box.style.background = 'var(--theme-popup-header)'
box.addEventListener('mousedown', (e) => {
e.stopPropagation()
})
box.addEventListener('click', (e) => {
e.stopPropagation()
editor.focus()
})
const editor = document.createElement('div')
editor.style.cursor = 'text'
editor.style.padding = '0'
editor.contentEditable = 'true'
editor.style.outline = 'none'
editor.style.minWidth = '2rem'
editor.style.whiteSpace = 'nowrap'
if (existingCmd !== undefined) {
editor.innerText = existingCmd.text
}
editor.addEventListener('input', (e) => {
if (editor.innerText.length > maxTextLength) {
e.preventDefault()
editor.innerText = editor.innerText.substring(0, maxTextLength)
moveCaretToEnd()
}
})
editor.addEventListener('paste', (e) => {
e.preventDefault()
const selection = window.getSelection()
const text = (e.clipboardData?.getData('text/plain') ?? '').trim()
if (text.length === 0 || selection === null || selection.rangeCount === 0) {
return
}
let selectedLen = 0
const range = selection.getRangeAt(0)
if (editor.contains(range.commonAncestorContainer)) {
selectedLen = range.endOffset - range.startOffset
}
const availableLen = maxTextLength - (selectedLen > 0 ? selectedLen : editor.innerText.length)
if (availableLen > 0) {
const lines = text.slice(0, availableLen).split('\n')
const pastedNode = document.createDocumentFragment()
for (let i = 0; i < lines.length; i++) {
pastedNode.appendChild(document.createTextNode(lines[i]))
if (i < lines.length - 1) {
pastedNode.appendChild(document.createElement('br'))
}
}
range.deleteContents()
range.insertNode(pastedNode)
// move caret to the end of pasted node
range.collapse(false)
selection.addRange(range)
}
})
editor.addEventListener('keydown', (e) => {
if (e.key === 'Escape') {
e.preventDefault()
if (liveTextBox !== undefined) {
const cmdIndex = liveTextBox.cmdIndex
if (cmdIndex >= 0) {
// reset changingCmdIndex in clients
setTimeout(() => {
props.cmdUnchanged?.(cmdIndex)
}, 0)
}
}
closeLiveTextBox()
replayCommands()
} else if (e.key === 'Enter' && e.ctrlKey) {
e.preventDefault()
storeTextCommand()
}
})
box.appendChild(editor)
const moveCaretToEnd = (): void => {
const selection = window.getSelection()
const range = document.createRange()
range.setStartAfter(editor.lastChild ?? editor)
range.collapse(true)
selection?.removeAllRanges()
selection?.addRange(range)
}
const selectAll = (): void => {
const selection = window.getSelection()
const range = document.createRange()
range.selectNodeContents(editor)
selection?.removeAllRanges()
selection?.addRange(range)
}
const makeHandle = (): HTMLDivElement => {
const handle = document.createElement('div')
handle.style.position = 'absolute'
handle.style.top = `-${handleSize / 2}px`
handle.style.width = `${handleSize}px`
handle.style.height = `${handleSize}px`
handle.style.color = 'var(--global-on-accent-TextColor)'
handle.style.background = 'var(--global-accent-IconColor)'
handle.style.borderRadius = '50%'
handle.style.display = 'flex'
handle.style.alignItems = 'center'
handle.style.justifyContent = 'center'
return handle
}
const dragHandle = makeHandle()
dragHandle.style.left = `-${handleSize / 2}px`
dragHandle.style.cursor = 'grab'
dragHandle.addEventListener('pointerdown', (e) => {
e.preventDefault()
dragHandle.style.cursor = 'grabbing'
dragHandle.setPointerCapture(e.pointerId)
const x = e.clientX
const y = e.clientY
const dragStart = { x, y }
const pointerMove = (e: PointerEvent): void => {
e.preventDefault()
const x = e.clientX
const y = e.clientY
const dx = x - dragStart.x
const dy = y - dragStart.y
dragStart.x = x
dragStart.y = y
let newX = box.offsetLeft + dx
let newY = box.offsetTop + dy
// For screenshots the canvas always has the same size as the underlying image
// and we should not be able to drag the text box outside of the screenshot
if (props.autoSize !== true) {
newX = Math.max(0, newX)
newY = Math.max(0, newY)
if (newX + box.offsetWidth > node.clientWidth) {
newX = node.clientWidth - box.offsetWidth
}
if (newY + box.offsetHeight > node.clientHeight) {
newY = node.clientHeight - box.offsetHeight
}
}
box.style.left = `${newX}px`
box.style.top = `${newY}px`
if (liveTextBox !== undefined) {
liveTextBox.pos.x = newX + padding
liveTextBox.pos.y = newY + padding
}
}
const pointerUp = (e: PointerEvent): void => {
setTimeout(() => {
editor.focus()
}, 100)
e.preventDefault()
dragHandle.style.cursor = 'grab'
dragHandle.releasePointerCapture(e.pointerId)
dragHandle.removeEventListener('pointermove', pointerMove)
dragHandle.removeEventListener('pointerup', pointerUp)
}
dragHandle.addEventListener('pointermove', pointerMove)
dragHandle.addEventListener('pointerup', pointerUp)
})
box.appendChild(dragHandle)
const deleteButton = makeHandle()
deleteButton.style.right = `-${handleSize / 2}px`
deleteButton.style.cursor = 'pointer'
deleteButton.innerHTML = crossSvg
deleteButton.addEventListener('click', () => {
node.removeChild(box)
if (liveTextBox?.cmdIndex !== undefined) {
props.cmdDeleted?.(liveTextBox.cmdIndex)
}
liveTextBox = undefined
})
box.appendChild(deleteButton)
node.appendChild(box)
liveTextBox = { box, editor, pos, cmdIndex }
updateLiveTextBox()
setTimeout(() => {
editor.focus()
}, 100)
selectAll()
}
function updateLiveTextBox (): void {
if (liveTextBox !== undefined) {
liveTextBox.editor.style.color = draw.penColor
liveTextBox.editor.style.lineHeight = `${draw.fontSize / draw.lineScale()}px`
liveTextBox.editor.style.fontSize = `${draw.fontSize / draw.lineScale()}px`
liveTextBox.editor.style.fontFamily = draw.fontFace
}
}
function closeLiveTextBox (): void {
if (liveTextBox !== undefined) {
node.removeChild(liveTextBox.box)
liveTextBox = undefined
}
}
function storeTextCommand (defer = false): void {
if (liveTextBox !== undefined) {
const text = (liveTextBox.editor.innerText ?? '').trim()
if (text !== '') {
const cmd: DrawTextCmd = {
type: 'text',
text,
pos: draw.mouseToCanvasPoint(liveTextBox.pos),
fontSize: draw.fontSize,
fontFace: draw.fontFace,
color: draw.penColor
}
const cmdIndex = liveTextBox.cmdIndex
const notify = (): void => {
if (cmdIndex >= 0) {
props.cmdChanged?.(cmdIndex, cmd)
} else {
props.cmdAdded?.(cmd)
}
}
if (defer) {
setTimeout(notify, 0)
} else {
notify()
}
} else {
props.cmdUnchanged?.(liveTextBox.cmdIndex)
}
}
}
function storeLineCommand (): void {
if (draw.points.length > 0) {
const erasing = draw.tool === 'erase'
const cmd: DrawingCmd = {
lineWidth: (erasing ? draw.eraserWidth : draw.penWidth) * draw.lineScale(),
const cmd: DrawLineCmd = {
type: 'line',
lineWidth: erasing ? draw.eraserWidth : draw.penWidth,
erasing,
penColor: draw.penColor,
points: draw.points
}
commands.push(cmd)
props.cmdAdded?.(cmd)
}
}
@@ -360,28 +717,23 @@ export function drawing (
} else if (draw.tool === 'pan') {
canvas.style.cursor = 'move'
canvasCursor.style.visibility = 'hidden'
} else if (draw.tool === 'text') {
canvas.style.cursor = 'text'
canvasCursor.style.visibility = 'hidden'
} else {
canvas.style.cursor = 'default'
canvasCursor.style.visibility = 'hidden'
}
}
function clearCanvas (): void {
function replayCommands (): void {
draw.ctx.reset()
draw.offset = { x: 0, y: 0 }
}
function replayCommands ({ offset, startIndex }: { offset?: Point, startIndex?: number } = {}): void {
if (startIndex === undefined || startIndex === 0) {
clearCanvas()
}
if (offset !== undefined) {
draw.offset = offset
}
for (let i = startIndex ?? 0; i < commands.length; i++) {
for (let i = 0; i < commands.length; i++) {
if (liveTextBox?.cmdIndex === i) {
continue
}
draw.drawCommand(commands[i])
}
commandCount = commands.length
}
return {
@@ -389,45 +741,65 @@ export function drawing (
update (props: DrawingProps) {
let replay = false
let offset: Point | undefined
let startIndex: number | undefined
if (props.offset?.x !== draw.offset.x || props.offset?.y !== draw.offset.y) {
offset = props.offset
if (props.offset !== undefined && (props.offset.x !== draw.offset.x || props.offset.y !== draw.offset.y)) {
draw.offset = props.offset
replay = true
}
if (commands !== props.commands) {
commands = props.commands
replay = true
} else if (props.commandCount !== undefined && props.commandCount !== commandCount) {
startIndex = commandCount
replay = true
}
let updateCursor = false
let updateTextBox = false
if (draw.tool !== props.tool) {
draw.tool = props.tool ?? 'pen'
updateCursor = true
}
if (draw.penColor !== props.penColor) {
draw.penColor = props.penColor ?? 'blue'
updateTextBox = true
updateCursor = true
}
if (draw.penWidth !== props.penWidth) {
draw.penWidth = props.penWidth ?? 5
draw.penWidth = props.penWidth ?? draw.penWidth
updateCursor = true
}
if (draw.eraserWidth !== props.eraserWidth) {
draw.eraserWidth = props.eraserWidth ?? 5
draw.eraserWidth = props.eraserWidth ?? draw.eraserWidth
updateCursor = true
}
if (draw.fontSize !== props.fontSize) {
draw.fontSize = props.fontSize ?? draw.fontSize
updateTextBox = true
}
if (props.readonly !== readonly) {
readonly = props.readonly ?? false
updateCursor = true
}
if (props.changingCmdIndex === undefined) {
if (liveTextBox !== undefined) {
storeTextCommand(true)
closeLiveTextBox()
replay = true
}
} else {
if (liveTextBox === undefined) {
makeLiveTextBox(props.changingCmdIndex)
replay = true
} else if (liveTextBox.cmdIndex !== props.changingCmdIndex) {
storeTextCommand(true)
closeLiveTextBox()
replay = true
}
}
if (updateCursor) {
updateCanvasCursor()
}
if (updateTextBox) {
updateLiveTextBox()
}
if (replay) {
replayCommands({ offset, startIndex })
replayCommands()
}
}
}
+1 -1
View File
@@ -82,7 +82,7 @@ FulltextStorage & {
return {}
},
close: async () => {},
loadChunk: async (domain: Domain, idx?: number, recheck?: boolean) => ({
loadChunk: async (domain: Domain, idx?: number) => ({
idx: -1,
index: -1,
docs: [],
+2 -2
View File
@@ -663,8 +663,8 @@ class Connection implements ClientConnection {
})
}
loadChunk (domain: Domain, idx?: number, recheck?: boolean): Promise<DocChunk> {
return this.sendRequest({ method: 'loadChunk', params: [domain, idx, recheck] })
loadChunk (domain: Domain, idx?: number): Promise<DocChunk> {
return this.sendRequest({ method: 'loadChunk', params: [domain, idx] })
}
closeChunk (idx: number): Promise<void> {
@@ -16,17 +16,30 @@
import type { Asset, IntlString } from '@hcengineering/platform'
import type { AnySvelteComponent } from '@hcengineering/ui'
import { AppItem } from '@hcengineering/workbench-resources'
import { RoomType } from '@hcengineering/love'
import { currentRoom } from '../../stores'
import { isConnected, isSharingEnabled, isCameraEnabled, isMicEnabled } from '../../utils'
import love from '../../plugin'
export let label: IntlString
export let icon: Asset | AnySvelteComponent
export let selected: boolean = false
export let size: 'small' | 'medium' | 'large' = 'small'
$: allowCam = $currentRoom?.type === RoomType.Video
</script>
<AppItem
{label}
{icon}
icon={$isSharingEnabled
? love.icon.SharingDisabled
: $isConnected && allowCam && !$isCameraEnabled && !$isMicEnabled
? love.icon.CamDisabled
: $isConnected && !allowCam && !$isMicEnabled
? love.icon.MicDisabled
: !allowCam || (!$isCameraEnabled && $isMicEnabled)
? love.icon.Mic
: icon}
{selected}
{size}
kind={$isSharingEnabled
@@ -13,10 +13,11 @@
// limitations under the License.
-->
<script lang="ts">
import { Doc, DocumentQuery, Ref, Space } from '@hcengineering/core'
import { Doc, DocumentQuery, Ref, Space, mergeQueries } from '@hcengineering/core'
import { Button } from '@hcengineering/ui'
import { selectionStore } from '@hcengineering/view-resources'
import type { TestCase, TestProject } from '@hcengineering/test-management'
import { createQuery } from '@hcengineering/presentation'
import testManagement from '../../plugin'
import { showCreateTestRunPopup } from '../../utils'
@@ -24,6 +25,19 @@
export let query: DocumentQuery<Doc> = {}
export let space: Ref<Space>
const docQuery = createQuery()
let haveTestCases = false
$: resultQuery = mergeQueries(query, { space })
$: docQuery.query(
testManagement.class.TestCase,
resultQuery,
(res) => {
haveTestCases = res.length > 0
},
{ limit: 1 }
)
const project: Ref<TestProject> = space as any
const handleRun = async (): Promise<void> => {
@@ -42,5 +56,6 @@
justify={'left'}
kind={'primary'}
label={testManagement.string.RunTestCases}
disabled={!haveTestCases}
on:click={handleRun}
/>
@@ -13,7 +13,7 @@
// limitations under the License.
-->
<script lang="ts">
import { DrawingBoardToolbar, DrawingCmd, DrawingTool, drawing } from '@hcengineering/presentation'
import { DrawingBoardToolbar, DrawingCmd, DrawingTool, DrawTextCmd, drawing } from '@hcengineering/presentation'
import { Loading } from '@hcengineering/ui'
import { onMount, onDestroy } from 'svelte'
import { Array as YArray, Map as YMap } from 'yjs'
@@ -31,29 +31,47 @@
let penColor: string
let penWidth: number
let eraserWidth: number
let commandCount: number
let fontSize: number
let commands: DrawingCmd[] = []
let offset: { x: number, y: number }
let offset: { x: number, y: number } = { x: 0, y: 0 }
let changingCmdIndex: number | undefined
let toolbar: HTMLDivElement
let oldSelected = false
$: onSelectedChanged(selected)
function listenSavedCommands (): void {
if (savedCmds.length === 0) {
commands = []
} else {
for (let i = commands.length; i < savedCmds.length; i++) {
commands.push(savedCmds.get(i))
}
}
commandCount = savedCmds.length
commands = savedCmds.toArray()
}
function listenSavedProps (): void {
offset = savedProps.get('offset')
// We have only local offset for now
// A global offset should be implemented as a "Follow" feature
// offset = savedProps.get('offset')
}
function showCommandProps (index: number): void {
changingCmdIndex = index
const anyCmd = commands[index]
if (anyCmd?.type === 'text') {
const cmd = anyCmd as DrawTextCmd
penColor = cmd.color
fontSize = cmd.fontSize
}
}
function onSelectedChanged (selected: boolean): void {
if (oldSelected !== selected) {
if (oldSelected && !selected && changingCmdIndex !== undefined) {
changingCmdIndex = undefined
}
oldSelected = selected
}
}
onMount(() => {
commands = savedCmds.toArray()
offset = savedProps.get('offset')
// offset = savedProps.get('offset')
savedCmds.observe(listenSavedCommands)
savedProps.observe(listenSavedProps)
})
@@ -84,18 +102,34 @@
use:drawing={{
readonly,
autoSize: true,
commandCount,
commands,
offset,
tool,
penColor,
penWidth,
eraserWidth,
fontSize,
changingCmdIndex,
cmdAdded: (cmd) => {
savedCmds.push([cmd])
changingCmdIndex = undefined
},
panned: (offset) => {
savedProps.set('offset', offset)
cmdChanging: showCommandProps,
cmdChanged: (index, cmd) => {
savedCmds.delete(index)
savedCmds.insert(index, [cmd])
changingCmdIndex = undefined
},
cmdUnchanged: () => {
changingCmdIndex = undefined
},
cmdDeleted: (index) => {
savedCmds.delete(index)
changingCmdIndex = undefined
},
panned: (newOffset) => {
offset = newOffset
// savedProps.set('offset', offset)
}
}}
>
@@ -113,9 +147,11 @@
bind:penColor
bind:penWidth
bind:eraserWidth
bind:fontSize
on:clear={() => {
savedCmds.delete(0, savedCmds.length)
savedProps.set('offset', { x: 0, y: 0 })
offset = { x: 0, y: 0 }
// savedProps.set('offset', { x: 0, y: 0 })
}}
/>
{/if}
@@ -103,6 +103,7 @@
kind={selected ? 'primary' : 'ghost'}
icon={IconScribble}
disabled={loading}
noFocus
on:click={() => {
showBoardPopup(savedBoard, editor)
}}
@@ -119,7 +119,7 @@
}
&.accented,
&.accented.selected {
background-color: var(--global-disabled-PriorityColor);
background-color: var(--button-secondary-active-BackgroundColor);
}
&.positive .icon-container,
&.negative .icon-container,
+1 -1
View File
@@ -35,7 +35,7 @@ describe.skip('test-backup-find', () => {
const docs: Doc[] = []
while (true) {
const chunk = await client.loadChunk(DOMAIN_TX, 0, true)
const chunk = await client.loadChunk(DOMAIN_TX, 0)
const part = await client.loadDocs(
DOMAIN_TX,
chunk.docs.map((doc) => doc.id as Ref<Doc>)
@@ -605,7 +605,7 @@ export async function createPushNotification (
const limiter = new RateLimiter(5)
for (const subscription of userSubscriptions) {
await limiter.exec(async () => {
await limiter.add(async () => {
await sendPushToSubscription(control, target, subscription, data)
})
}
+4 -2
View File
@@ -91,7 +91,8 @@ export async function backupWorkspace (
externalStorage: StorageAdapter
) => DbConfiguration,
region: string,
recheck: boolean = false,
freshBackup: boolean = false,
clean: boolean = false,
downloadLimit: number,
onFinish?: (backupStorage: StorageAdapter, workspaceStorage: StorageAdapter) => Promise<void>
@@ -126,7 +127,8 @@ export async function backupWorkspace (
return getConfig(ctx, mainDbUrl, workspace, branding, externalStorage)
},
region,
recheck,
freshBackup,
clean,
downloadLimit
)
if (result && onFinish !== undefined) {
+2 -1
View File
@@ -51,6 +51,7 @@
"@hcengineering/server-tool": "^0.6.0",
"@hcengineering/server-client": "^0.6.0",
"@hcengineering/server-token": "^0.6.11",
"@hcengineering/server-core": "^0.6.1"
"@hcengineering/server-core": "^0.6.1",
"fast-equals": "^5.0.1"
}
}
+110 -63
View File
@@ -32,6 +32,7 @@ import core, {
Ref,
SortingOrder,
systemAccountEmail,
toIdMap,
TxProcessor,
WorkspaceId,
type BackupStatus,
@@ -41,9 +42,10 @@ import core, {
type TxCUD
} from '@hcengineering/core'
import { BlobClient, createClient } from '@hcengineering/server-client'
import { type StorageAdapter } from '@hcengineering/server-core'
import { estimateDocSize, type StorageAdapter } from '@hcengineering/server-core'
import { generateToken } from '@hcengineering/server-token'
import { connect } from '@hcengineering/server-tool'
import { deepEqual } from 'fast-equals'
import { createReadStream, createWriteStream, existsSync, mkdirSync, statSync } from 'node:fs'
import { rm } from 'node:fs/promises'
import { basename, dirname } from 'node:path'
@@ -58,7 +60,6 @@ export * from './storage'
const dataBlobSize = 50 * 1024 * 1024
const dataUploadSize = 2 * 1024 * 1024
const retrieveChunkSize = 2 * 1024 * 1024
const defaultLevel = 9
@@ -134,7 +135,7 @@ async function loadDigest (
date?: number
): Promise<Map<Ref<Doc>, string>> {
ctx = ctx.newChild('load digest', { domain, count: snapshots.length })
ctx.info('load-digest', { domain, count: snapshots.length })
ctx.info('loading-digest', { domain, snapshots: snapshots.length })
const result = new Map<Ref<Doc>, string>()
for (const s of snapshots) {
const d = s.domains[domain]
@@ -186,6 +187,7 @@ async function loadDigest (
}
}
ctx.end()
ctx.info('load-digest', { domain, snapshots: snapshots.length, documents: result.size })
return result
}
async function verifyDigest (
@@ -477,9 +479,8 @@ export async function cloneWorkspace (
idx = it.idx
let needRetrieve: Ref<Doc>[] = []
let needRetrieveSize = 0
for (const { id, hash, size } of it.docs) {
for (const { id, hash } of it.docs) {
processed++
if (Date.now() - st > 2500) {
ctx.info('processed', { processed, time: Date.now() - st, workspace: targetWorkspaceId.name })
@@ -488,11 +489,9 @@ export async function cloneWorkspace (
changes.added.set(id as Ref<Doc>, hash)
needRetrieve.push(id as Ref<Doc>)
needRetrieveSize += size
if (needRetrieveSize > retrieveChunkSize) {
if (needRetrieve.length > 200) {
needRetrieveChunks.push(needRetrieve)
needRetrieveSize = 0
needRetrieve = []
}
}
@@ -532,7 +531,7 @@ export async function cloneWorkspace (
for (const d of docs) {
if (d._class === core.class.Blob) {
const blob = d as Blob
await executor.exec(async () => {
await executor.add(async () => {
try {
ctx.info('clone blob', { name: blob._id, contentType: blob.contentType })
const readable = await storageAdapter.get(ctx, sourceWorkspaceId, blob._id)
@@ -662,7 +661,8 @@ export async function backup (
include?: Set<string>
skipDomains: string[]
force: boolean
recheck: boolean
freshBackup: boolean // If passed as true, will download all documents except blobs as new backup
clean: boolean // If set will perform a clena of old backup files
timeout: number
connectTimeout: number
skipBlobContentTypes: string[]
@@ -676,7 +676,8 @@ export async function backup (
token?: string
} = {
force: false,
recheck: false,
freshBackup: false,
clean: false,
timeout: 0,
skipDomains: [],
connectTimeout: 30000,
@@ -693,7 +694,7 @@ export async function backup (
ctx = ctx.newChild('backup', {
workspaceId: workspaceId.name,
force: options.force,
recheck: options.recheck,
recheck: options.freshBackup,
timeout: options.timeout
})
@@ -741,7 +742,7 @@ export async function backup (
let lastTxChecked = false
// Skip backup if there is no transaction changes.
if (options.getLastTx !== undefined) {
if (options.getLastTx !== undefined && !options.freshBackup) {
lastTx = await options.getLastTx()
if (lastTx !== undefined) {
if (lastTx._id === backupInfo.lastTxId && !options.force) {
@@ -771,7 +772,7 @@ export async function backup (
options.connectTimeout
)) as CoreClient & BackupClient)
if (!lastTxChecked) {
if (!lastTxChecked && !options.freshBackup) {
lastTx = await connection.findOne(
core.class.Tx,
{ objectSpace: { $ne: core.space.Model } },
@@ -888,19 +889,13 @@ export async function backup (
}
while (true) {
try {
const currentChunk = await ctx.with('loadChunk', {}, () => connection.loadChunk(domain, idx, options.recheck))
const currentChunk = await ctx.with('loadChunk', {}, () => connection.loadChunk(domain, idx))
idx = currentChunk.idx
ops++
let needRetrieve: Ref<Doc>[] = []
let currentNeedRetrieveSize = 0
for (const { id, hash, size } of currentChunk.docs) {
if (domain === DOMAIN_BLOB) {
result.blobsSize += size
} else {
result.dataSize += size
}
for (const { id, hash } of currentChunk.docs) {
processed++
if (Date.now() - st > 2500) {
ctx.info('processed', {
@@ -911,19 +906,18 @@ export async function backup (
})
st = Date.now()
}
const _hash = doTrimHash(hash) as string
const kHash = doTrimHash(digest.get(id as Ref<Doc>) ?? oldHash.get(id as Ref<Doc>))
if (kHash !== undefined) {
const serverDocHash = doTrimHash(hash) as string
const currentHash = doTrimHash(digest.get(id as Ref<Doc>) ?? oldHash.get(id as Ref<Doc>))
if (currentHash !== undefined) {
if (digest.delete(id as Ref<Doc>)) {
oldHash.set(id as Ref<Doc>, kHash)
oldHash.set(id as Ref<Doc>, currentHash)
}
if (kHash !== _hash) {
if (currentHash !== serverDocHash || (options.freshBackup && domain !== DOMAIN_BLOB)) {
if (changes.updated.has(id as Ref<Doc>)) {
removeFromNeedRetrieve(needRetrieve, id as Ref<Doc>)
}
changes.updated.set(id as Ref<Doc>, _hash)
changes.updated.set(id as Ref<Doc>, serverDocHash)
needRetrieve.push(id as Ref<Doc>)
currentNeedRetrieveSize += size
changed++
} else if (changes.updated.has(id as Ref<Doc>)) {
// We have same
@@ -936,22 +930,19 @@ export async function backup (
// We need to clean old need retrieve in case of duplicates.
removeFromNeedRetrieve(needRetrieve, id)
}
changes.added.set(id as Ref<Doc>, _hash)
changes.added.set(id as Ref<Doc>, serverDocHash)
needRetrieve.push(id as Ref<Doc>)
changed++
currentNeedRetrieveSize += size
}
if (currentNeedRetrieveSize > retrieveChunkSize) {
if (needRetrieve.length > 0) {
needRetrieveChunks.push(needRetrieve)
}
currentNeedRetrieveSize = 0
if (needRetrieve.length > 200) {
needRetrieveChunks.push(needRetrieve)
needRetrieve = []
}
}
if (needRetrieve.length > 0) {
needRetrieveChunks.push(needRetrieve)
needRetrieve = []
}
if (currentChunk.finished) {
ctx.info('processed-end', {
@@ -1036,6 +1027,8 @@ export async function backup (
global.gc?.()
} catch (err) {}
let lastSize = 0
while (needRetrieveChunks.length > 0) {
if (canceled()) {
return
@@ -1048,11 +1041,13 @@ export async function backup (
ctx.info('Retrieve chunk', {
needRetrieve: needRetrieveChunks.reduce((v, docs) => v + docs.length, 0),
toLoad: needRetrieve.length,
workspace: workspaceId.name
workspace: workspaceId.name,
lastSize: Math.round((lastSize * 100) / (1024 * 1024)) / 100
})
let docs: Doc[] = []
try {
docs = await ctx.with('load-docs', {}, async (ctx) => await connection.loadDocs(domain, needRetrieve))
lastSize = docs.reduce((p, it) => p + estimateDocSize(it), 0)
if (docs.length !== needRetrieve.length) {
const nr = new Set(docs.map((it) => it._id))
ctx.error('failed to retrieve all documents', { missing: needRetrieve.filter((it) => !nr.has(it)) })
@@ -1146,6 +1141,12 @@ export async function backup (
break
}
if (domain === DOMAIN_BLOB) {
result.blobsSize += (d as Blob).size
} else {
result.dataSize += JSON.stringify(d).length
}
function processChanges (d: Doc, error: boolean = false): void {
processed++
progress(10 + (processed / totalChunks) * 90)
@@ -1319,6 +1320,41 @@ export async function backup (
}
let processed = 0
if (!canceled()) {
backupInfo.lastTxId = lastTx?._id ?? '0' // We could store last tx, since full backup is complete
await storage.writeFile(infoFile, gzipSync(JSON.stringify(backupInfo, undefined, 2), { level: defaultLevel }))
if (options.freshBackup && options.clean) {
// Preparing a list of files to clean
ctx.info('Cleaning old backup files...')
for (const sn of backupInfo.snapshots.slice(0, backupInfo.snapshots.length - 1)) {
const filesToDelete: string[] = []
for (const [domain, dsn] of [...Object.entries(sn.domains)]) {
if (domain === DOMAIN_BLOB) {
continue
}
// eslint-disable-next-line @typescript-eslint/no-dynamic-delete
delete (sn.domains as any)[domain]
filesToDelete.push(...(dsn.snapshots ?? []))
filesToDelete.push(...(dsn.storage ?? []))
if (dsn.snapshot !== undefined) {
filesToDelete.push(dsn.snapshot)
}
}
for (const file of filesToDelete) {
ctx.info('Removing file...', { file })
await storage.delete(file)
// eslint-disable-next-line @typescript-eslint/no-dynamic-delete
delete sizeInfo[file]
}
}
ctx.info('Cleaning complete...')
await storage.writeFile(infoFile, gzipSync(JSON.stringify(backupInfo, undefined, 2), { level: defaultLevel }))
}
}
const addFileSize = async (file: string | undefined | null): Promise<void> => {
if (file != null) {
const sz = sizeInfo[file]
@@ -1588,6 +1624,8 @@ export async function backupFind (storage: BackupStorage, id: Ref<Doc>, domain?:
/**
* @public
* Restore state of DB to specified point.
*
* Recheck mean we download and compare every document on our side and if found difference upload changed version to server.
*/
export async function restore (
ctx: MeasureContext,
@@ -1689,7 +1727,7 @@ export async function restore (
try {
while (true) {
const st = Date.now()
const it = await connection.loadChunk(c, idx, opt.recheck)
const it = await connection.loadChunk(c, idx)
chunks++
idx = it.idx
@@ -1723,11 +1761,13 @@ export async function restore (
// Let's find difference
const docsToAdd = new Map(
Array.from(changeset.entries()).filter(
([it]) =>
!serverChangeset.has(it) ||
(serverChangeset.has(it) && doTrimHash(serverChangeset.get(it)) !== doTrimHash(changeset.get(it)))
)
opt.recheck === true // If recheck we check all documents.
? Array.from(changeset.entries())
: Array.from(changeset.entries()).filter(
([it]) =>
!serverChangeset.has(it) ||
(serverChangeset.has(it) && doTrimHash(serverChangeset.get(it)) !== doTrimHash(changeset.get(it)))
)
)
const docsToRemove = Array.from(serverChangeset.keys()).filter((it) => !changeset.has(it))
@@ -1738,15 +1778,12 @@ export async function restore (
async function sendChunk (doc: Doc | undefined, len: number): Promise<void> {
if (doc !== undefined) {
docsToAdd.delete(doc._id)
if (opt.recheck === true) {
// We need to clear %hash% in case our is wrong.
delete (doc as any)['%hash%']
}
docs.push(doc)
}
sendSize = sendSize + len
if (sendSize > dataUploadSize || (doc === undefined && docs.length > 0)) {
let docsToSend = docs
totalSend += docs.length
ctx.info('upload-' + c, {
docs: docs.length,
@@ -1757,32 +1794,42 @@ export async function restore (
})
// Correct docs without space
for (const d of docs) {
if (d._class === core.class.DocIndexState) {
// We need to clean old stuff from restored document.
if ('stages' in d) {
delete (d as any).stages
delete (d as any).attributes
;(d as any).needIndex = true
;(d as any)['%hash%'] = ''
}
}
if (d.space == null) {
d.space = core.space.Workspace
;(d as any)['%hash%'] = ''
}
if (TxProcessor.isExtendsCUD(d._class)) {
const tx = d as TxCUD<Doc>
if (tx.objectSpace == null) {
tx.objectSpace = core.space.Workspace
;(tx as any)['%hash%'] = ''
}
}
}
try {
await connection.upload(c, docs)
} catch (err: any) {
ctx.error('error during upload', { err, docs: JSON.stringify(docs) })
if (opt.recheck === true) {
// We need to download all documents and compare them.
const serverDocs = toIdMap(
await connection.loadDocs(
c,
docs.map((it) => it._id)
)
)
docsToSend = docs.filter((doc) => {
const serverDoc = serverDocs.get(doc._id)
if (serverDoc !== undefined) {
const { '%hash%': _h1, ...dData } = doc as any
const { '%hash%': _h2, ...sData } = serverDoc as any
return !deepEqual(dData, sData)
}
return true
})
} else {
try {
await connection.upload(c, docsToSend)
} catch (err: any) {
ctx.error('error during upload', { err, docs: JSON.stringify(docs) })
}
}
docs.length = 0
sendSize = 0
@@ -1964,7 +2011,7 @@ export async function restore (
if (opt.skip?.has(c) === true) {
continue
}
await limiter.exec(async () => {
await limiter.add(async () => {
ctx.info('processing domain', { domain: c, workspaceId: workspaceId.name })
let retry = 5
let delay = 1
+8 -4
View File
@@ -65,7 +65,8 @@ class BackupWorker {
externalStorage: StorageAdapter
) => DbConfiguration,
readonly region: string,
readonly recheck: boolean = false
readonly freshWorkspace: boolean = false,
readonly clean: boolean = false
) {}
canceled = false
@@ -186,7 +187,8 @@ class BackupWorker {
backup(ctx, '', getWorkspaceId(ws.workspace), storage, {
skipDomains: [],
force: true,
recheck: this.recheck,
freshBackup: this.freshWorkspace,
clean: this.clean,
timeout: this.config.Timeout * 1000,
connectTimeout: 5 * 60 * 1000, // 5 minutes to,
blobDownloadLimit: this.downloadLimit,
@@ -303,7 +305,8 @@ export async function doBackupWorkspace (
externalStorage: StorageAdapter
) => DbConfiguration,
region: string,
recheck: boolean,
freshWorkspace: boolean,
clean: boolean,
downloadLimit: number
): Promise<boolean> {
const backupWorker = new BackupWorker(
@@ -313,7 +316,8 @@ export async function doBackupWorkspace (
workspaceStorageAdapter,
getConfig,
region,
recheck
freshWorkspace,
clean
)
backupWorker.downloadLimit = downloadLimit
const { processed } = await backupWorker.doBackup(ctx, [workspace], Number.MAX_VALUE)
+15 -10
View File
@@ -74,9 +74,17 @@ export class StorageExtension implements Extension {
return await this.loadDocument(documentName, context)
}
async afterLoadDocument ({ documentName, document }: withContext<afterLoadDocumentPayload>): Promise<any> {
// remember the markup for the document
this.markups.set(documentName, this.configuration.transformer.fromYdoc(document))
async afterLoadDocument ({ context, documentName, document }: withContext<afterLoadDocumentPayload>): Promise<any> {
const { ctx } = this.configuration
const { connectionId } = context
try {
// remember the markup for the document
this.markups.set(documentName, this.configuration.transformer.fromYdoc(document))
} catch {
ctx.warn('document is not of a markup type', { documentName, connectionId })
this.markups.set(documentName, {})
}
}
async onStoreDocument ({ context, documentName, document }: withContext<onStoreDocumentPayload>): Promise<void> {
@@ -141,17 +149,14 @@ export class StorageExtension implements Extension {
const { ctx, adapter } = this.configuration
try {
const prevMarkup = this.markups.get(documentName) ?? {}
const currMarkup = this.configuration.transformer.fromYdoc(document)
await ctx.with('save-document', {}, (ctx) =>
const currMarkup = await ctx.with('save-document', {}, (ctx) =>
adapter.saveDocument(ctx, documentName, document, context, {
prev: prevMarkup,
curr: currMarkup
prev: () => this.markups.get(documentName) ?? {},
curr: () => this.configuration.transformer.fromYdoc(document)
})
)
this.markups.set(documentName, currMarkup)
this.markups.set(documentName, currMarkup ?? {})
} catch (err) {
ctx.error('failed to save document', { documentName, error: err })
throw new Error('Failed to save document')
+4 -4
View File
@@ -24,9 +24,9 @@ export interface CollabStorageAdapter {
documentId: string,
document: YDoc,
context: Context,
markup: {
prev: Record<string, string>
curr: Record<string, string>
getMarkup: {
prev: () => Record<string, string>
curr: () => Record<string, string>
}
) => Promise<void>
) => Promise<Record<string, string> | undefined>
}
+23 -16
View File
@@ -84,11 +84,11 @@ export class PlatformStorageAdapter implements CollabStorageAdapter {
documentName: string,
document: YDoc,
context: Context,
markup: {
prev: Record<string, string>
curr: Record<string, string>
getMarkup: {
prev: () => Record<string, string>
curr: () => Record<string, string>
}
): Promise<void> {
): Promise<Record<string, string> | undefined> {
const { clientFactory } = context
const { documentId } = decodeDocumentId(documentName)
@@ -110,8 +110,8 @@ export class PlatformStorageAdapter implements CollabStorageAdapter {
}
ctx.info('save document content to platform', { documentName })
await ctx.with('save-to-platform', {}, (ctx) => {
return this.saveDocumentToPlatform(ctx, client, documentName, markup)
return await ctx.with('save-to-platform', {}, (ctx) => {
return this.saveDocumentToPlatform(ctx, client, documentName, getMarkup)
})
} finally {
await client.close()
@@ -122,14 +122,25 @@ export class PlatformStorageAdapter implements CollabStorageAdapter {
ctx: MeasureContext,
client: Omit<TxOperations, 'close'>,
documentName: string,
markup: {
prev: Record<string, string>
curr: Record<string, string>
getMarkup: {
prev: () => Record<string, string>
curr: () => Record<string, string>
}
): Promise<void> {
): Promise<Record<string, string> | undefined> {
const { documentId, workspaceId } = decodeDocumentId(documentName)
const { objectAttr, objectClass, objectId } = documentId
const attribute = client.getHierarchy().findAttribute(objectClass, objectAttr)
if (attribute === undefined) {
ctx.warn('attribute not found', { documentName, objectClass, objectAttr })
return
}
const markup = {
prev: getMarkup.prev(),
curr: getMarkup.curr()
}
const currMarkup = markup.curr[objectAttr]
const prevMarkup = markup.prev[objectAttr]
@@ -138,12 +149,6 @@ export class PlatformStorageAdapter implements CollabStorageAdapter {
return
}
const attribute = client.getHierarchy().findAttribute(objectClass, objectAttr)
if (attribute === undefined) {
ctx.warn('attribute not found', { documentName, objectClass, objectAttr })
return
}
const current = await ctx.with('query', {}, () => {
return client.findOne(objectClass, { _id: objectId })
})
@@ -191,6 +196,8 @@ export class PlatformStorageAdapter implements CollabStorageAdapter {
data
)
})
return markup.curr
}
}
+11 -10
View File
@@ -101,6 +101,15 @@ export function initStatisticsContext (
let oldMetricsValue = ''
const serviceId = encodeURIComponent(os.hostname() + '-' + serviceName)
const handleError = (err: any): void => {
errorToSend++
if (errorToSend % 2 === 0) {
if (err.code !== 'UND_ERR_SOCKET') {
console.error(err)
}
}
}
const intTimer = setInterval(() => {
try {
if (metricsFile !== undefined || ops?.logConsole === true) {
@@ -138,18 +147,10 @@ export function initStatisticsContext (
},
body: statData
}
).catch((err) => {
errorToSend++
if (errorToSend % 2 === 0) {
console.error(err)
}
})
).catch(handleError)
}
} catch (err: any) {
errorToSend++
if (errorToSend % 20 === 0) {
console.error(err)
}
handleError(err)
}
}, METRICS_UPDATE_INTERVAL)
+5 -11
View File
@@ -9,7 +9,6 @@ import {
type StorageIterator,
type WorkspaceId
} from '@hcengineering/core'
import { estimateDocSize } from './utils'
export * from '@hcengineering/storage'
@@ -20,7 +19,7 @@ export function getBucketId (workspaceId: WorkspaceId): string {
return toWorkspaceString(workspaceId)
}
const chunkSize = 512 * 1024
const chunkSize = 200
/**
* @public
@@ -44,8 +43,7 @@ export class BackupClientOps {
loadChunk (
ctx: MeasureContext,
domain: Domain,
idx?: number,
recheck?: boolean
idx?: number
): Promise<{
idx: number
docs: DocInfo[]
@@ -64,22 +62,18 @@ export class BackupClientOps {
}
}
} else {
chunk = { idx, iterator: this.storage.find(ctx, domain, recheck), finished: false, index: 0 }
chunk = { idx, iterator: this.storage.find(ctx, domain), finished: false, index: 0 }
this.chunkInfo.set(idx, chunk)
}
let size = 0
const docs: DocInfo[] = []
while (size < chunkSize) {
while (docs.length < chunkSize) {
const _docs = await chunk.iterator.next(ctx)
if (_docs.length === 0) {
chunk.finished = true
break
}
for (const doc of _docs) {
size += estimateDocSize(doc)
docs.push(doc)
}
docs.push(..._docs)
}
return {
+6 -37
View File
@@ -12,7 +12,6 @@ import core, {
type Class,
type Client,
type Doc,
type DocInfo,
type MeasureContext,
type ModelDb,
type Ref,
@@ -21,7 +20,7 @@ import core, {
type WorkspaceIdWithUrl
} from '@hcengineering/core'
import platform, { PlatformError, Severity, Status, unknownError } from '@hcengineering/platform'
import { createHash, type Hash } from 'crypto'
import { type Hash } from 'crypto'
import fs from 'fs'
import { BackupClientOps } from './storage'
import type { Pipeline } from './types'
@@ -86,7 +85,7 @@ export function estimateDocSize (_obj: any): number {
return result
}
/**
* Return some estimation for object size
* Calculate hash for object
*/
export function updateHashForDoc (hash: Hash, _obj: any): void {
const toProcess = [_obj]
@@ -98,7 +97,9 @@ export function updateHashForDoc (hash: Hash, _obj: any): void {
if (typeof obj === 'function') {
continue
}
for (const key in obj) {
const keys = Object.keys(obj).sort()
// We need sorted list of keys to make it consistent
for (const key of keys) {
// include prototype properties
const value = obj[key]
const type = getTypeOf(value)
@@ -220,38 +221,6 @@ export function loadBrandingMap (brandingPath?: string): BrandingMap {
return brandings
}
export function toDocInfo (d: Doc, bulkUpdate: Map<Ref<Doc>, string>, recheck?: boolean): DocInfo {
let digest: string | null = (d as any)['%hash%']
if ('%hash%' in d) {
delete d['%hash%']
}
const pos = (digest ?? '').indexOf('|')
const oldDigest = digest
if (digest == null || digest === '' || recheck === true) {
const size = estimateDocSize(d)
const hash = createHash('sha256')
updateHashForDoc(hash, d)
digest = hash.digest('base64')
const newDigest = `${digest}|${size.toString(16)}`
if (recheck !== true || oldDigest !== newDigest) {
bulkUpdate.set(d._id, `${digest}|${size.toString(16)}`)
}
return {
id: d._id,
hash: digest,
size
}
} else {
return {
id: d._id,
hash: pos >= 0 ? digest.slice(0, pos) : digest,
size: pos >= 0 ? parseInt(digest.slice(pos + 1), 16) : 0
}
}
}
export function wrapPipeline (
ctx: MeasureContext,
pipeline: Pipeline,
@@ -284,7 +253,7 @@ export function wrapPipeline (
closeChunk: (idx) => backupOps.closeChunk(ctx, idx),
getHierarchy: () => pipeline.context.hierarchy,
getModel: () => pipeline.context.modelDb,
loadChunk: (domain, idx, recheck) => backupOps.loadChunk(ctx, domain, idx, recheck),
loadChunk: (domain, idx) => backupOps.loadChunk(ctx, domain, idx),
loadDocs: (domain, docs) => backupOps.loadDocs(ctx, domain, docs),
upload: (domain, docs) => backupOps.upload(ctx, domain, docs),
searchFulltext: async (query, options) => ({ docs: [], total: 0 }),
+3 -3
View File
@@ -402,7 +402,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
): Promise<{ classUpdate: Ref<Class<Doc>>[], processed: number }> {
const _classUpdate = new Set<Ref<Class<Doc>>>()
let processed = 0
await rateLimiter.exec(async () => {
await rateLimiter.add(async () => {
let st = Date.now()
let groupBy = await this.storage.groupBy(ctx, DOMAIN_DOC_INDEX_STATE, 'objectClass', { needIndex: true })
@@ -581,7 +581,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
if (docs.length === 0) {
return
}
await pushQueue.exec(async () => {
await pushQueue.add(async () => {
try {
try {
await ctx.with('push-elastic', {}, () => this.fulltextAdapter.updateMany(ctx, this.workspace, docs))
@@ -646,7 +646,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
const docState = valueIds.get(doc._id as Ref<DocIndexState>) as WithLookup<DocIndexState>
const indexedDoc = createIndexedDoc(doc, this.hierarchy.findAllMixins(doc), doc.space)
await rateLimit.exec(async () => {
await rateLimit.add(async () => {
await ctx.with('process-document', { _class: doc._class }, async (ctx) => {
try {
// Copy content attributes as well.
+2 -2
View File
@@ -43,8 +43,8 @@ export class LowLevelMiddleware extends BaseMiddleware implements Middleware {
}
const adapterManager = context.adapterManager
context.lowLevelStorage = {
find (ctx: MeasureContext, domain: Domain, recheck?: boolean): StorageIterator {
return adapterManager.getAdapter(domain, false).find(ctx, domain, recheck)
find (ctx: MeasureContext, domain: Domain): StorageIterator {
return adapterManager.getAdapter(domain, false).find(ctx, domain)
},
load (ctx: MeasureContext, domain: Domain, docs: Ref<Doc>[]): Promise<Doc[]> {
+31 -30
View File
@@ -118,30 +118,33 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
}
if (this.wasInit === false) {
this.wasInit = (async () => {
const spaces: SpaceWithMembers[] =
(await this.next?.findAll(
ctx,
core.class.Space,
{},
{
projection: {
private: 1,
_class: 1,
_id: 1,
members: 1
await ctx.with('init-space-security', {}, async (ctx) => {
ctx.contextData = undefined
const spaces: SpaceWithMembers[] =
(await this.next?.findAll(
ctx,
core.class.Space,
{},
{
projection: {
private: 1,
_class: 1,
_id: 1,
members: 1
}
}
)) ?? []
this.spacesMap.clear()
this.publicSpaces.clear()
this.systemSpaces.clear()
for (const space of spaces) {
if (space._class === core.class.SystemSpace) {
this.systemSpaces.add(space._id)
} else {
this.addSpace(space)
}
)) ?? []
this.spacesMap.clear()
this.publicSpaces.clear()
this.systemSpaces.clear()
for (const space of spaces) {
if (space._class === core.class.SystemSpace) {
this.systemSpaces.add(space._id)
} else {
this.addSpace(space)
}
}
})
})()
}
if (this.wasInit instanceof Promise) {
@@ -559,7 +562,7 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
if (options?.lookup !== undefined) {
for (const object of findResult) {
if (object.$lookup !== undefined) {
await this.filterLookup(ctx, object.$lookup)
this.filterLookup(ctx, object.$lookup)
}
}
}
@@ -600,25 +603,23 @@ export class SpaceSecurityMiddleware extends BaseMiddleware implements Middlewar
return result
}
async isUnavailable (ctx: MeasureContext<SessionData>, space: Ref<Space>): Promise<boolean> {
filterLookup<T extends Doc>(ctx: MeasureContext, lookup: LookupData<T>): void {
if (Object.keys(lookup).length === 0) return
const account = ctx.contextData.account
if (isSystem(account, ctx)) return false
return !this.getAllAllowedSpaces(account, true).includes(space)
}
async filterLookup<T extends Doc>(ctx: MeasureContext, lookup: LookupData<T>): Promise<void> {
if (isSystem(account, ctx)) return
const allowedSpaces = this.getAllAllowedSpaces(account, true)
for (const key in lookup) {
const val = lookup[key]
if (Array.isArray(val)) {
const arr: AttachedDoc[] = []
for (const value of val) {
if (!(await this.isUnavailable(ctx, value.space))) {
if (allowedSpaces.includes(value.space)) {
arr.push(value)
}
}
lookup[key] = arr as any
} else if (val !== undefined) {
if (await this.isUnavailable(ctx, val.space)) {
if (!allowedSpaces.includes(val.space)) {
lookup[key] = undefined
}
}
+57 -65
View File
@@ -62,8 +62,6 @@ import core, {
type WorkspaceId
} from '@hcengineering/core'
import {
estimateDocSize,
toDocInfo,
type DbAdapter,
type DbAdapterHandler,
type DomainHelperOperations,
@@ -86,8 +84,8 @@ import {
} from 'mongodb'
import { DBCollectionHelper, getMongoClient, getWorkspaceMongoDB, type MongoClientReference } from './utils'
function translateDoc (doc: Doc): Doc {
return { ...doc, '%hash%': null } as any
function translateDoc (doc: Doc, hash: string): Doc {
return { ...doc, '%hash%': hash } as any
}
function isLookupQuery<T extends Doc> (query: DocumentQuery<T>): boolean {
@@ -253,14 +251,14 @@ abstract class MongoAdapterBase implements DbAdapter {
operations: DocumentUpdate<T>
): Promise<void> {
if (isOperator(operations)) {
await this.db.collection(domain).updateMany(this.translateRawQuery(query), { $set: { '%hash%': null } })
await this.db.collection(domain).updateMany(this.translateRawQuery(query), { $set: { '%hash%': this.curHash() } })
await this.db
.collection(domain)
.updateMany(this.translateRawQuery(query), { ...operations } as unknown as UpdateFilter<Document>)
} else {
await this.db
.collection(domain)
.updateMany(this.translateRawQuery(query), { $set: { ...operations, '%hash%': null } })
.updateMany(this.translateRawQuery(query), { $set: { ...operations, '%hash%': this.curHash() } })
}
}
@@ -1023,66 +1021,57 @@ abstract class MongoAdapterBase implements DbAdapter {
return docs
}
find (_ctx: MeasureContext, domain: Domain, recheck?: boolean): StorageIterator {
curHash (): string {
return Date.now().toString(16) // Current hash value
}
strimSize (str: string): string {
const pos = str.indexOf('|')
if (pos > 0) {
return str.substring(0, pos)
}
return str
}
find (_ctx: MeasureContext, domain: Domain): StorageIterator {
const ctx = _ctx.newChild('find', { domain })
const coll = this.db.collection<Doc>(domain)
let mode: 'hashed' | 'non-hashed' = 'hashed'
let iterator: FindCursor<Doc>
const bulkUpdate = new Map<Ref<Doc>, string>()
const flush = async (flush = false): Promise<void> => {
if (bulkUpdate.size > 1000 || flush) {
if (bulkUpdate.size > 0) {
await ctx.with('bulk-write-find', {}, () =>
coll.bulkWrite(
Array.from(bulkUpdate.entries()).map((it) => ({
updateOne: {
filter: { _id: it[0], '%hash%': null },
update: { $set: { '%hash%': it[1] } }
}
})),
{ ordered: false }
)
)
}
bulkUpdate.clear()
}
}
return {
next: async () => {
if (iterator === undefined) {
await coll.updateMany({ '%hash%': { $in: [null, ''] } }, { $set: { '%hash%': this.curHash() } })
iterator = coll.find(
recheck === true ? {} : { '%hash%': { $nin: ['', null] } },
recheck === true
? {}
: {
projection: {
'%hash%': 1,
_id: 1
}
}
{},
{
projection: {
'%hash%': 1,
_id: 1
}
}
)
}
let d = await ctx.with('next', { mode }, () => iterator.next())
if (d == null && mode === 'hashed' && recheck !== true) {
mode = 'non-hashed'
await iterator.close()
await flush(true) // We need to flush, so wrong id documents will be updated.
iterator = coll.find({ '%hash%': { $in: ['', null] } })
d = await ctx.with('next', { mode }, () => iterator.next())
}
const d = await ctx.with('next', {}, () => iterator.next())
const result: DocInfo[] = []
if (d != null) {
result.push(toDocInfo(d, bulkUpdate, recheck))
result.push({
id: d._id,
hash: this.strimSize((d as any)['%hash%'])
})
}
if (iterator.bufferedCount() > 0) {
result.push(...iterator.readBufferedDocuments().map((it) => toDocInfo(it, bulkUpdate, recheck)))
result.push(
...iterator.readBufferedDocuments().map((it) => ({
id: it._id,
hash: this.strimSize((it as any)['%hash%'])
}))
)
}
await ctx.with('flush', {}, () => flush())
return result
},
close: async () => {
await ctx.with('flush', {}, () => flush(true))
await ctx.with('close', {}, () => iterator.close())
ctx.end()
}
@@ -1104,7 +1093,7 @@ abstract class MongoAdapterBase implements DbAdapter {
return ctx.with('upload', { domain }, (ctx) => {
const coll = this.collection(domain)
return uploadDocuments(ctx, docs, coll)
return uploadDocuments(ctx, docs, coll, this.curHash())
})
}
@@ -1137,7 +1126,7 @@ abstract class MongoAdapterBase implements DbAdapter {
updateOne: {
filter: { _id: it[0] },
update: {
$set: { ...set, '%hash%': null },
$set: { ...set, '%hash%': this.curHash() },
...($unset !== undefined ? { $unset } : {})
}
}
@@ -1419,7 +1408,7 @@ class MongoAdapter extends MongoAdapterBase {
protected txCreateDoc (bulk: OperationBulk, tx: TxCreateDoc<Doc>): void {
const doc = TxProcessor.createDoc2Doc(tx)
bulk.add.push(translateDoc(doc))
bulk.add.push(translateDoc(doc, this.curHash()))
}
protected txUpdateDoc (bulk: OperationBulk, tx: TxUpdateDoc<Doc>): void {
@@ -1439,7 +1428,7 @@ class MongoAdapter extends MongoAdapterBase {
update: {
$set: {
...Object.fromEntries(Object.entries(desc.$update).map((it) => [arr + '.$.' + it[0], it[1]])),
'%hash%': null
'%hash%': this.curHash()
}
}
}
@@ -1451,7 +1440,7 @@ class MongoAdapter extends MongoAdapterBase {
$set: {
modifiedBy: tx.modifiedBy,
modifiedOn: tx.modifiedOn,
'%hash%': null
'%hash%': this.curHash()
}
}
}
@@ -1470,7 +1459,7 @@ class MongoAdapter extends MongoAdapterBase {
$set: {
modifiedBy: tx.modifiedBy,
modifiedOn: tx.modifiedOn,
'%hash%': null
'%hash%': this.curHash()
}
} as unknown as UpdateFilter<Document>,
{ returnDocument: 'after', includeResultMetadata: true }
@@ -1488,7 +1477,7 @@ class MongoAdapter extends MongoAdapterBase {
$set: {
modifiedBy: tx.modifiedBy,
modifiedOn: tx.modifiedOn,
'%hash%': null
'%hash%': this.curHash()
}
}
}
@@ -1506,7 +1495,7 @@ class MongoAdapter extends MongoAdapterBase {
...tx.operations,
modifiedBy: tx.modifiedBy,
modifiedOn: tx.modifiedOn,
'%hash%': null
'%hash%': this.curHash()
})) {
;(upd as any)[k] = v
}
@@ -1552,8 +1541,9 @@ class MongoTxAdapter extends MongoAdapterBase implements TxAdapter {
{ domain: 'tx' },
async () => {
try {
const hash = this.curHash()
await this.txCollection().insertMany(
baseTxes.map((it) => translateDoc(it)),
baseTxes.map((it) => translateDoc(it, hash)),
{
ordered: false
}
@@ -1582,8 +1572,9 @@ class MongoTxAdapter extends MongoAdapterBase implements TxAdapter {
{ domain: DOMAIN_MODEL_TX },
async () => {
try {
const hash = this.curHash()
await this.db.collection<Doc>(DOMAIN_MODEL_TX).insertMany(
modelTxes.map((it) => translateDoc(it)),
modelTxes.map((it) => translateDoc(it, hash)),
{
ordered: false
}
@@ -1637,22 +1628,23 @@ class MongoTxAdapter extends MongoAdapterBase implements TxAdapter {
}
}
export async function uploadDocuments (ctx: MeasureContext, docs: Doc[], coll: Collection<Document>): Promise<void> {
export async function uploadDocuments (
ctx: MeasureContext,
docs: Doc[],
coll: Collection<Document>,
curHash: string
): Promise<void> {
const ops = Array.from(docs)
while (ops.length > 0) {
const part = ops.splice(0, 500)
await coll.bulkWrite(
part.map((it) => {
const digest: string | null = (it as any)['%hash%']
if ('%hash%' in it) {
delete it['%hash%']
}
const size = digest != null ? estimateDocSize(it) : 0
const digest: string = (it as any)['%hash%'] ?? curHash
return {
replaceOne: {
filter: { _id: it._id },
replacement: { ...it, '%hash%': digest == null ? null : `${digest}|${size.toString(16)}` },
replacement: { ...it, '%hash%': digest },
upsert: true
}
}
+1 -1
View File
@@ -31,7 +31,7 @@ BEGIN
EXECUTE format('
ALTER TABLE %I ADD COLUMN "%%hash%%" text;', tbl_name);
EXECUTE format('
UPDATE %I SET "%%hash%%" = data->>''%%data%%'';', tbl_name);
UPDATE %I SET "%%hash%%" = data->>''%%hash%%'';', tbl_name);
END IF;
END LOOP;
END $$;
+111 -164
View File
@@ -59,9 +59,7 @@ import {
type DbAdapter,
type DbAdapterHandler,
type DomainHelperOperations,
estimateDocSize,
type ServerFindOptions,
toDocInfo,
type TxAdapter
} from '@hcengineering/server-core'
import type postgres from 'postgres'
@@ -96,7 +94,7 @@ async function * createCursorGenerator (
client: postgres.ReservedSql,
sql: string,
schema: Schema,
bulkSize = 50
bulkSize = 1000
): AsyncGenerator<Doc[]> {
const cursor = client.unsafe(sql).cursor(bulkSize)
try {
@@ -447,8 +445,8 @@ abstract class PostgresAdapterBase implements DbAdapter {
;(operations as any) = { ...(operations as any).$set }
}
const isOps = isOperator(operations)
if ((operations as any)['%hash%'] === undefined) {
;(operations as any)['%hash%'] = null
if ((operations as any)['%hash%'] == null) {
;(operations as any)['%hash%'] = this.curHash()
}
const schemaFields = getSchemaAndFields(domain)
if (isOps) {
@@ -459,7 +457,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
if (doc === undefined) continue
const prevAttachedTo = (doc as any).attachedTo
TxProcessor.applyUpdate(doc, operations)
;(doc as any)['%hash%'] = null
;(doc as any)['%hash%'] = this.curHash()
const converted = convertDoc(domain, doc, this.workspaceId.name, schemaFields)
const params: any[] = [doc._id, this.workspaceId.name]
let paramsIndex = params.length + 1
@@ -542,73 +540,78 @@ abstract class PostgresAdapterBase implements DbAdapter {
query: DocumentQuery<T>,
options?: ServerFindOptions<T>
): Promise<FindResult<T>> {
return ctx.with('findAll', { _class }, async () => {
try {
const domain = translateDomain(options?.domain ?? this.hierarchy.getDomain(_class))
const sqlChunks: string[] = []
const joins = this.buildJoin(_class, options?.lookup)
if (options?.domainLookup !== undefined) {
const baseDomain = translateDomain(this.hierarchy.getDomain(_class))
let fquery = ''
return ctx.with(
'findAll',
{},
async () => {
try {
const domain = translateDomain(options?.domain ?? this.hierarchy.getDomain(_class))
const sqlChunks: string[] = []
const joins = this.buildJoin(_class, options?.lookup)
if (options?.domainLookup !== undefined) {
const baseDomain = translateDomain(this.hierarchy.getDomain(_class))
const domain = translateDomain(options.domainLookup.domain)
const key = options.domainLookup.field
const as = `dl_lookup_${domain}_${key}`
joins.push({
isReverse: false,
table: domain,
path: options.domainLookup.field,
toAlias: as,
toField: '_id',
fromField: key,
fromAlias: baseDomain,
toClass: undefined
})
const domain = translateDomain(options.domainLookup.domain)
const key = options.domainLookup.field
const as = `dl_lookup_${domain}_${key}`
joins.push({
isReverse: false,
table: domain,
path: options.domainLookup.field,
toAlias: as,
toField: '_id',
fromField: key,
fromAlias: baseDomain,
toClass: undefined
})
}
const select = `SELECT ${this.getProjection(domain, options?.projection, joins)} FROM ${domain}`
const secJoin = this.addSecurity(query, domain, ctx.contextData)
if (secJoin !== undefined) {
sqlChunks.push(secJoin)
}
if (joins.length > 0) {
sqlChunks.push(this.buildJoinString(joins))
}
sqlChunks.push(`WHERE ${this.buildQuery(_class, domain, query, joins, options)}`)
return (await this.mgr.read(ctx.id, async (connection) => {
let total = options?.total === true ? 0 : -1
if (options?.total === true) {
const totalReq = `SELECT COUNT(${domain}._id) as count FROM ${domain}`
const totalSql = [totalReq, ...sqlChunks].join(' ')
const totalResult = await connection.unsafe(totalSql)
const parsed = Number.parseInt(totalResult[0].count)
total = Number.isNaN(parsed) ? 0 : parsed
}
if (options?.sort !== undefined) {
sqlChunks.push(this.buildOrder(_class, domain, options.sort, joins))
}
if (options?.limit !== undefined) {
sqlChunks.push(`LIMIT ${options.limit}`)
}
const finalSql: string = [select, ...sqlChunks].join(' ')
fquery = finalSql
const result = await connection.unsafe(finalSql)
if (options?.lookup === undefined && options?.domainLookup === undefined) {
return toFindResult(
result.map((p) => parseDocWithProjection(p as any, domain, options?.projection)),
total
)
} else {
const res = this.parseLookup<T>(result, joins, options?.projection, domain)
return toFindResult(res, total)
}
})) as FindResult<T>
} catch (err) {
ctx.error('Error in findAll', { err })
throw err
}
const select = `SELECT ${this.getProjection(domain, options?.projection, joins)} FROM ${domain}`
const secJoin = this.addSecurity(query, domain, ctx.contextData)
if (secJoin !== undefined) {
sqlChunks.push(secJoin)
}
if (joins.length > 0) {
sqlChunks.push(this.buildJoinString(joins))
}
sqlChunks.push(`WHERE ${this.buildQuery(_class, domain, query, joins, options)}`)
const findId = ctx.id ?? generateId()
return (await this.mgr.read(findId, async (connection) => {
let total = options?.total === true ? 0 : -1
if (options?.total === true) {
const totalReq = `SELECT COUNT(${domain}._id) as count FROM ${domain}`
const totalSql = [totalReq, ...sqlChunks].join(' ')
const totalResult = await connection.unsafe(totalSql)
const parsed = Number.parseInt(totalResult[0].count)
total = Number.isNaN(parsed) ? 0 : parsed
}
if (options?.sort !== undefined) {
sqlChunks.push(this.buildOrder(_class, domain, options.sort, joins))
}
if (options?.limit !== undefined) {
sqlChunks.push(`LIMIT ${options.limit}`)
}
const finalSql: string = [select, ...sqlChunks].join(' ')
const result = await connection.unsafe(finalSql)
if (options?.lookup === undefined && options?.domainLookup === undefined) {
return toFindResult(
result.map((p) => parseDocWithProjection(p as any, domain, options?.projection)),
total
)
} else {
const res = this.parseLookup<T>(result, joins, options?.projection, domain)
return toFindResult(res, total)
}
})) as FindResult<T>
} catch (err) {
ctx.error('Error in findAll', { err })
throw err
}
})
},
() => ({ fquery })
)
}
addSecurity<T extends Doc>(query: DocumentQuery<T>, domain: string, sessionContext: SessionData): string | undefined {
@@ -1231,6 +1234,10 @@ abstract class PostgresAdapterBase implements DbAdapter {
return res
}
curHash (): string {
return Date.now().toString(16) // Current hash value
}
private getProjection<T extends Doc>(
baseDomain: string,
projection: Projection<T> | undefined,
@@ -1269,60 +1276,31 @@ abstract class PostgresAdapterBase implements DbAdapter {
return []
}
find (_ctx: MeasureContext, domain: Domain, recheck?: boolean): StorageIterator {
strimSize (str: string): string {
const pos = str.indexOf('|')
if (pos > 0) {
return str.substring(0, pos)
}
return str
}
find (_ctx: MeasureContext, domain: Domain): StorageIterator {
const ctx = _ctx.newChild('find', { domain })
let initialized: boolean = false
let client: postgres.ReservedSql
let mode: 'hashed' | 'non_hashed' = 'hashed'
const bulkUpdate = new Map<Ref<Doc>, string>()
const tdomain = translateDomain(domain)
const schema = getSchema(domain)
const findId = generateId()
const flush = async (flush = false): Promise<void> => {
if (bulkUpdate.size > 1000 || flush) {
if (bulkUpdate.size > 0) {
const entries = Array.from(bulkUpdate.entries())
bulkUpdate.clear()
try {
while (entries.length > 0) {
const part = entries.splice(0, 200)
const data: string[] = part.flat()
const indexes = part.map((val, idx) => `($${2 * idx + 1}::text, $${2 * idx + 2}::text)`).join(', ')
await ctx.with('bulk-write-find', {}, () => {
return this.mgr.write(
findId,
async (client) =>
await client.unsafe(
`
UPDATE ${tdomain} SET "%hash%" = update_data.hash
FROM (values ${indexes}) AS update_data(_id, hash)
WHERE ${tdomain}."workspaceId" = '${this.workspaceId.name}' AND ${tdomain}."_id" = update_data._id
`,
data
)
)
})
}
} catch (err: any) {
ctx.error('failed to update hash', { err })
}
}
}
}
const workspaceId = this.workspaceId
function createBulk (projection: string, query: string, limit = 50): AsyncGenerator<Doc[]> {
const sql = `SELECT ${projection} FROM ${tdomain} WHERE "workspaceId" = '${workspaceId.name}' AND ${query}`
function createBulk (projection: string, limit = 50000): AsyncGenerator<Doc[]> {
const sql = `SELECT ${projection} FROM ${tdomain} WHERE "workspaceId" = '${workspaceId.name}'`
return createCursorGenerator(client, sql, schema, limit)
}
let bulk: AsyncGenerator<Doc[]>
let forcedRecheck = false
return {
next: async () => {
@@ -1331,60 +1309,29 @@ abstract class PostgresAdapterBase implements DbAdapter {
client = await this.client.reserve()
}
if (recheck === true) {
await this.mgr.write(
findId,
async (client) =>
await client`UPDATE ${client(tdomain)} SET "%hash%" = NULL WHERE "workspaceId" = ${this.workspaceId.name} AND "%hash%" IS NOT NULL`
)
}
// We need update hash to be set properly
await client.unsafe(
`UPDATE ${tdomain} SET "%hash%" = '${this.curHash()}' WHERE "workspaceId" = '${this.workspaceId.name}' AND "%hash%" IS NULL OR "%hash%" = ''`
)
initialized = true
await flush(true) // We need to flush, so wrong id documents will be updated.
bulk = createBulk('_id, "%hash%"', '"%hash%" IS NOT NULL AND "%hash%" <> \'\'')
bulk = createBulk('_id, "%hash%"')
}
let docs = await ctx.with('next', { mode }, () => bulk.next())
if (!forcedRecheck && docs.done !== true && docs.value?.length > 0) {
// Check if we have wrong hash stored, and update all of them.
forcedRecheck = true
for (const d of docs.value) {
const digest: string | null = (d as any)['%hash%']
const pos = (digest ?? '').indexOf('|')
if (pos === -1) {
await bulk.return([]) // We need to close generator
docs = { done: true, value: undefined }
await this.mgr.write(
findId,
async (client) =>
await client`UPDATE ${client(tdomain)} SET "%hash%" = NULL WHERE "workspaceId" = ${this.workspaceId.name} AND "%hash%" IS NOT NULL`
)
break
}
}
}
if ((docs.done === true || docs.value.length === 0) && mode === 'hashed') {
forcedRecheck = true
mode = 'non_hashed'
bulk = createBulk('*', '"%hash%" IS NULL OR "%hash%" = \'\'')
docs = await ctx.with('next', { mode }, () => bulk.next())
}
const docs = await ctx.with('next', {}, () => bulk.next())
if (docs.done === true || docs.value.length === 0) {
return []
}
const result: DocInfo[] = []
for (const d of docs.value) {
result.push(toDocInfo(d, bulkUpdate))
result.push({
id: d._id,
hash: this.strimSize((d as any)['%hash%'])
})
}
await ctx.with('flush', {}, () => flush())
return result
},
close: async () => {
await ctx.with('flush', {}, () => flush(true))
await bulk.return([]) // We need to close generator, just in case
client?.release()
ctx.end()
@@ -1431,12 +1378,9 @@ abstract class PostgresAdapterBase implements DbAdapter {
const doc = part[i]
const variables: string[] = []
const digest: string | null = (doc as any)['%hash%']
if ('%hash%' in doc) {
delete doc['%hash%']
if (!('%hash%' in doc) || doc['%hash%'] === '' || doc['%hash%'] == null) {
;(doc as any)['%hash%'] = this.curHash() // We need to set current hash
}
const size = digest != null ? estimateDocSize(doc) : 0
;(doc as any)['%hash%'] = digest == null ? null : `${digest}|${size.toString(16)}`
const d = convertDoc(domain, doc, this.workspaceId.name, schemaFields)
values.push(d.workspaceId)
@@ -1472,7 +1416,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
const tdomain = translateDomain(domain)
const toClean = [...docs]
while (toClean.length > 0) {
const part = toClean.splice(0, 200)
const part = toClean.splice(0, 2500)
await ctx.with('clean', {}, () => {
return this.mgr.write(
ctx.id,
@@ -1492,7 +1436,7 @@ abstract class PostgresAdapterBase implements DbAdapter {
const key = isDataField(domain, field) ? `data ->> '${field}'` : `"${field}"`
return ctx.with('groupBy', { domain }, async (ctx) => {
try {
return await this.mgr.read(ctx.id ?? generateId(), async (connection) => {
return await this.mgr.read(ctx.id, async (connection) => {
const result = await connection.unsafe(
`SELECT DISTINCT ${key} as ${field}, Count(*) AS count FROM ${translateDomain(domain)} WHERE ${this.buildRawQuery(domain, query ?? {})} GROUP BY ${key}`
)
@@ -1519,8 +1463,8 @@ abstract class PostgresAdapterBase implements DbAdapter {
const doc = map.get(_id)
if (doc === undefined) continue
const op = { ...ops }
if ((op as any)['%hash%'] === undefined) {
;(op as any)['%hash%'] = null
if ((op as any)['%hash%'] == null) {
;(op as any)['%hash%'] = this.curHash()
}
TxProcessor.applyUpdate(doc, op)
const converted = convertDoc(domain, doc, this.workspaceId.name, schemaFields)
@@ -1560,6 +1504,9 @@ abstract class PostgresAdapterBase implements DbAdapter {
const values: DBDoc[] = []
for (let i = 0; i < part.length; i++) {
const doc = part[i]
if ((doc as any)['%hash%'] == null) {
;(doc as any)['%hash%'] = this.curHash()
}
const d = convertDoc(domain, doc, this.workspaceId.name, schemaFields)
values.push(d)
}
@@ -1622,7 +1569,7 @@ class PostgresAdapter extends PostgresAdapterBase {
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
;(doc as any)['%hash%'] = this.curHash()
const domain = this.hierarchy.getDomain(tx.objectClass)
const converted = convertDoc(domain, doc, this.workspaceId.name, schemaFields)
const { extractedFields } = parseUpdate(tx.attributes as Partial<Doc>, schemaFields)
@@ -1714,7 +1661,7 @@ class PostgresAdapter extends PostgresAdapterBase {
for (const tx of withOperator ?? []) {
let doc: Doc | undefined
const ops: any = { '%hash%': null, ...tx.operations }
const ops: any = { '%hash%': this.curHash(), ...tx.operations }
result.push(
await ctx.with('tx-update-doc', { _class: tx.objectClass }, async (ctx) => {
await this.mgr.write(ctx.id, async (client) => {
@@ -1723,7 +1670,7 @@ class PostgresAdapter extends PostgresAdapterBase {
ops.modifiedBy = tx.modifiedBy
ops.modifiedOn = tx.modifiedOn
TxProcessor.applyUpdate(doc, ops)
;(doc as any)['%hash%'] = null
;(doc as any)['%hash%'] = this.curHash()
const converted = convertDoc(domain, doc, this.workspaceId.name, schemaFields)
const columns: string[] = []
const { extractedFields, remainingData } = parseUpdate(ops, schemaFields)
+3 -3
View File
@@ -304,9 +304,9 @@ export function convertDoc<T extends Doc> (
// Check if some fields are missing
for (const [key, _type] of Object.entries(schemaAndFields.schema)) {
if (!(key in doc)) {
// We missing required field, and we need to add a dummy value for it.
if (_type.notNull) {
if (_type.notNull) {
if (!(key in doc) || (doc as any)[key] == null) {
// We missing required field, and we need to add a dummy value for it.
// Null value is not allowed
switch (_type.type) {
case 'bigint':
+1 -1
View File
@@ -94,7 +94,7 @@ class StorageBlobAdapter implements DbAdapter {
async close (): Promise<void> {}
find (ctx: MeasureContext, domain: Domain, recheck?: boolean): StorageIterator {
find (ctx: MeasureContext, domain: Domain): StorageIterator {
return this.client.find(ctx, this.workspaceId)
}
+3 -3
View File
@@ -263,10 +263,10 @@ export class ClientSession implements Session {
return this.ops
}
async loadChunk (ctx: ClientSessionCtx, domain: Domain, idx?: number, recheck?: boolean): Promise<void> {
async loadChunk (ctx: ClientSessionCtx, domain: Domain, idx?: number): Promise<void> {
this.lastRequest = Date.now()
try {
const result = await this.getOps().loadChunk(ctx.ctx, domain, idx, recheck)
const result = await this.getOps().loadChunk(ctx.ctx, domain, idx)
await ctx.sendResponse(result)
} catch (err: any) {
await ctx.sendError('Failed to upload', unknownError(err))
@@ -326,7 +326,7 @@ export class ClientSession implements Session {
* @public
*/
export interface BackupSession extends Session {
loadChunk: (ctx: ClientSessionCtx, domain: Domain, idx?: number, recheck?: boolean) => Promise<void>
loadChunk: (ctx: ClientSessionCtx, domain: Domain, idx?: number) => Promise<void>
closeChunk: (ctx: ClientSessionCtx, idx: number) => Promise<void>
loadDocs: (ctx: ClientSessionCtx, domain: Domain, docs: Ref<Doc>[]) => Promise<void>
}