Merge branch 'develop' of https://github.com/hcengineering/platform into staging-new

Signed-off-by: Artem Savchenko <armisav@gmail.com>
This commit is contained in:
Artem Savchenko
2025-09-22 12:50:33 +07:00
83 changed files with 1551 additions and 1686 deletions
-19
View File
@@ -972,25 +972,6 @@
"runtimeArgs": ["--nolazy", "-r", "ts-node/register"],
"sourceMaps": true,
"cwd": "${workspaceRoot}/services/export/pod-export"
},
{
"name": "Msg2File",
"type": "node",
"request": "launch",
"args": ["src/index.ts"],
"env": {
"ACCOUNTS_URL": "http://localhost:3000",
"DB_URL": "postgresql://root@localhost:26257/defaultdb?sslmode=disable",
"PORT": "9087",
"SECRET": "secret",
"SERVICE_ID": "msg2file-service",
"STORAGE_CONFIG": "datalake|http://huly.local:4030"
},
"runtimeVersion": "20",
"runtimeArgs": ["--nolazy", "-r", "ts-node/register"],
"sourceMaps": true,
"outputCapture": "std",
"cwd": "${workspaceRoot}/services/msg2file"
}
]
}
+10 -105
View File
@@ -304,9 +304,6 @@ importers:
'@rush-temp/communication-types':
specifier: file:./projects/communication-types.tgz
version: file:projects/communication-types.tgz(@types/node@22.15.29)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))
'@rush-temp/communication-yaml':
specifier: file:./projects/communication-yaml.tgz
version: file:projects/communication-yaml.tgz(@types/node@22.15.29)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))
'@rush-temp/contact':
specifier: file:./projects/contact.tgz
version: file:projects/contact.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(@types/node@22.15.29)(babel-jest@29.7.0(@babel/core@7.23.9))(esbuild@0.25.9)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))
@@ -928,9 +925,6 @@ importers:
'@rush-temp/pod-media':
specifier: file:./projects/pod-media.tgz
version: file:projects/pod-media.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(@swc/core@1.13.5)(babel-jest@29.7.0(@babel/core@7.23.9))
'@rush-temp/pod-msg2file':
specifier: file:./projects/pod-msg2file.tgz
version: file:projects/pod-msg2file.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(@swc/core@1.13.5)(babel-jest@29.7.0(@babel/core@7.23.9))
'@rush-temp/pod-notification':
specifier: file:./projects/pod-notification.tgz
version: file:projects/pod-notification.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(@swc/core@1.13.5)(babel-jest@29.7.0(@babel/core@7.23.9))
@@ -1732,9 +1726,6 @@ importers:
'@types/node':
specifier: ^22.15.29
version: 22.15.29
'@types/node-cron':
specifier: ^3.0.11
version: 3.0.11
'@types/nodemailer':
specifier: ^6.4.17
version: 6.4.17
@@ -2155,9 +2146,6 @@ importers:
msgpackr-extract:
specifier: ^3.0.3
version: 3.0.3
node-cron:
specifier: ^3.0.3
version: 3.0.3
node-forge:
specifier: ^1.3.1
version: 1.3.1
@@ -4580,7 +4568,7 @@ packages:
version: 0.0.0
'@rush-temp/communication-client-query@file:projects/communication-client-query.tgz':
resolution: {integrity: sha512-6OA/sAql1ZLhhDX9J2yZOAEgfXqSxYh6yxpH/YGrC3PfYyhVAoe7piCjpB+JpKR5wL1kzpl1eP7Bj0IEq2jNJQ==, tarball: file:projects/communication-client-query.tgz}
resolution: {integrity: sha512-+68bXbpIU6XTujL5Dj3pFCdFelV7mE+Ln6zZEvf2NPHhAcJqy3+KjHQyIbytnWhCFsGyk+HE3S08ps6Kl6TKtw==, tarball: file:projects/communication-client-query.tgz}
version: 0.0.0
'@rush-temp/communication-cockroach@file:projects/communication-cockroach.tgz':
@@ -4588,7 +4576,7 @@ packages:
version: 0.0.0
'@rush-temp/communication-query@file:projects/communication-query.tgz':
resolution: {integrity: sha512-sXcvSDoXmgZH233jFRePio23qQosvQZoACKqKgJ2Obr+oY4ZzJ2uX7DM0lyIKbhFUuLNFl56u2RfM5ra8SxRDg==, tarball: file:projects/communication-query.tgz}
resolution: {integrity: sha512-KMDSFpLNW1ekFoAf+BbRxPeRYnMUjYu+l80F+fgtvkKsX/6oKtb6xRhQyVh4EQiRlQBaMV//1b2d7cBwuF5HDw==, tarball: file:projects/communication-query.tgz}
version: 0.0.0
'@rush-temp/communication-resources@file:projects/communication-resources.tgz':
@@ -4604,21 +4592,17 @@ packages:
version: 0.0.0
'@rush-temp/communication-server@file:projects/communication-server.tgz':
resolution: {integrity: sha512-KVc53phC1b9CNptvJmMb706t74gsCiEQGWvCpvTlRWMd3q+pj7i2994gJ0/KzuCojcHxPYVn2QxNojpNSV66Eg==, tarball: file:projects/communication-server.tgz}
resolution: {integrity: sha512-HfKzdbndaS3opK1inqcTOrVnKvU7PS2asLmKuR4JHooiTTvFFCHcXXYO++t1QmdGmankZRn3Qsq59CIt/UPgsA==, tarball: file:projects/communication-server.tgz}
version: 0.0.0
'@rush-temp/communication-shared@file:projects/communication-shared.tgz':
resolution: {integrity: sha512-SustnNpr/eopCXmm0ZahvUZXZPcD7psSWM+CkWVwiPlLVPGWMVUURFOcnC94F3OhY+7qu9N0Rgr3gDvr4VKXWA==, tarball: file:projects/communication-shared.tgz}
resolution: {integrity: sha512-/ZUh5sviGQkQGJnOe6ImJHpKCcn7xNyyMeOcl4Tr3TqXjsH+uOQ0MAFI0Ul4f1LI+dGIrp1QsZWSVWvEVprRQQ==, tarball: file:projects/communication-shared.tgz}
version: 0.0.0
'@rush-temp/communication-types@file:projects/communication-types.tgz':
resolution: {integrity: sha512-uzZ82V1R+h+TmxBCrmKaAfUgoiX3Z8ciPGrkc42M++M5d5mlfIjTvPh0O+y50pHoFuQ5NQ3pNIjeDkrJE8Dp0w==, tarball: file:projects/communication-types.tgz}
version: 0.0.0
'@rush-temp/communication-yaml@file:projects/communication-yaml.tgz':
resolution: {integrity: sha512-D+fijwXQM2oMeT6QOmkRVLik5uyVCNsml07rz5RvKcPEpbqRh3YBtIvcUABoCStLvNTIzYuXnqbCHzSl/ibidw==, tarball: file:projects/communication-yaml.tgz}
version: 0.0.0
'@rush-temp/communication@file:projects/communication.tgz':
resolution: {integrity: sha512-kO/x+NBE4Nmtz5KhHO0kxty/0N1A2PH4dNVTdOviEV4mVAhFm7DkgWHu0SwWJvQrCzsB0xEga3/CGuGjcPyMjA==, tarball: file:projects/communication.tgz}
version: 0.0.0
@@ -5424,7 +5408,7 @@ packages:
version: 0.0.0
'@rush-temp/pod-fulltext@file:projects/pod-fulltext.tgz':
resolution: {integrity: sha512-0G3FKVxHWnyPuHWVBiR3Xtns3rLxGTC6ferQCX8NSC3dGOGOvk4TZgB8leohmJavwwxdZQKk1QFJqmXvHux10A==, tarball: file:projects/pod-fulltext.tgz}
resolution: {integrity: sha512-kPi4ermTxQtUKZ5xzbJgT07rLj2gG6m4YQgArWiONz47ooSoizi/Sg8Pu2FFyVV0klzQwqUOseuoVqgdySnFNg==, tarball: file:projects/pod-fulltext.tgz}
version: 0.0.0
'@rush-temp/pod-github@file:projects/pod-github.tgz':
@@ -5451,10 +5435,6 @@ packages:
resolution: {integrity: sha512-5LRHcx/zPNVCJbIBfjigt+of64jfGrH5fjoDyvhASTEJbtG4RGRj8kKZoUI+1JHibDZzcVn2NLc+3nFMhksreA==, tarball: file:projects/pod-media.tgz}
version: 0.0.0
'@rush-temp/pod-msg2file@file:projects/pod-msg2file.tgz':
resolution: {integrity: sha512-7b2GrlsmevpahNIADhqNUBpshocrdjlc2L0kf2BAFVSEFpjzUtRju2XqwwoZ2l7xA0Z/b+Un9PHt+XYMWcNzlg==, tarball: file:projects/pod-msg2file.tgz}
version: 0.0.0
'@rush-temp/pod-notification@file:projects/pod-notification.tgz':
resolution: {integrity: sha512-1EeImD9ZNRxWq0jDp6o50hOpuoB0sNI54V/V1pm5VONdRF0S2lVzPCr35YehJ5BJcvbhy22au/7rk9XadoYI9w==, tarball: file:projects/pod-notification.tgz}
version: 0.0.0
@@ -5524,7 +5504,7 @@ packages:
version: 0.0.0
'@rush-temp/presentation@file:projects/presentation.tgz':
resolution: {integrity: sha512-nMn7FUSF2F6YuJAstYbvV8Yi9upkKqeSccsU3EOrKjJm1c57sa7x8i4K3Z6l4RoSQk86/F5CrzZVgqIMKCjF0w==, tarball: file:projects/presentation.tgz}
resolution: {integrity: sha512-yUvZa29zEHKk183hweT+OqJYrQtt4WwtYd6X4Y5WknLTG4fnd0kDEyFyOT9LgUUFF0/vW2iVYWLcMU93+W8PUg==, tarball: file:projects/presentation.tgz}
version: 0.0.0
'@rush-temp/print-assets@file:projects/print-assets.tgz':
@@ -5804,7 +5784,7 @@ packages:
version: 0.0.0
'@rush-temp/server-indexer@file:projects/server-indexer.tgz':
resolution: {integrity: sha512-OYx2FWCD6jtTnsKmYnhhS3BsvhtCN9pIYkvIoSe0mzZi3xEAhnwu2T59+JpZgRO1fvD1yVxBy7nw0hyFtz03Hg==, tarball: file:projects/server-indexer.tgz}
resolution: {integrity: sha512-InjSApXF2jmClorioQIIV+JxbKtuYjl3a9i4eetEWuG2z7R+Ya7L6f98cqhEhYkdA4hM5TTZYbS+BcRLWlSqQg==, tarball: file:projects/server-indexer.tgz}
version: 0.0.0
'@rush-temp/server-inventory-resources@file:projects/server-inventory-resources.tgz':
@@ -6116,7 +6096,7 @@ packages:
version: 0.0.0
'@rush-temp/tool@file:projects/tool.tgz':
resolution: {integrity: sha512-bsqWzFX5qr02kfG1zQZOgHsGWzR/Z66VC16ZoQnXs6mkdHFSmQ2p2l8eNL/aYPu2qfrRkMNAnI8cu6aULfKuDg==, tarball: file:projects/tool.tgz}
resolution: {integrity: sha512-RQFbb0YGcuuAzNtGrbYDMyFE5VKDEsPA+bylES+hRUvKKBou2icVHwzDbFWrdSwGwfrr0ZhXVFQfTwSKOM18Bw==, tarball: file:projects/tool.tgz}
version: 0.0.0
'@rush-temp/tracker-assets@file:projects/tracker-assets.tgz':
@@ -7176,9 +7156,6 @@ packages:
'@types/mysql@2.15.27':
resolution: {integrity: sha512-YfWiV16IY0OeBfBCk8+hXKmdTKrKlwKN1MNKAPBu5JYxLwBEZl7QzeEpGnlZb3VMGJrrGmB84gXiH+ofs/TezA==}
'@types/node-cron@3.0.11':
resolution: {integrity: sha512-0ikrnug3/IyneSHqCBeslAhlK2aBfYek1fGo4bP4QnZPmiqSGRK+Oy7ZMisLWkesffJvQ1cqAcBnJC+8+nxIAg==}
'@types/node-fetch@2.6.12':
resolution: {integrity: sha512-8nneRWKCg3rMtF69nLQJnOYUcbafYeFSjqkw3jCRLsqkWFlHaoQrr5mXmofFGOx3DKn7UfmBMyov8ySvLRVldA==}
@@ -11808,10 +11785,6 @@ packages:
node-api-version@0.2.0:
resolution: {integrity: sha512-fthTTsi8CxaBXMaBAD7ST2uylwvsnYxh2PfaScwpMhos6KlSFajXQPcM4ogNE1q2s3Lbz9GCGqeIHC+C6OZnKg==}
node-cron@3.0.3:
resolution: {integrity: sha512-dOal67//nohNgYWb+nWmg5dkFdIwDm8EpeGYMekPMrngV3637lqnX0lbUcCtgibHTz6SEz7DAIjKvKDFYCnO1A==}
engines: {node: '>=6.0.0'}
node-domexception@1.0.0:
resolution: {integrity: sha512-/jKZoMpw0F8GRwl4/eLROPA3cfcXtLApP0QzLmUT/HuPCZWyB7IY9ZrMeKw2O/nFIqPQB3PVM9aYm0F312AXDQ==}
engines: {node: '>=10.5.0'}
@@ -19286,29 +19259,6 @@ snapshots:
- supports-color
- ts-node
'@rush-temp/communication-yaml@file:projects/communication-yaml.tgz(@types/node@22.15.29)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))':
dependencies:
'@types/js-yaml': 4.0.9
'@typescript-eslint/eslint-plugin': 6.21.0(@typescript-eslint/parser@6.21.0(eslint@8.56.0)(typescript@5.8.3))(eslint@8.56.0)(typescript@5.8.3)
'@typescript-eslint/parser': 6.21.0(eslint@8.56.0)(typescript@5.8.3)
esbuild: 0.25.9
esbuild-plugin-copy: 2.1.1(esbuild@0.25.9)
eslint: 8.56.0
eslint-config-standard-with-typescript: 40.0.0(@typescript-eslint/eslint-plugin@6.21.0(@typescript-eslint/parser@6.21.0(eslint@8.56.0)(typescript@5.8.3))(eslint@8.56.0)(typescript@5.8.3))(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))(eslint@8.56.0)(typescript@5.8.3)
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)
jest: 29.7.0(@types/node@22.15.29)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))
js-yaml: 4.1.0
prettier: 3.2.5
typescript: 5.8.3
transitivePeerDependencies:
- '@types/node'
- babel-plugin-macros
- node-notifier
- supports-color
- ts-node
'@rush-temp/communication@file:projects/communication.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(@types/node@22.15.29)(babel-jest@29.7.0(@babel/core@7.23.9))(esbuild@0.25.9)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))':
dependencies:
'@types/jest': 29.5.12
@@ -25435,47 +25385,6 @@ snapshots:
- node-notifier
- supports-color
'@rush-temp/pod-msg2file@file:projects/pod-msg2file.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(@swc/core@1.13.5)(babel-jest@29.7.0(@babel/core@7.23.9))':
dependencies:
'@tsconfig/node16': 1.0.4
'@types/cors': 2.8.17
'@types/express': 4.17.21
'@types/jest': 29.5.12
'@types/js-yaml': 4.0.9
'@types/node': 22.15.29
'@types/node-cron': 3.0.11
'@types/uuid': 8.3.4
'@typescript-eslint/eslint-plugin': 6.21.0(@typescript-eslint/parser@6.21.0(eslint@8.56.0)(typescript@5.8.3))(eslint@8.56.0)(typescript@5.8.3)
'@typescript-eslint/parser': 6.21.0(eslint@8.56.0)(typescript@5.8.3)
cors: 2.8.5
dotenv: 16.0.3
esbuild: 0.25.9
eslint: 8.56.0
eslint-config-standard-with-typescript: 40.0.0(@typescript-eslint/eslint-plugin@6.21.0(@typescript-eslint/parser@6.21.0(eslint@8.56.0)(typescript@5.8.3))(eslint@8.56.0)(typescript@5.8.3))(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))(eslint@8.56.0)(typescript@5.8.3)
eslint-plugin-import: 2.29.1(eslint@8.56.0)
eslint-plugin-n: 15.7.0(eslint@8.56.0)
eslint-plugin-node: 11.1.0(eslint@8.56.0)
eslint-plugin-promise: 6.1.1(eslint@8.56.0)
express: 4.21.2
jest: 29.7.0(@types/node@22.15.29)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))
js-yaml: 4.1.0
node-cron: 3.0.3
postgres: 3.4.7
prettier: 3.2.5
ts-jest: 29.1.2(@babel/core@7.23.9)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.23.9))(esbuild@0.25.9)(jest@29.7.0(@types/node@22.15.29)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3)))(typescript@5.8.3)
ts-node: 10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3)
typescript: 5.8.3
uuid: 8.3.2
transitivePeerDependencies:
- '@babel/core'
- '@jest/types'
- '@swc/core'
- '@swc/wasm'
- babel-jest
- babel-plugin-macros
- node-notifier
- supports-color
'@rush-temp/pod-notification@file:projects/pod-notification.tgz(@babel/core@7.23.9)(@jest/types@29.6.3)(@swc/core@1.13.5)(babel-jest@29.7.0(@babel/core@7.23.9))':
dependencies:
'@tsconfig/node16': 1.0.4
@@ -30348,6 +30257,7 @@ snapshots:
'@elastic/elasticsearch': 7.17.14
'@faker-js/faker': 8.4.1
'@types/jest': 29.5.12
'@types/js-yaml': 4.0.9
'@types/mime-types': 2.1.4
'@types/minio': 7.0.18
'@types/node': 22.15.29
@@ -30368,6 +30278,7 @@ snapshots:
eslint-plugin-promise: 6.1.1(eslint@8.56.0)
fast-equals: 5.2.2
jest: 29.7.0(@types/node@22.15.29)(ts-node@10.9.2(@swc/core@1.13.5)(@types/node@22.15.29)(typescript@5.8.3))
js-yaml: 4.1.0
libphonenumber-js: 1.10.56
mime-types: 2.1.35
mongodb: 6.16.0(gcp-metadata@5.3.0(encoding@0.1.13))(snappy@7.2.2)(socks@2.8.3)
@@ -32146,8 +32057,6 @@ snapshots:
dependencies:
'@types/node': 22.15.29
'@types/node-cron@3.0.11': {}
'@types/node-fetch@2.6.12':
dependencies:
'@types/node': 22.15.29
@@ -37734,10 +37643,6 @@ snapshots:
dependencies:
semver: 7.7.2
node-cron@3.0.3:
dependencies:
uuid: 8.3.2
node-domexception@1.0.0: {}
node-fetch@2.7.0(encoding@0.1.13):
-1
View File
@@ -19,7 +19,6 @@ rush docker:build -p 20 \
--to @hcengineering/pod-datalake \
--to @hcengineering/pod-mail-worker \
--to @hcengineering/pod-export \
--to @hcengineering/pod-msg2file \
--to @hcengineering/pod-media \
--to @hcengineering/pod-preview \
--to @hcengineering/pod-external \
+1
View File
@@ -298,6 +298,7 @@ export async function configurePlatform (onWorkbenchConnect?: () => Promise<void
setMetadata(presentation.metadata.MailUrl, config.MAIL_URL)
setMetadata(recorder.metadata.StreamUrl, config.STREAM_URL ?? '')
setMetadata(presentation.metadata.StatsUrl, config.STATS_URL)
setMetadata(presentation.metadata.HulylakeUrl, config.HULYLAKE_URL ?? '')
setMetadata(presentation.metadata.PulseUrl, config.PULSE_URL ?? '')
setMetadata(textEditor.metadata.Collaborator, config.COLLABORATOR ?? '')
+1
View File
@@ -56,6 +56,7 @@ export interface Config {
PULSE_URL?: string
PASSWORD_STRICTNESS?: 'very_strict' | 'strict' | 'normal' | 'none'
EXCLUDED_APPLICATIONS_FOR_ANONYMOUS?: string
HULYLAKE_URL?: string
}
export interface Branding {
-1
View File
@@ -236,7 +236,6 @@ services:
- LAST_NAME_FIRST=true
- BRANDING_PATH=/var/cfg/branding.json
- AI_BOT_URL=http://huly.local:4010
- MSG2FILE_URL=http://huly.local:9087
- COMMUNICATION_TIME_LOGGING_ENABLED=true
restart: unless-stopped
fulltext_cockroach:
+3 -16
View File
@@ -286,6 +286,7 @@ services:
# - DISABLE_SIGNUP=true
- OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4318/v1/traces
- PULSE_URL=ws://huly.local:8099/ws
- HULYLAKE_URL=http://huly.local:8096
restart: unless-stopped
transactor_cockroach:
image: hardcoreeng/transactor
@@ -321,13 +322,13 @@ services:
- LAST_NAME_FIRST=true
- BRANDING_PATH=/var/cfg/branding.json
- AI_BOT_URL=http://huly.local:4010
- MSG2FILE_URL=http://huly.local:9087
- COMMUNICATION_TIME_LOGGING_ENABLED=true
- RATE_LIMIT_MAX=250 # 250 requests per 30 seconds
- RATE_LIMIT_WINDOW=30000
- FILES_URL=http://huly.local:4030/blob/:workspace/:blobId/:filename
- COMMUNICATION_API_ENABLED=true
- OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4318/v1/traces
- HYLYLAKE_URL=http://huly.local:8096
restart: unless-stopped
rekoni:
image: hardcoreeng/rekoni-service
@@ -366,6 +367,7 @@ services:
- REKONI_URL=http://huly.local:4004
- ACCOUNTS_URL=http://huly.local:3000
- OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4318/v1/traces
- HULYLAKE_URL=http://huly.local:8096
print:
image: hardcoreeng/print
extra_hosts:
@@ -440,21 +442,6 @@ services:
# - POSTHOG_HOST=${POSTHOG_HOST}
# - POSTHOG_API_KEY=${POSTHOG_API_KEY}
# - MAX_PAYLOAD_SIZE=10mb
msg2file:
image: hardcoreeng/msg2file
ports:
- 9087:9087
extra_hosts:
- 'huly.local:host-gateway'
restart: unless-stopped
environment:
- ACCOUNTS_URL=http://huly.local:3000
- DB_URL=postgresql://root@huly.local:26257/defaultdb?sslmode=disable
- PORT=9087
- SECRET=secret
- SERVICE_ID=msg2file-service
- STORAGE_CONFIG=${STORAGE_CONFIG}
- OTEL_EXPORTER_OTLP_ENDPOINT=http://jaeger:4318/v1/traces
export:
image: hardcoreeng/export
extra_hosts:
+1
View File
@@ -28,5 +28,6 @@
"COMMUNICATION_API_ENABLED": "true",
"BACKUP_URL": "http://huly.local:4039/api/backup",
"PULSE_URL": "ws://huly.local:8099/ws",
"HULYLAKE_URL": "http://huly.local:8096",
"EXCLUDED_APPLICATIONS_FOR_ANONYMOUS": "[\"chunter\", \"notification\"]"
}
+3 -1
View File
@@ -201,7 +201,8 @@ export interface Config {
COMMUNICATION_API_ENABLED?: string
BILLING_URL?: string,
EXCLUDED_APPLICATIONS_FOR_ANONYMOUS?: string,
PULSE_URL?: string
PULSE_URL?: string,
HULYLAKE_URL?: string
}
export interface Branding {
@@ -491,6 +492,7 @@ export async function configurePlatform() {
setMetadata(billingPlugin.metadata.BillingURL, config.BILLING_URL ?? '')
setMetadata(presentation.metadata.PulseUrl, config.PULSE_URL)
setMetadata(presentation.metadata.HulylakeUrl, config.HULYLAKE_URL ?? '')
const languages = myBranding.languages
? (myBranding.languages as string).split(',').map((l) => l.trim())
+7 -3
View File
@@ -21,7 +21,7 @@
"docker:staging": "../../common/scripts/docker_tag.sh hardcoreeng/tool staging",
"docker:push": "../../common/scripts/docker_tag.sh hardcoreeng/tool",
"run-local-mongo": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret FULLTEXT_URL=http://localhost:4700 ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3333 MINIO_ACCESS_KEY=minioadmin MINIO_SECRET_KEY=minioadmin MINIO_ENDPOINT=localhost ACCOUNT_DB_URL=mongodb://localhost:27017 DB_URL=mongodb://localhost:27017 TELEGRAM_DATABASE=telegram-service REKONI_URL=http://localhost:4004 MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) QUEUE_CONFIG='localhost:19092' node --expose-gc --max-old-space-size=18000 ./bundle/bundle.js",
"run-local": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret FULLTEXT_URL=http://localhost:4702 ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3332 STORAGE_CONFIG='datalake|http://huly.local:4030' ACCOUNT_DB_URL=postgresql://root@huly.local:26257/defaultdb?sslmode=disable DB_URL=postgresql://root@huly.local:26257/defaultdb?sslmode=disable TELEGRAM_DATABASE=telegram-service REKONI_URL=http://localhost:4004 REGION_INFO='cockroach|CockroachDB' MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) QUEUE_CONFIG='localhost:19092' node --expose-gc --max-old-space-size=18000 $TOOL_OPT ./bundle/bundle.js",
"run-local": "rush bundle --to @hcengineering/tool >/dev/null && cross-env SERVER_SECRET=secret FULLTEXT_URL=http://localhost:4702 ACCOUNTS_URL=http://localhost:3000 TRANSACTOR_URL=ws://localhost:3332 STORAGE_CONFIG='datalake|http://huly.local:4030' HULYLAKE_URL=http://huly.local:8096 ACCOUNT_DB_URL=postgresql://root@huly.local:26257/defaultdb?sslmode=disable DB_URL=postgresql://root@huly.local:26257/defaultdb?sslmode=disable TELEGRAM_DATABASE=telegram-service REKONI_URL=http://localhost:4004 REGION_INFO='cockroach|CockroachDB' MODEL_VERSION=$(node ../../common/scripts/show_version.js) GIT_REVISION=$(git describe --all --long) QUEUE_CONFIG='localhost:19092' node --expose-gc --max-old-space-size=18000 $TOOL_OPT ./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 ACCOUNT_DB_URL=mongodb://localhost:27017 DB_URL=mongodb://localhost:27017 TELEGRAM_DATABASE=telegram-service 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-mongo": "rushx run-local-mongo upgrade-workspace -- $1",
@@ -55,7 +55,8 @@
"@types/request": "~2.48.8",
"jest": "^29.7.0",
"ts-jest": "^29.1.1",
"@types/jest": "^29.5.5"
"@types/jest": "^29.5.5",
"@types/js-yaml": "^4.0.9"
},
"dependencies": {
"@elastic/elasticsearch": "^7.17.14",
@@ -180,6 +181,9 @@
"msgpackr-extract": "^3.0.3",
"@hcengineering/kafka": "^0.6.0",
"@hcengineering/api-client": "^0.6.0",
"@faker-js/faker": "^8.4.1"
"@faker-js/faker": "^8.4.1",
"@hcengineering/hulylake-client": "^0.6.0",
"js-yaml": "^4.1.0",
"@hcengineering/communication-types": "^0.1.0"
}
}
File diff suppressed because it is too large Load Diff
+58 -7
View File
@@ -56,6 +56,14 @@ import {
type Account as OldAccount,
type Workspace as OldWorkspace
} from '@hcengineering/account-service'
import { getClient as getHulylakeClient } from '@hcengineering/hulylake-client'
import {
getDBClient,
createPostgreeDestroyAdapter,
createPostgresAdapter,
createPostgresTxAdapter,
shutdownPostgres
} from '@hcengineering/postgres'
import { faker } from '@faker-js/faker'
import { getPlatformQueue } from '@hcengineering/kafka'
@@ -69,6 +77,7 @@ import {
isDeletingMode,
MeasureMetricsContext,
metricsToString,
SocialId,
SocialIdType,
systemAccountEmail,
systemAccountUuid,
@@ -80,7 +89,8 @@ import {
type Tx,
type Version,
type WorkspaceDataId,
type WorkspaceUuid
type WorkspaceUuid,
type PersonUuid
} from '@hcengineering/core'
import { consoleModelLogger, type MigrateOperation } from '@hcengineering/model'
import {
@@ -92,12 +102,6 @@ import {
} from '@hcengineering/mongo'
import { getModelVersion } from '@hcengineering/model-all'
import {
createPostgreeDestroyAdapter,
createPostgresAdapter,
createPostgresTxAdapter,
shutdownPostgres
} from '@hcengineering/postgres'
import {
QueueTopic,
workspaceEvents,
@@ -131,6 +135,7 @@ import { mkdir, writeFile } from 'fs/promises'
import { dirname } from 'path'
import { restoreMarkupRefs } from './markup'
import { restoreGithubIntegrations } from './restoreGithub'
import { migrateWorkspaceMessages } from './communication'
const colorConstants = {
colorRed: '\u001b[31m',
@@ -2845,6 +2850,52 @@ export function devTool (
}, dbUrl)
})
program
.command('migrate-communication-to-hulylake')
.description('Migrate communication messages to hulylake')
.action(async () => {
const { dbUrl } = prepareTools()
const storageConfig = storageConfigFromEnv()
const hulylakeUrl = process.env.HULYLAKE_URL ?? ''
if (hulylakeUrl === '') {
throw new Error('HULYLAKE_URL should be specified')
}
const storage: StorageAdapter = buildStorageFromConfig(storageConfig)
const token = generateToken(systemAccountUuid, undefined, {
service: 'tool'
})
const db = getDBClient(dbUrl, undefined, 'tool')
const dbClient = await db.getClient()
const accountClient = getAccountClient(token)
const personUuidBySocialId = new Map<PersonId, PersonUuid>()
await withAccountDatabase(async (accountDb) => {
const workspaces = await accountDb.workspace.find({})
for (const ws of workspaces) {
try {
const hulylake = getHulylakeClient(hulylakeUrl, ws.uuid, token)
console.log('start workspace migration', ws.name)
await migrateWorkspaceMessages(
toolCtx.newChild(ws.name, {}),
ws,
dbClient,
storage,
hulylake,
accountClient,
personUuidBySocialId
)
console.log('done workspace migration', ws.name)
} catch (err: any) {
console.error('failed to migrate workspace', ws.name)
console.error(err)
}
}
db.close()
}, dbUrl)
})
extendProgram?.(program)
process.on('unhandledRejection', (reason, promise) => {
+1
View File
@@ -61,6 +61,7 @@
"@hcengineering/emoji": "^0.6.0",
"@hcengineering/theme": "^0.6.5",
"@hcengineering/retry": "^0.6.0",
"@hcengineering/hulylake-client": "^0.6.0",
"fast-equals": "^5.2.2",
"png-chunks-extract": "^1.0.0",
"svelte": "^4.2.20",
+54 -41
View File
@@ -33,26 +33,24 @@ import {
type UpdateAttachmentsOperation,
type UpdateNotificationContextEvent,
type UpdateNotificationEvent,
type UpdateNotificationQuery,
type NotificationQuery,
type UpdatePatchEvent
} from '@hcengineering/communication-sdk-types'
import {
type AccountID,
type AccountUuid,
type CardID,
type CardType,
type Collaborator,
type ContextID,
type FindCollaboratorsParams,
type FindLabelsParams,
type FindMessagesGroupsParams,
type FindMessagesParams,
type FindNotificationContextParams,
type FindNotificationsParams,
type FindMessagesMetaParams,
type Label,
type Markdown,
type Message,
type MessageID,
type MessagesGroup,
MessageType,
type Notification,
type NotificationContext,
@@ -61,7 +59,10 @@ import {
type AttachmentData,
type AttachmentParams,
type AttachmentUpdateData,
type WithTotal
type WithTotal,
type NotificationID,
type Emoji,
type MessageMeta
} from '@hcengineering/communication-types'
import core, {
generateId,
@@ -71,17 +72,19 @@ import core, {
SocialIdType,
type Tx,
type TxDomainEvent,
AccountRole
AccountRole,
generateUuid
} from '@hcengineering/core'
import { onDestroy } from 'svelte'
import { addNotification, NotificationSeverity, languageStore } from '@hcengineering/ui'
import { translate } from '@hcengineering/platform'
import { getMetadata, translate } from '@hcengineering/platform'
import view from '@hcengineering/view'
import { v4 as uuid } from 'uuid'
import { getCurrentWorkspaceUuid, getFilesUrl } from './file'
import { addTxListener, removeTxListener, type TxListener } from './utils'
import { get } from 'svelte/store'
import { getClient as getHulylakeClient } from '@hcengineering/hulylake-client'
import { getCurrentWorkspaceUuid } from './file'
import { addTxListener, removeTxListener, type TxListener } from './utils'
import presentation from './plugin'
export {
createCollaboratorsQuery,
@@ -107,7 +110,12 @@ export async function setCommunicationClient (platformClient: PlatformClient): P
client.close()
}
const _client = new Client(platformClient)
initLiveQueries(_client, getCurrentWorkspaceUuid(), getFilesUrl(), onDestroy)
const token = getMetadata(presentation.metadata.Token) ?? ''
const hulylakeUrl = getMetadata(presentation.metadata.HulylakeUrl) ?? ''
const hulylake = getHulylakeClient(hulylakeUrl, getCurrentWorkspaceUuid(), token)
initLiveQueries(_client, hulylake, onDestroy)
client = _client
onClientListeners.forEach((fn) => {
fn()
@@ -161,7 +169,7 @@ class Client {
async createMessage (cardId: CardID, cardType: CardType, content: Markdown): Promise<CreateMessageResult> {
const event: CreateMessageEvent = {
type: MessageEventType.CreateMessage,
messageType: MessageType.Message,
messageType: MessageType.Text,
cardId,
cardType,
content,
@@ -198,28 +206,28 @@ class Client {
await this.sendEvent(event)
}
async addReaction (cardId: CardID, messageId: MessageID, reaction: string): Promise<void> {
async addReaction (cardId: CardID, messageId: MessageID, emoji: Emoji): Promise<void> {
const event: ReactionPatchEvent = {
type: MessageEventType.ReactionPatch,
cardId,
messageId,
operation: {
opcode: 'add',
reaction
reaction: emoji
},
socialId: this.getSocialId()
}
await this.sendEvent(event)
}
async removeReaction (cardId: CardID, messageId: MessageID, reaction: string): Promise<void> {
async removeReaction (cardId: CardID, messageId: MessageID, emoji: Emoji): Promise<void> {
const event: ReactionPatchEvent = {
type: MessageEventType.ReactionPatch,
cardId,
messageId,
operation: {
opcode: 'remove',
reaction
reaction: emoji
},
socialId: this.getSocialId()
}
@@ -245,7 +253,7 @@ class Client {
opcode: 'add',
attachments: ops.add.map((it) => ({
...it,
id: it.id ?? (uuid() as AttachmentID)
id: it.id ?? (generateUuid() as AttachmentID)
}))
})
}
@@ -262,7 +270,7 @@ class Client {
opcode: 'set',
attachments: ops.set.map((it) => ({
...it,
id: it.id ?? (uuid() as AttachmentID)
id: it.id ?? (generateUuid() as AttachmentID)
}))
})
}
@@ -286,7 +294,7 @@ class Client {
await this.sendEvent(event)
}
async addCollaborators (cardId: CardID, cardType: CardType, collaborators: AccountID[]): Promise<void> {
async addCollaborators (cardId: CardID, cardType: CardType, collaborators: AccountUuid[]): Promise<void> {
const event: AddCollaboratorsEvent = {
type: NotificationEventType.AddCollaborators,
cardId,
@@ -297,7 +305,7 @@ class Client {
await this.sendEvent(event)
}
async removeCollaborators (cardId: CardID, cardType: CardType, collaborators: AccountID[]): Promise<void> {
async removeCollaborators (cardId: CardID, cardType: CardType, collaborators: AccountUuid[]): Promise<void> {
const event: RemoveCollaboratorsEvent = {
type: NotificationEventType.RemoveCollaborators,
cardId,
@@ -329,7 +337,11 @@ class Client {
await this.sendEvent(event)
}
async updateNotifications (contextId: ContextID, query: UpdateNotificationQuery, read: boolean): Promise<void> {
async updateNotifications (
contextId: ContextID,
query: Pick<NotificationQuery, 'type' | 'untilDate'> & { id?: NotificationID },
read: boolean
): Promise<void> {
const event: UpdateNotificationEvent = {
type: NotificationEventType.UpdateNotification,
contextId,
@@ -342,37 +354,32 @@ class Client {
await this.sendEvent(event)
}
async findMessages (params: FindMessagesParams, queryId?: number): Promise<Message[]> {
async findMessagesMeta (params: FindMessagesMetaParams): Promise<MessageMeta[]> {
return (
await this.connection.domainRequest<Message[]>(COMMUNICATION, {
findMessages: { params, queryId }
})
).value
}
async findMessagesGroups (params: FindMessagesGroupsParams): Promise<MessagesGroup[]> {
return (
await this.connection.domainRequest<MessagesGroup[]>(COMMUNICATION, {
findMessagesGroups: { params }
await this.connection.domainRequest<MessageMeta[]>(COMMUNICATION, {
findMessagesMeta: { params }
})
).value
}
async findNotificationContexts (
params: FindNotificationContextParams,
queryId?: number
subscription?: number | string
): Promise<NotificationContext[]> {
return (
await this.connection.domainRequest<NotificationContext[]>(COMMUNICATION, {
findNotificationContexts: { params, queryId }
findNotificationContexts: { params, subscription }
})
).value
}
async findNotifications (params: FindNotificationsParams, queryId?: number): Promise<WithTotal<Notification>> {
async findNotifications (
params: FindNotificationsParams,
subscription?: number | string
): Promise<WithTotal<Notification>> {
return (
await this.connection.domainRequest<WithTotal<Notification>>(COMMUNICATION, {
findNotifications: { params, queryId }
findNotifications: { params, subscription }
})
).value
}
@@ -393,9 +400,15 @@ class Client {
).value
}
async unsubscribeQuery (id: number): Promise<void> {
async subscribeCard (cardId: CardID, subscription: string | number): Promise<void> {
await this.connection.domainRequest<Message[]>(COMMUNICATION, {
unsubscribeQuery: id
subscribeCard: { cardId, subscription }
})
}
async unsubscribeCard (cardId: CardID, subscription: string | number): Promise<void> {
await this.connection.domainRequest<Message[]>(COMMUNICATION, {
unsubscribeCard: { cardId, subscription }
})
}
@@ -438,7 +451,7 @@ class Client {
return id
}
private getAccount (): AccountID {
private getAccount (): AccountUuid {
return getCurrentAccount().uuid
}
}
+2 -1
View File
@@ -180,7 +180,8 @@ export default plugin(presentationId, {
StatsUrl: '' as Metadata<string>,
MailUrl: '' as Metadata<string>,
PreviewUrl: '' as Metadata<string>,
PulseUrl: '' as Metadata<string>
PulseUrl: '' as Metadata<string>,
HulylakeUrl: '' as Metadata<string>
},
status: {
FileTooLarge: '' as StatusCode
@@ -65,7 +65,7 @@
)
$: tab?.id &&
contextsQuery.query({ card: tab.id as Ref<Card>, limit: 1 }, (res) => {
contextsQuery.query({ cardId: tab.id as Ref<Card>, limit: 1 }, (res) => {
context = res.getResult()[0]
isContextLoaded = true
})
@@ -48,7 +48,7 @@
})
$: notificationsQuery.query(
{ card: cardId, limit: 1, read: false, order: SortingOrder.Descending, type: NotificationType.Message },
{ cardId, limit: 1, read: false, order: SortingOrder.Descending, type: NotificationType.Message },
(res) => {
count = res.getResult().length
}
@@ -84,7 +84,7 @@
}
})
$: contextsQuery.query({ card: _id, limit: 1 }, (res) => {
$: contextsQuery.query({ cardId: _id, limit: 1 }, (res) => {
context = res.getResult()[0]
isContextLoaded = true
})
@@ -42,10 +42,14 @@
let person: Person | undefined = undefined
$: messagesQuery.query(
{ card: card._id, strict: true, attachments: true, reactions: true, limit: 1, order: SortingOrder.Descending },
{ cardId: card._id, limit: 1, order: SortingOrder.Descending },
(res) => {
const msgs = res.getResult().reverse()
message = msgs[msgs.length - 1]
},
{
attachments: true,
reactions: true
}
)
@@ -103,7 +107,7 @@
{#if isCompact}
<MessagePreview {card} {message} colorInherit />
{:else}
<MessagePresenter {card} {message} hideHeader hideAvatar readonly padding="0" thread={false} />
<MessagePresenter {card} {message} hideHeader hideAvatar readonly padding="0" showThreads={false} />
{/if}
{/if}
</div>
@@ -93,7 +93,7 @@
console.error('Failed to create thread card')
return
}
const blobs: BlobParams[] = descriptionBox.getAttachments().map((attachment) => ({
const blobs: (BlobParams & { mimeType: string })[] = descriptionBox.getAttachments().map((attachment) => ({
blobId: attachment.file,
mimeType: attachment.type,
fileName: attachment.name,
@@ -113,7 +113,11 @@
}
}
async function createMessage (card: Ref<Card>, markup: Markup, blobs: BlobParams[]): Promise<void> {
async function createMessage (
card: Ref<Card>,
markup: Markup,
blobs: (BlobParams & { mimeType: string })[]
): Promise<void> {
const markdown = markupToMarkdown(markupToJSON(markup))
const { messageId } = await communicationClient.createMessage(card, type, markdown)
@@ -121,7 +125,7 @@
void communicationClient.attachmentPatch<BlobParams>(card, messageId, {
add: blobs.map((it) => ({
id: it.blobId as any as AttachmentID,
type: it.mimeType,
mimeType: it.mimeType,
params: it
}))
})
@@ -70,7 +70,7 @@
$: if (favorites.length > 0) {
contextsQuery.query(
{
card: favorites.map((it) => it.attachedTo),
cardId: favorites.map((it) => it.attachedTo),
notifications: {
read: false,
type: NotificationType.Message,
@@ -105,7 +105,7 @@
$: if (cards.length > 0) {
notificationContextsQuery.query(
{
card: cards.map((it) => it._id),
cardId: cards.map((it) => it._id),
notifications: {
type: NotificationType.Message,
order: SortingOrder.Descending,
@@ -35,7 +35,7 @@
let total = 0
notificationsQuery.query({ limit: 1, total: true, read: false, strict: true, context: context.id }, (res) => {
notificationsQuery.query({ limit: 1, total: true, read: false, strict: true, contextId: context.id }, (res) => {
total = res.getTotal()
})
@@ -48,7 +48,6 @@
query.query(
{
notifications: {
message: true,
order: SortingOrder.Descending,
limit: 3
},
@@ -63,7 +62,8 @@
if (contexts.length < limit && window.hasPrevPage()) {
void window.loadPrevPage()
}
}
},
{ message: true }
)
$: cardsQuery.query(cardPlugin.class.Card, { _id: { $in: contexts.map((c) => c.cardId) } }, (res) => {
@@ -16,12 +16,11 @@
<script lang="ts">
import { Notification } from '@hcengineering/communication-types'
import { Card } from '@hcengineering/card'
import NotificationPreview from './preview/NotificationPreview.svelte'
export let notification: Notification
export let card: Card
</script>
{#if notification.message}
<NotificationPreview {card} message={notification.message} date={notification.created} />
{/if}
<NotificationPreview {card} message={notification.message} creator={notification.creator} date={notification.created} />
@@ -13,7 +13,7 @@
-->
<script lang="ts">
import { Notification, ReactionNotificationContent, SocialID } from '@hcengineering/communication-types'
import { Message, Notification, ReactionNotificationContent, SocialID } from '@hcengineering/communication-types'
import { EmojiPresenter } from '@hcengineering/emoji-resources'
import { Card } from '@hcengineering/card'
import { Label } from '@hcengineering/ui'
@@ -31,7 +31,7 @@
$: content = notification.content as ReactionNotificationContent
let author: Person | undefined
$: void updateAuthor(content.creator)
$: void updateAuthor(notification.creator)
async function updateAuthor (socialId: SocialID): Promise<void> {
author = $employeeByPersonIdStore.get(socialId)
@@ -40,6 +40,9 @@
author = (await getPersonByPersonId(socialId)) ?? undefined
}
}
let message: Message | undefined = undefined
$: message = notification.message
</script>
{#if notification.message}
@@ -57,7 +60,8 @@
</div>
<NotificationPreview
{card}
message={notification.message}
{message}
creator={notification.creator}
date={notification.created}
kind="column"
padding="0"
@@ -26,7 +26,8 @@
import PreviewTemplate from './PreviewTemplate.svelte'
export let card: Card
export let message: Message
export let message: Message | undefined = undefined
export let creator: SocialID
export let date: Date
export let color: 'primary' | 'secondary' = 'primary'
export let kind: 'default' | 'column' = 'default'
@@ -36,9 +37,10 @@
let person: WithLookup<Person> | undefined = undefined
$: void updatePerson(message.creator)
$: void updatePerson(creator)
function getTooltipLabel (message: Message): IntlString {
function getTooltipLabel (message: Message | undefined): IntlString {
if (message == null) return getEmbeddedLabel('')
const text = markupToText(jsonToMarkup(markdownToMarkup(message.content)))
if (text.length > tooltipLimit) {
return getEmbeddedLabel(text.substring(0, tooltipLimit) + '...')
@@ -56,20 +58,24 @@
{color}
{person}
{padding}
socialId={message.creator}
socialId={creator}
{date}
fixHeight={message.type !== MessageType.Activity}
fixHeight={message == null || message.type !== MessageType.Activity}
tooltipLabel={getTooltipLabel(message)}
>
<svelte:fragment slot="content">
{#if isActivityMessage(message)}
<ActivityMessageViewer {message} {card} author={person} />
{:else}
<LiteMessageViewer message={markdownToMarkup(message.content)} />
{#if message}
{#if isActivityMessage(message)}
<ActivityMessageViewer {message} {card} author={person} />
{:else}
<LiteMessageViewer message={markdownToMarkup(message.content)} />
{/if}
{/if}
</svelte:fragment>
<svelte:fragment slot="after">
<AttachmentsPreview {message} />
{#if message}
<AttachmentsPreview {message} />
{/if}
</svelte:fragment>
</PreviewTemplate>
+16 -14
View File
@@ -72,20 +72,27 @@ export const replyInThread: MessageActionFunction = async (message: Message, par
await showForbidden()
return
}
await attachCardToMessage(message, parentCard, createThreadTitle(message, parentCard), chat.masterTag.Thread)
await attachCardToMessage(
message,
parentCard,
createThreadTitle(message, parentCard),
chat.masterTag.Thread,
`${parentCard._id}_${message.id}` as Ref<Card>
)
}
export async function attachCardToMessage (
message: Message,
parentCard: Card,
title: string,
type: Ref<MasterTag>
type: Ref<MasterTag>,
_id?: Ref<Card>
): Promise<void> {
const client = getClient()
const communicationClient = getCommunicationClient()
const hierarchy = client.getHierarchy()
const thread = message.thread
const thread = _id != null ? message.threads.find((it) => it.threadId === _id) : undefined
if (thread != null) {
const _id = thread.threadId
const card = await client.findOne(cardPlugin.class.Card, { _id: _id as Ref<Card> })
@@ -95,7 +102,7 @@ export async function attachCardToMessage (
return
}
const threadCardID = generateId<Card>()
const threadCardID = _id ?? generateId<Card>()
await communicationClient.attachThread(parentCard._id, message.id, threadCardID, type)
@@ -140,11 +147,7 @@ function createThreadTitle (message: Message, parent: Card): string {
}
export const canReplyInThread: MessageActionVisibilityTester = (message: Message): boolean => {
return (
message.type === MessageType.Message &&
message.extra?.threadRoot !== true &&
(!message.removed || message.thread != null)
)
return message.type === MessageType.Text && message.extra?.threadRoot !== true
}
export const translateMessage: MessageActionFunction = async (message: Message): Promise<void> => {
@@ -183,7 +186,7 @@ export const translateMessage: MessageActionFunction = async (message: Message):
export const canTranslateMessage: MessageActionVisibilityTester = (message: Message): boolean => {
const url = getMetadata(aiBot.metadata.EndpointURL) ?? ''
if (url === '') return false
return message.type === MessageType.Message && !message.removed
return message.type === MessageType.Text
}
export const showOriginalMessage: MessageActionFunction = async (message: Message): Promise<void> => {
@@ -205,19 +208,18 @@ export const editMessage: MessageActionFunction = async (message: Message): Prom
}
export const canEditMessage: MessageActionVisibilityTester = (message: Message): boolean => {
if (message.type !== MessageType.Message || message.removed) return false
if (message.type !== MessageType.Text) return false
const me = getCurrentAccount()
return me.socialIds.includes(message.creator)
}
export const removeMessage: MessageActionFunction = async (message: Message): Promise<void> => {
const communicationClient = getCommunicationClient()
message.removed = true
await communicationClient.removeMessage(message.cardId, message.id)
}
export const canRemoveMessage: MessageActionVisibilityTester = (message: Message): boolean => {
if (message.type !== MessageType.Message || message.removed) return false
if (message.type !== MessageType.Text) return false
const me = getCurrentAccount()
return me.socialIds.includes(message.creator)
}
@@ -234,7 +236,7 @@ export const createCard: MessageActionFunction = async (message: Message, card:
}
export const canCreateCard: MessageActionVisibilityTester = (message: Message): boolean => {
return canReplyInThread(message) && message.thread == null
return canReplyInThread(message)
}
let allMessageActions: MessageAction[] | undefined
@@ -31,7 +31,7 @@
{/if}
{#if isAppletAttachment(attachment)}
{@const applet = applets.find((it) => it.type === attachment.type)}
{@const applet = applets.find((it) => it.type === attachment.mimeType)}
{#if applet}
<span class:lower>
<Label label={applet.label} />:
@@ -107,7 +107,7 @@
type="type-popup"
okLabel={presentation.string.Create}
okAction={attachCard}
canSave={selectedType !== undefined && title.trim() !== '' && _message.thread == null}
canSave={selectedType !== undefined && title.trim() !== ''}
onCancel={() => dispatch('close')}
on:close
>
@@ -127,10 +127,10 @@
/>
</div>
<div class="mt-4" />
<MessagePresenter {card} message={{ ..._message, reactions: [], thread: undefined }} readonly={true} padding="0" />
<MessagePresenter {card} message={{ ..._message, reactions: {}, threads: [] }} readonly={true} padding="0" />
</div>
<svelte:fragment slot="footer">
{#if _message.thread != null && !inProgress}
{#if !inProgress}
<div class="footer-error">
<Label label={communication.string.MessageAlreadyHasCardAttached} />
</div>
@@ -29,7 +29,7 @@
import { getCurrentAccount, SortingOrder } from '@hcengineering/core'
import { createEventDispatcher, onDestroy, onMount, tick } from 'svelte'
import { MessagesNavigationAnchors } from '@hcengineering/communication'
import { isAppFocusedStore, deviceOptionsStore as deviceInfo } from '@hcengineering/ui'
import { deviceOptionsStore as deviceInfo, isAppFocusedStore } from '@hcengineering/ui'
import { createMessagesObserver, getGroupDay, groupMessagesByDay, MessagesGroup } from '../messages'
import MessagesGroupPresenter from './message/MessagesGroupPresenter.svelte'
@@ -101,27 +101,34 @@
$: reinit(position)
$: query.query(queryDef, (res: Window<Message>) => {
window = res
messages = (queryDef.order === SortingOrder.Ascending ? res.getResult() : res.getResult().reverse()).filter(
(it) => !it.removed || it.thread != null
)
$: query.query(
queryDef,
(res: Window<Message>) => {
window = res
messages = queryDef.order === SortingOrder.Ascending ? res.getResult() : res.getResult().reverse()
if (messages.length < limit && res.hasNextPage()) {
void window.loadNextPage()
} else if (messages.length < limit && res.hasPrevPage()) {
void window.loadPrevPage()
if (messages.length < limit && res.hasNextPage()) {
void window.loadNextPage()
} else if (messages.length < limit && res.hasPrevPage()) {
void window.loadPrevPage()
}
groups = groupMessagesByDay(messages)
isLoading = messages.length < limit && (res.hasNextPage() || res.hasPrevPage())
void onUpdate(messages)
},
{
autoExpand: true,
threads: true,
attachments: true,
reactions: true
}
groups = groupMessagesByDay(messages)
isLoading = messages.length < limit && (res.hasNextPage() || res.hasPrevPage())
void onUpdate(messages)
})
)
$: if (context !== undefined) {
void notificationsQuery.query(
{
context: context.id,
contextId: context.id,
read: false
},
(res) => {
@@ -193,10 +200,7 @@
function getBaseQuery (): MessageQueryParams {
if (position === MessagesNavigationAnchors.ConversationStart) {
return {
card: card._id,
replies: true,
attachments: true,
reactions: true,
cardId: card._id,
order: SortingOrder.Ascending,
limit
}
@@ -206,10 +210,7 @@
const unread = initialLastView != null && initialLastUpdate != null && initialLastUpdate > initialLastView
const order = unread && !shouldScrollToEnd ? SortingOrder.Ascending : SortingOrder.Descending
return {
card: card._id,
replies: true,
attachments: true,
reactions: true,
cardId: card._id,
order,
limit,
from: unread && !shouldScrollToEnd && initialLastView != null ? initialLastView : undefined
@@ -13,7 +13,7 @@
<script lang="ts">
import { Icon, IconComponent, IconSize, tooltip } from '@hcengineering/ui'
import { PersonId } from '@hcengineering/core'
import { PersonUuid } from '@hcengineering/core'
import { EmojiPresenter } from '@hcengineering/emoji-resources'
import ReactionsTooltip from './ReactionsTooltip.svelte'
@@ -24,7 +24,7 @@
export let count: number | undefined = undefined
export let selected: boolean = false
export let active: boolean = false
export let socialIds: PersonId[] = []
export let persons: PersonUuid[] = []
</script>
<!-- svelte-ignore a11y-click-events-have-key-events -->
@@ -34,7 +34,7 @@
class:selected
class:active
on:click
use:tooltip={(count ?? 0) > 0 ? { component: ReactionsTooltip, props: { socialIds } } : undefined}
use:tooltip={(count ?? 0) > 0 ? { component: ReactionsTooltip, props: { persons } } : undefined}
>
<div class="reaction__emoji" class:foreground={icon != null}>
{#if icon}
@@ -14,13 +14,13 @@
<script lang="ts">
import { createEventDispatcher } from 'svelte'
import { showPopup } from '@hcengineering/ui'
import { getCurrentAccount, groupByArray } from '@hcengineering/core'
import { Reaction } from '@hcengineering/communication-types'
import { getCurrentAccount } from '@hcengineering/core'
import { Emoji, EmojiData } from '@hcengineering/communication-types'
import emojiPlugin from '@hcengineering/emoji'
import ReactionPresenter from './ReactionPresenter.svelte'
export let reactions: Reaction[] = []
export let reactions: Record<Emoji, EmojiData[]> = {}
const dispatch = createEventDispatcher()
const me = getCurrentAccount()
@@ -46,18 +46,15 @@
() => {}
)
}
let reactionsByEmoji = new Map<string, Reaction[]>()
$: reactionsByEmoji = groupByArray(reactions, (it) => it.reaction)
</script>
<div class="reactions">
{#each reactionsByEmoji as [emoji, reactions] (emoji)}
{#each Object.entries(reactions) as [emoji, data] (emoji)}
<ReactionPresenter
{emoji}
selected={reactions.some((it) => me.socialIds.includes(it.creator))}
socialIds={reactions.map((it) => it.creator)}
count={reactions.length}
selected={data.some((it) => it.person === me.uuid)}
persons={data.map((it) => it.person)}
count={data.length}
on:click={() => dispatch('click', emoji)}
/>
{/each}
@@ -12,21 +12,20 @@
<!-- limitations under the License. -->
<script lang="ts">
import { notEmpty, PersonId, Ref } from '@hcengineering/core'
import { AccountUuid, notEmpty, PersonUuid } from '@hcengineering/core'
import { ObjectPresenter } from '@hcengineering/view-resources'
import contact, { Person } from '@hcengineering/contact'
import { getPersonRefsByPersonIdsCb } from '@hcengineering/contact-resources'
import { Employee } from '@hcengineering/contact'
import { employeeByAccountStore } from '@hcengineering/contact-resources'
export let socialIds: PersonId[] = []
export let persons: PersonUuid[] = []
let persons: Set<Ref<Person>>
$: getPersonRefsByPersonIdsCb(socialIds, (personsMap) => {
persons = new Set(personsMap.values().filter(notEmpty))
})
let employees: Employee[] = []
$: employees = persons.map((it) => $employeeByAccountStore.get(it as AccountUuid)).filter(notEmpty)
</script>
<div class="m-2 flex-col flex-gap-2">
{#each persons as person (person)}
<ObjectPresenter objectId={person} _class={contact.class.Person} disabled />
{#each employees as emp (emp._id)}
<ObjectPresenter objectId={emp._id} _class={emp._class} value={emp} disabled />
{/each}
</div>
@@ -18,7 +18,7 @@
import { formatName, Person } from '@hcengineering/contact'
import { Message } from '@hcengineering/communication-types'
import { Card } from '@hcengineering/card'
import { IconDelete, Label } from '@hcengineering/ui'
import { Label } from '@hcengineering/ui'
import communication from '../../plugin'
import MessageInput from './MessageInput.svelte'
@@ -34,7 +34,7 @@
export let compact: boolean = false
export let hideAvatar: boolean = false
export let hideHeader: boolean = false
export let thread: boolean = true
export let showThreads: boolean = true
function formatDate (date: Date): string {
return date.toLocaleTimeString('default', {
@@ -74,7 +74,7 @@
/>
{/if}
{#if !isEditing}
<MessageFooter {message} {thread} />
<MessageFooter {message} {showThreads} />
{/if}
</div>
</div>
@@ -82,38 +82,32 @@
<div class="message__body">
{#if !hideAvatar}
<div class="message__avatar">
{#if !message.removed}
<PersonPreviewProvider value={author}>
<Avatar name={author?.name} person={author} size="medium" />
</PersonPreviewProvider>
{:else}
<Avatar icon={IconDelete} size="medium" />
{/if}
<PersonPreviewProvider value={author}>
<Avatar name={author?.name} person={author} size="medium" />
</PersonPreviewProvider>
</div>
{/if}
<div class="message__content">
<div class="message__header">
{#if !message.removed}
<PersonPreviewProvider value={author}>
<div class="message__username">
{formatName(author?.name ?? '')}
</div>
</PersonPreviewProvider>
{/if}
<PersonPreviewProvider value={author}>
<div class="message__username">
{formatName(author?.name ?? '')}
</div>
</PersonPreviewProvider>
<div class="message__date">
{formatDate(message.created)}
</div>
{#if message.edited && !message.removed}
{#if message.modified}
<div class="message__edited-marker">
(<Label label={communication.string.Edited} />)
</div>
{/if}
{#if !message.removed && $translateMessagesStore.get(message.id)?.inProgress === true}
{#if $translateMessagesStore.get(message.id)?.inProgress === true}
<div class="message__translating">
<Label label={communication.string.Translating} />
</div>
{/if}
{#if !message.removed && $translateMessagesStore.get(message.id)?.shown === true}
{#if $translateMessagesStore.get(message.id)?.shown === true}
<div class="message__show-original" on:click={() => showOriginalMessage(message, card)}>
<Label label={communication.string.ShowOriginal} />
</div>
@@ -136,7 +130,7 @@
/>
{/if}
{#if !isEditing}
<MessageFooter {message} {thread} />
<MessageFooter {message} {showThreads} />
{/if}
</div>
</div>
@@ -61,10 +61,6 @@
{#if isActivityMessage(message)}
<ActivityMessageViewer {message} {card} {author} />
{:else if message.removed}
<span class="overflow-label removed-label">
<Label label={communication.string.MessageWasRemoved} />
</span>
{:else}
<MarkupMessageViewer message={displayMarkup} />
{/if}
@@ -15,10 +15,10 @@
<script lang="ts">
import { getClient, getCommunicationClient } from '@hcengineering/presentation'
import cardPlugin from '@hcengineering/card'
import { getCurrentAccount } from '@hcengineering/core'
import cardPlugin, { Card } from '@hcengineering/card'
import { getCurrentAccount, Ref } from '@hcengineering/core'
import { AttachmentPreview, LinkPreview } from '@hcengineering/attachment-resources'
import { AttachmentID, Message, MessageType } from '@hcengineering/communication-types'
import { AttachmentID, Emoji, Message, MessageType } from '@hcengineering/communication-types'
import { getResource } from '@hcengineering/platform'
import { isAppletAttachment, isBlobAttachment, isLinkPreviewAttachment } from '@hcengineering/communication-shared'
import { Component } from '@hcengineering/ui'
@@ -29,7 +29,7 @@
import communication from '../../plugin'
export let message: Message
export let thread: boolean = true
export let showThreads: boolean = true
const me = getCurrentAccount()
const communicationClient = getCommunicationClient()
@@ -39,16 +39,18 @@
return message.type !== MessageType.Activity && message.extra?.threadRoot !== true
}
async function handleReaction (event: CustomEvent<string>): Promise<void> {
async function handleReaction (event: CustomEvent<Emoji>): Promise<void> {
event.preventDefault()
event.stopPropagation()
const emoji = event.detail
await toggleReaction(message, emoji)
}
async function handleReply (): Promise<void> {
async function handleReply (event: CustomEvent<Ref<Card> | undefined>): Promise<void> {
if (!canReply()) return
const t = message.thread
const threadID = event.detail
if (threadID === undefined) return
const t = message.threads.find((it) => it.threadId === threadID)
if (t === undefined) return
const _id = t.threadId
const client = getClient()
@@ -70,10 +72,10 @@
$: applets = message.attachments.filter(isAppletAttachment) ?? []
</script>
{#if applets.length > 0 && !message.removed}
{#if applets.length > 0}
<div class="message__applets">
{#each applets as applet (applet.id)}
{@const appletModel = appletsModels.find((it) => it.type === applet.type)}
{@const appletModel = appletsModels.find((it) => it.type === applet.mimeType)}
{#if appletModel}
<Component
is={appletModel.component}
@@ -87,13 +89,13 @@
</div>
{/if}
{#if blobs.length > 0 && !message.removed}
{#if blobs.length > 0}
<div class="message__files">
{#each blobs as blob (blob.id)}
<AttachmentPreview
value={{
file: blob.params.blobId,
type: blob.params.mimeType,
type: blob.mimeType,
name: blob.params.fileName,
size: blob.params.size,
metadata: blob.params.metadata
@@ -103,7 +105,7 @@
{/each}
</div>
{/if}
{#if links.length > 0 && !message.removed}
{#if links.length > 0}
<div class="message__links">
{#each links as link (link.id)}
<LinkPreview
@@ -126,14 +128,17 @@
{/each}
</div>
{/if}
{#if message.reactions.length > 0 && !message.removed}
{#if Object.keys(message.reactions).length > 0}
<div class="message__reactions">
<ReactionsList reactions={message.reactions} on:click={handleReaction} />
</div>
{/if}
{#if thread && message.thread && message.thread.threadId}
{#if showThreads && message.threads.length > 0}
<div class="message__replies overflow-label">
<MessageThread thread={message.thread} on:click={handleReply} />
{#each message.threads as thread (thread.threadId)}
<MessageThread {thread} on:click={handleReply} />
{/each}
</div>
{/if}
@@ -147,9 +152,10 @@
margin-left: -0.5rem;
padding-bottom: 0;
display: flex;
align-items: flex-start;
align-self: stretch;
flex-direction: column;
gap: 0.5rem;
overflow: hidden;
width: 100%;
}
.message__files {
@@ -163,7 +163,7 @@
if (toAttach.length > 0) {
void communicationClient.attachmentPatch<AppletParams>(card._id, messageId, {
add: toAttach.map((it) => ({
type: it.type,
mimeType: it.mimeType,
params: it.params
}))
})
@@ -172,7 +172,7 @@
async function createMessage (
markdown: string,
blobs: BlobParams[],
blobs: (BlobParams & { mimeType: string })[],
links: LinkPreviewParams[],
urlsToLoad: string[],
appletDrafts: AppletDraft[]
@@ -186,7 +186,7 @@
void communicationClient.attachmentPatch<BlobParams>(card._id, messageId, {
add: blobs.map((it) => ({
id: it.blobId as any as AttachmentID,
type: it.mimeType,
mimeType: it.mimeType,
params: it
}))
})
@@ -195,7 +195,7 @@
if (links.length > 0) {
void communicationClient.attachmentPatch<LinkPreviewParams>(card._id, messageId, {
add: links.map((it) => ({
type: linkPreviewType,
mimeType: linkPreviewType,
params: it
}))
})
@@ -210,7 +210,7 @@
void communicationClient.attachmentPatch<LinkPreviewParams>(card._id, messageId, {
add: [
{
type: linkPreviewType,
mimeType: linkPreviewType,
params
}
]
@@ -221,7 +221,7 @@
async function editMessage (
message: Message,
markdown: string,
blobs: BlobParams[],
blobs: (BlobParams & { mimeType: string })[],
links: LinkPreviewParams[],
appletDrafts: AppletDraft[]
): Promise<void> {
@@ -240,7 +240,7 @@
void communicationClient.attachmentPatch<BlobParams>(card._id, message.id, {
add: attachBlobs.map((it) => ({
id: it.blobId as any as AttachmentID,
type: it.mimeType,
mimeType: it.mimeType,
params: it
}))
})
@@ -274,7 +274,7 @@
void communicationClient.attachmentPatch(card._id, message.id, {
add: attachLinks.map((it) => ({
type: linkPreviewType,
mimeType: linkPreviewType,
params: it
}))
})
@@ -495,7 +495,10 @@
if (result != null) {
draft = {
...draft,
applets: [...draft.applets, { id: generateId(), type: applet.type, appletId: applet._id, params: result }]
applets: [
...draft.applets,
{ id: generateId(), mimeType: applet.type, appletId: applet._id, params: result }
]
}
}
})
@@ -43,13 +43,11 @@
export let hideAvatar: boolean = false
export let hideHeader: boolean = false
export let readonly: boolean = false
export let thread: boolean = true
export let showThreads: boolean = true
let isEditing = false
let isDeleted = false
let author: Person | undefined
$: isDeleted = message.removed
$: isEditing = $messageEditingStore === message.id
$: void updateAuthor(message.creator)
@@ -138,8 +136,7 @@
let isActionsPanelOpened = false
$: showActions = !isEditing && !isDeleted && !readonly
$: isThread = message.thread != null
$: showActions = !isEditing && !readonly
</script>
<!-- svelte-ignore a11y-click-events-have-key-events -->
@@ -152,7 +149,7 @@
class:noHover={readonly}
style:padding
>
{#if message.type === MessageType.Activity || (message.removed && message.thread?.threadId === undefined)}
{#if message.type === MessageType.Activity}
<OneRowMessageBody {message} {card} {author} {hideAvatar} {hideHeader} />
{:else}
<MessageBody
@@ -160,10 +157,10 @@
{card}
{author}
{isEditing}
compact={compact && !isThread}
compact={compact && message.threads.length === 0}
{hideAvatar}
{hideHeader}
{thread}
{showThreads}
/>
{/if}
@@ -51,8 +51,6 @@
{#each messages as message, index (message.id)}
{@const previousMessage = messages[index - 1]}
{@const compact =
!message.removed &&
!previousMessage?.removed &&
previousMessage !== undefined &&
previousMessage.creator === message.creator &&
previousMessage.type === message.type &&
@@ -18,7 +18,6 @@
import { formatName, Person } from '@hcengineering/contact'
import { Message } from '@hcengineering/communication-types'
import { Card } from '@hcengineering/card'
import { IconDelete } from '@hcengineering/ui'
import MessageContentViewer from './MessageContentViewer.svelte'
import MessageFooter from './MessageFooter.svelte'
@@ -35,24 +34,17 @@
minute: 'numeric'
})
}
let isDeleted = false
$: isDeleted = message.removed
</script>
<div class="message__body">
{#if !hideAvatar}
<div class="message__avatar">
{#if !isDeleted}
<PersonPreviewProvider value={author}>
<Avatar name={author?.name} person={author} size="x-small" />
</PersonPreviewProvider>
{:else}
<Avatar icon={IconDelete} size="x-small" />
{/if}
<PersonPreviewProvider value={author}>
<Avatar name={author?.name} person={author} size="x-small" />
</PersonPreviewProvider>
</div>
{/if}
{#if !isDeleted && !hideHeader}
{#if !hideHeader}
<div class="message__header">
<PersonPreviewProvider value={author}>
<div class="message__username">
@@ -69,11 +61,9 @@
<MessageContentViewer {message} {card} {author} />
</div>
</div>
{#if !isDeleted}
<div class="message__footer">
<MessageFooter {message} />
</div>
{/if}
<div class="message__footer">
<MessageFooter {message} />
</div>
<style lang="scss">
.message__body {
@@ -17,6 +17,7 @@
import { createQuery } from '@hcengineering/presentation'
import { Thread } from '@hcengineering/communication-types'
import cardPlugin, { Card } from '@hcengineering/card'
import { createEventDispatcher } from 'svelte'
import ThreadCollaborators from './ThreadCollaborators.svelte'
import ThreadRepliesCount from './ThreadRepliesCount.svelte'
@@ -25,6 +26,8 @@
import ThreadTitle from './ThreadTitle.svelte'
export let thread: Thread
const dispatch = createEventDispatcher()
const threadCardQuery = createQuery()
let threadCard: Card | undefined
@@ -52,13 +55,13 @@
<div class="replies-container flex-grow" bind:clientWidth>
<!-- svelte-ignore a11y-click-events-have-key-events -->
<!-- svelte-ignore a11y-no-static-element-interactions -->
<div class="replies" on:click style:max-width={`${clientWidth}px`}>
<ThreadCollaborators threadId={thread.threadId} />
<div class="replies" on:click={() => dispatch('click', threadCard?._id)} style:max-width={`${clientWidth}px`}>
<ThreadCollaborators persons={thread.repliedPersons} />
{#if thread.repliesCount > 0}
<span class="text overflow-label">
<ThreadRepliesCount count={thread.repliesCount} />
<ThreadLastReply lastReply={thread.lastReply} />
<ThreadLastReply lastReply={thread.lastReplyDate} />
</span>
{/if}
@@ -12,32 +12,23 @@
<!-- limitations under the License. -->
<script lang="ts">
import { AccountUuid, Ref } from '@hcengineering/core'
import { Card } from '@hcengineering/card'
import { AccountUuid, PersonUuid } from '@hcengineering/core'
import { Avatar, employeeByAccountStore } from '@hcengineering/contact-resources'
import { createCollaboratorsQuery } from '@hcengineering/presentation'
import { Collaborator } from '@hcengineering/communication-types'
import { Person } from '@hcengineering/contact'
export let threadId: Ref<Card>
export let persons: Record<PersonUuid, number> = {}
const displayPersonsNumber = 4
const collaboratorsQuery = createCollaboratorsQuery()
let _persons: Person[] = []
let collaborators: Collaborator[] = []
let persons: Person[] = []
$: updatePersons(persons, $employeeByAccountStore)
$: collaboratorsQuery.query({ card: threadId }, (res) => {
collaborators = res
})
$: updatePersons(collaborators, $employeeByAccountStore)
function updatePersons (collaborators: Collaborator[], employeeByAccount: Map<AccountUuid, Person>): void {
function updatePersons (persons: Record<PersonUuid, number>, employeeByAccount: Map<AccountUuid, Person>): void {
const newPersons: Person[] = []
for (const collaborator of collaborators) {
const person = employeeByAccount.get(collaborator.account)
for (const [personUuid, count] of Object.entries(persons)) {
if (count < 1) continue
const person = employeeByAccount.get(personUuid as AccountUuid)
if (person !== undefined) {
newPersons.push(person)
}
@@ -47,19 +38,21 @@
}
}
persons = newPersons
_persons = newPersons
}
$: count = Object.entries(persons).filter(([, count]) => count > 0).length
</script>
{#if persons.length > 0}
{#if _persons.length > 0}
<div class="thread__avatars">
{#each persons as person}
{#each _persons as person}
<Avatar size="x-small" {person} name={person.name} />
{/each}
</div>
{#if collaborators.length > displayPersonsNumber}
+{collaborators.length - displayPersonsNumber}
{#if count > displayPersonsNumber}
+{count - displayPersonsNumber}
{/if}
{/if}
+3 -3
View File
@@ -99,14 +99,14 @@ export function messageToDraft (message: Message): MessageDraft {
return {
_id: message.id,
content: toMarkup(message.content),
blobs: message.attachments.filter(isBlobAttachment).map((it) => it.params),
blobs: message.attachments.filter(isBlobAttachment).map((it) => ({ ...it.params, mimeType: it.mimeType })),
links: message.attachments.filter(isLinkPreviewAttachment).map((it) => it.params),
applets: message.attachments
.filter(isAppletAttachment)
.map((it) => ({
id: it.id,
type: it.type,
appletId: applets.find((a) => a.type === it.type)?._id as any,
mimeType: it.mimeType,
appletId: applets.find((a) => a.type === it.mimeType)?._id as any,
params: it.params
}))
.filter((it) => it.appletId !== undefined)
@@ -13,12 +13,11 @@
import { get, writable } from 'svelte/store'
import { createLabelsQuery, createQuery, onClient, onCommunicationClient } from '@hcengineering/presentation'
import type { Label, Message, MessageID } from '@hcengineering/communication-types'
import type { Markup, Ref } from '@hcengineering/core'
import { type Label, type LabelID, type Message, type MessageID } from '@hcengineering/communication-types'
import core, { getCurrentAccount, type Markup, type Ref } from '@hcengineering/core'
import { languageStore } from '@hcengineering/ui'
import { type Card } from '@hcengineering/card'
import cardPlugin, { type Card } from '@hcengineering/card'
import communication from '@hcengineering/communication'
import core from '@hcengineering/core'
export const labelsStore = writable<Label[]>([])
export const messageEditingStore = writable<MessageID | undefined>(undefined)
@@ -53,6 +52,11 @@ export function isShownTranslatedMessage (messageId: MessageID): boolean {
return result?.shown === true && result?.result != null
}
export function isCardSubscribed (cardId: Ref<Card>): boolean {
const me = getCurrentAccount()
const labelId = cardPlugin.label.Subscribed as string as LabelID
return get(labelsStore).some((it) => it.account === me.uuid && it.cardId === cardId && it.labelId === labelId)
}
const query = createLabelsQuery(true)
onCommunicationClient(() => {
+2 -2
View File
@@ -41,7 +41,7 @@ export interface Action {
export interface AppletDraft {
id: string
type: AppletType
mimeType: AppletType
appletId: Ref<Applet>
params: Record<string, any>
}
@@ -49,7 +49,7 @@ export interface AppletDraft {
export interface MessageDraft {
_id: string
content: Markup
blobs: BlobParams[]
blobs: Array<BlobParams & { mimeType: string }>
links: LinkPreviewParams[]
applets: AppletDraft[]
}
+8 -14
View File
@@ -21,16 +21,16 @@ import { type Card } from '@hcengineering/card'
import { AccountRole, type Data, getCurrentAccount, type Ref, type Space, type Markup } from '@hcengineering/core'
import { getMetadata, translate } from '@hcengineering/platform'
import { addNotification, languageStore, NotificationSeverity, showPopup } from '@hcengineering/ui'
import { type LinkPreviewParams, type Message } from '@hcengineering/communication-types'
import { type Emoji, type LinkPreviewParams, type Message } from '@hcengineering/communication-types'
import emoji from '@hcengineering/emoji'
import { markdownToMarkup, markupToMarkdown } from '@hcengineering/text-markdown'
import { jsonToMarkup, markupToJSON } from '@hcengineering/text'
import { isCardSubscribed, guestCommunicationAllowedCards } from './stores'
import IconAt from './components/icons/At.svelte'
import communication from './plugin'
import { type TextInputAction } from './types'
import { guestCommunicationAllowedCards } from './stores'
import { get } from 'svelte/store'
import view from '@hcengineering/view'
import { type Direct } from '@hcengineering/communication'
@@ -51,19 +51,14 @@ export async function subscribe (card: Card): Promise<void> {
export async function canSubscribe (card: Card): Promise<boolean> {
const isEnabled = getMetadata(communication.metadata.Enabled) === true
if (!isEnabled) return false
const client = getCommunicationClient()
const me = getCurrentAccount()
const collaborator = (await client.findCollaborators({ card: card._id, account: me.uuid, limit: 1 }))[0]
return collaborator === undefined
return !isCardSubscribed(card._id)
}
export async function canUnsubscribe (card: Card): Promise<boolean> {
const isEnabled = getMetadata(communication.metadata.Enabled) === true
if (!isEnabled) return false
const client = getCommunicationClient()
const me = getCurrentAccount()
const collaborator = (await client.findCollaborators({ card: card._id, account: me.uuid, limit: 1 }))[0]
return collaborator !== undefined
return isCardSubscribed(card._id)
}
export const defaultMessageInputActions: TextInputAction[] = [
@@ -107,12 +102,11 @@ export function toMarkup (markdown: string): Markup {
return jsonToMarkup(markdownToMarkup(markdown))
}
export async function toggleReaction (message: Message, emoji: string): Promise<void> {
export async function toggleReaction (message: Message, emoji: Emoji): Promise<void> {
const me = getCurrentAccount()
const communicationClient = getCommunicationClient()
const { socialIds } = me
const reaction = message.reactions.find((it) => it.reaction === emoji && socialIds.includes(it.creator))
if (reaction !== undefined) {
const alreadyAdded = (message.reactions[emoji] ?? []).some((it) => it.person === me.uuid)
if (alreadyAdded) {
await communicationClient.removeReaction(message.cardId, message.id, emoji)
} else {
await communicationClient.addReaction(message.cardId, message.id, emoji)
+3 -3
View File
@@ -11,9 +11,9 @@
// See the License for the specific language governing permissions and
// limitations under the License.
import { AttachedDoc, Ref } from '@hcengineering/core'
import { AttachedDoc, Ref, AccountUuid } from '@hcengineering/core'
import { PersonSpace } from '@hcengineering/contact'
import { AccountID, MessageID } from '@hcengineering/communication-types'
import { MessageID } from '@hcengineering/communication-types'
import { Card } from '@hcengineering/card'
export interface PollAnswer extends AttachedDoc<Poll> {
@@ -22,7 +22,7 @@ export interface PollAnswer extends AttachedDoc<Poll> {
}
export interface UserVote {
account: AccountID
account: AccountUuid
options: { id: string, label: string, votedAt: Date }[]
}
+3 -10
View File
@@ -11,15 +11,8 @@
// See the License for the specific language governing permissions and
// limitations under the License.
import {
AccountID,
AppletAttachment,
AppletParams,
AppletType,
Message,
MessageID
} from '@hcengineering/communication-types'
import { AttachedDoc, Configuration, Doc, Ref } from '@hcengineering/core'
import { AppletAttachment, AppletParams, AppletType, Message, MessageID } from '@hcengineering/communication-types'
import { AttachedDoc, Configuration, Doc, Ref, AccountUuid } from '@hcengineering/core'
import { Asset, IntlString, Resource } from '@hcengineering/platform'
import { Card, MasterTag } from '@hcengineering/card'
import { AnyComponent } from '@hcengineering/ui'
@@ -85,7 +78,7 @@ export interface PollAnswer extends AttachedDoc<Poll> {
}
export interface UserVote {
account: AccountID
account: AccountUuid
options: { id: string, label: string, votedAt: Date }[]
}
+2 -1
View File
@@ -77,6 +77,7 @@
"@hcengineering/mongo": "^0.6.1",
"@hcengineering/kafka": "^0.6.0",
"@hcengineering/communication-server": "^0.1.0",
"@hcengineering/communication-sdk-types": "^0.1.0"
"@hcengineering/communication-sdk-types": "^0.1.0",
"@hcengineering/hulylake-client": "^0.6.0"
}
}
@@ -82,6 +82,7 @@ class TestQueue {
elasticIndexName,
serverSecret: 'secret',
dbURL: dbUrl,
hulylakeUrl: 'http://localhost:8096',
config: dbConfig,
externalStorage: createDummyStorageAdapter(),
listener: {
+3
View File
@@ -104,6 +104,8 @@ if (accountsUrl === undefined) {
process.exit(1)
}
const hulylakeUrl = process.env.HULYLAKE_URL ?? ''
const storageConfig: StorageConfiguration = storageConfigFromEnv()
const externalStorage = buildStorageFromConfig(storageConfig)
@@ -116,6 +118,7 @@ const onClose = startIndexer(metricsContext, {
externalStorage,
elasticIndexName,
dbURL,
hulylakeUrl,
port: servicePort,
serverSecret,
accountsUrl
+3
View File
@@ -33,6 +33,7 @@ import {
import { type QueueSourced, type FulltextDBConfiguration } from '@hcengineering/server-indexer'
import { generateToken } from '@hcengineering/server-token'
import { type Event } from '@hcengineering/communication-sdk-types'
import { getClient as getHulylakeClient } from '@hcengineering/hulylake-client'
import { WorkspaceIndexer } from './workspace'
@@ -55,6 +56,7 @@ export class WorkspaceManager {
private readonly opt: {
queue: PlatformQueue
dbURL: string
hulylakeUrl: string
config: FulltextDBConfiguration
externalStorage: StorageAdapter
elasticIndexName: string
@@ -332,6 +334,7 @@ export class WorkspaceManager {
this.opt.externalStorage,
this.fulltextAdapter,
this.contentAdapter,
getHulylakeClient(this.opt.hulylakeUrl, workspace, token ?? ''),
(token) => this.getTransactorAPIEndpoint(token),
this.opt.listener
)
+1
View File
@@ -84,6 +84,7 @@ export async function startIndexer (
queue: PlatformQueue
model: Tx[]
dbURL: string
hulylakeUrl: string
config: FulltextDBConfiguration
externalStorage: StorageAdapter
elasticIndexName: string
+3
View File
@@ -39,6 +39,7 @@ import { FullTextIndexPipeline } from '@hcengineering/server-indexer'
import { getConfig } from '@hcengineering/server-pipeline'
import { generateToken } from '@hcengineering/server-token'
import { Api as CommunicationApi } from '@hcengineering/communication-server'
import { type HulylakeClient } from '@hcengineering/hulylake-client'
import { fulltextModelFilter } from './utils'
@@ -61,6 +62,7 @@ export class WorkspaceIndexer {
externalStorage: StorageAdapter,
ftadapter: FullTextAdapter,
contentAdapter: ContentTextAdapter,
hulylake: HulylakeClient,
endpointProvider: (token: string) => Promise<string | undefined>,
listener?: FulltextListener
): Promise<WorkspaceIndexer> {
@@ -147,6 +149,7 @@ export class WorkspaceIndexer {
})
}
},
hulylake,
communicationApi,
listener
)
+2 -2
View File
@@ -108,10 +108,10 @@ async function handleCommunicationTx (
.flatMap((it) => it.attachments)
.filter((it): it is BlobAttachment => 'blobId' in it.params)
const messages: VideoTranscodeRequest[] = attachments.map(({ params }) => ({
const messages: VideoTranscodeRequest[] = attachments.map(({ mimeType, params }) => ({
workspaceUuid,
blobId: params.blobId,
contentType: params.mimeType,
contentType: mimeType,
source
}))
+3 -4
View File
@@ -113,15 +113,14 @@ export function start (
const communicationApiFactory: CommunicationApiFactory = async (ctx, workspace, broadcastSessions) => {
if (dbUrl.startsWith('mongodb') || !opt.communicationApiEnabled) {
return {
findMessages: async () => [],
findMessagesGroups: async () => [],
findMessagesMeta: async () => [],
findNotificationContexts: async () => [],
findCollaborators: async () => [],
findNotifications: async () => [],
findLabels: async () => [],
findThreads: async () => [],
findPeers: async () => [],
unsubscribeQuery: async () => {},
subscribeCard: () => {},
unsubscribeCard: () => {},
event: async () => {
return {}
},
-10
View File
@@ -2366,11 +2366,6 @@
"projectFolder": "models/chat",
"shouldPublish": false
},
{
"packageName": "@hcengineering/pod-msg2file",
"projectFolder": "services/msg2file",
"shouldPublish": false
},
{
"packageName": "@hcengineering/process",
"projectFolder": "plugins/process",
@@ -2511,11 +2506,6 @@
"projectFolder": "communication/packages/shared",
"shouldPublish": false
},
{
"packageName": "@hcengineering/communication-yaml",
"projectFolder": "communication/packages/yaml",
"shouldPublish": false
},
{
"packageName": "@hcengineering/communication-rest-client",
"projectFolder": "communication/packages/rest-client",
+2
View File
@@ -279,6 +279,7 @@ export function start (
mailUrl?: string
billingUrl?: string
pulseUrl?: string
hulylakeUrl?: string
},
port: number,
extraConfig?: Record<string, string | undefined>
@@ -357,6 +358,7 @@ export function start (
MAIL_URL: config.mailUrl,
BILLING_URL: config.billingUrl,
PULSE_URL: config.pulseUrl,
HULYLAKE_URL: config.hulylakeUrl,
...(extraConfig ?? {})
}
res.status(200)
+4 -1
View File
@@ -130,6 +130,8 @@ export function startFront (ctx: MeasureContext, extraConfig?: Record<string, st
const billingUrl = process.env.BILLING_URL
const hulylakeUrl = process.env.HULYLAKE_URL
setMetadata(serverToken.metadata.Secret, serverSecret)
setMetadata(serverToken.metadata.Service, 'front')
@@ -161,7 +163,8 @@ export function startFront (ctx: MeasureContext, extraConfig?: Record<string, st
streamUrl,
mailUrl,
billingUrl,
pulseUrl
pulseUrl,
hulylakeUrl
}
console.log('Starting Front service with', config)
const shutdown = start(ctx, config, SERVER_PORT, extraConfig)
+1 -1
View File
@@ -53,6 +53,6 @@
"@hcengineering/communication-sdk-types": "^0.1.0",
"@hcengineering/communication-shared": "^0.1.0",
"@hcengineering/communication-types": "^0.1.0",
"@hcengineering/communication-yaml": "^0.1.0"
"@hcengineering/hulylake-client": "^0.6.0"
}
}
+75 -81
View File
@@ -38,7 +38,6 @@ import core, {
type ModelDb,
platformNow,
type Ref,
SortingOrder,
type Space,
systemAccount,
toIdMap,
@@ -88,9 +87,14 @@ import {
type Message,
type MessageID
} from '@hcengineering/communication-types'
import { parseYaml } from '@hcengineering/communication-yaml'
import { applyPatches, isBlobAttachment, isLinkPreviewAttachment } from '@hcengineering/communication-shared'
import {
isBlobAttachment,
isLinkPreviewAttachment,
loadMessages,
loadMessagesGroups
} from '@hcengineering/communication-shared'
import { markdownToMarkup } from '@hcengineering/text-markdown'
import { type HulylakeClient } from '@hcengineering/hulylake-client'
export * from './types'
export * from './utils'
@@ -99,9 +103,6 @@ const printThresholdMs = 2500
const textLimit = 500 * 1024
const messageGroupsLimit = 100
const messagesLimit = 1000
// Inner presentation in message queue differs from sdk-types,
// also date is always filled at the output queue
export type QueueSourced<T extends Event> = Omit<T, 'date'> & { date: string }
@@ -237,6 +238,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
readonly storageAdapter: StorageAdapter,
readonly contentAdapter: ContentTextAdapter,
readonly broadcastUpdate: (ctx: MeasureContext, classes: Ref<Class<Doc>>[]) => void,
readonly hulylake: HulylakeClient,
readonly communicationApi?: CommunicationApi,
readonly listener?: FulltextListener
) {
@@ -586,10 +588,11 @@ export class FullTextIndexPipeline implements FullTextPipeline {
const rateLimit = new RateLimiter(10)
let lastPrint = platformNow()
await ctx.with('process-message-groups', {}, async (ctx) => {
let groups = await communicationApi.findMessagesGroups(this.communicationSession, {
limit: messageGroupsLimit,
order: SortingOrder.Ascending
})
// let groups = await communicationApi.findMessagesGroups(this.communicationSession, {
// limit: messageGroupsLimit,
// order: SortingOrder.Ascending
// })
let groups = [] as any[]
while (groups.length > 0) {
if (this.cancelling) {
return processed
@@ -608,27 +611,13 @@ export class FullTextIndexPipeline implements FullTextPipeline {
cardInfo = { space: cardDoc[0].space, _class: cardDoc[0]._class }
cardsInfo.set(group.cardId, cardInfo)
}
const blob = await this.storageAdapter.read(ctx, this.workspace, group.blobId)
const messagesFile = Buffer.concat(blob as any).toString()
const messagesParsedFile = parseYaml(messagesFile)
let patchedMessages
if (group.patches !== undefined && group.patches.length > 0) {
const patchesByMessage = groupByArray(group.patches, (it) => it.messageId)
patchedMessages = messagesParsedFile.messages.map((message) => {
const patches = patchesByMessage.get(message.id) ?? []
if (patches.length === 0) {
return message
} else {
return applyPatches(message, patches)
}
})
} else {
patchedMessages = messagesParsedFile.messages
}
for (const message of patchedMessages) {
if (message.removed) {
continue
}
// const blob = await this.storageAdapter.read(ctx, this.workspace, group.blobId)
// const messagesFile = Buffer.concat(blob as any).toString()
// const messagesParsedFile = parseYaml(messagesFile)
// const messages = messagesParsedFile.messages
const messages = [] as Message[]
for (const message of messages) {
await rateLimit.add(async () => {
await this.processCommunicationMessage(
ctx,
@@ -662,21 +651,22 @@ export class FullTextIndexPipeline implements FullTextPipeline {
if (this.cancelling) {
return processed
}
groups = await communicationApi.findMessagesGroups(this.communicationSession, {
limit: messageGroupsLimit,
order: SortingOrder.Ascending,
fromDate: {
greater: groups[groups.length - 1].toDate
}
})
// groups = await communicationApi.findMessagesGroups(this.communicationSession, {
// limit: messageGroupsLimit,
// order: SortingOrder.Ascending,
// fromDate: {
// greater: groups[groups.length - 1].toDate
// }
// })
groups = []
}
})
await ctx.with('process-messages', {}, async (ctx) => {
let messages = await communicationApi.findMessages(this.communicationSession, {
attachments: true,
limit: messagesLimit,
order: SortingOrder.Ascending
})
// let messages = await communicationApi.findMessages(this.communicationSession, {
// limit: messagesLimit,
// order: SortingOrder.Ascending
// })
let messages = [] as any[]
while (messages.length > 0) {
for (const message of messages) {
if (control !== undefined) {
@@ -723,14 +713,14 @@ export class FullTextIndexPipeline implements FullTextPipeline {
lastPrint = now
}
}
messages = await communicationApi.findMessages(this.communicationSession, {
attachments: true,
limit: messagesLimit,
order: SortingOrder.Ascending,
created: {
greater: messages[messages.length - 1].created
}
})
// messages = await communicationApi.findMessages(this.communicationSession, {
// limit: messagesLimit,
// order: SortingOrder.Ascending,
// created: {
// greater: messages[messages.length - 1].created
// }
// })
messages = []
}
})
await rateLimit.waitProcessing()
@@ -842,35 +832,38 @@ export class FullTextIndexPipeline implements FullTextPipeline {
return
}
const getMessage = async (cardId: CardID, msgId: MessageID): Promise<Message | undefined> => {
const messages = await communicationApi.findMessages(this.communicationSession, {
card: cardId,
id: msgId,
attachments: true
})
if (messages.length === 1) {
return messages[0]
}
const messagesGroups = await communicationApi.findMessagesGroups(this.communicationSession, {
card: cardId,
messageId: msgId
})
if (messagesGroups.length !== 1) {
const meta = (
await communicationApi.findMessagesMeta(this.communicationSession, {
cardId,
id: msgId,
limit: 1
})
)[0]
if (meta === undefined) {
return undefined
}
const group = messagesGroups[0]
const blob = await this.storageAdapter.read(ctx, this.workspace, group.blobId)
const messagesFile = Buffer.concat(blob as any).toString()
const messagesParsedFile = parseYaml(messagesFile)
const message = messagesParsedFile.messages.find((m) => m.id === msgId)
if (group.patches === undefined || message === undefined) {
return message
}
const relevantPatches = group.patches.filter((p) => p.messageId === msgId)
if (relevantPatches.length === 0) {
return message
} else {
return applyPatches(message, relevantPatches)
const messagesGroups = await loadMessagesGroups(this.hulylake, cardId)
const group = messagesGroups.find((it) => it.blobId === meta.blobId)
if (group === undefined) {
return undefined
}
return (
await loadMessages(
this.hulylake,
group.blobId,
{
cardId,
id: msgId
},
{
attachments: true,
reactions: true,
threads: true
}
)
)[0]
}
const cardDoc = (await this.storage.findAll(ctx, card.class.Card, { _id: cardId }))[0]
// If message was already fully replaced, other transactions can skip the message
@@ -898,8 +891,9 @@ export class FullTextIndexPipeline implements FullTextPipeline {
for (const operation of event.operations) {
if (operation.opcode === 'attach' || operation.opcode === 'set' || operation.opcode === 'update') {
for (const blobData of operation.blobs) {
const blobAttachment: Omit<BlobAttachment, 'type'> = {
const blobAttachment: BlobAttachment = {
id: blobData.blobId as any as AttachmentID,
mimeType: blobData.mimeType ?? '',
params: blobData as BlobParams,
creator: event.socialId,
created: new Date(Date.parse(event.date))
@@ -1090,7 +1084,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
cardId: CardID,
cardSpace: Ref<Space>,
cardClass: Ref<Class<Card>>,
message: Pick<Message, 'id' | 'edited' | 'created' | 'creator' | 'content' | 'extra' | 'thread' | 'attachments'>
message: Pick<Message, 'id' | 'modified' | 'created' | 'creator' | 'content' | 'extra' | 'threads' | 'attachments'>
): Promise<void> {
const indexedDoc = createIndexedDocFromMessage(cardId, cardSpace, cardClass, message)
const markup = markdownToMarkup(message.content)
@@ -1130,7 +1124,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
ctx: MeasureContext<any>,
pushQueue: ElasticPushQueue,
parentDoc: { id: Ref<Doc>, _class: Ref<Class<Doc>>[], space: Ref<Space>, attachedTo?: Ref<Doc> },
blobAttachment: Omit<BlobAttachment, 'type'>
blobAttachment: BlobAttachment
): Promise<void> {
try {
const indexedDoc: IndexedDoc = {
@@ -1148,7 +1142,7 @@ export class FullTextIndexPipeline implements FullTextPipeline {
attachedToCard: parentDoc.attachedTo
}
indexedDoc.fulltextSummary = ''
await this.handleBlobRef(ctx, blobAttachment.params.blobId, indexedDoc, blobAttachment.params.mimeType)
await this.handleBlobRef(ctx, blobAttachment.params.blobId, indexedDoc, blobAttachment.mimeType)
if (this.listener?.onIndexing !== undefined) {
await this.listener.onIndexing(indexedDoc)
}
+2 -2
View File
@@ -117,9 +117,9 @@ export function createIndexedDocFromMessage (
cardId: Ref<Card>,
cardSpace: Ref<Space>,
cardClass: Ref<Class<Card>>,
message: Pick<Message, 'id' | 'edited' | 'created' | 'creator'>
message: Pick<Message, 'id' | 'modified' | 'created' | 'creator'>
): IndexedDoc {
const modifiedDate = message.edited ?? message.created
const modifiedDate = message.modified ?? message.created
const modifiedOn = modifiedDate.getTime()
const indexedDoc = {
id: `${message.id}@${cardId}` as any,
+16 -18
View File
@@ -109,30 +109,23 @@ export class CommunicationMiddleware extends BaseMiddleware implements Middlewar
async handleCommand (_ctx: MeasureContext<SessionData>, args: DomainParams): Promise<any> {
const ctx = this.getCommunicationCtx(_ctx)
if (args.findMessages !== undefined) {
const { params, queryId } = args.findMessages
return await this.communicationApi.findMessages(ctx, params, queryId)
}
if (args.findMessagesGroups !== undefined) {
const { params } = args.findMessagesGroups
return await this.communicationApi.findMessagesGroups(ctx, params)
if (args.findMessagesMeta !== undefined) {
const { params } = args.findMessagesMeta
return await this.communicationApi.findMessagesMeta(ctx, params)
}
if (args.findNotificationContexts !== undefined) {
const { params, queryId } = args.findNotificationContexts
return await this.communicationApi.findNotificationContexts(ctx, params, queryId)
const { params, subscription } = args.findNotificationContexts
return await this.communicationApi.findNotificationContexts(ctx, params, subscription)
}
if (args.findNotifications !== undefined) {
const { params, queryId } = args.findNotifications
return await this.communicationApi.findNotifications(ctx, params, queryId)
const { params, subscription } = args.findNotifications
return await this.communicationApi.findNotifications(ctx, params, subscription)
}
if (args.findLabels !== undefined) {
const { params } = args.findLabels
return await this.communicationApi.findLabels(ctx, params)
}
if (args.findThreads !== undefined) {
const { params } = args.findThreads
return await this.communicationApi.findThreads(ctx, params)
}
if (args.findCollaborators !== undefined) {
const { params } = args.findCollaborators
return await this.communicationApi.findCollaborators(ctx, params)
@@ -141,9 +134,14 @@ export class CommunicationMiddleware extends BaseMiddleware implements Middlewar
const { params } = args.findPeers
return await this.communicationApi.findPeers(ctx, params)
}
if (args.unsubscribeQuery !== undefined) {
const { id } = args.unsubscribeQuery
await this.communicationApi.unsubscribeQuery(ctx, id)
if (args.subscribeCard !== undefined) {
const { cardId, subscription } = args.subscribeCard
this.communicationApi.subscribeCard(ctx, cardId, subscription)
return
}
if (args.unsubscribeCard !== undefined) {
const { cardId, subscription } = args.unsubscribeCard
this.communicationApi.unsubscribeCard(ctx, cardId, subscription)
return
}
if (args.event !== undefined) {
+2 -2
View File
@@ -281,7 +281,7 @@ async function createMailThread (
const subjectId = generateMessageId()
const createSubjectEvent: CreateMessageEvent = {
type: MessageEventType.CreateMessage,
messageType: MessageType.Message,
messageType: MessageType.Text,
cardId: data.channel,
cardType: chat.masterTag.Thread,
content: data.subject,
@@ -322,7 +322,7 @@ async function createMailMessage (
const messageId = getHulyIdFromEmailMessageId(data.mailId, data.from) ?? generateMessageId()
const createMessageEvent: CreateMessageEvent = {
type: MessageEventType.CreateMessage,
messageType: MessageType.Message,
messageType: MessageType.Text,
cardId: threadId,
cardType: chat.masterTag.Thread,
content: data.content,
+1 -1
View File
@@ -44,7 +44,7 @@ export function toMessageEvent (tx: Tx): CreateMessageEvent | undefined {
return undefined
}
const event: CreateMessageEvent = domainTx.event
const isMessage = event.cardType === chat.masterTag.Thread && event.messageType === MessageType.Message
const isMessage = event.cardType === chat.masterTag.Thread && event.messageType === MessageType.Text
if (!isMessage) {
return undefined
}
-7
View File
@@ -1,7 +0,0 @@
module.exports = {
extends: ['./node_modules/@hcengineering/platform-rig/profiles/default/eslint.config.json'],
parserOptions: {
tsconfigRootDir: __dirname,
project: './tsconfig.json'
}
}
-4
View File
@@ -1,4 +0,0 @@
*
!/lib/**
!CHANGELOG.md
/lib/**/__tests__/
-7
View File
@@ -1,7 +0,0 @@
FROM hardcoreeng/base:v20250113a
WORKDIR /usr/src/app
COPY bundle/bundle.js ./
EXPOSE 9001
CMD [ "node", "bundle.js" ]
-4
View File
@@ -1,4 +0,0 @@
{
"$schema": "https://developer.microsoft.com/json-schemas/rig-package/rig.schema.json",
"rigPackageName": "@hcengineering/platform-rig"
}
-7
View File
@@ -1,7 +0,0 @@
module.exports = {
preset: 'ts-jest',
testEnvironment: 'node',
testMatch: ['**/?(*.)+(spec|test).[jt]s?(x)'],
roots: ["./src"],
coverageReporters: ["text-summary", "html"]
}
-79
View File
@@ -1,79 +0,0 @@
{
"name": "@hcengineering/pod-msg2file",
"version": "0.6.0",
"main": "lib/index.js",
"svelte": "src/index.ts",
"types": "types/index.d.ts",
"files": [
"lib/**/*",
"types/**/*",
"tsconfig.json"
],
"author": "Hardcore Engineering Inc.",
"scripts": {
"build": "compile",
"build:watch": "compile",
"test": "jest --passWithNoTests --silent",
"_phase:bundle": "rushx bundle",
"_phase:docker-build": "rushx docker:build",
"_phase:docker-staging": "rushx docker:staging",
"bundle": "node ../../common/scripts/esbuild.js --external=ws",
"docker:build": "../../common/scripts/docker_build.sh hardcoreeng/msg2file .",
"docker:tbuild": "docker build -t hardcoreeng/msg2file . --platform=linux/amd64 && ../../common/scripts/docker_tag_push.sh hardcoreeng/msg2file",
"docker:staging": "../../common/scripts/docker_tag.sh hardcoreeng/msg2file staging",
"docker:push": "../../common/scripts/docker_tag.sh hardcoreeng/msg2file",
"run-local": "cross-env ts-node src/index.ts",
"format": "format src",
"_phase:build": "compile transpile src",
"_phase:test": "jest --passWithNoTests --silent",
"_phase:format": "format src",
"_phase:validate": "compile validate"
},
"devDependencies": {
"@hcengineering/platform-rig": "^0.6.0",
"@tsconfig/node16": "^1.0.4",
"@types/cors": "^2.8.12",
"@types/express": "^4.17.13",
"@types/jest": "^29.5.5",
"@types/js-yaml": "^4.0.9",
"@types/node": "^22.15.29",
"@types/node-cron": "^3.0.11",
"@types/uuid": "^8.3.1",
"@typescript-eslint/eslint-plugin": "^6.11.0",
"@typescript-eslint/parser": "^6.11.0",
"esbuild": "^0.25.9",
"eslint": "^8.54.0",
"eslint-config-standard-with-typescript": "^40.0.0",
"eslint-plugin-import": "^2.26.0",
"eslint-plugin-n": "^15.4.0",
"eslint-plugin-node": "^11.1.0",
"eslint-plugin-promise": "^6.1.1",
"jest": "^29.7.0",
"prettier": "^3.1.0",
"ts-jest": "^29.1.1",
"ts-node": "^10.8.0",
"typescript": "^5.8.3"
},
"dependencies": {
"@hcengineering/api-client": "^0.6.0",
"@hcengineering/card": "^0.6.0",
"@hcengineering/communication-types": "^0.1.0",
"@hcengineering/communication-shared": "^0.1.0",
"@hcengineering/communication-yaml": "^0.1.0",
"@hcengineering/communication-sdk-types": "^0.1.0",
"@hcengineering/communication-rest-client": "^0.1.0",
"@hcengineering/core": "^0.6.32",
"@hcengineering/platform": "^0.6.11",
"@hcengineering/server-client": "^0.6.0",
"@hcengineering/server-core": "^0.6.1",
"@hcengineering/server-storage": "^0.6.0",
"@hcengineering/server-token": "^0.6.11",
"cors": "^2.8.5",
"dotenv": "~16.0.0",
"express": "^4.21.2",
"js-yaml": "^4.1.0",
"node-cron": "^3.0.3",
"postgres": "^3.4.7",
"uuid": "^8.3.2"
}
}
-50
View File
@@ -1,50 +0,0 @@
//
// Copyright © 2025 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.
//
interface Config {
AccountsURL: string
DbUrl: string
MaxSyncAttempts: number
MinSyncMessagesCount: number
MessagesPerFile: number
Port: number
Secret: string
ServiceID: string
}
const parseNumber = (str: string | undefined): number | undefined => (str !== undefined ? Number(str) : undefined)
const config: Config = (() => {
const params: Partial<Config> = {
AccountsURL: process.env.ACCOUNTS_URL,
DbUrl: process.env.DB_URL,
MaxSyncAttempts: parseNumber(process.env.MAX_SYNC_ATTEMPTS) ?? 3,
MessagesPerFile: parseNumber(process.env.MESSAGES_PER_FILE) ?? 500,
MinSyncMessagesCount: parseNumber(process.env.MIN_SYNC_MESSAGES_COUNT) ?? 60,
Port: parseNumber(process.env.PORT),
Secret: process.env.SECRET ?? 'secret',
ServiceID: process.env.SERVICE_ID ?? 'msg2file-service'
}
const missingEnv = (Object.keys(params) as Array<keyof Config>).filter((key) => params[key] === undefined)
if (missingEnv.length > 0) {
throw Error(`Missing env variables: ${missingEnv.join(', ')}`)
}
return params as Config
})()
export default config
-139
View File
@@ -1,139 +0,0 @@
//
// Copyright © 2025 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 postgres from 'postgres'
import { MessageID, type CardID, type WorkspaceID } from '@hcengineering/communication-types'
import { Domain } from '@hcengineering/communication-sdk-types'
import config from './config'
export interface SyncRecord {
workspace: WorkspaceID
card: CardID
attempt: number
created: Date
}
export async function getDb (): Promise<PostgresDB> {
const sql = postgres(config.DbUrl, {
connection: {
application_name: config.ServiceID
},
fetch_types: true,
prepare: true
})
return await PostgresDB.create(sql)
}
export class PostgresDB {
private readonly syncTable = 'msg2file.sync_record'
constructor (private readonly client: postgres.Sql) {}
static async create (client: postgres.Sql): Promise<PostgresDB> {
await this.init(client)
return new PostgresDB(client)
}
static async init (client: postgres.Sql): Promise<void> {
const sql = `
CREATE SCHEMA IF NOT EXISTS msg2file;
CREATE TABLE IF NOT EXISTS msg2file.sync_record
(
workspace UUID NOT NULL,
card VARCHAR(255) NOT NULL,
attempt INT NOT NULL,
created TIMESTAMPTZ NOT NULL DEFAULT NOW(),
PRIMARY KEY (workspace, card)
);
`
await client.unsafe(sql)
}
async getRecords (limit: number, date: Date): Promise<SyncRecord[]> {
const sql = `
SELECT workspace, card
FROM ${this.syncTable}
WHERE created < $1::timestamptz
ORDER BY created
LIMIT ${limit};`
const result = await this.client.unsafe(sql, [date])
return result.map((raw) => ({
workspace: raw.workspace,
card: raw.card
})) as SyncRecord[]
}
async createRecord (workspace: WorkspaceID, card: CardID): Promise<void> {
const sql = `
INSERT INTO ${this.syncTable} (workspace, card, attempt)
VALUES ($1::uuid, $2::varchar, 0)
ON CONFLICT DO NOTHING;`
await this.client.unsafe(sql, [workspace, card])
}
async removeRecord (workspace: WorkspaceID, card: CardID): Promise<void> {
const sql = `
DELETE
FROM ${this.syncTable}
WHERE workspace = $1::uuid
AND card = $2::varchar;`
await this.client.unsafe(sql, [workspace, card])
}
async increaseAttempt (workspace: WorkspaceID, card: CardID): Promise<void> {
const sql = `
UPDATE ${this.syncTable}
SET attempt = attempt + 1
WHERE workspace = $1::uuid
AND card = $2::varchar;`
await this.client.unsafe(sql, [workspace, card])
}
async removeMessages (workspace: WorkspaceID, card: CardID, ids: MessageID[]): Promise<void> {
if (ids.length === 0) return
const sql = `
DELETE
FROM ${Domain.Message}
WHERE workspace_id = $1::uuid
AND card_id = $2::varchar
AND id = ANY ($3::varchar[]);`
await this.client.unsafe(sql, [workspace, card, ids])
}
async removePatches (workspace: WorkspaceID, card: CardID, ids: MessageID[]): Promise<void> {
if (ids.length === 0) return
const sql = `
DELETE
FROM ${Domain.Patch}
WHERE workspace_id = $1::uuid
AND card_id = $2::varchar
AND message_id = ANY ($3::varchar[]);`
await this.client.unsafe(sql, [workspace, card, ids])
}
async close (): Promise<void> {
await this.client.end({ timeout: 0 })
}
}
-21
View File
@@ -1,21 +0,0 @@
//
// Copyright © 2025 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 { config } from 'dotenv'
import { main } from './main'
config()
void main()
-56
View File
@@ -1,56 +0,0 @@
//
// Copyright © 2025 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 { setMetadata } from '@hcengineering/platform'
import serverClient from '@hcengineering/server-client'
import { initStatisticsContext } from '@hcengineering/server-core'
import serverToken from '@hcengineering/server-token'
import cron from 'node-cron'
import { buildStorageFromConfig, storageConfigFromEnv } from '@hcengineering/server-storage'
import config from './config'
import { job } from './worker'
import { startServer } from './server'
import { getDb } from './db'
export const main = async (): Promise<void> => {
setMetadata(serverClient.metadata.Endpoint, config.AccountsURL)
setMetadata(serverClient.metadata.UserAgent, config.ServiceID)
setMetadata(serverToken.metadata.Secret, config.Secret)
setMetadata(serverToken.metadata.Service, 'msg2file')
const ctx = initStatisticsContext(config.ServiceID, {})
const storage = buildStorageFromConfig(storageConfigFromEnv())
const db = await getDb()
const server = startServer(ctx, db)
const task = cron.schedule('0 0 * * *', () => {
void job(ctx, storage, db)
})
const shutdown = (): void => {
task.stop()
server.close(() => process.exit())
}
process.on('SIGINT', shutdown)
process.on('SIGTERM', shutdown)
process.on('uncaughtException', (e) => {
console.error(e)
})
process.on('unhandledRejection', (e) => {
console.error(e)
})
}
-40
View File
@@ -1,40 +0,0 @@
//
// Copyright © 2025 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 { Readable } from 'stream'
import { ParsedFile } from '@hcengineering/communication-types'
import { parseYaml } from '@hcengineering/communication-yaml'
export async function parseFileStream (stream: Readable): Promise<ParsedFile> {
return await new Promise((resolve, reject) => {
let yamlData = ''
stream.on('data', (chunk) => {
yamlData += chunk.toString()
})
stream.on('end', () => {
try {
resolve(parseYaml(yamlData))
} catch (error) {
reject(error)
}
})
stream.on('error', (error) => {
reject(error)
})
})
}
-48
View File
@@ -1,48 +0,0 @@
//
// Copyright © 2025 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 { generateToken } from '@hcengineering/server-token'
import { systemAccountUuid } from '@hcengineering/core'
import { getTransactorEndpoint } from '@hcengineering/server-client'
import { createRestClient, RestClient } from '@hcengineering/api-client'
import { WorkspaceID } from '@hcengineering/communication-types'
import {
createRestClient as createCommunicationRestClient,
RestClient as CommunicationRestClient
} from '@hcengineering/communication-rest-client'
import config from './config'
export async function connectPlatform (workspace: WorkspaceID): Promise<RestClient> {
const token = generateToken(systemAccountUuid, workspace, { service: config.ServiceID })
const endpoint = toHttpUrl(await getTransactorEndpoint(token))
return createRestClient(endpoint, workspace, token)
}
export async function connectCommunication (workspace: WorkspaceID): Promise<CommunicationRestClient> {
const token = generateToken(systemAccountUuid, workspace, { service: config.ServiceID })
const endpoint = toHttpUrl(await getTransactorEndpoint(token))
return createCommunicationRestClient(endpoint, workspace, token)
}
function toHttpUrl (url: string): string {
if (url.startsWith('ws://')) {
return url.replace('ws://', 'http://')
}
if (url.startsWith('wss://')) {
return url.replace('wss://', 'https://')
}
return url
}
-118
View File
@@ -1,118 +0,0 @@
//
// Copyright © 2025 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 { Token } from '@hcengineering/server-token'
import cors from 'cors'
import express, { type Express, type NextFunction, type Request, type Response } from 'express'
import { type Server } from 'http'
import { extractToken } from '@hcengineering/server-client'
import { MeasureContext } from '@hcengineering/core'
import { CardID } from '@hcengineering/communication-types'
import config from './config'
import { register } from './worker'
import { PostgresDB } from './db'
export function startServer (ctx: MeasureContext, db: PostgresDB): Server {
const app = createServer(ctx, db)
const port = config.Port
return app.listen(port, (): void => {
ctx.info(`Msg2file service has been started at :${port}`)
})
}
type AsyncRequestHandler = (req: Request, res: Response, token: Token, next: NextFunction) => Promise<void>
const handleRequest = async (
fn: AsyncRequestHandler,
req: Request,
res: Response,
next: NextFunction
): Promise<void> => {
const token = extractToken(req.headers)
if (token === undefined) {
throw new ApiError(401)
}
try {
await fn(req, res, token, next)
} catch (err: unknown) {
next(err)
}
}
const wrapRequest = (fn: AsyncRequestHandler) => (req: Request, res: Response, next: NextFunction) => {
void handleRequest(fn, req, res, next)
}
function createServer (ctx: MeasureContext, db: PostgresDB): Express {
const app = express()
app.use(cors())
app.use(express.json())
app.post(
'/register/:card',
wrapRequest(async (req, res, token) => {
const { card } = req.params
const { workspace } = token
if (card == null || card === '') {
throw new ApiError(400)
}
ctx.info('Register card', { workspace, card })
await register(workspace, card as CardID, db)
res.status(200)
res.json({})
})
)
app.use((err: any, _req: any, res: any, _next: any) => {
console.error(err)
if (err instanceof ApiError) {
res.status(err.code).send({ code: err.code, message: err.message })
return
}
res.status(500).send(err.message?.length > 0 ? { message: err.message } : err)
})
return app
}
class ApiError extends Error {
readonly code: number
constructor (code: number, message?: string) {
super(message ?? generateErrorMessage(code))
this.code = code
}
}
const generateErrorMessage = (code: number): string => {
if (code === 401) {
return 'Unauthorized'
}
if (code === 404) {
return 'Not Found'
}
if (code === 400) {
return 'Bad Request'
}
return 'Error'
}
-50
View File
@@ -1,50 +0,0 @@
//
// Copyright © 2025 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 { retry } from '@hcengineering/communication-shared'
import { StorageAdapter, UploadedObjectInfo } from '@hcengineering/server-core'
import { MeasureContext } from '@hcengineering/core'
import { BlobID, WorkspaceID } from '@hcengineering/communication-types'
import { type Readable } from 'stream'
export async function getFile (
storage: StorageAdapter,
ctx: MeasureContext,
workspace: WorkspaceID,
blob: BlobID
): Promise<Readable> {
return await retry(() => storage.get(ctx, { uuid: workspace } as any, blob), { retries: 3 })
}
export async function removeFile (
storage: StorageAdapter,
ctx: MeasureContext,
workspace: WorkspaceID,
blob: BlobID
): Promise<void> {
await retry(() => storage.remove(ctx, { uuid: workspace } as any, [blob]), { retries: 3 })
}
export async function uploadFile (
storage: StorageAdapter,
ctx: MeasureContext,
workspace: WorkspaceID,
blob: BlobID,
content: Readable | string
): Promise<UploadedObjectInfo> {
return await retry(async () => await storage.put(ctx, { uuid: workspace } as any, blob, content, 'text/yaml'), {
retries: 3
})
}
-492
View File
@@ -1,492 +0,0 @@
//
// Copyright © 2025 Hardcore Engineering Inc.
//
// Licensed under the Eclipse Public License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License. You may
// obtain a copy of the License at https://www.eclipse.org/legal/epl-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
//
// See the License for the specific language governing permissions and
// limitations under the License.
//
import core, { MeasureContext, RateLimiter, Ref, Blob, groupByArray } from '@hcengineering/core'
import {
type CardID,
type FileMessage,
type FileMetadata,
type MessageID,
type WorkspaceID,
SortingOrder,
MessagesGroup,
BlobID
} from '@hcengineering/communication-types'
import cardPlugin, { type Card } from '@hcengineering/card'
import yaml from 'js-yaml'
import { v4 as uuid } from 'uuid'
import { MessageEventType } from '@hcengineering/communication-sdk-types'
import { applyPatches, retry } from '@hcengineering/communication-shared'
import { StorageAdapter } from '@hcengineering/server-core'
import { deserializeMessage } from '@hcengineering/communication-yaml'
import { RestClient as CommunicationRestClient } from '@hcengineering/communication-rest-client'
import { PostgresDB, SyncRecord } from './db'
import config from './config'
import { parseFileStream } from './parser'
import { connectCommunication, connectPlatform } from './platform'
import { getFile, removeFile, uploadFile } from './storage'
export async function register (workspace: WorkspaceID, card: CardID, db: PostgresDB): Promise<void> {
await db.createRecord(workspace, card)
}
export async function job (ctx: MeasureContext, storage: StorageAdapter, db: PostgresDB): Promise<void> {
ctx.info('Job started', { date: new Date() })
const start = new Date()
try {
const rateLimiter = new RateLimiter(10)
while (true) {
const records = await db.getRecords(100, start)
if (records.length === 0) {
break
}
for (const record of records) {
await rateLimiter.add(async () => {
await ctx.with(
'process',
{},
async () => {
await processRecord(ctx, record, db, storage)
},
{
workspace: record.workspace,
card: record.card,
attempt: record.attempt
}
)
})
}
await rateLimiter.waitProcessing()
}
} finally {
ctx.info('Job finished', { date: new Date() })
}
}
async function processRecord (
ctx: MeasureContext,
record: SyncRecord,
db: PostgresDB,
storage: StorageAdapter
): Promise<void> {
try {
ctx.info('Start processing record', { workspace: record.workspace, card: record.card, attempt: record.attempt })
await msg2file(ctx, record.workspace, record.card, storage, db)
await db.removeRecord(record.workspace, record.card)
} catch (e) {
ctx.error('Failed to process record', { workspace: record.workspace, card: record.card, err: e })
if (record.attempt < config.MaxSyncAttempts) {
await db.increaseAttempt(record.workspace, record.card)
} else {
await db.removeRecord(record.workspace, record.card)
}
}
}
async function msg2file (
ctx: MeasureContext,
workspace: WorkspaceID,
cardId: CardID,
storage: StorageAdapter,
db: PostgresDB
): Promise<void> {
ctx.info('Processing card', { card: cardId })
const platformClient = await connectPlatform(workspace)
const card = await platformClient.findOne<Card>(cardPlugin.class.Card, { _id: cardId as Ref<Card> })
if (card === undefined) {
ctx.error('Card not found, skip processing', { workspace, card: cardId })
return
}
const client = await connectCommunication(workspace)
await applyPatchesToGroups(ctx, client, workspace, card, storage, db)
await newMessages2file(ctx, client, workspace, card, storage, db)
}
async function applyPatchesToGroups (
ctx: MeasureContext,
client: CommunicationRestClient,
workspace: WorkspaceID,
card: Card,
storage: StorageAdapter,
db: PostgresDB
): Promise<void> {
let groups = await client.findMessagesGroups({
card: card._id,
patches: true,
limit: 100,
order: SortingOrder.Ascending
})
let firstBlobId: BlobID | undefined
while (groups.length > 0) {
if (firstBlobId === groups[0].blobId) {
ctx.error('Group is repeated', { group: groups[0].blobId })
throw new Error(`Group is repeated ${groups[0].blobId}`)
}
firstBlobId = groups[0].blobId
for (const group of groups) {
await applyPatchesToGroup(ctx, client, storage, workspace, group, db, card._id)
}
groups = await client.findMessagesGroups({
card: card._id,
patches: true,
limit: 100,
order: SortingOrder.Ascending,
fromDate: {
greater: groups[groups.length - 1].toDate
}
})
}
}
async function applyPatchesToGroup (
ctx: MeasureContext,
client: CommunicationRestClient,
storage: StorageAdapter,
workspace: WorkspaceID,
group: MessagesGroup,
db: PostgresDB,
card: CardID
): Promise<void> {
if (group.patches == null || group.patches.length === 0) {
return
}
ctx.info('Start apply patches to group', { group: group.blobId, patches: group.patches.length })
try {
const file = await getFile(storage, ctx, workspace, group.blobId)
const parsedFile = await parseFileStream(file)
const patchesByMessage = groupByArray(group.patches, (it) => it.messageId)
const updatedMessages = parsedFile.messages.map((message) => {
const patches = patchesByMessage.get(message.id) ?? []
if (patches.length === 0) {
return message
} else {
return applyPatches(message, patches)
}
})
const blob = await uploadGroupFile(
ctx,
storage,
workspace,
parsedFile.cardId,
parsedFile.title,
parsedFile.fromDate,
parsedFile.toDate,
updatedMessages.map(deserializeMessage)
)
await createGroup(client, group.cardId, blob, parsedFile.fromDate, parsedFile.toDate, updatedMessages.length)
await removeGroup(client, group.cardId, group.blobId)
await removePatches(db, workspace, card, Array.from(patchesByMessage.keys()))
await removeFile(storage, ctx, workspace, group.blobId)
} catch (error) {
ctx.error('Failed to apply patches', { group, error })
}
}
async function newMessages2file (
ctx: MeasureContext,
client: CommunicationRestClient,
workspace: WorkspaceID,
card: Card,
storage: StorageAdapter,
db: PostgresDB
): Promise<void> {
const lastGroup = (
await client.findMessagesGroups({
card: card._id,
orderBy: 'toDate',
order: SortingOrder.Descending,
limit: 1
})
)[0]
let firstMessageId: MessageID | undefined
while (true) {
const messages = (
await client.findMessages({
card: card._id,
order: SortingOrder.Ascending,
limit: config.MessagesPerFile,
reactions: true,
replies: true,
attachments: true
})
).map(deserializeMessage)
if (messages.length === 0) {
break
}
if (firstMessageId === messages[0].id) {
ctx.error('Message repeated', { card: card._id, message: firstMessageId })
throw new Error(`Message repeated ${firstMessageId}`)
}
firstMessageId = messages[0].id
const firstMessage = messages[0]
const fromDate = firstMessage.created
const messagesFromExistingGroup =
lastGroup == null || fromDate > lastGroup.toDate
? []
: messages.filter(({ created }) => created <= lastGroup.toDate)
const newMessages =
messagesFromExistingGroup.length > 0 ? messages.slice(messagesFromExistingGroup.length) : messages
await pushMessagesToExistingGroup(client, card, messagesFromExistingGroup, ctx, storage, workspace)
if (messages.length < config.MinSyncMessagesCount) {
const ids = [...messagesFromExistingGroup].map((it) => it.id)
await removeMessages(db, workspace, card._id, [...ids])
await removePatches(db, workspace, card._id, [...ids])
break
}
const savedMessages = await createNewGroup(client, card, newMessages, ctx, storage, workspace)
const ids = [...messagesFromExistingGroup, ...savedMessages].map((it) => it.id)
await removeMessages(db, workspace, card._id, [...ids])
await removePatches(db, workspace, card._id, [...ids])
}
}
async function createNewGroup (
client: CommunicationRestClient,
card: Card,
messages: FileMessage[],
ctx: MeasureContext,
storage: StorageAdapter,
workspace: WorkspaceID
): Promise<FileMessage[]> {
if (messages.length === 0) return []
const lastMessage = messages[messages.length - 1]
const lastCreated = lastMessage.created
const messagesFromRange = (
await client.findMessages({
card: card._id,
order: SortingOrder.Ascending,
created: lastCreated,
reactions: true,
replies: true,
attachments: true
})
)
.filter((it) => it.id !== lastMessage.id)
.map(deserializeMessage)
const allMessages = messages.concat(messagesFromRange)
const fromDate = allMessages[0].created
const toDate = allMessages[allMessages.length - 1].created
const blob = await uploadGroupFile(ctx, storage, workspace, card._id, card.title, fromDate, toDate, allMessages)
await createGroup(client, card._id, blob, fromDate, toDate, allMessages.length)
return allMessages
}
async function pushMessagesToExistingGroup (
client: CommunicationRestClient,
card: Card,
messages: FileMessage[],
ctx: MeasureContext,
storage: StorageAdapter,
workspace: WorkspaceID
): Promise<void> {
if (messages.length === 0) return
ctx.warn('Push messages to existing group', {
firstId: messages[0].id,
lastId: messages[messages.length - 1].id,
count: messages.length
})
while (true) {
const message = messages[0]
if (message == null) break
const group = await findTargetGroup(client, card._id, message.created)
if (group == null) {
ctx.error('Failed to find group for message', { card: card._id, message: message.id })
throw new Error('Failed to find group for message ' + message.id)
}
const messagesFromGroup = messages.filter((it) => {
const date = it.created
return date <= group.toDate && date >= group.fromDate
})
messages.splice(0, messagesFromGroup.length)
await pushMessagesToGroup(client, group, messagesFromGroup, ctx, storage, workspace)
}
}
async function pushMessagesToGroup (
client: CommunicationRestClient,
group: MessagesGroup,
messages: FileMessage[],
ctx: MeasureContext,
storage: StorageAdapter,
workspace: WorkspaceID
): Promise<void> {
try {
const file = await getFile(storage, ctx, workspace, group.blobId)
const parsedFile = await parseFileStream(file)
const newMessages = parsedFile.messages
.map(deserializeMessage)
.concat(messages)
.sort((a, b) => a.created.getTime() - b.created.getTime())
const blob = await uploadGroupFile(
ctx,
storage,
workspace,
parsedFile.cardId,
parsedFile.title,
parsedFile.fromDate,
parsedFile.toDate,
newMessages
)
await removeFile(storage, ctx, workspace, group.blobId)
await removeGroup(client, group.cardId, group.blobId)
await createGroup(client, group.cardId, blob, parsedFile.fromDate, parsedFile.toDate, newMessages.length)
} catch (err: any) {
ctx.error('Failed to push messages to group', { group: group.blobId, error: err })
throw err
}
}
async function createGroup (
client: CommunicationRestClient,
cardId: CardID,
blobId: Ref<Blob>,
fromDate: Date,
toDate: Date,
count: number
): Promise<void> {
await retry(
async () =>
await client.event(
{
type: MessageEventType.CreateMessagesGroup,
group: {
cardId,
blobId,
fromDate,
toDate,
count
},
socialId: core.account.System,
date: new Date()
},
core.account.System
),
{ retries: 3 }
)
}
async function removeGroup (client: CommunicationRestClient, cardId: CardID, blobId: Ref<Blob>): Promise<void> {
await retry(
async () =>
await client.event(
{
type: MessageEventType.RemoveMessagesGroup,
cardId,
blobId,
socialId: core.account.System,
date: new Date()
},
core.account.System
),
{ retries: 3 }
)
}
async function removeMessages (db: PostgresDB, workspace: WorkspaceID, card: CardID, ids: MessageID[]): Promise<void> {
while (ids.length > 0) {
const chunk = ids.splice(0, 100)
await retry(
async () => {
await db.removeMessages(workspace, card, chunk)
},
{ retries: 3 }
)
}
}
async function removePatches (db: PostgresDB, workspace: WorkspaceID, card: CardID, ids: MessageID[]): Promise<void> {
while (ids.length > 0) {
const chunk = ids.splice(0, 100)
await retry(
async () => {
await db.removePatches(workspace, card, chunk)
},
{ retries: 3 }
)
}
}
async function uploadGroupFile (
ctx: MeasureContext,
storage: StorageAdapter,
workspace: WorkspaceID,
cardId: CardID,
title: string,
fromDate: Date,
toDate: Date,
messages: FileMessage[]
): Promise<Ref<Blob>> {
const metadata: FileMetadata = {
cardId,
title,
fromDate,
toDate
}
const yamlMetadata = yaml.dump(metadata, { noRefs: true }).trim()
const yamlMessages = yaml.dump(messages, { noRefs: true, indent: 0 }).trim()
const yamlContent = `---\n${yamlMetadata}\n---\n${yamlMessages}`
const blobId = uuid() as Ref<Blob>
await uploadFile(storage, ctx, workspace, blobId, yamlContent)
return blobId
}
async function findTargetGroup (
client: CommunicationRestClient,
card: CardID,
created: Date
): Promise<MessagesGroup | undefined> {
return (
await client.findMessagesGroups({
card,
fromDate: { lessOrEqual: created },
toDate: { greaterOrEqual: created },
limit: 1,
order: SortingOrder.Ascending,
orderBy: 'fromDate'
})
)[0]
}
-12
View File
@@ -1,12 +0,0 @@
{
"extends": "./node_modules/@hcengineering/platform-rig/profiles/default/tsconfig.json",
"compilerOptions": {
"rootDir": "./src",
"outDir": "./lib",
"declarationDir": "./types",
"tsBuildInfoFile": ".build/build.tsbuildinfo"
},
"include": ["src/**/*"],
"exclude": ["node_modules", "lib", "dist", "types", "bundle"]
}