diff --git a/.vscode/launch.json b/.vscode/launch.json index 120bfe5e16..0cb4fdfd4c 100644 --- a/.vscode/launch.json +++ b/.vscode/launch.json @@ -95,6 +95,7 @@ "ACCOUNT_PORT": "3000", "FRONT_URL": "http://localhost:8080", "SES_URL": "", + // "DB_NS": "account-2", // "WS_LIVENESS_DAYS": "1", "MINIO_ACCESS_KEY": "minioadmin", "MINIO_SECRET_KEY": "minioadmin", @@ -381,9 +382,7 @@ "PORT": "4005", "SECRET": "secret", "MONGO_URL": "mongodb://localhost:27017", - "MINIO_ENDPOINT": "localhost", - "MINIO_ACCESS_KEY": "minioadmin", - "MINIO_SECRET_KEY": "minioadmin" + "STORAGE_CONFIG": "minio|localhost?accessKey=minioadmin&secretKey=minioadmin" }, "runtimeArgs": ["--nolazy", "-r", "ts-node/register"], "runtimeVersion": "20", diff --git a/common/config/rush/pnpm-lock.yaml b/common/config/rush/pnpm-lock.yaml index aba68922b4..f07ee37bdb 100644 --- a/common/config/rush/pnpm-lock.yaml +++ b/common/config/rush/pnpm-lock.yaml @@ -155,6 +155,9 @@ dependencies: '@rush-temp/cloud-branding': specifier: file:./projects/cloud-branding.tgz version: file:projects/cloud-branding.tgz(@types/node@20.11.19)(bufferutil@4.0.8)(esbuild@0.20.1)(ts-node@10.9.2)(utf-8-validate@6.0.4) + '@rush-temp/cloud-datalake': + specifier: file:./projects/cloud-datalake.tgz + version: file:projects/cloud-datalake.tgz(@types/node@20.11.19)(bufferutil@4.0.8)(esbuild@0.20.1)(ts-node@10.9.2)(utf-8-validate@6.0.4) '@rush-temp/collaboration': specifier: file:./projects/collaboration.tgz version: file:projects/collaboration.tgz(esbuild@0.20.1)(ts-node@10.9.2) @@ -1394,6 +1397,9 @@ dependencies: aws-sdk: specifier: ^2.1423.0 version: 2.1664.0 + aws4fetch: + specifier: ^1.0.20 + version: 1.0.20 base64-js: specifier: ^1.5.1 version: 1.5.1 @@ -1477,7 +1483,7 @@ dependencies: version: 22.8.8 electron: specifier: ^32.1.1 - version: 32.1.1 + version: 32.2.1 electron-builder: specifier: ^25.0.5 version: 25.0.5 @@ -1748,6 +1754,9 @@ dependencies: postcss-loader: specifier: ^7.0.2 version: 7.3.4(postcss@8.4.35)(typescript@5.3.3)(webpack@5.90.3) + postgres: + specifier: ^3.4.4 + version: 3.4.4 posthog-js: specifier: ~1.122.0 version: 1.122.0 @@ -11374,6 +11383,10 @@ packages: xml2js: 0.6.2 dev: false + /aws4fetch@1.0.20: + resolution: {integrity: sha512-/djoAN709iY65ETD6LKCtyyEI04XIBP5xVvfmNxsEP0uJB5tyaGBztSryRr4HqMStr9R06PisQE7m9zDTXKu6g==} + dev: false + /axobject-query@4.0.0: resolution: {integrity: sha512-+60uv1hiVFhHZeO+Lz0RYzsVHy5Wr1ayX0mwda9KPDVLNJgZ1T9Ny7VmFbLDzxsH0D87I86vgj3gFrjTJUYznw==} dependencies: @@ -11548,6 +11561,13 @@ packages: dev: false optional: true + /base32-encode@2.0.0: + resolution: {integrity: sha512-mlmkfc2WqdDtMl/id4qm3A7RjW6jxcbAoMjdRmsPiwQP0ufD4oXItYMnPgVHe80lnAIy+1xwzhHE1s4FoIceSw==} + engines: {node: ^12.20.0 || ^14.13.1 || >=16.0.0} + dependencies: + to-data-view: 2.0.0 + dev: false + /base64-js@1.5.1: resolution: {integrity: sha512-AKpaYlHn8t4SVbOHCy+b5+KKgvR4vrsD8vbvrbiQJps7fKDTkjkDry6ji0rUJjC0kzbNePLwzxq8iypo41qeWA==} dev: false @@ -13734,8 +13754,8 @@ packages: resolution: {integrity: sha512-hWFbUk9u3fQHcKzTAcjZAN7XH9bL9oH9g20RRDU/DVDNqdMI03GzlBZfR/R8R1krYu9AT4biLqSCAxnt9LMAfA==} dev: false - /electron@32.1.1: - resolution: {integrity: sha512-NlWvG6kXOJbZbELmzP3oV7u50I3NHYbCeh+AkUQ9vGyP7b74cFMx9HdTzejODeztW1jhr3SjIBbUZzZ45zflfQ==} + /electron@32.2.1: + resolution: {integrity: sha512-GCPI/5hU34pPcNltNpz+uylhhuTm9BM0N8RmrbVgaWBodLSmmcCkvpgN0BseKhO6IwQOPzWaovrcZ/nPIpfGaQ==} engines: {node: '>= 12.20.55'} hasBin: true requiresBuild: true @@ -20790,6 +20810,11 @@ packages: resolution: {integrity: sha512-i/hbxIE9803Alj/6ytL7UHQxRvZkI9O4Sy+J3HGc4F4oo/2eQAjTSNJ0bfxyse3bH0nuVesCk+3IRLaMtG3H6w==} dev: false + /postgres@3.4.4: + resolution: {integrity: sha512-IbyN+9KslkqcXa8AO9fxpk97PA4pzewvpi2B3Dwy9u4zpV32QicaEdgmF3eSQUzdRk7ttDHQejNgAEr4XoeH4A==} + engines: {node: '>=12'} + dev: false + /posthog-js@1.122.0: resolution: {integrity: sha512-+8R2/nLaWyI5Jp2Ly7L52qcgDFU3xryyoNG52DPJ8dlGnagphxIc0mLNGurgyKeeTGycsOsuOIP4dtofv3ZoBA==} deprecated: This version of posthog-js is deprecated, please update posthog-js, and do not use this version! Check out our JS docs at https://posthog.com/docs/libraries/js @@ -23505,6 +23530,11 @@ packages: resolution: {integrity: sha512-3f0uOEAQwIqGuWW2MVzYg8fV/QNnc/IpuJNG837rLuczAaLVHslWHZQj4IGiEl5Hs3kkbhwL9Ab7Hrsmuj+Smw==} dev: false + /to-data-view@2.0.0: + resolution: {integrity: sha512-RGEM5KqlPHr+WVTPmGNAXNeFEmsBnlkxXaIfEpUYV0AST2Z5W1EGq9L/MENFrMMmL2WQr1wjkmZy/M92eKhjYA==} + engines: {node: ^12.20.0 || ^14.13.1 || >=16.0.0} + dev: false + /to-fast-properties@2.0.0: resolution: {integrity: sha512-/OaKK0xYrs3DmxRYqL/yDc+FxFUVYhDlXMhRmv3z915w2HF1tnN1omB354j8VUGO/hbRzyD6Y3sA7v7GS/ceog==} engines: {node: '>=4'} @@ -26627,6 +26657,44 @@ packages: - utf-8-validate dev: false + file:projects/cloud-datalake.tgz(@types/node@20.11.19)(bufferutil@4.0.8)(esbuild@0.20.1)(ts-node@10.9.2)(utf-8-validate@6.0.4): + resolution: {integrity: sha512-AA2lTsmPKPeYA1MTwIscZFRO40m9Ctc59Er2x8VRLNBBt4mQ01b1CCay4VFVPWYxAzh+Ru9RoUIB7lS+m8sj9Q==, tarball: file:projects/cloud-datalake.tgz} + id: file:projects/cloud-datalake.tgz + name: '@rush-temp/cloud-datalake' + version: 0.0.0 + dependencies: + '@cloudflare/workers-types': 4.20241004.0 + '@types/jest': 29.5.12 + '@typescript-eslint/eslint-plugin': 6.21.0(@typescript-eslint/parser@6.21.0)(eslint@8.56.0)(typescript@5.6.2) + '@typescript-eslint/parser': 6.21.0(eslint@8.56.0)(typescript@5.6.2) + aws4fetch: 1.0.20 + base32-encode: 2.0.0 + eslint: 8.56.0 + eslint-config-standard-with-typescript: 40.0.0(@typescript-eslint/eslint-plugin@6.21.0)(eslint-plugin-import@2.29.1)(eslint-plugin-n@15.7.0)(eslint-plugin-promise@6.1.1)(eslint@8.56.0)(typescript@5.6.2) + 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) + itty-router: 5.0.18 + jest: 29.7.0(@types/node@20.11.19)(ts-node@10.9.2) + postgres: 3.4.4 + prettier: 3.2.5 + ts-jest: 29.1.2(esbuild@0.20.1)(jest@29.7.0)(typescript@5.6.2) + typescript: 5.6.2 + wrangler: 3.80.2(@cloudflare/workers-types@4.20241004.0)(bufferutil@4.0.8)(utf-8-validate@6.0.4) + transitivePeerDependencies: + - '@babel/core' + - '@jest/types' + - '@types/node' + - babel-jest + - babel-plugin-macros + - bufferutil + - esbuild + - node-notifier + - supports-color + - ts-node + - utf-8-validate + 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} id: file:projects/collaboration.tgz @@ -27018,6 +27086,7 @@ packages: name: '@rush-temp/datalake' version: 0.0.0 dependencies: + '@aws-sdk/client-s3': 3.577.0 '@types/jest': 29.5.12 '@types/node': 20.11.19 '@types/node-fetch': 2.6.11 @@ -27057,7 +27126,7 @@ packages: '@vercel/webpack-asset-relocator-loader': 1.7.4 cross-env: 7.0.3 dotenv: 16.0.3 - electron: 32.1.1 + electron: 32.2.1 electron-builder: 25.0.5 electron-squirrel-startup: 1.0.1 node-loader: 2.0.0(webpack@5.90.3) @@ -27193,7 +27262,7 @@ packages: css-loader: 5.2.7(webpack@5.90.3) dotenv: 16.0.3 dotenv-webpack: 8.0.1(webpack@5.90.3) - electron: 32.1.1 + electron: 32.2.1 electron-context-menu: 4.0.4 electron-log: 5.1.7 electron-squirrel-startup: 1.0.1 diff --git a/desktop/src/main/start.ts b/desktop/src/main/start.ts index 731aa1022a..b4e96f7ac0 100644 --- a/desktop/src/main/start.ts +++ b/desktop/src/main/start.ts @@ -124,6 +124,21 @@ function hookOpenWindow (window: BrowserWindow): void { }) } +function handleAuthRedirects (window: BrowserWindow): void { + window.webContents.on('will-redirect', (event) => { + if (event?.url.startsWith(`${FRONT_URL}/login/auth`)) { + console.log('Auth happened, redirecting to local index') + const urlObj = new URL(decodeURIComponent(event.url)) + event.preventDefault() + + void (async (): Promise => { + await window.loadFile(path.join('dist', 'ui', 'index.html')) + window.webContents.send('handle-auth', urlObj.searchParams.get('token')) + })() + } + }) +} + const createWindow = async (): Promise => { mainWindow = new BrowserWindow({ width: defaultWidth, @@ -146,6 +161,7 @@ const createWindow = async (): Promise => { } await mainWindow.loadFile(path.join('dist', 'ui', 'index.html')) addPermissionHandlers(mainWindow.webContents.session) + handleAuthRedirects(mainWindow) // In this example, only windows with the `about:blank` url will be created. // All other urls will be blocked. diff --git a/desktop/src/ui/index.ts b/desktop/src/ui/index.ts index 58b0443aaf..1a393b089a 100644 --- a/desktop/src/ui/index.ts +++ b/desktop/src/ui/index.ts @@ -92,6 +92,15 @@ window.addEventListener('DOMContentLoaded', () => { setDownloadProgress(progress) }) + ipcMain.handleAuth((token) => { + const authLoc = { + path: ['login', 'auth'], + query: { token } + } + + navigate(authLoc) + }) + ipcMain.on('start-backup', () => { // We need to obtain current token and endpoint and trigger backup const token = getMetadata(presentation.metadata.Token) diff --git a/desktop/src/ui/platform.ts b/desktop/src/ui/platform.ts index 359874b35f..c43c95eeed 100644 --- a/desktop/src/ui/platform.ts +++ b/desktop/src/ui/platform.ts @@ -39,7 +39,7 @@ import { taskId } from '@hcengineering/task' import telegram, { telegramId } from '@hcengineering/telegram' import { templatesId } from '@hcengineering/templates' import tracker, { trackerId } from '@hcengineering/tracker' -import uiPlugin, { getCurrentLocation, locationStorageKeyId, navigate, setLocationStorageKey } from '@hcengineering/ui' +import uiPlugin, { getCurrentLocation, locationStorageKeyId, locationToUrl, navigate, parseLocation, setLocationStorageKey } from '@hcengineering/ui' import { uploaderId } from '@hcengineering/uploader' import { viewId } from '@hcengineering/view' import workbench, { workbenchId } from '@hcengineering/workbench' @@ -96,7 +96,7 @@ import '@hcengineering/analytics-collector-assets' import '@hcengineering/text-editor-assets' import { coreId } from '@hcengineering/core' -import presentation, { parsePreviewConfig, presentationId } from '@hcengineering/presentation' +import presentation, { parsePreviewConfig, parseUploadConfig, presentationId } from '@hcengineering/presentation' import textEditor, { textEditorId } from '@hcengineering/text-editor' import love, { loveId } from '@hcengineering/love' import print, { printId } from '@hcengineering/print' @@ -205,6 +205,7 @@ export async function configurePlatform (): Promise { setMetadata(presentation.metadata.FilesURL, config.FILES_URL) setMetadata(presentation.metadata.CollaboratorUrl, config.COLLABORATOR_URL) setMetadata(presentation.metadata.PreviewConfig, parsePreviewConfig(config.PREVIEW_CONFIG)) + setMetadata(presentation.metadata.UploadConfig, parseUploadConfig(config.UPLOAD_CONFIG, config.UPLOAD_URL)) setMetadata(presentation.metadata.FrontUrl, config.FRONT_URL) setMetadata(textEditor.metadata.Collaborator, config.COLLABORATOR ?? '') @@ -336,6 +337,7 @@ export async function configurePlatform (): Promise { } const last = localStorage.getItem(locationStorageKeyId) + if (config.INITIAL_URL !== '') { console.log('NAVIGATE', config.INITIAL_URL, getCurrentLocation()) // NavigationExpandedDefault=false fills buggy: @@ -352,5 +354,6 @@ export async function configurePlatform (): Promise { } else { navigate({ path: [] }) } + console.log('Initial location is: ', getCurrentLocation()) } diff --git a/desktop/src/ui/preload.ts b/desktop/src/ui/preload.ts index afd497a98f..a3984a4d9d 100644 --- a/desktop/src/ui/preload.ts +++ b/desktop/src/ui/preload.ts @@ -130,6 +130,12 @@ const expose: IPCMainExposed = { }) }, + handleAuth: (callback) => { + ipcRenderer.on('handle-auth', (event, value) => { + callback(value) + }) + }, + async setFrontCookie (host: string, name: string, value: string): Promise { ipcRenderer.send('set-front-cookie', host, name, value) }, diff --git a/desktop/src/ui/types.ts b/desktop/src/ui/types.ts index 55e8e39d66..a8d8e16450 100644 --- a/desktop/src/ui/types.ts +++ b/desktop/src/ui/types.ts @@ -30,6 +30,7 @@ export interface Config { AI_URL?:string BRANDING_URL?: string PREVIEW_CONFIG: string + UPLOAD_CONFIG: string DESKTOP_UPDATES_URL?: string DESKTOP_UPDATES_CHANNEL?: string TELEGRAM_BOT_URL?: string @@ -80,6 +81,7 @@ export interface IPCMainExposed { sendNotification: (notififationParams: NotificationParams) => void getScreenAccess: () => Promise getScreenSources: () => Promise + handleAuth: (callback: (token: string) => void) => void cancelBackup: () => void startBackup: (token: string, endpoint: string, workspace: string) => void diff --git a/dev/docker-compose.yaml b/dev/docker-compose.yaml index 4dbab23640..c82f014d1d 100644 --- a/dev/docker-compose.yaml +++ b/dev/docker-compose.yaml @@ -75,6 +75,7 @@ services: - SERVER_SECRET=secret # - DB_URL=postgresql://postgres:example@postgres:5432 - DB_URL=${MONGO_URL} + # - DB_NS=account-2 - REGION_INFO=|Mongo;pg|Postgres - TRANSACTOR_URL=ws://host.docker.internal:3333,ws://host.docker.internal:3331;;pg - SES_URL= diff --git a/dev/prod/src/analytics/posthog.ts b/dev/prod/src/analytics/posthog.ts index dff093d8d7..90fa34d5ec 100644 --- a/dev/prod/src/analytics/posthog.ts +++ b/dev/prod/src/analytics/posthog.ts @@ -2,7 +2,7 @@ import { AnalyticProvider } from "@hcengineering/analytics" import posthog from 'posthog-js' export class PosthogAnalyticProvider implements AnalyticProvider { - init (config: Record): boolean { + init(config: Record): boolean { if (config.POSTHOG_API_KEY !== undefined && config.POSTHOG_API_KEY !== '' && config.POSTHOG_HOST !== null) { posthog.init(config.POSTHOG_API_KEY, { api_host: config.POSTHOG_HOST }) return true @@ -11,15 +11,17 @@ export class PosthogAnalyticProvider implements AnalyticProvider { } setUser(email: string): void { - posthog.identify(email, { email: email }) + if (!posthog._isIdentified()) { + posthog.identify(email, { email: email }) + } } setTag(key: string, value: string): void { posthog.setPersonProperties({ [key]: value }) } setWorkspace(ws: string): void { this.setTag('workspace', ws) - posthog.group('workspace', ws, { - name: `${ws}` + posthog.group('workspace', ws, { + name: `${ws}` }) } logout(): void { diff --git a/dev/prod/src/platform.ts b/dev/prod/src/platform.ts index a43cf316fc..68031ac2e4 100644 --- a/dev/prod/src/platform.ts +++ b/dev/prod/src/platform.ts @@ -108,7 +108,12 @@ import github, { githubId } from '@hcengineering/github' import '@hcengineering/github-assets' import { coreId } from '@hcengineering/core' -import presentation, { loadServerConfig, parsePreviewConfig, presentationId } from '@hcengineering/presentation' +import presentation, { + loadServerConfig, + parsePreviewConfig, + parseUploadConfig, + presentationId +} from '@hcengineering/presentation' import { setMetadata } from '@hcengineering/platform' import { setDefaultLanguage, initThemeStore } from '@hcengineering/theme' @@ -150,6 +155,7 @@ export interface Config { // Could be defined for dev environment FRONT_URL?: string PREVIEW_CONFIG?: string + UPLOAD_CONFIG?: string } export interface Branding { @@ -292,6 +298,7 @@ export async function configurePlatform() { setMetadata(presentation.metadata.FrontUrl, config.FRONT_URL) setMetadata(presentation.metadata.PreviewConfig, parsePreviewConfig(config.PREVIEW_CONFIG)) + setMetadata(presentation.metadata.UploadConfig, parseUploadConfig(config.UPLOAD_CONFIG, config.UPLOAD_URL)) setMetadata(textEditor.metadata.Collaborator, config.COLLABORATOR) diff --git a/dev/tool/package.json b/dev/tool/package.json index aa54749dd4..48324f1b94 100644 --- a/dev/tool/package.json +++ b/dev/tool/package.json @@ -19,8 +19,8 @@ "docker:tbuild": "docker build -t hardcoreeng/tool . --platform=linux/amd64 && ../../common/scripts/docker_tag_push.sh hardcoreeng/tool", "docker:staging": "../../common/scripts/docker_tag.sh hardcoreeng/tool staging", "docker:push": "../../common/scripts/docker_tag.sh hardcoreeng/tool", - "run-local": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3333 MINIO_ACCESS_KEY=minioadmin MINIO_SECRET_KEY=minioadmin MINIO_ENDPOINT=localhost MONGO_URL=mongodb://localhost:27017 DB_URL=mongodb://localhost:27017 TELEGRAM_DATABASE=telegram-service ELASTIC_URL=http://localhost:9200 REKONI_URL=http://localhost:4004 MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) node --max-old-space-size=18000 ./bundle/bundle.js", - "run-local-pg": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3333 MINIO_ACCESS_KEY=minioadmin MINIO_SECRET_KEY=minioadmin MINIO_ENDPOINT=localhost MONGO_URL=mongodb://localhost:27017 DB_URL=postgresql://postgres:example@localhost:5432 TELEGRAM_DATABASE=telegram-service ELASTIC_URL=http://localhost:9200 REKONI_URL=http://localhost:4004 MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) node --max-old-space-size=18000 ./bundle/bundle.js", + "run-local": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3333 MINIO_ACCESS_KEY=minioadmin MINIO_SECRET_KEY=minioadmin MINIO_ENDPOINT=localhost MONGO_URL=mongodb://localhost:27017 DB_URL=mongodb://localhost:27017 TELEGRAM_DATABASE=telegram-service ELASTIC_URL=http://localhost:9200 REKONI_URL=http://localhost:4004 MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) node --expose-gc --max-old-space-size=18000 ./bundle/bundle.js", + "run-local-pg": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3333 MINIO_ACCESS_KEY=minioadmin MINIO_SECRET_KEY=minioadmin MINIO_ENDPOINT=localhost MONGO_URL=mongodb://localhost:27017 DB_URL=postgresql://postgres:example@localhost:5432 TELEGRAM_DATABASE=telegram-service ELASTIC_URL=http://localhost:9200 REKONI_URL=http://localhost:4004 MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) node --expose-gc --max-old-space-size=18000 ./bundle/bundle.js", "run-local-brk": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3333 MINIO_ACCESS_KEY=minioadmin MINIO_SECRET_KEY=minioadmin MINIO_ENDPOINT=localhost MONGO_URL=mongodb://localhost:27017 DB_URL=mongodb://localhost:27017 TELEGRAM_DATABASE=telegram-service ELASTIC_URL=http://localhost:9200 REKONI_URL=http://localhost:4004 MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) node --inspect-brk --enable-source-maps --max-old-space-size=18000 ./bundle/bundle.js", "run": "rush bundle --to @hcengineering/tool >/dev/null && cross-env node --max-old-space-size=8000 ./bundle/bundle.js", "upgrade": "rushx run-local upgrade", diff --git a/dev/tool/src/index.ts b/dev/tool/src/index.ts index e470b9a64e..821bcf97f9 100644 --- a/dev/tool/src/index.ts +++ b/dev/tool/src/index.ts @@ -15,7 +15,6 @@ // import accountPlugin, { - type AccountDB, assignWorkspace, confirmEmail, createAcc, @@ -34,6 +33,7 @@ import accountPlugin, { setAccountAdmin, setRole, updateWorkspace, + type AccountDB, type Workspace } from '@hcengineering/account' import { setMetadata } from '@hcengineering/platform' @@ -41,13 +41,21 @@ import { backup, backupFind, backupList, + backupRemoveLast, backupSize, + checkBackupIntegrity, compactBackup, createFileBackupStorage, createStorageBackupStorage, restore } from '@hcengineering/server-backup' -import serverClientPlugin, { BlobClient, createClient, getTransactorEndpoint } from '@hcengineering/server-client' +import serverClientPlugin, { + BlobClient, + createClient, + getTransactorEndpoint, + listAccountWorkspaces, + updateBackupInfo +} from '@hcengineering/server-client' import { getServerPipeline, registerServerPlugins, registerStringLoaders } from '@hcengineering/server-pipeline' import serverToken, { decodeToken, generateToken } from '@hcengineering/server-token' import toolPlugin, { FileModelLogger } from '@hcengineering/server-tool' @@ -65,6 +73,7 @@ import core, { getWorkspaceId, MeasureMetricsContext, metricsToString, + RateLimiter, systemAccountEmail, versionToString, type Data, @@ -77,8 +86,8 @@ import core, { } from '@hcengineering/core' import { consoleModelLogger, type MigrateOperation } from '@hcengineering/model' import contact from '@hcengineering/model-contact' -import { backupDownload } from '@hcengineering/server-backup/src/backup' import { getMongoClient, getWorkspaceMongoDB, shutdown } from '@hcengineering/mongo' +import { backupDownload } from '@hcengineering/server-backup/src/backup' import type { StorageAdapter, StorageAdapterEx } from '@hcengineering/server-core' import { deepEqual } from 'fast-equals' @@ -104,7 +113,7 @@ import { restoreRecruitingTaskTypes } from './clean' import { changeConfiguration } from './configuration' -import { moveFromMongoToPG, moveWorkspaceFromMongoToPG, moveAccountDbFromMongoToPG } from './db' +import { moveAccountDbFromMongoToPG, moveFromMongoToPG, moveWorkspaceFromMongoToPG } from './db' import { fixJsonMarkup, migrateMarkup, restoreLostMarkup } from './markup' import { fixMixinForeignAttributes, showMixinForeignAttributes } from './mixin' import { fixAccountEmails, renameAccount } from './renameAccount' @@ -814,6 +823,13 @@ export function devTool ( const storage = await createFileBackupStorage(dirName) await compactBackup(toolCtx, storage, cmd.force) }) + program + .command('backup-check ') + .description('Compact a given backup, will create one snapshot clean unused resources') + .action(async (dirName: string, cmd: any) => { + const storage = await createFileBackupStorage(dirName) + await checkBackupIntegrity(toolCtx, storage) + }) program .command('backup-restore [date]') @@ -864,6 +880,61 @@ export function devTool ( await backup(toolCtx, endpoint, wsid, storage) }) }) + program + .command('backup-s3-clean ') + .description('dump workspace transactions and minio resources') + .action(async (bucketName: string, days: string, cmd) => { + const backupStorageConfig = storageConfigFromEnv(process.env.STORAGE) + const storageAdapter = createStorageFromConfig(backupStorageConfig.storages[0]) + + const daysInterval = Date.now() - parseInt(days) * 24 * 60 * 60 * 1000 + try { + const token = generateToken(systemAccountEmail, { name: 'any' }) + const workspaces = (await listAccountWorkspaces(token)).filter((it) => { + const lastBackup = it.backupInfo?.lastBackup ?? 0 + if (lastBackup > daysInterval) { + // No backup required, interval not elapsed + return true + } + + if (it.lastVisit == null) { + return false + } + + return false + }) + workspaces.sort((a, b) => { + return (b.backupInfo?.backupSize ?? 0) - (a.backupInfo?.backupSize ?? 0) + }) + + for (const ws of workspaces) { + const storage = await createStorageBackupStorage( + toolCtx, + storageAdapter, + getWorkspaceId(bucketName), + ws.workspace + ) + await backupRemoveLast(storage, daysInterval) + await updateBackupInfo(generateToken(systemAccountEmail, { name: 'any' }), { + backups: ws.backupInfo?.backups ?? 0, + backupSize: ws.backupInfo?.backupSize ?? 0, + blobsSize: ws.backupInfo?.blobsSize ?? 0, + dataSize: ws.backupInfo?.dataSize ?? 0, + lastBackup: daysInterval + }) + } + } finally { + await storageAdapter.close() + } + }) + program + .command('backup-clean ') + .description('dump workspace transactions and minio resources') + .action(async (dirName: string, days: string, cmd) => { + const daysInterval = Date.now() - parseInt(days) * 24 * 60 * 60 * 1000 + const storage = await createFileBackupStorage(dirName) + await backupRemoveLast(storage, daysInterval) + }) program .command('backup-s3-compact ') @@ -880,6 +951,20 @@ export function devTool ( } await storageAdapter.close() }) + program + .command('backup-s3-check ') + .description('Compact a given backup to just one snapshot') + .action(async (bucketName: string, dirName: string, cmd: any) => { + const backupStorageConfig = storageConfigFromEnv(process.env.STORAGE) + const storageAdapter = createStorageFromConfig(backupStorageConfig.storages[0]) + try { + const storage = await createStorageBackupStorage(toolCtx, storageAdapter, getWorkspaceId(bucketName), dirName) + await checkBackupIntegrity(toolCtx, storage) + } catch (err: any) { + toolCtx.error('failed to size backup', { err }) + } + await storageAdapter.close() + }) program .command('backup-s3-restore [date]') @@ -1100,7 +1185,7 @@ export function devTool ( .command('move-files') .option('-w, --workspace ', 'Selected workspace only', '') .option('-m, --move ', 'When set to true, the files will be moved, otherwise copied', 'false') - .option('-bl, --blobLimit ', 'A blob size limit in megabytes (default 50mb)', '50') + .option('-bl, --blobLimit ', 'A blob size limit in megabytes (default 50mb)', '999999') .option('-c, --concurrency ', 'Number of files being processed concurrently', '10') .option('--disabled', 'Include disabled workspaces', false) .action( @@ -1125,6 +1210,7 @@ export function devTool ( const workspaces = await listWorkspacesPure(db) workspaces.sort((a, b) => b.lastVisit - a.lastVisit) + const rateLimit = new RateLimiter(10) for (const workspace of workspaces) { if (cmd.workspace !== '' && workspace.workspace !== cmd.workspace) { continue @@ -1134,12 +1220,14 @@ export function devTool ( continue } - console.log('start', workspace.workspace, index, '/', workspaces.length) - await moveFiles(toolCtx, getWorkspaceId(workspace.workspace), exAdapter, params) - console.log('done', workspace.workspace) - - index += 1 + await rateLimit.exec(async () => { + console.log('start', workspace.workspace, index, '/', workspaces.length) + await moveFiles(toolCtx, getWorkspaceId(workspace.workspace), exAdapter, params) + console.log('done', workspace.workspace) + index += 1 + }) } + await rateLimit.waitProcessing() } catch (err: any) { console.error(err) } diff --git a/dev/tool/src/storage.ts b/dev/tool/src/storage.ts index 3a76b25ea2..a1ec1a7719 100644 --- a/dev/tool/src/storage.ts +++ b/dev/tool/src/storage.ts @@ -144,6 +144,23 @@ async function processAdapter ( let movedBytes = 0 let batchBytes = 0 + function printStats (): void { + const duration = Date.now() - time + console.log( + '...processed', + processedCnt, + Math.round(processedBytes / 1024 / 1024) + 'MB', + 'moved', + movedCnt, + Math.round(movedBytes / 1024 / 1024) + 'MB', + '+' + Math.round(batchBytes / 1024 / 1024) + 'MB', + Math.round(duration / 1000) + 's' + ) + + batchBytes = 0 + time = Date.now() + } + const rateLimiter = new RateLimiter(params.concurrency) const iterator = await source.listStream(ctx, workspaceId) @@ -152,15 +169,7 @@ async function processAdapter ( const targetBlobs = new Map, ListBlobResult>() - while (true) { - const part = await targetIterator.next() - for (const p of part) { - targetBlobs.set(p._id, p) - } - if (part.length === 0) { - break - } - } + let targetFilled = false const toRemove: string[] = [] try { @@ -168,6 +177,20 @@ async function processAdapter ( const dataBulk = await iterator.next() if (dataBulk.length === 0) break + if (!targetFilled) { + // Only fill target if have something to move. + targetFilled = true + while (true) { + const part = await targetIterator.next() + for (const p of part) { + targetBlobs.set(p._id, p) + } + if (part.length === 0) { + break + } + } + } + for (const data of dataBulk) { let targetBlob: Blob | ListBlobResult | undefined = targetBlobs.get(data._id) if (targetBlob !== undefined) { @@ -219,22 +242,7 @@ async function processAdapter ( if (processedCnt % 100 === 0) { await rateLimiter.waitProcessing() - - const duration = Date.now() - time - - console.log( - '...processed', - processedCnt, - Math.round(processedBytes / 1024 / 1024) + 'MB', - 'moved', - movedCnt, - Math.round(movedBytes / 1024 / 1024) + 'MB', - '+' + Math.round(batchBytes / 1024 / 1024) + 'MB', - Math.round(duration / 1000) + 's' - ) - - batchBytes = 0 - time = Date.now() + printStats() } } } @@ -246,6 +254,7 @@ async function processAdapter ( await source.remove(ctx, workspaceId, part) } } + printStats() } finally { await iterator.close() } diff --git a/models/contact/src/index.ts b/models/contact/src/index.ts index e7772d1147..08e06f05d8 100644 --- a/models/contact/src/index.ts +++ b/models/contact/src/index.ts @@ -860,7 +860,8 @@ export function createModel (builder: Builder): void { title: contact.string.Employees, query: contact.completion.EmployeeQuery, context: ['search', 'mention'], - classToSearch: contact.mixin.Employee + classToSearch: contact.mixin.Employee, + priority: 1000 }, contact.completion.EmployeeCategory ) @@ -874,7 +875,8 @@ export function createModel (builder: Builder): void { title: contact.string.People, query: contact.completion.PersonQuery, context: ['search', 'spotlight'], - classToSearch: contact.class.Person + classToSearch: contact.class.Person, + priority: 900 }, contact.completion.PersonCategory ) @@ -888,7 +890,8 @@ export function createModel (builder: Builder): void { title: contact.string.Organizations, query: contact.completion.OrganizationQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: contact.class.Organization + classToSearch: contact.class.Organization, + priority: 800 }, contact.completion.OrganizationCategory ) diff --git a/models/controlled-documents/src/index.ts b/models/controlled-documents/src/index.ts index 528dfa2fcd..14f81efa24 100644 --- a/models/controlled-documents/src/index.ts +++ b/models/controlled-documents/src/index.ts @@ -908,7 +908,8 @@ export function defineSearch (builder: Builder): void { label: documents.string.SearchDocument, query: documents.completion.DocumentMetaQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: documents.class.DocumentMeta + classToSearch: documents.class.DocumentMeta, + priority: 800 }, documents.completion.DocumentMetaCategory ) diff --git a/models/document/src/index.ts b/models/document/src/index.ts index 6fe9c9db38..6a51aaca3c 100644 --- a/models/document/src/index.ts +++ b/models/document/src/index.ts @@ -492,7 +492,8 @@ function defineDocument (builder: Builder): void { label: document.string.SearchDocument, query: document.completion.DocumentQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: document.class.Document + classToSearch: document.class.Document, + priority: 800 }, document.completion.DocumentQueryCategory ) diff --git a/models/drive/src/index.ts b/models/drive/src/index.ts index c262462c2e..c4d35d18b8 100644 --- a/models/drive/src/index.ts +++ b/models/drive/src/index.ts @@ -444,7 +444,8 @@ function defineFolder (builder: Builder): void { label: presentation.string.Search, query: drive.completion.FolderQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: drive.class.Folder + classToSearch: drive.class.Folder, + priority: 700 }, drive.completion.FolderCategory ) @@ -594,7 +595,8 @@ function defineFile (builder: Builder): void { label: presentation.string.Search, query: drive.completion.FileQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: drive.class.File + classToSearch: drive.class.File, + priority: 600 }, drive.completion.FileCategory ) diff --git a/models/recruit/src/index.ts b/models/recruit/src/index.ts index d0ec049ac2..8d0a97046f 100644 --- a/models/recruit/src/index.ts +++ b/models/recruit/src/index.ts @@ -1048,7 +1048,8 @@ export function createModel (builder: Builder): void { title: recruit.string.Applications, query: recruit.completion.ApplicationQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: recruit.class.Applicant + classToSearch: recruit.class.Applicant, + priority: 500 }, recruit.completion.ApplicationCategory ) @@ -1062,7 +1063,8 @@ export function createModel (builder: Builder): void { title: recruit.string.Vacancies, query: recruit.completion.VacancyQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: recruit.class.Vacancy + classToSearch: recruit.class.Vacancy, + priority: 550 }, recruit.completion.VacancyCategory ) diff --git a/models/tracker/src/index.ts b/models/tracker/src/index.ts index cf5e11bc30..227d794db4 100644 --- a/models/tracker/src/index.ts +++ b/models/tracker/src/index.ts @@ -602,7 +602,8 @@ export function createModel (builder: Builder): void { title: tracker.string.Issues, query: tracker.completion.IssueQuery, context: ['search', 'mention', 'spotlight'], - classToSearch: tracker.class.Issue + classToSearch: tracker.class.Issue, + priority: 300 }, tracker.completion.IssueCategory ) diff --git a/packages/presentation/src/components/markup/NodeContent.svelte b/packages/presentation/src/components/markup/NodeContent.svelte index 80eb244986..35442b8143 100644 --- a/packages/presentation/src/components/markup/NodeContent.svelte +++ b/packages/presentation/src/components/markup/NodeContent.svelte @@ -112,7 +112,8 @@ {:else if node.type === MarkupNodeType.hard_break}
{:else if node.type === MarkupNodeType.ordered_list} -
    + {@const start = toNumber(attrs.start) ?? 1} +
      {#if nodes.length > 0} {#each nodes as node} diff --git a/packages/presentation/src/file.ts b/packages/presentation/src/file.ts index 84e435ec17..5d1fa55756 100644 --- a/packages/presentation/src/file.ts +++ b/packages/presentation/src/file.ts @@ -13,12 +13,30 @@ // limitations under the License. // -import { concatLink, type Blob, type Ref } from '@hcengineering/core' +import { concatLink, type Blob as PlatformBlob, type Ref } from '@hcengineering/core' import { PlatformError, Severity, Status, getMetadata } from '@hcengineering/platform' import { v4 as uuid } from 'uuid' import plugin from './plugin' +export type FileUploadMethod = 'form-data' | 'signed-url' + +export interface UploadConfig { + 'form-data': { + url: string + } + 'signed-url'?: { + url: string + size: number + } +} + +export interface FileUploadParams { + method: FileUploadMethod + url: string + headers: Record +} + interface FileUploadError { key: string error: string @@ -34,6 +52,46 @@ type FileUploadResult = FileUploadSuccess | FileUploadError const defaultUploadUrl = '/files' const defaultFilesUrl = '/files/:workspace/:filename?file=:blobId&workspace=:workspace' +function parseInt (value: string, fallback: number): number { + const number = Number.parseInt(value) + return Number.isInteger(number) ? number : fallback +} + +export function parseUploadConfig (config: string, uploadUrl: string): UploadConfig { + const uploadConfig: UploadConfig = { + 'form-data': { url: uploadUrl }, + 'signed-url': undefined + } + + if (config !== undefined) { + const configs = config.split(';') + for (const c of configs) { + if (c === '') { + continue + } + + const [key, size, url] = c.split('|') + + if (url === undefined || url === '') { + throw new Error(`Bad upload config: ${c}`) + } + + if (key === 'form-data') { + uploadConfig['form-data'] = { url } + } else if (key === 'signed-url') { + uploadConfig['signed-url'] = { + url, + size: parseInt(size, 0) * 1024 * 1024 + } + } else { + throw new Error(`Unknown upload config key: ${key}`) + } + } + } + + return uploadConfig +} + function getFilesUrl (): string { const filesUrl = getMetadata(plugin.metadata.FilesURL) ?? defaultFilesUrl const frontUrl = getMetadata(plugin.metadata.FrontUrl) ?? window.location.origin @@ -61,6 +119,42 @@ export function getUploadUrl (): string { return template.replaceAll(':workspace', encodeURIComponent(getCurrentWorkspaceId())) } +function getUploadConfig (): UploadConfig { + return getMetadata(plugin.metadata.UploadConfig) ?? { 'form-data': { url: getUploadUrl() } } +} + +function getFileUploadMethod (blob: Blob): { method: FileUploadMethod, url: string } { + const config = getUploadConfig() + + const signedUrl = config['signed-url'] + if (signedUrl !== undefined && signedUrl.size < blob.size) { + return { method: 'signed-url', url: signedUrl.url } + } + + return { method: 'form-data', url: config['form-data'].url } +} + +/** + * @public + */ +export function getFileUploadParams (blobId: string, blob: Blob): FileUploadParams { + const workspaceId = encodeURIComponent(getCurrentWorkspaceId()) + const fileId = encodeURIComponent(blobId) + + const { method, url: urlTemplate } = getFileUploadMethod(blob) + + const url = urlTemplate.replaceAll(':workspace', workspaceId).replaceAll(':blobId', fileId) + + const headers: Record = + method !== 'signed-url' + ? { + Authorization: 'Bearer ' + (getMetadata(plugin.metadata.Token) as string) + } + : {} + + return { method, url, headers } +} + /** * @public */ @@ -79,12 +173,40 @@ export function getFileUrl (file: string, filename?: string): string { /** * @public */ -export async function uploadFile (file: File): Promise> { - const uploadUrl = getUploadUrl() - +export async function uploadFile (file: File): Promise> { const id = generateFileId() + const params = getFileUploadParams(id, file) + + if (params.method === 'signed-url') { + await uploadFileWithSignedUrl(file, id, params.url) + } else { + await uploadFileWithFormData(file, id, params.url) + } + + return id as Ref +} + +/** + * @public + */ +export async function deleteFile (id: string): Promise { + const fileUrl = getFileUrl(id) + + const resp = await fetch(fileUrl, { + method: 'DELETE', + headers: { + Authorization: 'Bearer ' + (getMetadata(plugin.metadata.Token) as string) + } + }) + + if (resp.status !== 200) { + throw new Error('Failed to delete file') + } +} + +async function uploadFileWithFormData (file: File, uuid: string, uploadUrl: string): Promise { const data = new FormData() - data.append('file', file, id) + data.append('file', file, uuid) const resp = await fetch(uploadUrl, { method: 'POST', @@ -110,24 +232,54 @@ export async function uploadFile (file: File): Promise> { if ('error' in result[0]) { throw Error(`Failed to upload file: ${result[0].error}`) } - - return id as Ref } -/** - * @public - */ -export async function deleteFile (id: string): Promise { - const fileUrl = getFileUrl(id) - - const resp = await fetch(fileUrl, { - method: 'DELETE', +async function uploadFileWithSignedUrl (file: File, uuid: string, uploadUrl: string): Promise { + const response = await fetch(uploadUrl, { + method: 'POST', headers: { Authorization: 'Bearer ' + (getMetadata(plugin.metadata.Token) as string) } }) - if (resp.status !== 200) { - throw new Error('Failed to delete file') + if (response.ok) { + throw Error(`Failed to genearte signed upload URL: ${response.statusText}`) + } + + const signedUrl = await response.text() + if (signedUrl === undefined || signedUrl === '') { + throw Error('Missing signed upload URL') + } + + try { + const response = await fetch(signedUrl, { + body: file, + method: 'PUT', + headers: { + 'Content-Type': file.type, + 'Content-Length': file.size.toString(), + 'x-amz-meta-last-modified': file.lastModified.toString() + } + }) + + if (!response.ok) { + throw Error(`Failed to upload file: ${response.statusText}`) + } + + // confirm we uploaded file + await fetch(uploadUrl, { + method: 'PUT', + headers: { + Authorization: 'Bearer ' + (getMetadata(plugin.metadata.Token) as string) + } + }) + } catch (err) { + // abort the upload + await fetch(uploadUrl, { + method: 'DELETE', + headers: { + Authorization: 'Bearer ' + (getMetadata(plugin.metadata.Token) as string) + } + }) } } diff --git a/packages/presentation/src/plugin.ts b/packages/presentation/src/plugin.ts index 067ccd235a..0f826224c4 100644 --- a/packages/presentation/src/plugin.ts +++ b/packages/presentation/src/plugin.ts @@ -43,6 +43,7 @@ import { type InstantTransactions, type ObjectSearchCategory } from './types' +import { type UploadConfig } from './file' /** * @public @@ -138,6 +139,7 @@ export default plugin(presentationId, { Workspace: '' as Metadata, WorkspaceId: '' as Metadata, FrontUrl: '' as Asset, + UploadConfig: '' as Metadata, PreviewConfig: '' as Metadata, ClientHook: '' as Metadata, SessionId: '' as Metadata diff --git a/packages/presentation/src/search.ts b/packages/presentation/src/search.ts index 113882b907..0b49f5009e 100644 --- a/packages/presentation/src/search.ts +++ b/packages/presentation/src/search.ts @@ -98,9 +98,14 @@ async function doFulltextSearch ( } return sections.sort((a, b) => { - const maxScoreA = Math.max(...(a?.items ?? []).map((obj) => obj?.score ?? 0)) - const maxScoreB = Math.max(...(b?.items ?? []).map((obj) => obj?.score ?? 0)) - return maxScoreB - maxScoreA + const ac = categories.indexOf(a.category) + const bc = categories.indexOf(b.category) + if (ac === bc) { + const maxScoreA = Math.max(...(a?.items ?? []).map((obj) => obj?.score ?? 0)) + const maxScoreB = Math.max(...(b?.items ?? []).map((obj) => obj?.score ?? 0)) + return maxScoreB - maxScoreA + } + return ac - bc }) } @@ -114,6 +119,8 @@ export async function searchFor ( let categories = categoriesByContext.get(context) if (categories === undefined) { categories = await client.findAll(plugin.class.ObjectSearchCategory, { context }) + + categories.sort((a, b) => (b.priority ?? 0) - (a.priority ?? 0)) categoriesByContext.set(context, categories) } diff --git a/packages/presentation/src/types.ts b/packages/presentation/src/types.ts index fad908bd4b..9e9a09b5f2 100644 --- a/packages/presentation/src/types.ts +++ b/packages/presentation/src/types.ts @@ -72,6 +72,8 @@ export interface ObjectSearchCategory extends Doc { // Query for documents with pattern query: Resource classToSearch?: Ref> + + priority?: number } export interface ComponentExt { diff --git a/packages/ui/src/components/Scroller.svelte b/packages/ui/src/components/Scroller.svelte index a19a6cee4a..8b9548a68d 100644 --- a/packages/ui/src/components/Scroller.svelte +++ b/packages/ui/src/components/Scroller.svelte @@ -194,8 +194,13 @@ const topBar = Y - rectScroll.y - shiftTop - 2 const heightScroll = rectScroll.height - 4 - divBar.clientHeight - shiftTop - shiftBottom const procBar = topBar / heightScroll - divScroll.scrollTop = (divScroll.scrollHeight - divScroll.clientHeight) * procBar - } else { + + if (scrollDirection === 'vertical-reverse') { + divScroll.scrollTop = (divScroll.scrollHeight - divScroll.clientHeight) * (procBar - 1) + } else { + divScroll.scrollTop = (divScroll.scrollHeight - divScroll.clientHeight) * procBar + } + } else if (isScrolling === 'horizontal') { let X = 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) { @@ -513,7 +518,11 @@ divBar.style.top = `${topBar}px` const heightScroll = rectScroll.height - 4 - barHeight - shiftTop - shiftBottom const procBar = (topBar - rectScroll.top - shiftTop - 2) / heightScroll - divScroll.scrollTop = (divScroll.scrollHeight - divScroll.clientHeight) * procBar + if (scrollDirection === 'vertical-reverse') { + divScroll.scrollTop = (divScroll.scrollHeight - divScroll.clientHeight) * (procBar - 1) + } else { + divScroll.scrollTop = (divScroll.scrollHeight - divScroll.clientHeight) * procBar + } } } @@ -655,6 +664,7 @@
      { onScrollStart(ev, 'vertical') @@ -895,6 +905,11 @@ max-height: calc(100% - 12px); transform: scaleX(0.5); + &.reverse { + top: auto; + bottom: 2px; + } + &:hover, &.hovered { transform: scaleX(1); diff --git a/plugins/attachment-resources/src/components/AttachmentStyleBoxCollabEditor.svelte b/plugins/attachment-resources/src/components/AttachmentStyleBoxCollabEditor.svelte index a65b770905..84e415dff8 100644 --- a/plugins/attachment-resources/src/components/AttachmentStyleBoxCollabEditor.svelte +++ b/plugins/attachment-resources/src/components/AttachmentStyleBoxCollabEditor.svelte @@ -35,7 +35,7 @@ getModelRefActions } from '@hcengineering/text-editor-resources' import { AnySvelteComponent, getEventPositionElement, getPopupPositionElement, navigate } from '@hcengineering/ui' - import { uploadFiles } from '@hcengineering/uploader' + import { type FileUploadCallbackParams, uploadFiles } from '@hcengineering/uploader' import view from '@hcengineering/view' import { getCollaborationUser, getObjectId, getObjectLinkFragment } from '@hcengineering/view-resources' import { Analytics } from '@hcengineering/analytics' @@ -135,14 +135,7 @@ progress = true - await uploadFiles( - list, - { objectId: object._id, objectClass: object._class }, - {}, - async (uuid, name, file, path, metadata) => { - await createAttachment(uuid, name, file, metadata) - } - ) + await uploadFiles(list, { onFileUploaded }) inputFile.value = '' progress = false @@ -151,14 +144,7 @@ async function attachFiles (files: File[] | FileList): Promise { progress = true if (files.length > 0) { - await uploadFiles( - files, - { objectId: object._id, objectClass: object._class }, - {}, - async (uuid, name, file, path, metadata) => { - await createAttachment(uuid, name, file, metadata) - } - ) + await uploadFiles(files, { onFileUploaded }) } progress = false } @@ -174,6 +160,10 @@ } } + async function onFileUploaded ({ uuid, name, file, metadata }: FileUploadCallbackParams): Promise { + await createAttachment(uuid, name, file, metadata) + } + async function createAttachment ( uuid: Ref, name: string, diff --git a/plugins/attachment-resources/src/components/AttachmentStyledBox.svelte b/plugins/attachment-resources/src/components/AttachmentStyledBox.svelte index 2f16f8ad50..f8e7c0fe7b 100644 --- a/plugins/attachment-resources/src/components/AttachmentStyledBox.svelte +++ b/plugins/attachment-resources/src/components/AttachmentStyledBox.svelte @@ -13,7 +13,7 @@ // limitations under the License. --> diff --git a/plugins/drive-resources/src/components/FilePanel.svelte b/plugins/drive-resources/src/components/FilePanel.svelte index 9c460c2c7c..f5d84e8d4e 100644 --- a/plugins/drive-resources/src/components/FilePanel.svelte +++ b/plugins/drive-resources/src/components/FilePanel.svelte @@ -18,7 +18,7 @@ import { Panel } from '@hcengineering/panel' import { createQuery, getClient, getFileUrl } from '@hcengineering/presentation' import { Button, IconMoreH } from '@hcengineering/ui' - import { showFilesUploadPopup } from '@hcengineering/uploader' + import { FileUploadCallbackParams, showFilesUploadPopup } from '@hcengineering/uploader' import view from '@hcengineering/view' import { showMenu } from '@hcengineering/view-resources' @@ -68,27 +68,30 @@ function handleUploadFile (): void { if (object != null) { void showFilesUploadPopup( - { objectId: object._id, objectClass: object._class }, { - maxNumberOfFiles: 1, - hideProgress: true + onFileUploaded, + showProgress: { + target: { objectId: object._id, objectClass: object._class } + }, + maxNumberOfFiles: 1 }, - {}, - async (uuid, name, file, path, metadata) => { - const data = { - file: uuid, - title: name, - size: file.size, - type: file.type, - lastModified: file instanceof File ? file.lastModified : Date.now(), - metadata - } - - await createFileVersion(client, _id, data) - } + {} ) } } + + async function onFileUploaded ({ uuid, name, file, metadata }: FileUploadCallbackParams): Promise { + const data = { + file: uuid, + title: name, + size: file.size, + type: file.type, + lastModified: file instanceof File ? file.lastModified : Date.now(), + metadata + } + + await createFileVersion(client, _id, data) + } {#if object && version} diff --git a/plugins/drive-resources/src/utils.ts b/plugins/drive-resources/src/utils.ts index 52473b0281..38d87e1ef2 100644 --- a/plugins/drive-resources/src/utils.ts +++ b/plugins/drive-resources/src/utils.ts @@ -177,7 +177,14 @@ export async function uploadFilesToDrive (dt: DataTransfer, space: Ref, p ? { objectId: parent, objectClass: drive.class.Folder } : { objectId: space, objectClass: drive.class.Drive } - await uploadFiles(files, target, {}, onFileUploaded) + const options = { + onFileUploaded, + showProgress: { + target + } + } + + await uploadFiles(files, options) } export async function uploadFilesToDrivePopup (space: Ref, parent: Ref): Promise { @@ -189,12 +196,15 @@ export async function uploadFilesToDrivePopup (space: Ref, parent: Ref, parent: Ref): Prom return current } - const callback: FileUploadCallback = async (uuid, name, file, path, metadata) => { + const callback: FileUploadCallback = async ({ uuid, name, file, path, metadata }) => { const folder = await findParent(path) try { const data = { diff --git a/plugins/login-resources/src/components/BottomAction.svelte b/plugins/login-resources/src/components/BottomAction.svelte index 1d8d5e02c1..173b2e7105 100644 --- a/plugins/login-resources/src/components/BottomAction.svelte +++ b/plugins/login-resources/src/components/BottomAction.svelte @@ -18,7 +18,7 @@ import { NavLink } from '@hcengineering/presentation' import { getHref } from '../utils' - import { BottomAction } from '../index' + import { BottomAction, goTo } from '../index' export let action: BottomAction @@ -28,7 +28,16 @@ {/if} {#if action.page} - + { + if (action.func !== undefined) { + action.func() + } else if (action.page !== undefined) { + goTo(action.page) + } + }}> {:else} {/if} diff --git a/plugins/login-resources/src/components/SelectWorkspace.svelte b/plugins/login-resources/src/components/SelectWorkspace.svelte index 9e2b447fe6..fcb70f0180 100644 --- a/plugins/login-resources/src/components/SelectWorkspace.svelte +++ b/plugins/login-resources/src/components/SelectWorkspace.svelte @@ -169,7 +169,12 @@ {#if workspaces.length}
      - + { + goTo('createWorkspace') + }}>
      {/if}
      diff --git a/plugins/uploader-resources/src/types.ts b/plugins/uploader-resources/src/types.ts new file mode 100644 index 0000000000..5997123f75 --- /dev/null +++ b/plugins/uploader-resources/src/types.ts @@ -0,0 +1,28 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// +import type { IndexedObject } from '@uppy/core' + +// For Uppy 4.0 compatibility +export type Meta = IndexedObject +export type Body = IndexedObject + +/** @public */ +export type UppyMeta = Meta & { + uuid: string + relativePath?: string +} + +/** @public */ +export type UppyBody = Body diff --git a/plugins/uploader-resources/src/uppy.ts b/plugins/uploader-resources/src/uppy.ts index 84a7316e2d..ce175d3a40 100644 --- a/plugins/uploader-resources/src/uppy.ts +++ b/plugins/uploader-resources/src/uppy.ts @@ -14,12 +14,18 @@ // import { type Blob, type Ref, generateId } from '@hcengineering/core' -import { getMetadata } from '@hcengineering/platform' -import presentation, { generateFileId, getFileMetadata, getUploadUrl } from '@hcengineering/presentation' +import { getMetadata, PlatformError, unknownError } from '@hcengineering/platform' +import presentation, { + type FileUploadParams, + generateFileId, + getFileMetadata, + getFileUploadParams, + getUploadUrl +} from '@hcengineering/presentation' import { getCurrentLanguage } from '@hcengineering/theme' -import type { FileUploadCallback, FileUploadOptions } from '@hcengineering/uploader' +import type { FileUploadOptions } from '@hcengineering/uploader' -import Uppy, { type IndexedObject, type UppyOptions } from '@uppy/core' +import Uppy, { type UppyFile, type UppyOptions } from '@uppy/core' import XHR from '@uppy/xhr-upload' import En from '@uppy/locales/lib/en_US' @@ -29,6 +35,8 @@ import Pt from '@uppy/locales/lib/pt_PT' import Ru from '@uppy/locales/lib/ru_RU' import Zh from '@uppy/locales/lib/zh_CN' +import type { UppyBody, UppyMeta } from './types' + type Locale = UppyOptions['locale'] const locales: Record = { @@ -44,22 +52,65 @@ function getUppyLocale (lang: string): Locale { return locales[lang] ?? En } -// For Uppy 4.0 compatibility -type Meta = IndexedObject -type Body = IndexedObject +interface XHRFileProcessor { + name: string + onBeforeUpload: (uppy: Uppy, file: UppyFile, params: FileUploadParams) => Promise + onAfterUpload: (uppy: Uppy, file: UppyFile, params: FileUploadParams) => Promise +} -/** @public */ -export type UppyMeta = Meta & { - relativePath?: string +const FormDataFileProcessor: XHRFileProcessor = { + name: 'form-data', + + onBeforeUpload: async (uppy: Uppy, file: UppyFile, { url, headers }: FileUploadParams): Promise => { + const xhrUpload = 'xhrUpload' in file && typeof file.xhrUpload === 'object' ? file.xhrUpload : {} + const state = { + xhrUpload: { + ...xhrUpload, + endpoint: url, + method: 'POST', + formData: true, + headers + } + } + uppy.setFileState(file.id, state) + }, + + onAfterUpload: async (uppy: Uppy, file: UppyFile, params: FileUploadParams): Promise => {} +} + +const SignedURLFileProcessor: XHRFileProcessor = { + name: 'signed-url', + + onBeforeUpload: async (uppy: Uppy, file: UppyFile, { url, headers }: FileUploadParams): Promise => { + const xhrUpload = 'xhrUpload' in file && typeof file.xhrUpload === 'object' ? file.xhrUpload : {} + const signedUrl = await getSignedUploadUrl(file, url) + const state = { + xhrUpload: { + ...xhrUpload, + formData: false, + method: 'PUT', + endpoint: signedUrl, + headers: { + ...headers, + // S3 direct upload does not require authorization + Authorization: '', + 'Content-Type': file.type + } + } + } + uppy.setFileState(file.id, state) + }, + + onAfterUpload: async (uppy: Uppy, file: UppyFile, params: FileUploadParams): Promise => { + const error = 'error' in file && file.error != null + await fetch(params.url, { method: error ? 'DELETE' : 'PUT' }) + } } /** @public */ -export type UppyBody = Body & { - uuid: string -} +export function getUppy (options: FileUploadOptions): Uppy { + const { onFileUploaded } = options -/** @public */ -export function getUppy (options: FileUploadOptions, onFileUploaded?: FileUploadCallback): Uppy { const uppyOptions: Partial = { id: generateId(), locale: getUppyLocale(getCurrentLanguage()), @@ -71,31 +122,17 @@ export function getUppy (options: FileUploadOptions, onFileUploaded?: FileUpload } } - const uppy = new Uppy(uppyOptions).use(XHR, { - endpoint: getUploadUrl(), - method: 'POST', - headers: { - Authorization: 'Bearer ' + (getMetadata(presentation.metadata.Token) as string) - }, - getResponseError: (_, response) => { - return new Error((response as Response).statusText) - } - }) - - // Hack to setup shouldRetry callback on xhrUpload that is not exposed in options - const xhrUpload = uppy.getState().xhrUpload ?? {} - uppy.getState().xhrUpload = { - ...xhrUpload, - shouldRetry: (response: Response) => response.status !== 413 - } + const uppy = new Uppy(uppyOptions) + // Ensure we always have UUID uppy.addPreProcessor(async (fileIds: string[]) => { for (const fileId of fileIds) { const file = uppy.getFile(fileId) if (file != null) { - const uuid = generateFileId() - file.meta.uuid = uuid - file.meta.name = uuid + // It may seem weird that we modify file name here + // but we need a way to pass file UUID to Datalake via form data + const uuid = file.meta.uuid ?? generateFileId() + uppy.setFileMeta(fileId, { uuid, name: uuid }) } } }) @@ -111,11 +148,78 @@ export function getUppy (options: FileUploadOptions, onFileUploaded?: FileUpload const uuid = file.meta.uuid as Ref if (uuid !== undefined) { const metadata = await getFileMetadata(file.data, uuid) - await onFileUploaded(uuid, file.name, file.data, file.meta.relativePath, metadata) + await onFileUploaded({ + uuid, + name: file.name, + file: file.data, + path: file.meta.relativePath, + metadata + }) + } else { + console.warn('missing file metadata uuid', file) } } }) } + configureXHR(uppy) + return uppy } + +function configureXHR (uppy: Uppy): Uppy { + uppy.use(XHR, { + endpoint: getUploadUrl(), + method: 'POST', + headers: { + Authorization: 'Bearer ' + (getMetadata(presentation.metadata.Token) as string) + }, + getResponseError: (_, response) => { + return new Error((response as Response).statusText) + } + }) + + // Hack to setup shouldRetry callback on xhrUpload that is not exposed in options + const xhrUpload = uppy.getState().xhrUpload ?? {} + uppy.getState().xhrUpload = { + ...xhrUpload, + shouldRetry: (response: Response) => !(response.status in [401, 403, 413]) + } + + uppy.addPreProcessor(async (fileIds: string[]) => { + for (const fileId of fileIds) { + const file = uppy.getFile(fileId) + if (file != null) { + const params = getFileUploadParams(file.meta.uuid, file.data) + const processor = getXHRProcessor(file, params) + await processor.onBeforeUpload(uppy, file, params) + } + } + }) + + uppy.addPostProcessor(async (fileIds: string[]) => { + for (const fileId of fileIds) { + const file = uppy.getFile(fileId) + if (file != null) { + const params = getFileUploadParams(file.meta.uuid, file.data) + const processor = getXHRProcessor(file, params) + await processor.onAfterUpload(uppy, file, params) + } + } + }) + + return uppy +} + +function getXHRProcessor (file: UppyFile, params: FileUploadParams): XHRFileProcessor { + return params.method === 'form-data' ? FormDataFileProcessor : SignedURLFileProcessor +} + +async function getSignedUploadUrl (file: UppyFile, signUrl: string): Promise { + const response = await fetch(signUrl, { method: 'POST' }) + if (!response.ok) { + throw new PlatformError(unknownError('Failed to get signed upload url')) + } + + return await response.text() +} diff --git a/plugins/uploader-resources/src/utils.ts b/plugins/uploader-resources/src/utils.ts index 41c625aaea..d6c667e714 100644 --- a/plugins/uploader-resources/src/utils.ts +++ b/plugins/uploader-resources/src/utils.ts @@ -14,54 +14,45 @@ // import { showPopup } from '@hcengineering/ui' -import { - type FileUploadCallback, - type FileUploadOptions, - type FileUploadPopupOptions, - type FileUploadTarget, - toFileWithPath -} from '@hcengineering/uploader' +import { type FileUploadOptions, type FileUploadPopupOptions, toFileWithPath } from '@hcengineering/uploader' import FileUploadPopup from './components/FileUploadPopup.svelte' import { dockFileUpload } from './store' import { getUppy } from './uppy' +import { generateFileId } from '@hcengineering/presentation' /** @public */ export async function showFilesUploadPopup ( - target: FileUploadTarget, options: FileUploadOptions, - popupOptions: FileUploadPopupOptions, - onFileUploaded: FileUploadCallback + popupOptions: FileUploadPopupOptions ): Promise { - const uppy = getUppy(options, onFileUploaded) + const uppy = getUppy(options) - showPopup(FileUploadPopup, { uppy, target, options: popupOptions }, undefined, (res) => { - if (res === true && options.hideProgress !== true) { + showPopup(FileUploadPopup, { uppy, options: popupOptions }, undefined, (res) => { + if (res === true && options.showProgress !== undefined) { + const { target } = options.showProgress dockFileUpload(target, uppy) } }) } /** @public */ -export async function uploadFiles ( - files: File[] | FileList, - target: FileUploadTarget, - options: FileUploadOptions, - onFileUploaded: FileUploadCallback -): Promise { +export async function uploadFiles (files: File[] | FileList, options: FileUploadOptions): Promise { const items = Array.from(files, (p) => toFileWithPath(p)) if (items.length === 0) return - const uppy = getUppy(options, onFileUploaded) + const uppy = getUppy(options) for (const data of items) { const { name, type, relativePath } = data - uppy.addFile({ name, type, data, meta: { relativePath } }) + const uuid = generateFileId() + uppy.addFile({ name, type, data, meta: { name: uuid, uuid, relativePath } }) } - if (options.hideProgress !== true) { + if (options.showProgress !== undefined) { + const { target } = options.showProgress dockFileUpload(target, uppy) } diff --git a/plugins/uploader/src/types.ts b/plugins/uploader/src/types.ts index 6d81bec926..b2315c8f3e 100644 --- a/plugins/uploader/src/types.ts +++ b/plugins/uploader/src/types.ts @@ -21,20 +21,10 @@ export interface FileWithPath extends File { } /** @public */ -export type UploadFilesPopupFn = ( - target: FileUploadTarget, - options: FileUploadOptions, - popupOptions: FileUploadPopupOptions, - onFileUploaded: FileUploadCallback -) => Promise +export type UploadFilesPopupFn = (options: FileUploadOptions, popupOptions: FileUploadPopupOptions) => Promise /** @public */ -export type UploadFilesFn = ( - files: File[] | FileList, - target: FileUploadTarget, - options: FileUploadOptions, - onFileUploaded: FileUploadCallback -) => Promise +export type UploadFilesFn = (files: File[] | FileList, options: FileUploadOptions) => Promise /** @public */ export interface FileUploadTarget { @@ -42,12 +32,20 @@ export interface FileUploadTarget { objectClass: Ref> } +/** @public */ +export interface FileUploadProgressOptions { + target: FileUploadTarget +} + /** @public */ export interface FileUploadOptions { + // Uppy options maxFileSize?: number maxNumberOfFiles?: number allowedFileTypes?: string[] | null - hideProgress?: boolean + + onFileUploaded?: FileUploadCallback + showProgress?: FileUploadProgressOptions } /** @public */ @@ -56,10 +54,13 @@ export interface FileUploadPopupOptions { } /** @public */ -export type FileUploadCallback = ( - uuid: Ref, - name: string, - file: FileWithPath | Blob, - path: string | undefined, +export interface FileUploadCallbackParams { + uuid: Ref + name: string + file: FileWithPath | Blob + path: string | undefined metadata: Record | undefined -) => Promise +} + +/** @public */ +export type FileUploadCallback = (params: FileUploadCallbackParams) => Promise diff --git a/plugins/uploader/src/utils.ts b/plugins/uploader/src/utils.ts index cadbf0b618..9199b5f40b 100644 --- a/plugins/uploader/src/utils.ts +++ b/plugins/uploader/src/utils.ts @@ -16,34 +16,27 @@ import { getResource } from '@hcengineering/platform' import uploader from './plugin' -import type { - FileUploadCallback, - FileUploadOptions, - FileUploadPopupOptions, - FileUploadTarget, - FileWithPath -} from './types' +import type { FileUploadOptions, FileUploadPopupOptions, FileWithPath } from './types' /** @public */ export async function showFilesUploadPopup ( - target: FileUploadTarget, options: FileUploadOptions, - popupOptions: FileUploadPopupOptions, - onFileUploaded: FileUploadCallback + popupOptions: FileUploadPopupOptions ): Promise { const fn = await getResource(uploader.function.ShowFilesUploadPopup) - await fn(target, options, popupOptions, onFileUploaded) + await fn(options, popupOptions) } /** @public */ -export async function uploadFiles ( - files: File[] | FileList, - target: FileUploadTarget, - options: FileUploadOptions, - onFileUploaded: FileUploadCallback -): Promise { +export async function uploadFile (file: File, options: FileUploadOptions): Promise { const fn = await getResource(uploader.function.UploadFiles) - await fn(files, target, options, onFileUploaded) + await fn([file], options) +} + +/** @public */ +export async function uploadFiles (files: File[] | FileList, options: FileUploadOptions): Promise { + const fn = await getResource(uploader.function.UploadFiles) + await fn(files, options) } /** @public */ diff --git a/pods/authProviders/src/openid.ts b/pods/authProviders/src/openid.ts index 297c1a5b0c..4029f2838a 100644 --- a/pods/authProviders/src/openid.ts +++ b/pods/authProviders/src/openid.ts @@ -38,21 +38,29 @@ export function registerOpenid ( const redirectURL = '/auth/openid/callback' if (openidClientId === undefined || openidClientSecret === undefined || issuer === undefined) return - void Issuer.discover(issuer).then((issuerObj) => { - const client = new issuerObj.Client({ - client_id: openidClientId, - client_secret: openidClientSecret, - redirect_uris: [concatLink(accountsUrl, redirectURL)], - response_types: ['code'] - }) + Issuer.discover(issuer) + .then((issuerObj) => { + measureCtx.info('Discovered issuer', { issuer: issuerObj }) - passport.use( - 'oidc', - new Strategy({ client, passReqToCallback: true }, (req: any, tokenSet: any, userinfo: any, done: any) => { - return done(null, userinfo) + const client = new issuerObj.Client({ + client_id: openidClientId, + client_secret: openidClientSecret, + redirect_uris: [concatLink(accountsUrl, redirectURL)], + response_types: ['code'] }) - ) - }) + measureCtx.info('Created OIDC client') + + passport.use( + 'oidc', + new Strategy({ client, passReqToCallback: true }, (req: any, tokenSet: any, userinfo: any, done: any) => { + return done(null, userinfo) + }) + ) + measureCtx.info('Registered OIDC strategy') + }) + .catch((err) => { + measureCtx.error('Failed to create OIDC client for IdP with the provided configuration', { err }) + }) router.get('/auth/openid', async (ctx, next) => { measureCtx.info('try auth via', { provider: 'openid' }) diff --git a/pods/backup/Dockerfile b/pods/backup/Dockerfile index 2450bc6e2c..1bcbcc317c 100644 --- a/pods/backup/Dockerfile +++ b/pods/backup/Dockerfile @@ -17,4 +17,4 @@ COPY bundle/bundle.js.map ./ COPY bundle/model.json ./ EXPOSE 3000 -CMD [ "node", "bundle.js" ] +CMD [ "node", "--expose-gc", "bundle.js" ] diff --git a/rush.json b/rush.json index e73a75304b..34ddb321b7 100644 --- a/rush.json +++ b/rush.json @@ -2115,6 +2115,11 @@ "packageName": "@hcengineering/cloud-branding", "projectFolder": "workers/branding", "shouldPublish": false + }, + { + "packageName": "@hcengineering/cloud-datalake", + "projectFolder": "workers/datalake", + "shouldPublish": false } ] } diff --git a/server/account-service/src/index.ts b/server/account-service/src/index.ts index f3e92f4922..9e0cb5f76c 100644 --- a/server/account-service/src/index.ts +++ b/server/account-service/src/index.ts @@ -95,7 +95,8 @@ export function serveAccount (measureCtx: MeasureContext, brandings: BrandingMap const hasSignUp = process.env.DISABLE_SIGNUP !== 'true' const methods = getMethods(hasSignUp) - const accountsDb = getAccountDB(dbUrl) + const dbNs = process.env.DB_NS + const accountsDb = getAccountDB(dbUrl, dbNs) const app = new Koa() const router = new Router() diff --git a/server/account/src/operations.ts b/server/account/src/operations.ts index ba27c43627..16af7c3ffb 100644 --- a/server/account/src/operations.ts +++ b/server/account/src/operations.ts @@ -1290,7 +1290,8 @@ export async function getPendingWorkspace ( workspaceName: result.workspaceName, operation, region, - version + workspaceVersion: result.version, + requestedVersion: version }) } diff --git a/server/account/src/utils.ts b/server/account/src/utils.ts index e572034cb7..f5067b6d3f 100644 --- a/server/account/src/utils.ts +++ b/server/account/src/utils.ts @@ -24,17 +24,12 @@ import { PostgresAccountDB } from './collections/postgres' import { accountPlugin } from './plugin' import type { Account, AccountDB, AccountInfo, RegionInfo, WorkspaceInfo } from './types' -/** - * @public - */ -export const ACCOUNT_DB = 'account' - -export async function getAccountDB (uri: string, db: string = ACCOUNT_DB): Promise<[AccountDB, () => void]> { +export async function getAccountDB (uri: string, dbNs?: string): Promise<[AccountDB, () => void]> { const isMongo = uri.startsWith('mongodb://') if (isMongo) { const client = getMongoClient(uri) - const db = (await client.getClient()).db(ACCOUNT_DB) + const db = (await client.getClient()).db(dbNs ?? 'account') const mongoAccount = new MongoAccountDB(db) await mongoAccount.init() @@ -48,6 +43,7 @@ export async function getAccountDB (uri: string, db: string = ACCOUNT_DB): Promi } else { const client = getDBClient(uri) const pgClient = await client.getClient() + // TODO: if dbNs is provided put tables in that schema const pgAccount = new PostgresAccountDB(pgClient) let error = false diff --git a/server/backup/src/backup.ts b/server/backup/src/backup.ts index e3ab794c66..1aef026c39 100644 --- a/server/backup/src/backup.ts +++ b/server/backup/src/backup.ts @@ -26,6 +26,7 @@ import core, { DOMAIN_FULLTEXT_BLOB, DOMAIN_MODEL, DOMAIN_TRANSIENT, + DOMAIN_TX, MeasureContext, MeasureMetricsContext, RateLimiter, @@ -44,8 +45,9 @@ import { type StorageAdapter } from '@hcengineering/server-core' import { fullTextPushStagePrefix } from '@hcengineering/server-indexer' import { generateToken } from '@hcengineering/server-token' import { connect } from '@hcengineering/server-tool' -import { createWriteStream, existsSync, mkdirSync } from 'node:fs' -import { dirname } from 'node:path' +import { createReadStream, createWriteStream, existsSync, mkdirSync } from 'node:fs' +import { rm } from 'node:fs/promises' +import { basename, dirname } from 'node:path' import { PassThrough } from 'node:stream' import { createGzip } from 'node:zlib' import { join } from 'path' @@ -187,6 +189,190 @@ async function loadDigest ( ctx.end() return result } +async function verifyDigest ( + ctx: MeasureContext, + storage: BackupStorage, + snapshots: BackupSnapshot[], + domain: Domain +): Promise { + ctx = ctx.newChild('verify digest', { domain, count: snapshots.length }) + ctx.info('verify-digest', { domain, count: snapshots.length }) + let modified = false + for (const s of snapshots) { + const d = s.domains[domain] + if (d === undefined) { + continue + } + + const storageToRemove = new Set() + // We need to verify storage has all necessary resources + ctx.info('checking', { domain }) + // We have required documents here. + const validDocs = new Set>() + + for (const sf of d.storage ?? []) { + const blobs = new Map() + try { + ctx.info('checking storage', { sf }) + const readStream = await storage.load(sf) + const ex = extract() + + ex.on('entry', (headers, stream, next) => { + const name = headers.name ?? '' + // We found blob data + if (name.endsWith('.json')) { + const chunks: Buffer[] = [] + const bname = name.substring(0, name.length - 5) + stream.on('data', (chunk) => { + chunks.push(chunk) + }) + stream.on('end', () => { + const bf = Buffer.concat(chunks as any) + const doc = JSON.parse(bf.toString()) as Doc + if (doc._class === core.class.Blob || doc._class === 'core:class:BlobData') { + const data = migradeBlobData(doc as Blob, '') + const d = blobs.get(bname) ?? (data !== '' ? Buffer.from(data, 'base64') : undefined) + if (d === undefined) { + blobs.set(bname, { doc, buffer: undefined }) + } else { + blobs.delete(bname) + const blob = doc as Blob + + if (blob.size === bf.length) { + validDocs.add(name as Ref) + } + } + } else { + validDocs.add(name as Ref) + } + next() + }) + } else { + const chunks: Buffer[] = [] + stream.on('data', (chunk) => { + chunks.push(chunk) + }) + stream.on('end', () => { + const bf = Buffer.concat(chunks as any) + const d = blobs.get(name) + if (d === undefined) { + blobs.set(name, { doc: undefined, buffer: bf }) + } else { + blobs.delete(name) + const doc = d?.doc as Blob + let sz = doc.size + if (Number.isNaN(sz) || sz !== bf.length) { + sz = bf.length + } + + // If blob size matches doc size, remove from requiredDocs + if (sz === bf.length) { + validDocs.add(name as Ref) + } + } + next() + }) + } + stream.resume() // just auto drain the stream + }) + + const unzip = createGunzip({ level: defaultLevel }) + const endPromise = new Promise((resolve) => { + ex.on('finish', () => { + resolve(null) + }) + unzip.on('error', (err) => { + ctx.error('error during reading of', { sf, err }) + modified = true + storageToRemove.add(sf) + resolve(null) + }) + }) + + readStream.on('end', () => { + readStream.destroy() + }) + readStream.pipe(unzip) + unzip.pipe(ex) + + await endPromise + } catch (err: any) { + ctx.error('error during reading of', { sf, err }) + // In case of invalid archive, we need to + // We need to remove broken storage file + modified = true + storageToRemove.add(sf) + } + } + if (storageToRemove.size > 0) { + modified = true + d.storage = (d.storage ?? []).filter((it) => !storageToRemove.has(it)) + } + + // if (d?.snapshot !== undefined) { + // Will not check old format + // } + const digestToRemove = new Set() + for (const snapshot of d?.snapshots ?? []) { + try { + ctx.info('checking', { snapshot }) + const changes: Snapshot = { + added: new Map(), + removed: [], + updated: new Map() + } + let lmodified = false + try { + const dataBlob = gunzipSync(await storage.loadFile(snapshot)) + .toString() + .split('\n') + const addedCount = parseInt(dataBlob.shift() ?? '0') + const added = dataBlob.splice(0, addedCount) + for (const it of added) { + const [k, v] = it.split(';') + if (validDocs.has(k as any)) { + changes.added.set(k as Ref, v) + } else { + lmodified = true + } + } + + const updatedCount = parseInt(dataBlob.shift() ?? '0') + const updated = dataBlob.splice(0, updatedCount) + for (const it of updated) { + const [k, v] = it.split(';') + if (validDocs.has(k as any)) { + changes.updated.set(k as Ref, v) + } else { + lmodified = true + } + } + + const removedCount = parseInt(dataBlob.shift() ?? '0') + const removed = dataBlob.splice(0, removedCount) + changes.removed = removed as Ref[] + } catch (err: any) { + ctx.warn('failed during processing of snapshot file, it will be skipped', { snapshot }) + digestToRemove.add(snapshot) + modified = true + } + + if (lmodified) { + modified = true + // Store changes without missing files + await writeChanges(storage, snapshot, changes) + } + } catch (err: any) { + digestToRemove.add(snapshot) + ctx.error('digest is broken, will do full backup for', { domain }) + modified = true + } + } + d.snapshots = (d.snapshots ?? []).filter((it) => !digestToRemove.has(it)) + } + ctx.end() + return modified +} async function write (chunk: any, stream: Writable): Promise { let needDrain = false @@ -662,6 +848,13 @@ export async function backup ( (options.include === undefined || options.include.has(it)) ) ] + domains.sort((a, b) => { + if (a === DOMAIN_TX) { + return -1 + } + + return a.localeCompare(b) + }) ctx.info('domains for dump', { domains: domains.length }) @@ -863,12 +1056,15 @@ export async function backup ( const digest = await ctx.with('load-digest', {}, (ctx) => loadDigest(ctx, storage, backupInfo.snapshots, domain)) let _pack: Pack | undefined + let _packClose = async (): Promise => {} let addedDocuments = (): number => 0 progress(0) let { changed, needRetrieveChunks } = await ctx.with('load-chunks', { domain }, (ctx) => loadChangesFromServer(ctx, domain, digest, changes) ) + processedChanges.removed = Array.from(digest.keys()) + digest.clear() progress(10) if (needRetrieveChunks.length > 0) { @@ -879,6 +1075,10 @@ export async function backup ( let processed = 0 let blobs = 0 + try { + global.gc?.() + } catch (err) {} + while (needRetrieveChunks.length > 0) { if (canceled()) { return @@ -910,11 +1110,16 @@ export async function backup ( while (docs.length > 0) { // Chunk data into small pieces - if (addedDocuments() > dataBlobSize && _pack !== undefined) { - _pack.finalize() - _pack = undefined + if ( + (addedDocuments() > dataBlobSize || processedChanges.added.size + processedChanges.updated.size > 500000) && + _pack !== undefined + ) { + await _packClose() if (changed > 0) { + try { + global.gc?.() + } catch (err) {} snapshot.domains[domain] = domainInfo domainInfo.added += processedChanges.added.size domainInfo.updated += processedChanges.updated.size @@ -940,7 +1145,9 @@ export async function backup ( const storageFile = join(backupIndex, `${domain}-data-${snapshot.date}-${stIndex}.tar.gz`) ctx.info('storing from domain', { domain, storageFile, workspace: workspaceId.name }) domainInfo.storage = [...(domainInfo.storage ?? []), storageFile] - const dataStream = await storage.write(storageFile) + const tmpFile = basename(storageFile) + '.tmp' + const tempFile = createWriteStream(tmpFile) + // const dataStream = await storage.write(storageFile) const sizePass = new PassThrough() let sz = 0 @@ -951,12 +1158,28 @@ export async function backup ( cb() } - sizePass.pipe(dataStream) + sizePass.pipe(tempFile) const storageZip = createGzip({ level: defaultLevel, memLevel: 9 }) addedDocuments = () => sz _pack.pipe(storageZip) storageZip.pipe(sizePass) + + _packClose = async () => { + await new Promise((resolve) => { + tempFile.on('close', () => { + resolve() + }) + _pack?.finalize() + }) + + // We need to upload file to storage + ctx.info('Upload pack file', { storageFile, size: sz, workspace: workspaceId.name }) + await storage.writeFile(storageFile, createReadStream(tmpFile)) + await rm(tmpFile) + + _pack = undefined + } } if (canceled()) { return @@ -1025,7 +1248,7 @@ export async function backup ( } }) - const finalBuffer = Buffer.concat(buffers) + const finalBuffer = Buffer.concat(buffers as any) if (finalBuffer.length !== blob.size) { ctx.error('download blob size mismatch', { _id: blob._id, @@ -1078,7 +1301,7 @@ export async function backup ( } } } - processedChanges.removed = Array.from(digest.keys()) + if (processedChanges.removed.length > 0) { changed++ } @@ -1097,7 +1320,7 @@ export async function backup ( processedChanges.added.clear() processedChanges.removed = [] processedChanges.updated.clear() - _pack?.finalize() + await _packClose() // This will allow to retry in case of critical error. await storage.writeFile(infoFile, gzipSync(JSON.stringify(backupInfo, undefined, 2), { level: defaultLevel })) } @@ -1108,6 +1331,14 @@ export async function backup ( if (canceled()) { break } + const oldUsed = process.memoryUsage().heapUsed + try { + global.gc?.() + } catch (err) {} + ctx.info('memory-stats', { + old: Math.round(oldUsed / (1024 * 1024)), + current: Math.round(process.memoryUsage().heapUsed / (1024 * 1024)) + }) await ctx.with('process-domain', { domain }, async (ctx) => { await processDomain(ctx, domain, (value) => { options.progress?.(Math.round(((domainProgress + value / 100) / domains.length) * 100)) @@ -1196,6 +1427,26 @@ export async function backupList (storage: BackupStorage): Promise { } } +/** + * @public + */ +export async function backupRemoveLast (storage: BackupStorage, date: number): Promise { + const infoFile = 'backup.json.gz' + + if (!(await storage.exists(infoFile))) { + throw new Error(`${infoFile} should present to restore`) + } + const backupInfo: BackupInfo = JSON.parse(gunzipSync(await storage.loadFile(infoFile)).toString()) + console.log('workspace:', backupInfo.workspace ?? '', backupInfo.version) + const old = backupInfo.snapshots.length + backupInfo.snapshots = backupInfo.snapshots.filter((it) => it.date < date) + if (old !== backupInfo.snapshots.length) { + console.log('removed snapshots: id:', old - backupInfo.snapshots.length) + + await storage.writeFile(infoFile, gzipSync(JSON.stringify(backupInfo, undefined, 2), { level: defaultLevel })) + } +} + /** * @public */ @@ -1458,6 +1709,12 @@ export async function restore ( // We need to load full changeset from server const serverChangeset = new Map, string>() + const oldUsed = process.memoryUsage().heapUsed + try { + global.gc?.() + } catch (err) {} + ctx.info('memory-stats', { old: oldUsed / (1024 * 1024), current: process.memoryUsage().heapUsed / (1024 * 1024) }) + let idx: number | undefined let loaded = 0 let el = 0 @@ -1929,7 +2186,7 @@ export async function compactBackup ( chunks.push(chunk) }) stream.on('end', () => { - const bf = Buffer.concat(chunks) + const bf = Buffer.concat(chunks as any) const d = blobs.get(name) if (d === undefined) { blobs.set(name, { doc: undefined, buffer: bf }) @@ -1980,12 +2237,16 @@ export async function compactBackup ( stream.resume() // just auto drain the stream }) + const unzip = createGunzip({ level: defaultLevel }) const endPromise = new Promise((resolve) => { ex.on('finish', () => { resolve(null) }) + unzip.on('error', (err) => { + ctx.error('error during processing', { snapshot, err }) + resolve(null) + }) }) - const unzip = createGunzip({ level: defaultLevel }) readStream.on('end', () => { readStream.destroy() @@ -2060,3 +2321,53 @@ function migradeBlobData (blob: Blob, etag: string): string { } return '' } + +/** + * Will check backup integrity, and in case of some missing resources, will update digest files, so next backup will backup all missing parts. + * @public + */ +export async function checkBackupIntegrity (ctx: MeasureContext, storage: BackupStorage): Promise { + console.log('starting backup compaction') + try { + let backupInfo: BackupInfo + + // Version 0.6.2, format of digest file is changed to + + const infoFile = 'backup.json.gz' + + if (await storage.exists(infoFile)) { + backupInfo = JSON.parse(gunzipSync(await storage.loadFile(infoFile)).toString()) + } else { + console.log('No backup found') + return + } + if (backupInfo.version !== '0.6.2') { + console.log('Invalid backup version') + return + } + + const domains: Domain[] = [] + for (const sn of backupInfo.snapshots) { + for (const d of Object.keys(sn.domains)) { + if (!domains.includes(d as Domain)) { + domains.push(d as Domain) + } + } + } + let modified = false + + for (const domain of domains) { + console.log('checking domain...', domain) + if (await verifyDigest(ctx, storage, backupInfo.snapshots, domain)) { + modified = true + } + } + if (modified) { + await storage.writeFile(infoFile, gzipSync(JSON.stringify(backupInfo, undefined, 2), { level: defaultLevel })) + } + } catch (err: any) { + console.error(err) + } finally { + console.log('end compacting') + } +} diff --git a/server/backup/src/service.ts b/server/backup/src/service.ts index 20355f0bb6..55ef156a8c 100644 --- a/server/backup/src/service.ts +++ b/server/backup/src/service.ts @@ -125,7 +125,9 @@ class BackupWorker { } return !workspacesIgnore.has(it.workspace) }) - workspaces.sort((a, b) => b.lastVisit - a.lastVisit) + workspaces.sort((a, b) => { + return (b.backupInfo?.backupSize ?? 0) - (a.backupInfo?.backupSize ?? 0) + }) ctx.info('Preparing for BACKUP', { total: workspaces.length, diff --git a/server/backup/src/storage.ts b/server/backup/src/storage.ts index a4070d1582..18e845c7c5 100644 --- a/server/backup/src/storage.ts +++ b/server/backup/src/storage.ts @@ -12,7 +12,8 @@ export interface BackupStorage { loadFile: (name: string) => Promise load: (name: string) => Promise write: (name: string) => Promise - writeFile: (name: string, data: string | Buffer) => Promise + + writeFile: (name: string, data: string | Buffer | Readable) => Promise exists: (name: string) => Promise stat: (name: string) => Promise @@ -51,14 +52,14 @@ class FileStorage implements BackupStorage { await rm(join(this.root, name)) } - async writeFile (name: string, data: string | Buffer): Promise { + async writeFile (name: string, data: string | Buffer | Readable): Promise { const fileName = join(this.root, name) const dir = dirname(fileName) if (!existsSync(dir)) { await mkdir(dir, { recursive: true }) } - await writeFile(fileName, data) + await writeFile(fileName, data as any) } } @@ -72,7 +73,7 @@ class AdapterStorage implements BackupStorage { async loadFile (name: string): Promise { const data = await this.client.read(this.ctx, this.workspaceId, join(this.root, name)) - return Buffer.concat(data) + return Buffer.concat(data as any) } async write (name: string): Promise { @@ -106,16 +107,9 @@ class AdapterStorage implements BackupStorage { await this.client.remove(this.ctx, this.workspaceId, [join(this.root, name)]) } - async writeFile (name: string, data: string | Buffer): Promise { + async writeFile (name: string, data: string | Buffer | Readable): Promise { // TODO: add mime type detection here. - await this.client.put( - this.ctx, - this.workspaceId, - join(this.root, name), - data, - 'application/octet-stream', - data.length - ) + await this.client.put(this.ctx, this.workspaceId, join(this.root, name), data, 'application/octet-stream') } } diff --git a/server/datalake/src/client.ts b/server/datalake/src/client.ts index 5571afce5a..6990cb60f9 100644 --- a/server/datalake/src/client.ts +++ b/server/datalake/src/client.ts @@ -15,7 +15,7 @@ import { type MeasureContext, type WorkspaceId, concatLink } from '@hcengineering/core' import FormData from 'form-data' -import fetch from 'node-fetch' +import fetch, { type RequestInit, type Response } from 'node-fetch' import { Readable } from 'stream' /** @public */ @@ -34,11 +34,6 @@ export interface StatObjectOutput { size?: number } -/** @public */ -export interface PutObjectOutput { - id: string -} - interface BlobUploadError { key: string error: string @@ -54,7 +49,11 @@ type BlobUploadResult = BlobUploadSuccess | BlobUploadError /** @public */ export class Client { - constructor (private readonly endpoint: string) {} + private readonly endpoint: string + + constructor (host: string, port?: number) { + this.endpoint = port !== undefined ? `${host}:${port}` : host + } getObjectUrl (ctx: MeasureContext, workspace: WorkspaceId, objectName: string): string { const path = `/blob/${workspace.name}/${encodeURIComponent(objectName)}` @@ -63,21 +62,7 @@ export class Client { async getObject (ctx: MeasureContext, workspace: WorkspaceId, objectName: string): Promise { const url = this.getObjectUrl(ctx, workspace, objectName) - - let response - try { - response = await fetch(url) - } catch (err: any) { - ctx.error('network error', { error: err }) - throw new Error(`Network error ${err}`) - } - - if (!response.ok) { - if (response.status === 404) { - throw new Error('Not Found') - } - throw new Error('HTTP error ' + response.status) - } + const response = await fetchSafe(ctx, url) if (response.body == null) { ctx.error('bad datalake response', { objectName }) @@ -99,20 +84,7 @@ export class Client { Range: `bytes=${offset}-${length ?? ''}` } - let response - try { - response = await fetch(url, { headers }) - } catch (err: any) { - ctx.error('network error', { error: err }) - throw new Error(`Network error ${err}`) - } - - if (!response.ok) { - if (response.status === 404) { - throw new Error('Not Found') - } - throw new Error('HTTP error ' + response.status) - } + const response = await fetchSafe(ctx, url, { headers }) if (response.body == null) { ctx.error('bad datalake response', { objectName }) @@ -129,20 +101,7 @@ export class Client { ): Promise { const url = this.getObjectUrl(ctx, workspace, objectName) - let response - try { - response = await fetch(url, { method: 'HEAD' }) - } catch (err: any) { - ctx.error('network error', { error: err }) - throw new Error(`Network error ${err}`) - } - - if (!response.ok) { - if (response.status === 404) { - return undefined - } - throw new Error('HTTP error ' + response.status) - } + const response = await fetchSafe(ctx, url, { method: 'HEAD' }) const headers = response.headers const lastModified = Date.parse(headers.get('Last-Modified') ?? '') @@ -158,30 +117,35 @@ export class Client { async deleteObject (ctx: MeasureContext, workspace: WorkspaceId, objectName: string): Promise { const url = this.getObjectUrl(ctx, workspace, objectName) - - let response - try { - response = await fetch(url, { method: 'DELETE' }) - } catch (err: any) { - ctx.error('network error', { error: err }) - throw new Error(`Network error ${err}`) - } - - if (!response.ok) { - if (response.status === 404) { - throw new Error('Not Found') - } - throw new Error('HTTP error ' + response.status) - } + await fetchSafe(ctx, url, { method: 'DELETE' }) } async putObject ( + ctx: MeasureContext, + workspace: WorkspaceId, + objectName: string, + stream: Readable | Buffer | string, + metadata: ObjectMetadata, + size?: number + ): Promise { + if (size === undefined || size < 64 * 1024 * 1024) { + await ctx.with('direct-upload', {}, async (ctx) => { + await this.uploadWithFormData(ctx, workspace, objectName, stream, metadata) + }) + } else { + await ctx.with('signed-url-upload', {}, async (ctx) => { + await this.uploadWithSignedURL(ctx, workspace, objectName, stream, metadata) + }) + } + } + + private async uploadWithFormData ( ctx: MeasureContext, workspace: WorkspaceId, objectName: string, stream: Readable | Buffer | string, metadata: ObjectMetadata - ): Promise { + ): Promise { const path = `/upload/form-data/${workspace.name}` const url = concatLink(this.endpoint, path) @@ -196,17 +160,7 @@ export class Client { } form.append('file', stream, options) - let response - try { - response = await fetch(url, { method: 'POST', body: form }) - } catch (err: any) { - ctx.error('network error', { error: err }) - throw new Error(`Network error ${err}`) - } - - if (!response.ok) { - throw new Error('HTTP error ' + response.status) - } + const response = await fetchSafe(ctx, url, { method: 'POST', body: form }) const result = (await response.json()) as BlobUploadResult[] if (result.length !== 1) { @@ -219,8 +173,68 @@ export class Client { if ('error' in uploadResult) { ctx.error('error during blob upload', { objectName, error: uploadResult.error }) throw new Error('Upload failed: ' + uploadResult.error) - } else { - return { id: uploadResult.id } } } + + private async uploadWithSignedURL ( + ctx: MeasureContext, + workspace: WorkspaceId, + objectName: string, + stream: Readable | Buffer | string, + metadata: ObjectMetadata + ): Promise { + const url = await this.signObjectSign(ctx, workspace, objectName) + + try { + await fetchSafe(ctx, url, { + body: stream, + method: 'PUT', + headers: { + 'Content-Type': metadata.type, + 'Content-Length': metadata.size?.toString() ?? '0', + 'x-amz-meta-last-modified': metadata.lastModified.toString() + } + }) + await this.signObjectComplete(ctx, workspace, objectName) + } catch { + await this.signObjectDelete(ctx, workspace, objectName) + } + } + + private async signObjectSign (ctx: MeasureContext, workspace: WorkspaceId, objectName: string): Promise { + const url = this.getSignObjectUrl(workspace, objectName) + const response = await fetchSafe(ctx, url, { method: 'POST' }) + return await response.text() + } + + private async signObjectComplete (ctx: MeasureContext, workspace: WorkspaceId, objectName: string): Promise { + const url = this.getSignObjectUrl(workspace, objectName) + await fetchSafe(ctx, url, { method: 'PUT' }) + } + + private async signObjectDelete (ctx: MeasureContext, workspace: WorkspaceId, objectName: string): Promise { + const url = this.getSignObjectUrl(workspace, objectName) + await fetchSafe(ctx, url, { method: 'DELETE' }) + } + + private getSignObjectUrl (workspace: WorkspaceId, objectName: string): string { + const path = `/upload/signed-url/${workspace.name}/${encodeURIComponent(objectName)}` + return concatLink(this.endpoint, path) + } +} + +async function fetchSafe (ctx: MeasureContext, url: string, init?: RequestInit): Promise { + let response + try { + response = await fetch(url, init) + } catch (err: any) { + ctx.error('network error', { error: err }) + throw new Error(`Network error ${err}`) + } + + if (!response.ok) { + throw new Error(response.status === 404 ? 'Not Found' : 'HTTP error ' + response.status) + } + + return response } diff --git a/server/datalake/src/index.ts b/server/datalake/src/index.ts index 071a063c50..3e58e548e5 100644 --- a/server/datalake/src/index.ts +++ b/server/datalake/src/index.ts @@ -37,7 +37,7 @@ export class DatalakeService implements StorageAdapter { static config = 'datalake' client: Client constructor (readonly opt: DatalakeConfig) { - this.client = new Client(opt.endpoint) + this.client = new Client(opt.endpoint, opt.port) } async initialize (ctx: MeasureContext, workspaceId: WorkspaceId): Promise {} @@ -129,7 +129,7 @@ export class DatalakeService implements StorageAdapter { await ctx.with('put', {}, async (ctx) => { await withRetry(ctx, 5, async () => { - return await this.client.putObject(ctx, workspaceId, objectName, stream, metadata) + await this.client.putObject(ctx, workspaceId, objectName, stream, metadata, size) }) }) diff --git a/server/front/readme.md b/server/front/readme.md index 20b4615959..0baf88ac2c 100644 --- a/server/front/readme.md +++ b/server/front/readme.md @@ -17,6 +17,7 @@ Front service is suited to deliver application bundles and resource assets, it a * MODEL_VERSION: Specifies the required model version. * SERVER_SECRET: Specifies the server secret. * PREVIEW_CONFIG: Specifies the preview configuration. +* UPLOAD_CONFIG: Specifies the upload configuration. * BRANDING_URL: Specifies the URL of the branding service. ## Preview service configuration diff --git a/server/front/src/index.ts b/server/front/src/index.ts index a174f93b00..0a87fae70c 100644 --- a/server/front/src/index.ts +++ b/server/front/src/index.ts @@ -256,6 +256,7 @@ export function start ( collaboratorUrl: string brandingUrl?: string previewConfig: string + uploadConfig: string pushPublicKey?: string disableSignUp?: string }, @@ -308,6 +309,7 @@ export function start ( COLLABORATOR_URL: config.collaboratorUrl, BRANDING_URL: config.brandingUrl, PREVIEW_CONFIG: config.previewConfig, + UPLOAD_CONFIG: config.uploadConfig, PUSH_PUBLIC_KEY: config.pushPublicKey, DISABLE_SIGNUP: config.disableSignUp, ...(extraConfig ?? {}) @@ -501,8 +503,15 @@ export function start ( void filesHandler(req, res) }) - // eslint-disable-next-line @typescript-eslint/no-misused-promises - app.post('/files', async (req, res) => { + app.post('/files', (req, res) => { + void handleUpload(req, res) + }) + + app.post('/files/*', (req, res) => { + void handleUpload(req, res) + }) + + const handleUpload = async (req: Request, res: Response): Promise => { await ctx.with( 'post-file', {}, @@ -538,7 +547,7 @@ export function start ( }, { url: req.path, query: req.query } ) - }) + } const handleDelete = async (req: Request, res: Response): Promise => { try { diff --git a/server/front/src/starter.ts b/server/front/src/starter.ts index e1301978f7..8823fb27f1 100644 --- a/server/front/src/starter.ts +++ b/server/front/src/starter.ts @@ -101,6 +101,11 @@ export function startFront (ctx: MeasureContext, extraConfig?: Record { + const sql = postgres(env.HYPERDRIVE.connectionString) + const { bucket } = selectStorage(env, workspace) + + const blob = await db.getBlob(sql, { workspace, name }) + if (blob === null || blob.deleted) { + return error(404) + } + + const cache = caches.default + const cached = await cache.match(request) + if (cached !== undefined) { + return cached + } + + const range = request.headers.has('Range') ? request.headers : undefined + const object = await bucket.get(blob.filename, { range }) + if (object === null) { + return error(404) + } + + const headers = r2MetadataHeaders(object) + if (range !== undefined && object?.range !== undefined) { + headers.set('Content-Range', rangeHeader(object.range, object.size)) + } + + const length = object?.range !== undefined && 'length' in object.range ? object?.range?.length : undefined + const status = length !== undefined && length < object.size ? 206 : 200 + + const response = new Response(object?.body, { headers, status }) + ctx.waitUntil(cache.put(request, response.clone())) + + return response +} + +export async function handleBlobHead ( + request: Request, + env: Env, + ctx: ExecutionContext, + workspace: string, + name: string +): Promise { + const sql = postgres(env.HYPERDRIVE.connectionString) + const { bucket } = selectStorage(env, workspace) + + const blob = await db.getBlob(sql, { workspace, name }) + if (blob === null) { + return error(404) + } + + const head = await bucket.head(blob.filename) + if (head?.httpMetadata === undefined) { + return error(404) + } + + const headers = r2MetadataHeaders(head) + return new Response(null, { headers, status: 200 }) +} + +export async function deleteBlob (env: Env, workspace: string, name: string): Promise { + const sql = postgres(env.HYPERDRIVE.connectionString) + + try { + await Promise.all([db.deleteBlob(sql, { workspace, name }), deleteVideo(env, workspace, name)]) + + return new Response(null, { status: 204 }) + } catch (err: any) { + const message = err instanceof Error ? err.message : String(err) + console.error({ error: 'failed to delete blob:' + message }) + return error(500) + } +} + +export async function postBlobFormData (request: Request, env: Env, workspace: string): Promise { + const sql = postgres(env.HYPERDRIVE.connectionString) + const formData = await request.formData() + + const files: [File, key: string][] = [] + formData.forEach((value: any, key: string) => { + if (typeof value === 'object') files.push([value, key]) + }) + + const result = await Promise.all( + files.map(async ([file, key]) => { + const { name, type, lastModified } = file + try { + const metadata = await saveBlob(env, sql, file, type, workspace, name, lastModified) + + // TODO this probably should happen via queue, let it be here for now + if (type.startsWith('video/')) { + const blobURL = getBlobURL(request, workspace, name) + await copyVideo(env, blobURL, workspace, name) + } + + return { key, metadata } + } catch (err: any) { + const error = err instanceof Error ? err.message : String(err) + console.error('failed to upload blob:', error) + return { key, error } + } + }) + ) + + return json(result) +} + +async function saveBlob ( + env: Env, + sql: postgres.Sql, + file: File, + type: string, + workspace: string, + name: string, + lastModified: number +): Promise { + const { location, bucket } = selectStorage(env, workspace) + + const size = file.size + const [mimetype, subtype] = type.split('/') + const httpMetadata = { contentType: type, cacheControl } + const filename = getUniqueFilename() + + const sha256hash = await getSha256(file) + + if (sha256hash !== null) { + // Lucky boy, nothing to upload, use existing blob + const hash = sha256hash + + const data = await db.getData(sql, { hash, location }) + if (data !== null) { + await db.createBlob(sql, { workspace, name, hash, location }) + } else { + await bucket.put(filename, file, { httpMetadata }) + await sql.begin((sql) => [ + db.createData(sql, { hash, location, filename, type: mimetype, subtype, size }), + db.createBlob(sql, { workspace, name, hash, location }) + ]) + } + + return { type, size, lastModified, name } + } else { + // For large files we cannot calculate checksum beforehead + // upload file with unique filename and then obtain checksum + const object = await bucket.put(filename, file, { httpMetadata }) + + const hash = + object.checksums.md5 !== undefined ? getMd5Checksum(object.checksums.md5) : (crypto.randomUUID() as UUID) + + const data = await db.getData(sql, { hash, location }) + if (data !== null) { + // We found an existing blob with the same hash + // we can safely remove the existing blob from storage + await Promise.all([bucket.delete(filename), db.createBlob(sql, { workspace, name, hash, location })]) + } else { + // Otherwise register a new hash and blob + await sql.begin((sql) => [ + db.createData(sql, { hash, location, filename, type: mimetype, subtype, size }), + db.createBlob(sql, { workspace, name, hash, location }) + ]) + } + + return { type, size, lastModified, name } + } +} + +export async function handleBlobUploaded (env: Env, workspace: string, name: string, filename: UUID): Promise { + const sql = postgres(env.HYPERDRIVE.connectionString) + const { location, bucket } = selectStorage(env, workspace) + + const object = await bucket.head(filename) + if (object?.httpMetadata === undefined) { + throw Error('blob not found') + } + + const hash = object.checksums.md5 !== undefined ? getMd5Checksum(object.checksums.md5) : (crypto.randomUUID() as UUID) + + const data = await db.getData(sql, { hash, location }) + if (data !== null) { + await Promise.all([bucket.delete(filename), db.createBlob(sql, { workspace, name, hash, location })]) + } else { + const size = object.size + const type = object.httpMetadata.contentType ?? 'application/octet-stream' + const [mimetype, subtype] = type.split('/') + + await db.createData(sql, { hash, location, filename, type: mimetype, subtype, size }) + await db.createBlob(sql, { workspace, name, hash, location }) + } +} + +function getUniqueFilename (): UUID { + return crypto.randomUUID() as UUID +} + +async function getSha256 (file: File): Promise { + if (file.size > HASH_LIMIT) { + return null + } + + const digestStream = new crypto.DigestStream('SHA-256') + await file.stream().pipeTo(digestStream) + const digest = await digestStream.digest + + return toUUID(new Uint8Array(digest)) +} + +function getMd5Checksum (digest: ArrayBuffer): UUID { + return toUUID(new Uint8Array(digest)) +} + +function rangeHeader (range: R2Range, size: number): string { + const offset = 'offset' in range ? range.offset : undefined + const length = 'length' in range ? range.length : undefined + const suffix = 'suffix' in range ? range.suffix : undefined + + const start = suffix !== undefined ? size - suffix : offset ?? 0 + const end = suffix !== undefined ? size : length !== undefined ? start + length : size + + return `bytes ${start}-${end - 1}/${size}` +} + +function r2MetadataHeaders (head: R2Object): Headers { + return head.httpMetadata !== undefined + ? new Headers({ + 'Accept-Ranges': 'bytes', + 'Content-Length': head.size.toString(), + 'Content-Type': head.httpMetadata.contentType ?? '', + 'Cache-Control': head.httpMetadata.cacheControl ?? cacheControl, + 'Last-Modified': head.uploaded.toUTCString(), + ETag: head.httpEtag + }) + : new Headers({ + 'Accept-Ranges': 'bytes', + 'Content-Length': head.size.toString(), + 'Cache-Control': cacheControl, + 'Last-Modified': head.uploaded.toUTCString(), + ETag: head.httpEtag + }) +} diff --git a/workers/datalake/src/cors.ts b/workers/datalake/src/cors.ts new file mode 100644 index 0000000000..978a2d31dd --- /dev/null +++ b/workers/datalake/src/cors.ts @@ -0,0 +1,103 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { type IRequest } from 'itty-router' + +// This is a copy of cors.ts from itty-router with following issues fixed: +// - https://github.com/kwhitley/itty-router/issues/242 +// - https://github.com/kwhitley/itty-router/issues/249 +export interface CorsOptions { + credentials?: true + origin?: boolean | string | string[] | RegExp | ((origin: string) => string | undefined) + maxAge?: number + allowMethods?: string | string[] + allowHeaders?: any + exposeHeaders?: string | string[] +} + +export type Preflight = (request: IRequest) => Response | undefined +export type Corsify = (response: Response, request?: IRequest) => Response | undefined + +export interface CorsPair { + preflight: Preflight + corsify: Corsify +} + +// Create CORS function with default options. +export const cors = (options: CorsOptions = {}): CorsPair => { + // Destructure and set defaults for options. + const { origin = '*', credentials = false, allowMethods = '*', allowHeaders, exposeHeaders, maxAge } = options + + const getAccessControlOrigin = (request?: Request): string | null | undefined => { + const requestOrigin = request?.headers.get('origin') // may be null if no request passed + if (requestOrigin === undefined || requestOrigin === null) return requestOrigin + + if (origin === true) return requestOrigin + if (origin instanceof RegExp) return origin.test(requestOrigin) ? requestOrigin : undefined + if (Array.isArray(origin)) return origin.includes(requestOrigin) ? requestOrigin : undefined + if (origin instanceof Function) return origin(requestOrigin) ?? undefined + + return origin === '*' && credentials ? requestOrigin : (origin as string) + } + + const appendHeadersAndReturn = (response: Response, headers: Record): Response => { + for (const [key, value] of Object.entries(headers)) { + if (value !== undefined && value !== null && value !== '') { + response.headers.append(key, value) + } + } + return response + } + + const preflight = (request: Request): Response | undefined => { + if (request.method === 'OPTIONS') { + const response = new Response(null, { status: 204 }) + + const allowMethodsHeader = Array.isArray(allowMethods) ? allowMethods.join(',') : allowMethods + const allowHeadersHeader = Array.isArray(allowHeaders) ? allowHeaders.join(',') : allowHeaders + const exposeHeadersHeader = Array.isArray(exposeHeaders) ? exposeHeaders.join(',') : exposeHeaders + + return appendHeadersAndReturn(response, { + 'access-control-allow-origin': getAccessControlOrigin(request), + 'access-control-allow-methods': allowMethodsHeader, + 'access-control-expose-headers': exposeHeadersHeader, + 'access-control-allow-headers': allowHeadersHeader ?? request.headers.get('access-control-request-headers'), + 'access-control-max-age': maxAge, + 'access-control-allow-credentials': credentials + }) + } // otherwise ignore + } + + const corsify = (response: Response, request?: Request): Response | undefined => { + // ignore if already has CORS headers + if (response?.headers?.has('access-control-allow-origin') || response.status === 101) { + return response + } + + const responseCopy = new Response(response.body, { + status: response.status, + statusText: response.statusText, + headers: response.headers + }) + + return appendHeadersAndReturn(responseCopy, { + 'access-control-allow-origin': getAccessControlOrigin(request), + 'access-control-allow-credentials': credentials + }) + } + + // Return corsify and preflight methods. + return { corsify, preflight } +} diff --git a/workers/datalake/src/db.ts b/workers/datalake/src/db.ts new file mode 100644 index 0000000000..ad7ea03f33 --- /dev/null +++ b/workers/datalake/src/db.ts @@ -0,0 +1,105 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import type postgres from 'postgres' +import { type Location, type UUID } from './types' + +export interface BlobDataId { + hash: UUID + location: Location +} + +export interface BlobDataRecord extends BlobDataId { + filename: UUID + size: number + type: string + subtype: string +} + +export interface BlobId { + workspace: string + name: string +} + +export interface BlobRecord extends BlobId { + hash: UUID + location: Location + deleted: boolean +} + +export interface BlobRecordWithFilename extends BlobRecord { + filename: string +} + +export async function getData (sql: postgres.Sql, dataId: BlobDataId): Promise { + const { hash, location } = dataId + + const rows = await sql` + SELECT hash, location, filename, size, type, subtype + FROM blob.data + WHERE hash = ${hash} AND location = ${location} + ` + + if (rows.length > 0) { + return rows[0] + } + + return null +} + +export async function createData (sql: postgres.Sql, data: BlobDataRecord): Promise { + const { hash, location, filename, size, type, subtype } = data + + await sql` + UPSERT INTO blob.data (hash, location, filename, size, type, subtype) + VALUES (${hash}, ${location}, ${filename}, ${size}, ${type}, ${subtype}) + ` +} + +export async function getBlob (sql: postgres.Sql, blobId: BlobId): Promise { + const { workspace, name } = blobId + + const rows = await sql` + SELECT b.workspace, b.name, b.hash, b.location, b.deleted, d.filename + FROM blob.blob AS b + JOIN blob.data AS d ON b.hash = d.hash AND b.location = d.location + WHERE b.workspace = ${workspace} AND b.name = ${name} + ` + + if (rows.length > 0) { + return rows[0] + } + + return null +} + +export async function createBlob (sql: postgres.Sql, blob: Omit): Promise { + const { workspace, name, hash, location } = blob + + await sql` + UPSERT INTO blob.blob (workspace, name, hash, location, deleted) + VALUES (${workspace}, ${name}, ${hash}, ${location}, false) + ` +} + +export async function deleteBlob (sql: postgres.Sql, blob: BlobId): Promise { + const { workspace, name } = blob + + await sql` + UPDATE blob.blob + SET deleted = true + WHERE workspace = ${workspace} AND name = ${name} + ` +} diff --git a/workers/datalake/src/encodings.ts b/workers/datalake/src/encodings.ts new file mode 100644 index 0000000000..7cf3b9347d --- /dev/null +++ b/workers/datalake/src/encodings.ts @@ -0,0 +1,37 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { type UUID } from './types' + +export const toUUID = (buffer: Uint8Array): UUID => { + const hex = toHex(buffer) + const hex32 = hex.slice(0, 32).padStart(32, '0') + return formatHexAsUUID(hex32) +} + +export const toHex = (buffer: Uint8Array): string => { + return Array.from(buffer) + .map((b) => b.toString(16).padStart(2, '0')) + .join('') +} + +export const etag = (id: string): string => `"${id}"` + +export function formatHexAsUUID (hexString: string): UUID { + if (hexString.length !== 32) { + throw new Error('Hex string must be exactly 32 characters long.') + } + return hexString.replace(/^(.{8})(.{4})(.{4})(.{4})(.{12})$/, '$1-$2-$3-$4-$5') as UUID +} diff --git a/workers/datalake/src/image.ts b/workers/datalake/src/image.ts new file mode 100644 index 0000000000..c309d114b2 --- /dev/null +++ b/workers/datalake/src/image.ts @@ -0,0 +1,50 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { getBlobURL } from './blob' + +const prefferedImageFormats = ['webp', 'avif', 'jpeg', 'png'] + +export async function getImage ( + request: Request, + workspace: string, + name: string, + transform: string +): Promise { + const Accept = request.headers.get('Accept') ?? 'image/*' + const image: Record = {} + + // select format based on Accept header + const formats = Accept.split(',') + for (const format of formats) { + const [type] = format.split(';') + const [clazz, kind] = type.split('/') + if (clazz === 'image' && prefferedImageFormats.includes(kind)) { + image.format = kind + break + } + } + + // apply transforms + transform.split(',').reduce((acc, param) => { + const [key, value] = param.split('=') + acc[key] = value + return acc + }, image) + + const blobURL = getBlobURL(request, workspace, name) + const imageRequest = new Request(blobURL, { headers: { Accept } }) + return await fetch(imageRequest, { cf: { image, cacheTtl: 3600 } }) +} diff --git a/workers/datalake/src/index.ts b/workers/datalake/src/index.ts new file mode 100644 index 0000000000..ffd0fed9e1 --- /dev/null +++ b/workers/datalake/src/index.ts @@ -0,0 +1,73 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { type IRequest, Router, error, html } from 'itty-router' +import { + deleteBlob as handleBlobDelete, + handleBlobGet, + handleBlobHead, + postBlobFormData as handleUploadFormData +} from './blob' +import { cors } from './cors' +import { getImage as handleImageGet } from './image' +import { getVideoMeta as handleVideoMetaGet } from './video' +import { handleSignAbort, handleSignComplete, handleSignCreate } from './sign' + +const { preflight, corsify } = cors({ + maxAge: 86400 +}) + +export default { + async fetch (request, env, ctx): Promise { + const router = Router({ + 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® Datalake™ https://huly.io + © 2024 Huly Labs` + ) + ) + .all('*', () => error(404)) + + return await router.fetch(request).catch(error) + } +} satisfies ExportedHandler diff --git a/workers/datalake/src/sign.ts b/workers/datalake/src/sign.ts new file mode 100644 index 0000000000..e1b2ed1aad --- /dev/null +++ b/workers/datalake/src/sign.ts @@ -0,0 +1,136 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { AwsClient } from 'aws4fetch' +import { error } from 'itty-router' + +import { handleBlobUploaded } from './blob' +import { type UUID } from './types' +import { selectStorage, type Storage } from './storage' + +const S3_SIGNED_LINK_TTL = 3600 + +interface SignBlobInfo { + uuid: UUID +} + +function signBlobKey (workspace: string, name: string): string { + return `s/${workspace}/${name}` +} + +function getS3Client (storage: Storage): AwsClient { + return new AwsClient({ + service: 's3', + region: 'auto', + accessKeyId: storage.bucketAccessKey, + secretAccessKey: storage.bucketSecretKey + }) +} + +export async function handleSignCreate ( + request: Request, + env: Env, + ctx: ExecutionContext, + workspace: string, + name: string +): Promise { + const storage = selectStorage(env, workspace) + const accountId = env.R2_ACCOUNT_ID + + const key = signBlobKey(workspace, name) + const uuid = crypto.randomUUID() as UUID + + // Generate R2 object link + const url = new URL(`https://${storage.bucketName}.${accountId}.r2.cloudflarestorage.com`) + url.pathname = uuid + url.searchParams.set('X-Amz-Expires', S3_SIGNED_LINK_TTL.toString()) + + // Sign R2 object link + let signed: Request + try { + const client = getS3Client(storage) + + signed = await client.sign(new Request(url, { method: 'PUT' }), { aws: { signQuery: true } }) + } catch (err: any) { + console.error({ error: 'failed to generate signed url', message: `${err}` }) + return error(500, 'failed to generate signed url') + } + + // Save upload details + const s3BlobInfo: SignBlobInfo = { uuid } + await env.datalake_blobs.put(key, JSON.stringify(s3BlobInfo), { expirationTtl: S3_SIGNED_LINK_TTL }) + + const headers = new Headers({ + Expires: new Date(Date.now() + S3_SIGNED_LINK_TTL * 1000).toISOString() + }) + return new Response(signed.url, { status: 200, headers }) +} + +export async function handleSignComplete ( + request: Request, + env: Env, + ctx: ExecutionContext, + workspace: string, + name: string +): Promise { + const { bucket } = selectStorage(env, workspace) + const key = signBlobKey(workspace, name) + + // Ensure we generated presigned URL earlier + // TODO what if we came after expiration date? + const signBlobInfo = await env.datalake_blobs.get(key, { type: 'json' }) + if (signBlobInfo === null) { + console.error({ error: 'blob sign info not found', workspace, name }) + return error(404) + } + + // Ensure the blob has been uploaded + const { uuid } = signBlobInfo + const head = await bucket.get(uuid) + if (head === null) { + console.error({ error: 'blob not found', workspace, name, uuid }) + return error(400) + } + + try { + await handleBlobUploaded(env, workspace, name, uuid) + } catch (err) { + const message = err instanceof Error ? err.message : String(err) + console.error({ error: message, workspace, name, uuid }) + return error(500, 'failed to upload blob') + } + + await env.datalake_blobs.delete(key) + + return new Response(null, { status: 201 }) +} + +export async function handleSignAbort ( + request: Request, + env: Env, + ctx: ExecutionContext, + workspace: string, + name: string +): Promise { + const key = signBlobKey(workspace, name) + + // Check if the blob has been uploaded + const s3BlobInfo = await env.datalake_blobs.get(key, { type: 'json' }) + if (s3BlobInfo !== null) { + await env.datalake_blobs.delete(key) + } + + return new Response(null, { status: 204 }) +} diff --git a/workers/datalake/src/storage.ts b/workers/datalake/src/storage.ts new file mode 100644 index 0000000000..777ab03317 --- /dev/null +++ b/workers/datalake/src/storage.ts @@ -0,0 +1,76 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the 'License'); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an 'AS IS' BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { type Location } from './types' + +export interface Storage { + location: Location + bucket: R2Bucket + + bucketName: string + bucketAccessKey: string + bucketSecretKey: string +} + +export function selectStorage (env: Env, workspace: string): Storage { + const location = selectLocation(env, workspace) + switch (location) { + case 'apac': + return { + location, + bucket: env.DATALAKE_APAC, + bucketName: env.DATALAKE_APAC_BUCKET_NAME, + bucketAccessKey: env.DATALAKE_APAC_ACCESS_KEY, + bucketSecretKey: env.DATALAKE_APAC_SECRET_KEY + } + case 'eeur': + return { + location, + bucket: env.DATALAKE_EEUR, + bucketName: env.DATALAKE_EEUR_BUCKET_NAME, + bucketAccessKey: env.DATALAKE_EEUR_ACCESS_KEY, + bucketSecretKey: env.DATALAKE_EEUR_SECRET_KEY + } + case 'weur': + return { + location, + bucket: env.DATALAKE_WEUR, + bucketName: env.DATALAKE_WEUR_BUCKET_NAME, + bucketAccessKey: env.DATALAKE_WEUR_ACCESS_KEY, + bucketSecretKey: env.DATALAKE_WEUR_SECRET_KEY + } + case 'enam': + return { + location, + bucket: env.DATALAKE_ENAM, + bucketName: env.DATALAKE_ENAM_BUCKET_NAME, + bucketAccessKey: env.DATALAKE_ENAM_ACCESS_KEY, + bucketSecretKey: env.DATALAKE_ENAM_SECRET_KEY + } + case 'wnam': + return { + location, + bucket: env.DATALAKE_WNAM, + bucketName: env.DATALAKE_WNAM_BUCKET_NAME, + bucketAccessKey: env.DATALAKE_WNAM_ACCESS_KEY, + bucketSecretKey: env.DATALAKE_WNAM_SECRET_KEY + } + } +} + +function selectLocation (env: Env, workspace: string): Location { + // TODO select location based on workspace + return 'weur' +} diff --git a/workers/datalake/src/types.ts b/workers/datalake/src/types.ts new file mode 100644 index 0000000000..5e5baad474 --- /dev/null +++ b/workers/datalake/src/types.ts @@ -0,0 +1,32 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +export type Location = 'weur' | 'eeur' | 'wnam' | 'enam' | 'apac' + +export type UUID = string & { __uuid: true } + +export interface CloudflareResponse { + success: boolean + errors: any + messages: any + result: any +} + +export interface StreamUploadResponse extends CloudflareResponse { + result: { + uid: string + uploadURL: string + } +} diff --git a/workers/datalake/src/video.ts b/workers/datalake/src/video.ts new file mode 100644 index 0000000000..229490c009 --- /dev/null +++ b/workers/datalake/src/video.ts @@ -0,0 +1,121 @@ +// +// Copyright © 2024 Hardcore Engineering Inc. +// +// Licensed under the Eclipse Public License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. You may +// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// +// See the License for the specific language governing permissions and +// limitations under the License. +// + +import { error, json } from 'itty-router' + +import { type CloudflareResponse, type StreamUploadResponse } from './types' + +export type StreamUploadState = 'ready' | 'error' | 'inprogress' | 'queued' | 'downloading' | 'pendingupload' + +// https://developers.cloudflare.com/api/operations/stream-videos-list-videos#response-body +export interface StreamDetailsResponse extends CloudflareResponse { + result: { + uid: string + thumbnail: string + status: { + state: StreamUploadState + } + playback: { + hls: string + dash: string + } + } +} + +interface StreamBlobInfo { + streamId: string +} + +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 { + const key = streamBlobKey(workspace, name) + + const streamInfo = await env.datalake_blobs.get(key, { type: 'json' }) + if (streamInfo === null) { + return error(404) + } + + const url = `https://api.cloudflare.com/client/v4/accounts/${env.STREAMS_ACCOUNT_ID}/stream/${streamInfo.streamId}` + const streamRequest = new Request(url, { + headers: { + Authorization: `Bearer ${env.STREAMS_AUTH_KEY}`, + 'Content-Type': 'application/json' + } + }) + + const streamResponse = await fetch(streamRequest) + const stream = await streamResponse.json() + + if (stream.success) { + return json({ + status: stream.result.status.state, + thumbnail: stream.result.thumbnail, + hls: stream.result.playback.hls + }) + } else { + return error(500, { errors: stream.errors }) + } +} + +export async function copyVideo (env: Env, source: string, workspace: string, name: string): Promise { + const key = streamBlobKey(workspace, name) + + const url = `https://api.cloudflare.com/client/v4/accounts/${env.STREAMS_ACCOUNT_ID}/stream/copy` + const request = new Request(url, { + method: 'POST', + headers: { + Authorization: `Bearer ${env.STREAMS_AUTH_KEY}`, + 'Content-Type': 'application/json' + }, + body: JSON.stringify({ url: source, meta: { name } }) + }) + + const response = await fetch(request) + const upload = await response.json() + + if (upload.success) { + const streamInfo: StreamBlobInfo = { + streamId: upload.result.uid + } + await env.datalake_blobs.put(key, JSON.stringify(streamInfo)) + } +} + +export async function deleteVideo (env: Env, workspace: string, name: string): Promise { + const key = streamBlobKey(workspace, name) + + const streamInfo = await env.datalake_blobs.get(key, { type: 'json' }) + if (streamInfo !== null) { + const url = `https://api.cloudflare.com/client/v4/accounts/${env.STREAMS_ACCOUNT_ID}/stream/${streamInfo.streamId}` + const request = new Request(url, { + method: 'DELETE', + headers: { + Authorization: `Bearer ${env.STREAMS_AUTH_KEY}`, + 'Content-Type': 'application/json' + } + }) + + await Promise.all([fetch(request), env.datalake_blobs.delete(key)]) + } +} diff --git a/workers/datalake/tsconfig.json b/workers/datalake/tsconfig.json new file mode 100644 index 0000000000..da8672e6cb --- /dev/null +++ b/workers/datalake/tsconfig.json @@ -0,0 +1,12 @@ +{ + "extends": "./node_modules/@hcengineering/platform-rig/profiles/default/tsconfig.json", + + "compilerOptions": { + "rootDir": "./src", + "outDir": "./lib", + "declarationDir": "./types", + "tsBuildInfoFile": ".build/build.tsbuildinfo", + "types": ["@cloudflare/workers-types", "jest"], + "lib": ["esnext"] + } +} \ No newline at end of file diff --git a/workers/datalake/worker-configuration.d.ts b/workers/datalake/worker-configuration.d.ts new file mode 100644 index 0000000000..dfc1e51e8e --- /dev/null +++ b/workers/datalake/worker-configuration.d.ts @@ -0,0 +1,30 @@ +// Generated by Wrangler on Sat Jul 06 2024 18:52:21 GMT+0200 (Central European Summer Time) +// by running `wrangler types` + +interface Env { + datalake_blobs: KVNamespace; + DATALAKE_APAC: R2Bucket; + DATALAKE_EEUR: R2Bucket; + DATALAKE_WEUR: R2Bucket; + DATALAKE_ENAM: R2Bucket; + DATALAKE_WNAM: R2Bucket; + HYPERDRIVE: Hyperdrive; + STREAMS_ACCOUNT_ID: string; + STREAMS_AUTH_KEY: string; + R2_ACCOUNT_ID: string; + DATALAKE_APAC_ACCESS_KEY: string; + DATALAKE_APAC_SECRET_KEY: string; + DATALAKE_APAC_BUCKET_NAME: string; + DATALAKE_EEUR_ACCESS_KEY: string; + DATALAKE_EEUR_SECRET_KEY: string; + DATALAKE_EEUR_BUCKET_NAME: string; + DATALAKE_WEUR_ACCESS_KEY: string; + DATALAKE_WEUR_SECRET_KEY: string; + DATALAKE_WEUR_BUCKET_NAME: string; + DATALAKE_ENAM_ACCESS_KEY: string; + DATALAKE_ENAM_SECRET_KEY: string; + DATALAKE_ENAM_BUCKET_NAME: string; + DATALAKE_WNAM_ACCESS_KEY: string; + DATALAKE_WNAM_SECRET_KEY: string; + DATALAKE_WNAM_BUCKET_NAME: string; +} diff --git a/workers/datalake/wrangler.toml b/workers/datalake/wrangler.toml new file mode 100644 index 0000000000..12f3afc6a3 --- /dev/null +++ b/workers/datalake/wrangler.toml @@ -0,0 +1,48 @@ +#:schema node_modules/wrangler/config-schema.json +name = "datalake-worker" +main = "src/index.ts" +compatibility_date = "2024-07-01" +compatibility_flags = ["nodejs_compat"] +keep_vars = true + +kv_namespaces = [ + { binding = "datalake_blobs", id = "64144eb146fd45febc928d44419ebb39", preview_id = "31c6f6e76e7e4524a59f87a4f381de82" } +] + +r2_buckets = [ + { binding = "DATALAKE_APAC", bucket_name = "datalake-apac", preview_bucket_name = "dev-datalake-eu-west" }, + { binding = "DATALAKE_EEUR", bucket_name = "datalake-eeur", preview_bucket_name = "dev-datalake-eu-west" }, + { binding = "DATALAKE_WEUR", bucket_name = "datalake-weur", preview_bucket_name = "dev-datalake-eu-west" }, + { binding = "DATALAKE_ENAM", bucket_name = "datalake-enam", preview_bucket_name = "dev-datalake-eu-west" }, + { binding = "DATALAKE_WNAM", bucket_name = "datalake-wnam", preview_bucket_name = "dev-datalake-eu-west" } +] + +[[hyperdrive]] +binding = "HYPERDRIVE" +id = "87259c3ae41e41a7b35e610d4282d85a" +localConnectionString = "postgresql://root:roach@localhost:26257/datalake" + +[observability] +enabled = true +head_sampling_rate = 1 + +[vars] +DATALAKE_EEUR_BUCKET_NAME = "datalake-eeur" +# DATALAKE_EEUR_ACCESS_KEY = "" +# DATALAKE_EEUR_SECRET_KEY = "" +DATALAKE_WEUR_BUCKET_NAME = "datalake-weur" +# DATALAKE_WEUR_ACCESS_KEY = "" +# DATALAKE_WEUR_SECRET_KEY = "" +DATALAKE_APAC_BUCKET_NAME = "datalake-apac" +# DATALAKE_APAC_ACCESS_KEY = "" +# DATALAKE_APAC_SECRET_KEY = "" +DATALAKE_ENAM_BUCKET_NAME = "datalake-enam" +# DATALAKE_ENAM_ACCESS_KEY = "" +# DATALAKE_ENAM_SECRET_KEY = "" +DATALAKE_WNAM_BUCKET_NAME = "datalake-wnam" +# DATALAKE_WNAM_ACCESS_KEY = "" +# DATALAKE_WNAM_SECRET_KEY = "" + +# STREAMS_ACCOUNT_ID = "" +# STREAMS_AUTH_KEY = "" +# R2_ACCOUNT_ID = ""