diff --git a/pods/fulltext/src/manager.ts b/pods/fulltext/src/manager.ts index c0ce6f3d23..6e2d823d1a 100644 --- a/pods/fulltext/src/manager.ts +++ b/pods/fulltext/src/manager.ts @@ -143,7 +143,7 @@ export class WorkspaceManager { private async processDocuments (msg: ConsumerMessage>>[], control: ConsumerControl): Promise { for (const m of msg) { - const ws = m.id as WorkspaceUuid + const ws = m.workspace let token: string try { @@ -164,7 +164,7 @@ export class WorkspaceManager { control: ConsumerControl ): Promise { for (const m of msg) { - const ws = m.id as WorkspaceUuid + const ws = m.workspace for (const mm of m.value) { this.ctx.info('workspace event', { type: mm.type, workspace: ws, restoring: Array.from(this.restoring) }) @@ -218,7 +218,7 @@ export class WorkspaceManager { control: ConsumerControl ): Promise { for (const m of msg) { - const ws = m.id as WorkspaceUuid + const ws = m.workspace for (const mm of m.value) { this.ctx.info('fulltext event', { type: mm.type, workspace: ws, restoring: Array.from(this.restoring) }) diff --git a/pods/media/src/handler.ts b/pods/media/src/handler.ts index 6835f34877..3babbc1abb 100644 --- a/pods/media/src/handler.ts +++ b/pods/media/src/handler.ts @@ -13,7 +13,14 @@ // limitations under the License. // -import core, { type Blob, type Doc, type MeasureContext, type TxCUD, type TxCreateDoc } from '@hcengineering/core' +import core, { + type Blob, + type Doc, + type MeasureContext, + type TxCUD, + type TxCreateDoc, + type WorkspaceUuid +} from '@hcengineering/core' import { PlatformQueueProducer } from '@hcengineering/server-core' import { VideoTranscodeRequest, VideoTranscodeResult } from './types' @@ -28,7 +35,7 @@ function shouldTranscode (contentType: string): boolean { export async function handleTx ( ctx: MeasureContext, - workspaceUuid: string, + workspaceUuid: WorkspaceUuid, tx: TxCUD, producer: PlatformQueueProducer ): Promise { diff --git a/pods/media/src/index.ts b/pods/media/src/index.ts index d116238b00..4683095ce7 100644 --- a/pods/media/src/index.ts +++ b/pods/media/src/index.ts @@ -15,7 +15,7 @@ import { Analytics } from '@hcengineering/analytics' import { configureAnalytics, SplitLogger } from '@hcengineering/analytics-service' -import { Doc, MeasureMetricsContext, TxCUD, WorkspaceUuid, newMetrics } from '@hcengineering/core' +import { Doc, MeasureMetricsContext, TxCUD, newMetrics } from '@hcengineering/core' import { getPlatformQueue } from '@hcengineering/kafka' import { setMetadata } from '@hcengineering/platform' import { initStatisticsContext, QueueTopic } from '@hcengineering/server-core' @@ -70,7 +70,7 @@ async function main (): Promise { queue.createConsumer>(ctx, QueueTopic.Tx, queue.getClientId(), async (msgs) => { for (const msg of msgs) { - const workspaceUuid = msg.id as WorkspaceUuid + const workspaceUuid = msg.workspace for (const tx of msg.value) { await handleTx(ctx, workspaceUuid, tx, transcodeProducer) } diff --git a/server/core/src/queue/dummyQueue.ts b/server/core/src/queue/dummyQueue.ts index 06cc48ded4..007760fb3a 100644 --- a/server/core/src/queue/dummyQueue.ts +++ b/server/core/src/queue/dummyQueue.ts @@ -39,7 +39,7 @@ export class DummyQueue implements PlatformQueue { topic: QueueTopic | string, groupId: string, onMessage: ( - msg: { id: WorkspaceUuid | string, value: T }[], + msg: { workspace: WorkspaceUuid, value: T }[], queue: { pause: () => void heartbeat: () => Promise diff --git a/server/core/src/queue/types.ts b/server/core/src/queue/types.ts index 7fb007ab92..c001fbae50 100644 --- a/server/core/src/queue/types.ts +++ b/server/core/src/queue/types.ts @@ -25,7 +25,7 @@ export interface ConsumerHandle { } export interface ConsumerMessage { - id: WorkspaceUuid | string + workspace: WorkspaceUuid value: T[] } @@ -68,7 +68,7 @@ export interface PlatformQueue { * Create a producer for a topic. */ export interface PlatformQueueProducer { - send: (id: WorkspaceUuid | string, msgs: T[]) => Promise + send: (workspace: WorkspaceUuid, msgs: T[], partitionKey?: string) => Promise close: () => Promise getQueue: () => PlatformQueue diff --git a/server/kafka/src/__test__/queue.spec.ts b/server/kafka/src/__test__/queue.spec.ts index ece022f836..a61cc336af 100644 --- a/server/kafka/src/__test__/queue.spec.ts +++ b/server/kafka/src/__test__/queue.spec.ts @@ -1,4 +1,4 @@ -import { generateId, MeasureMetricsContext } from '@hcengineering/core' +import { generateId, MeasureMetricsContext, type WorkspaceUuid } from '@hcengineering/core' import { createPlatformQueue, parseQueueConfig } from '..' jest.setTimeout(120000) @@ -26,7 +26,7 @@ describe('queue', () => { const producer = queue.getProducer(testCtx, 'qtest') for (let i = 0; i < docsCount; i++) { - await producer.send(genId, ['msg' + i]) + await producer.send(genId as any as WorkspaceUuid, ['msg' + i]) } await p1 @@ -55,7 +55,7 @@ describe('queue', () => { }) const producer = queue.getProducer(testCtx, 'test') - await producer.send(genId, ['msg']) + await producer.send(genId as any as WorkspaceUuid, ['msg']) await p } finally { diff --git a/server/kafka/src/index.ts b/server/kafka/src/index.ts index e8d97c3813..39686ee093 100644 --- a/server/kafka/src/index.ts +++ b/server/kafka/src/index.ts @@ -183,7 +183,7 @@ class PlatformQueueProducerImpl implements PlatformQueueProducer { return this.queue } - async send (id: WorkspaceUuid | string, msgs: any[]): Promise { + async send (workspace: WorkspaceUuid, msgs: any[], partitionKey?: string): Promise { if (this.connected !== undefined) { await this.connected this.connected = undefined @@ -192,8 +192,11 @@ class PlatformQueueProducerImpl implements PlatformQueueProducer { this.txProducer.send({ topic: this.topic, messages: msgs.map((m) => ({ - key: Buffer.from(`${id}`), - value: Buffer.from(JSON.stringify(m)) + key: Buffer.from(`${partitionKey ?? workspace}`), + value: Buffer.from(JSON.stringify(m)), + headers: { + workspace + } })) }) ) @@ -241,13 +244,15 @@ class PlatformQueueConsumerImpl implements ConsumerHandle { eachMessage: async ({ topic, message, pause, heartbeat }) => { const msgKey = message.key?.toString() ?? '' const msgData = JSON.parse(message.value?.toString() ?? '{}') + const workspace = (message.headers?.workspace?.toString() ?? msgKey) as WorkspaceUuid + let to = 1 while (true) { try { - await this.onMessage([{ id: msgKey, value: [msgData] }], { heartbeat, pause }) + await this.onMessage([{ workspace, value: [msgData] }], { heartbeat, pause }) break } catch (err: any) { - this.ctx.error('failed to process message', { err, msgKey, msgData }) + this.ctx.error('failed to process message', { err, msgKey, msgData, workspace }) await heartbeat() await new Promise((resolve) => setTimeout(resolve, to * 1000)) if (to < 10) { diff --git a/services/calendar/pod-calendar-mailer/src/index.ts b/services/calendar/pod-calendar-mailer/src/index.ts index 28a63a8766..de85be42b3 100644 --- a/services/calendar/pod-calendar-mailer/src/index.ts +++ b/services/calendar/pod-calendar-mailer/src/index.ts @@ -16,7 +16,7 @@ import { join } from 'path' import { Analytics } from '@hcengineering/analytics' import { SplitLogger, configureAnalytics } from '@hcengineering/analytics-service' -import { MeasureMetricsContext, newMetrics, WorkspaceUuid } from '@hcengineering/core' +import { MeasureMetricsContext, newMetrics } from '@hcengineering/core' import { setMetadata } from '@hcengineering/platform' import { initStatisticsContext, QueueTopic } from '@hcengineering/server-core' import serverToken from '@hcengineering/server-token' @@ -53,7 +53,7 @@ async function main (): Promise { queue.getClientId(), async (messages) => { for (const message of messages) { - const ws = message.id as WorkspaceUuid + const ws = message.workspace const records = message.value for (const record of records) { ctx.info('Processing event', { diff --git a/services/datalake/pod-datalake/src/datalake/datalake.ts b/services/datalake/pod-datalake/src/datalake/datalake.ts index f550aa67d2..c1e4cf599f 100644 --- a/services/datalake/pod-datalake/src/datalake/datalake.ts +++ b/services/datalake/pod-datalake/src/datalake/datalake.ts @@ -13,7 +13,7 @@ // limitations under the License. // -import { type MeasureContext, type Tx } from '@hcengineering/core' +import { type MeasureContext, type Tx, WorkspaceUuid } from '@hcengineering/core' import { PlatformQueueProducer } from '@hcengineering/server-core' import { Readable } from 'stream' @@ -35,7 +35,7 @@ export class DatalakeImpl implements Datalake { async list ( ctx: MeasureContext, - workspace: string, + workspace: WorkspaceUuid, options: { cursor?: string, limit?: number, derived?: boolean } ): Promise { const blobs = await this.db.listBlobs(ctx, workspace, options) @@ -49,7 +49,7 @@ export class DatalakeImpl implements Datalake { } } - async head (ctx: MeasureContext, workspace: string, name: string): Promise { + async head (ctx: MeasureContext, workspace: WorkspaceUuid, name: string): Promise { const blob = await this.db.getBlob(ctx, { workspace, name }) if (blob === null) { return null @@ -73,7 +73,7 @@ export class DatalakeImpl implements Datalake { async get ( ctx: MeasureContext, - workspace: string, + workspace: WorkspaceUuid, name: string, options: { range?: string } ): Promise { @@ -104,7 +104,7 @@ export class DatalakeImpl implements Datalake { } } - async delete (ctx: MeasureContext, workspace: string, name: string | string[]): Promise { + async delete (ctx: MeasureContext, workspace: WorkspaceUuid, name: string | string[]): Promise { if (Array.isArray(name)) { await this.db.deleteBlobList(ctx, { workspace, names: name }) } else { @@ -121,7 +121,7 @@ export class DatalakeImpl implements Datalake { async put ( ctx: MeasureContext, - workspace: string, + workspace: WorkspaceUuid, name: string, sha256: string, body: Buffer | Readable, @@ -183,7 +183,12 @@ export class DatalakeImpl implements Datalake { } } - async create (ctx: MeasureContext, workspace: string, name: string, filename: string): Promise { + async create ( + ctx: MeasureContext, + workspace: WorkspaceUuid, + name: string, + filename: string + ): Promise { const { location, bucket } = await this.selectStorage(ctx, workspace) const head = await bucket.head(ctx, filename) @@ -216,19 +221,19 @@ export class DatalakeImpl implements Datalake { return { name, size, contentType, lastModified, etag: hash } } - async getMeta (ctx: MeasureContext, workspace: string, name: string): Promise | null> { + async getMeta (ctx: MeasureContext, workspace: WorkspaceUuid, name: string): Promise | null> { return await this.db.getMeta(ctx, { workspace, name }) } - async setMeta (ctx: MeasureContext, workspace: string, name: string, meta: Record): Promise { + async setMeta (ctx: MeasureContext, workspace: WorkspaceUuid, name: string, meta: Record): Promise { await this.db.setMeta(ctx, { workspace, name }, meta) } - async setParent (ctx: MeasureContext, workspace: string, name: string, parent: string | null): Promise { + async setParent (ctx: MeasureContext, workspace: WorkspaceUuid, name: string, parent: string | null): Promise { await this.db.setParent(ctx, { workspace, name }, parent !== null ? { workspace, name: parent } : null) } - async selectStorage (ctx: MeasureContext, workspace: string, location?: Location): Promise { + async selectStorage (ctx: MeasureContext, workspace: WorkspaceUuid, location?: Location): Promise { location ??= this.selectLocation(workspace) const bucket = this.buckets.find((b) => b.location === location)?.bucket diff --git a/services/datalake/pod-datalake/src/datalake/types.ts b/services/datalake/pod-datalake/src/datalake/types.ts index 13bbe741df..c21444f6da 100644 --- a/services/datalake/pod-datalake/src/datalake/types.ts +++ b/services/datalake/pod-datalake/src/datalake/types.ts @@ -13,7 +13,7 @@ // limitations under the License. // -import { MeasureContext } from '@hcengineering/core' +import { MeasureContext, WorkspaceUuid } from '@hcengineering/core' import { type Readable } from 'stream' import { S3Bucket } from '../s3' import { WorkspaceStatsResult } from './db' @@ -51,27 +51,32 @@ export interface BlobStorage { export interface Datalake { list: ( ctx: MeasureContext, - workspace: string, + workspace: WorkspaceUuid, options: { cursor?: string, limit?: number, derived?: boolean } ) => Promise - head: (ctx: MeasureContext, workspace: string, name: string) => Promise - get: (ctx: MeasureContext, workspace: string, name: string, options: { range?: string }) => Promise - delete: (ctx: MeasureContext, workspace: string, name: string | string[]) => Promise + head: (ctx: MeasureContext, workspace: WorkspaceUuid, name: string) => Promise + get: ( + ctx: MeasureContext, + workspace: WorkspaceUuid, + name: string, + options: { range?: string } + ) => Promise + delete: (ctx: MeasureContext, workspace: WorkspaceUuid, name: string | string[]) => Promise put: ( ctx: MeasureContext, - workspace: string, + workspace: WorkspaceUuid, name: string, sha256: string, body: Buffer | Readable, options: Omit ) => Promise - create: (ctx: MeasureContext, workspace: string, name: string, filename: string) => Promise + create: (ctx: MeasureContext, workspace: WorkspaceUuid, name: string, filename: string) => Promise - getMeta: (ctx: MeasureContext, workspace: string, name: string) => Promise | null> - setMeta: (ctx: MeasureContext, workspace: string, name: string, meta: Record) => Promise + getMeta: (ctx: MeasureContext, workspace: WorkspaceUuid, name: string) => Promise | null> + setMeta: (ctx: MeasureContext, workspace: WorkspaceUuid, name: string, meta: Record) => Promise - setParent: (ctx: MeasureContext, workspace: string, name: string, parent: string | null) => Promise - selectStorage: (ctx: MeasureContext, workspace: string) => Promise + setParent: (ctx: MeasureContext, workspace: WorkspaceUuid, name: string, parent: string | null) => Promise + selectStorage: (ctx: MeasureContext, workspace: WorkspaceUuid) => Promise - getWorkspaceStats: (ctx: MeasureContext, workspace: string) => Promise + getWorkspaceStats: (ctx: MeasureContext, workspace: WorkspaceUuid) => Promise } diff --git a/services/datalake/pod-datalake/src/handlers/blob.ts b/services/datalake/pod-datalake/src/handlers/blob.ts index f190805b72..4ef864231c 100644 --- a/services/datalake/pod-datalake/src/handlers/blob.ts +++ b/services/datalake/pod-datalake/src/handlers/blob.ts @@ -14,7 +14,7 @@ // import { Analytics } from '@hcengineering/analytics' -import { MeasureContext } from '@hcengineering/core' +import { MeasureContext, type WorkspaceUuid } from '@hcengineering/core' import { type Request, type Response } from 'express' import { UploadedFile } from 'express-fileupload' import fs from 'fs' @@ -40,7 +40,7 @@ export async function handleWorkspaceStats ( datalake: Datalake ): Promise { const { workspace } = req.params - const stats = await datalake.getWorkspaceStats(ctx, workspace) + const stats = await datalake.getWorkspaceStats(ctx, workspace as WorkspaceUuid) res.status(200).json(stats) } @@ -50,7 +50,7 @@ export async function handleBlobList ( res: Response, datalake: Datalake ): Promise { - const { workspace } = req.params + const workspace = req.params.workspace as WorkspaceUuid const cursor = req.query.cursor as string const limit = extractIntParam(req.query.limit as string) const derived = req.query.derived === 'true' @@ -65,7 +65,8 @@ export async function handleBlobGet ( res: Response, datalake: Datalake ): Promise { - const { workspace, name, filename } = req.params + const { name, filename } = req.params + const workspace = req.params.workspace as WorkspaceUuid const range = req.headers.range @@ -121,7 +122,8 @@ export async function handleBlobHead ( res: Response, datalake: Datalake ): Promise { - const { workspace, name, filename } = req.params + const { name, filename } = req.params + const workspace = req.params.workspace as WorkspaceUuid const head = await datalake.head(ctx, workspace, name) if (head == null) { @@ -153,7 +155,7 @@ export async function handleBlobDelete ( const { workspace, name } = req.params try { - await datalake.delete(ctx, workspace, name) + await datalake.delete(ctx, workspace as WorkspaceUuid, name) ctx.info('deleted', { workspace, name }) res.status(204).send() @@ -174,7 +176,7 @@ export async function handleBlobDeleteList ( const body = req.body.names as DeleteBlobsRequest try { - await datalake.delete(ctx, workspace, body.names) + await datalake.delete(ctx, workspace as WorkspaceUuid, body.names) ctx.info('deleted', { workspace, names: body.names }) res.status(204).send() @@ -191,7 +193,8 @@ export async function handleBlobSetParent ( res: Response, datalake: Datalake ): Promise { - const { workspace, name } = req.params + const { name } = req.params + const workspace = req.params.workspace as WorkspaceUuid const { parent } = (await req.body) as BlobParentRequest if (parent != null) { @@ -233,7 +236,7 @@ export async function handleUploadFormData ( datalake: Datalake, tempDir: TemporaryDir ): Promise { - const { workspace } = req.params + const workspace = req.params.workspace as WorkspaceUuid if (req.files == null) { res.status(400).send('missing files') diff --git a/services/datalake/pod-datalake/src/handlers/image.ts b/services/datalake/pod-datalake/src/handlers/image.ts index 596b0687f4..5e56c60fa7 100644 --- a/services/datalake/pod-datalake/src/handlers/image.ts +++ b/services/datalake/pod-datalake/src/handlers/image.ts @@ -14,7 +14,7 @@ // import { Analytics } from '@hcengineering/analytics' -import { MeasureContext } from '@hcengineering/core' +import { MeasureContext, type WorkspaceUuid } from '@hcengineering/core' import { type Request, type Response } from 'express' import { createReadStream, createWriteStream } from 'fs' import sharp from 'sharp' @@ -105,7 +105,8 @@ export async function handleImageGet ( datalake: Datalake, tempDir: TemporaryDir ): Promise { - const { workspace, name, transform } = req.params + const { name, transform } = req.params + const workspace = req.params.workspace as WorkspaceUuid const accept = req.headers.accept ?? 'image/*' const { format, width, height, fit } = getImageTransformParams(accept, transform) diff --git a/services/datalake/pod-datalake/src/handlers/meta.ts b/services/datalake/pod-datalake/src/handlers/meta.ts index 192adf688d..b1f96e8e7a 100644 --- a/services/datalake/pod-datalake/src/handlers/meta.ts +++ b/services/datalake/pod-datalake/src/handlers/meta.ts @@ -13,7 +13,7 @@ // limitations under the License. // -import { type MeasureContext } from '@hcengineering/core' +import { type MeasureContext, type WorkspaceUuid } from '@hcengineering/core' import { type Request, type Response } from 'express' import { type Datalake } from '../datalake' @@ -24,7 +24,8 @@ export async function handleMetaGet ( res: Response, datalake: Datalake ): Promise { - const { workspace, name } = req.params + const { name } = req.params + const workspace = req.params.workspace as WorkspaceUuid const meta = await datalake.getMeta(ctx, workspace, name) if (meta == null) { @@ -41,7 +42,8 @@ export async function handleMetaPut ( res: Response, datalake: Datalake ): Promise { - const { workspace, name } = req.params + const { name } = req.params + const workspace = req.params.workspace as WorkspaceUuid const meta = req.body if (typeof meta !== 'object') { @@ -67,7 +69,8 @@ export async function handleMetaPatch ( res: Response, datalake: Datalake ): Promise { - const { workspace, name } = req.params + const { name } = req.params + const workspace = req.params.workspace as WorkspaceUuid if (typeof req.body !== 'object') { res.status(400).send() diff --git a/services/datalake/pod-datalake/src/handlers/multipart.ts b/services/datalake/pod-datalake/src/handlers/multipart.ts index 5bb64c3f21..f33b42fbb6 100644 --- a/services/datalake/pod-datalake/src/handlers/multipart.ts +++ b/services/datalake/pod-datalake/src/handlers/multipart.ts @@ -13,7 +13,7 @@ // limitations under the License. // -import { MeasureContext } from '@hcengineering/core' +import { MeasureContext, type WorkspaceUuid } from '@hcengineering/core' import { type Request, type Response } from 'express' import { cacheControl } from '../const' import { Datalake } from '../datalake' @@ -38,7 +38,7 @@ export async function handleMultipartUploadStart ( res: Response, datalake: Datalake ): Promise { - const { workspace } = req.params + const workspace = req.params.workspace as WorkspaceUuid const { bucket } = await datalake.selectStorage(ctx, workspace) @@ -63,7 +63,7 @@ export async function handleMultipartUploadPart ( res: Response, datalake: Datalake ): Promise { - const { workspace } = req.params + const workspace = req.params.workspace as WorkspaceUuid const partNumber = req.query.partNumber as string let uuid: string @@ -100,7 +100,8 @@ export async function handleMultipartUploadComplete ( res: Response, datalake: Datalake ): Promise { - const { workspace, name } = req.params + const { name } = req.params + const workspace = req.params.workspace as WorkspaceUuid const { bucket } = await datalake.selectStorage(ctx, workspace) let uuid: string @@ -129,7 +130,8 @@ export async function handleMultipartUploadAbort ( res: Response, datalake: Datalake ): Promise { - const { workspace, name } = req.params + const { name } = req.params + const workspace = req.params.workspace as WorkspaceUuid let uuid: string let uploadId: string diff --git a/services/datalake/pod-datalake/src/handlers/s3.ts b/services/datalake/pod-datalake/src/handlers/s3.ts index e60cf81c7c..bf899464c2 100644 --- a/services/datalake/pod-datalake/src/handlers/s3.ts +++ b/services/datalake/pod-datalake/src/handlers/s3.ts @@ -14,7 +14,7 @@ // import { Analytics } from '@hcengineering/analytics' -import { MeasureContext } from '@hcengineering/core' +import { MeasureContext, type WorkspaceUuid } from '@hcengineering/core' import { type Request, type Response } from 'express' import { type Datalake } from '../datalake' @@ -25,7 +25,7 @@ export async function handleS3CreateBlobParams ( res: Response, datalake: Datalake ): Promise { - const { workspace } = req.params + const workspace = req.params.workspace as WorkspaceUuid const { location, bucket } = await datalake.selectStorage(ctx, workspace) res.status(200).json({ location, bucket: bucket.bucket }) } @@ -36,7 +36,8 @@ export async function handleS3CreateBlob ( res: Response, datalake: Datalake ): Promise { - const { workspace, name } = req.params + const { name } = req.params + const workspace = req.params.workspace as WorkspaceUuid const filename = req.body.filename as string if (filename == null) { res.status(400).send('missing filename') diff --git a/services/telegram-bot/pod-telegram-bot/src/start.ts b/services/telegram-bot/pod-telegram-bot/src/start.ts index a0c1fe5314..5ab56d754f 100644 --- a/services/telegram-bot/pod-telegram-bot/src/start.ts +++ b/services/telegram-bot/pod-telegram-bot/src/start.ts @@ -15,7 +15,7 @@ import { Analytics } from '@hcengineering/analytics' import { SplitLogger, configureAnalytics } from '@hcengineering/analytics-service' -import { MeasureMetricsContext, WorkspaceUuid, newMetrics } from '@hcengineering/core' +import { MeasureMetricsContext, newMetrics } from '@hcengineering/core' import { setMetadata } from '@hcengineering/platform' import serverClient from '@hcengineering/server-client' import { initStatisticsContext, QueueTopic } from '@hcengineering/server-core' @@ -93,15 +93,15 @@ export const start = async (): Promise => { queue.getClientId(), async (messages) => { for (const message of messages) { - const id = message.id as WorkspaceUuid + const workspace = message.workspace const records = message.value for (const record of records) { switch (record.type) { case TelegramQueueMessageType.Notification: - await worker.processNotification(id, record, bot) + await worker.processNotification(workspace, record, bot) break case TelegramQueueMessageType.WorkspaceSubscription: - await worker.processWorkspaceSubscription(id, record) + await worker.processWorkspaceSubscription(workspace, record) break } }