Pass meta to tx and queue (#9695)

+ Pass ctx context via platform queue
+ Fix queue contracts
+ Add meta to tx, skip store.

Signed-off-by: Andrey Sobolev <haiodo@gmail.com>
This commit is contained in:
Andrey Sobolev
2025-08-20 12:13:43 +07:00
committed by GitHub
parent b438971926
commit c27ba3150f
26 changed files with 268 additions and 303 deletions
+2 -2
View File
@@ -1667,7 +1667,7 @@ export async function restoreFromv6All (
const queue = getPlatformQueue('tool', workspace.region)
const wsProducer = queue.getProducer<QueueWorkspaceMessage>(ctx, QueueTopic.Workspace)
await wsProducer.send(uuid, [workspaceEvents.restoring()])
await wsProducer.send(ctx, uuid, [workspaceEvents.restoring()])
const workspaceStorage: StorageAdapter = buildStorageFromConfig(storageConfig)
@@ -1693,7 +1693,7 @@ export async function restoreFromv6All (
await sendTransactorEvent(uuid, 'force-close')
ctx.info('workspace restored', { dataId })
await wsProducer.send(uuid, [workspaceEvents.restored()])
await wsProducer.send(ctx, uuid, [workspaceEvents.restored()])
if (activeWorkspaces.has(uuid)) {
// set workspace back to active
+6 -6
View File
@@ -415,7 +415,7 @@ export function devTool (
progress: 100
})
await wsProducer.send(res.workspaceUuid, [workspaceEvents.created()])
await wsProducer.send(measureCtx, res.workspaceUuid, [workspaceEvents.created()])
await queue.shutdown()
console.log(queue)
})
@@ -500,7 +500,7 @@ export function devTool (
console.log(metricsToString(measureCtx.metrics, 'upgrade', 60))
await wsProducer.send(info.uuid, [workspaceEvents.upgraded()])
await wsProducer.send(measureCtx, info.uuid, [workspaceEvents.upgraded()])
await queue.shutdown()
console.log('upgrade-workspace done')
})
@@ -1160,7 +1160,7 @@ export function devTool (
const queue = getPlatformQueue('tool', ws.region)
const wsProducer = queue.getProducer<QueueWorkspaceMessage>(toolCtx, QueueTopic.Workspace)
await wsProducer.send(ws.uuid, [workspaceEvents.restoring()])
await wsProducer.send(toolCtx, ws.uuid, [workspaceEvents.restoring()])
const workspaceStorage: StorageAdapter = buildStorageFromConfig(storageConfig)
@@ -1202,7 +1202,7 @@ export function devTool (
}
console.log('workspace restored')
await wsProducer.send(ws.uuid, [workspaceEvents.restored()])
await wsProducer.send(toolCtx, ws.uuid, [workspaceEvents.restored()])
} catch (err) {
toolCtx.error('failed to restore', { err })
}
@@ -2226,7 +2226,7 @@ export function devTool (
console.log('reindex workspace', workspace)
const queue = getPlatformQueue('tool', ws.region)
const wsProducer = queue.getProducer<QueueWorkspaceMessage>(toolCtx, QueueTopic.Workspace)
await wsProducer.send(ws.uuid, [workspaceEvents.fullReindex()])
await wsProducer.send(toolCtx, ws.uuid, [workspaceEvents.fullReindex()])
await queue.shutdown()
console.log('done', workspace)
})
@@ -2262,7 +2262,7 @@ export function devTool (
console.log('reindex workspace', ws)
const queue = getPlatformQueue('tool', ws.region)
const wsProducer = queue.getProducer<QueueWorkspaceMessage>(toolCtx, QueueTopic.Workspace)
await wsProducer.send(ws.uuid, [workspaceEvents.fullReindex()])
await wsProducer.send(toolCtx, ws.uuid, [workspaceEvents.fullReindex()])
await queue.shutdown()
}
console.log('done')
+1
View File
@@ -42,6 +42,7 @@ import { generateId } from './utils'
*/
export interface Tx extends Doc {
objectSpace: Ref<Space> // space where transaction will operate
meta?: Record<string, string | number | boolean> // meta information about transaction, non persisted to final DB's
}
/**
+2 -2
View File
@@ -138,7 +138,7 @@ describe('full-text-indexing', () => {
const dataId = generateId()
await queue.expectIndexingDoc(dataId, async () => {
await txProducer.send(wsId, [
await txProducer.send(toolCtx, wsId, [
createDoc(test.class.TestDocument, {
title: 'first doc',
description: dataId
@@ -211,7 +211,7 @@ describe('full-text-indexing', () => {
}
})
await wsProcessor.send(wsIds.uuid, [workspaceEvents.fullReindex()])
await wsProcessor.send(toolCtx, wsIds.uuid, [workspaceEvents.fullReindex()])
// Wait for reindex
await reindexAllP
+100 -107
View File
@@ -117,8 +117,8 @@ export class WorkspaceManager {
this.ctx,
QueueTopic.Workspace,
this.opt.queue.getClientId(),
async (msg, control) => {
await this.processWorkspaceEvent(msg, control)
async (ctx, msg, control) => {
await this.processWorkspaceEvent(ctx, msg, control)
}
)
@@ -126,7 +126,7 @@ export class WorkspaceManager {
this.ctx,
QueueTopic.Fulltext,
this.opt.queue.getClientId(),
async (msg, control) => {
async (ctx, msg, control) => {
await this.processFulltextEvent(msg, control)
}
)
@@ -136,14 +136,14 @@ export class WorkspaceManager {
this.ctx,
QueueTopic.Tx,
this.opt.queue.getClientId(),
async (msg, control) => {
async (ctx, msg, control) => {
clearTimeout(this.txInformer)
this.txInformer = setTimeout(() => {
this.ctx.info('tx message', { count: txMessages })
txMessages = 0
}, 5000)
txMessages += msg.length
txMessages += 1
await this.processTransactions(msg, control)
}
@@ -151,133 +151,126 @@ export class WorkspaceManager {
}
private async processTransactions (
msg: ConsumerMessage<TxCUD<Doc<Space>> | TxDomainEvent<QueueSourced<Event>>>[],
m: ConsumerMessage<TxCUD<Doc<Space>> | TxDomainEvent<QueueSourced<Event>>>,
control: ConsumerControl
): Promise<void> {
for (const m of msg) {
const ws = m.workspace
const ws = m.workspace
let token: string
try {
token = generateToken(systemAccountUuid, ws, { service: 'fulltext' })
} catch (err: any) {
this.ctx.error('Error generating token', { err, systemAccountUuid, ws })
continue
}
await this.withIndexer(this.ctx, ws, token, true, async (indexer) => {
await indexer.fulltext.processTransactions(this.ctx, m.value, control)
})
let token: string
try {
token = generateToken(systemAccountUuid, ws, { service: 'fulltext' })
} catch (err: any) {
this.ctx.error('Error generating token', { err, systemAccountUuid, ws })
throw err
}
await this.withIndexer(this.ctx, ws, token, true, async (indexer) => {
await indexer.fulltext.processTransactions(this.ctx, [m.value], control)
})
}
private async processWorkspaceEvent (
msg: ConsumerMessage<QueueWorkspaceMessage>[],
ctx: MeasureContext,
m: ConsumerMessage<QueueWorkspaceMessage>,
control: ConsumerControl
): Promise<void> {
for (const m of msg) {
const ws = m.workspace
const ws = m.workspace
const mm = m.value
for (const mm of m.value) {
this.ctx.info('workspace event', { type: mm.type, workspace: ws, restoring: Array.from(this.restoring) })
let token: string
try {
token = generateToken(systemAccountUuid, ws, { service: 'fulltext' })
} catch (err: any) {
this.ctx.error('Error generating token', { err, systemAccountUuid, ws })
continue
}
this.ctx.info('workspace event', { type: mm.type, workspace: ws, restoring: Array.from(this.restoring) })
let token: string
try {
token = generateToken(systemAccountUuid, ws, { service: 'fulltext' })
} catch (err: any) {
this.ctx.error('Error generating token', { err, systemAccountUuid, ws })
return
}
if (mm.type === QueueWorkspaceEvent.Restoring) {
this.restoring.add(ws)
await this.closeWorkspace(ws)
} else if (
mm.type === QueueWorkspaceEvent.Created ||
mm.type === QueueWorkspaceEvent.Restored ||
mm.type === QueueWorkspaceEvent.FullReindex
) {
if (mm.type === QueueWorkspaceEvent.Restored) {
this.restoring.delete(ws)
}
if (this.restoring.has(ws)) {
// Ignore fulltext in case of restoring
continue
}
await this.fulltextProducer.send(ws, [workspaceEvents.fullReindex()])
} else if (
mm.type === QueueWorkspaceEvent.Deleted ||
mm.type === QueueWorkspaceEvent.Archived ||
mm.type === QueueWorkspaceEvent.ClearIndex
) {
const workspaceInfo = await this.getWorkspaceInfo(this.ctx, token)
if (workspaceInfo !== undefined) {
await this.fulltextAdapter.clean(
this.ctx,
(workspaceInfo.dataId as unknown as WorkspaceUuid) ?? workspaceInfo.uuid
)
}
} else if (mm.type === QueueWorkspaceEvent.Upgraded) {
this.ctx.warn('Upgraded', this.supportedVersion)
await this.closeWorkspace(ws)
}
if (mm.type === QueueWorkspaceEvent.Restoring) {
this.restoring.add(ws)
await this.closeWorkspace(ws)
} else if (
mm.type === QueueWorkspaceEvent.Created ||
mm.type === QueueWorkspaceEvent.Restored ||
mm.type === QueueWorkspaceEvent.FullReindex
) {
if (mm.type === QueueWorkspaceEvent.Restored) {
this.restoring.delete(ws)
}
if (this.restoring.has(ws)) {
// Ignore fulltext in case of restoring
return
}
await this.fulltextProducer.send(ctx, ws, [workspaceEvents.fullReindex()])
} else if (
mm.type === QueueWorkspaceEvent.Deleted ||
mm.type === QueueWorkspaceEvent.Archived ||
mm.type === QueueWorkspaceEvent.ClearIndex
) {
const workspaceInfo = await this.getWorkspaceInfo(this.ctx, token)
if (workspaceInfo !== undefined) {
await this.fulltextAdapter.clean(
this.ctx,
(workspaceInfo.dataId as unknown as WorkspaceUuid) ?? workspaceInfo.uuid
)
}
} else if (mm.type === QueueWorkspaceEvent.Upgraded) {
this.ctx.warn('Upgraded', this.supportedVersion)
await this.closeWorkspace(ws)
}
}
private async processFulltextEvent (
msg: ConsumerMessage<QueueWorkspaceMessage>[],
m: ConsumerMessage<QueueWorkspaceMessage>,
control: ConsumerControl
): Promise<void> {
for (const m of msg) {
const ws = m.workspace
const ws = m.workspace
const mm = m.value
for (const mm of m.value) {
this.ctx.info('fulltext event', { type: mm.type, workspace: ws, restoring: Array.from(this.restoring) })
let token: string
try {
token = generateToken(systemAccountUuid, ws, { service: 'fulltext' })
} catch (err: any) {
this.ctx.error('Error generating token', { err, systemAccountUuid, ws })
continue
}
this.ctx.info('fulltext event', { type: mm.type, workspace: ws, restoring: Array.from(this.restoring) })
let token: string
try {
token = generateToken(systemAccountUuid, ws, { service: 'fulltext' })
} catch (err: any) {
this.ctx.error('Error generating token', { err, systemAccountUuid, ws })
return
}
if (mm.type === QueueWorkspaceEvent.FullReindex) {
await this.withIndexer(this.ctx, ws, token, true, async (indexer) => {
await indexer.dropWorkspace()
const toIndex = await indexer.getIndexClassess()
this.ctx.info('reindex starting full', { workspace: ws })
await this.ctx.with(
'reindex-workspace',
{},
async (ctx) => {
for (const { domain, classes } of toIndex) {
try {
await control.heartbeat()
await indexer.reindex(ctx, domain, classes, control)
} catch (err: any) {
ctx.error('failed to reindex domain', { workspace: ws })
throw err
}
}
},
{ workspace: ws }
)
this.ctx.info('reindex full done', { workspace: ws })
})
} else if (mm.type === QueueWorkspaceEvent.Reindex) {
const mmd = mm as QueueWorkspaceReindexMessage
if (!this.restoring.has(ws)) {
await this.withIndexer(this.ctx, ws, token, true, async (indexer) => {
if (mm.type === QueueWorkspaceEvent.FullReindex) {
await this.withIndexer(this.ctx, ws, token, true, async (indexer) => {
await indexer.dropWorkspace()
const toIndex = await indexer.getIndexClassess()
this.ctx.info('reindex starting full', { workspace: ws })
await this.ctx.with(
'reindex-workspace',
{},
async (ctx) => {
for (const { domain, classes } of toIndex) {
try {
await indexer.reindex(this.ctx, mmd.domain, mmd.classes, control)
await control.heartbeat()
await indexer.reindex(ctx, domain, classes, control)
} catch (err: any) {
this.ctx.error('failed to reindex domain', { workspace: ws })
ctx.error('failed to reindex domain', { workspace: ws })
throw err
}
})
}
},
{ workspace: ws }
)
this.ctx.info('reindex full done', { workspace: ws })
})
} else if (mm.type === QueueWorkspaceEvent.Reindex) {
const mmd = mm as QueueWorkspaceReindexMessage
if (!this.restoring.has(ws)) {
await this.withIndexer(this.ctx, ws, token, true, async (indexer) => {
try {
await indexer.reindex(this.ctx, mmd.domain, mmd.classes, control)
} catch (err: any) {
this.ctx.error('failed to reindex domain', { workspace: ws })
throw err
}
}
})
}
}
}
+17 -7
View File
@@ -218,14 +218,24 @@ export async function startIndexer (
req.body = {}
ctx.info('reindex', { workspace: decoded.workspace })
await manager.withIndexer(ctx, decoded.workspace, token, true, async (indexer) => {
indexer.lastUpdate = Date.now()
if (request?.onlyDrop ?? false) {
await manager.fulltextProducer.send(decoded.workspace, [workspaceEvents.clearIndex()])
} else {
await manager.fulltextProducer.send(decoded.workspace, [workspaceEvents.fullReindex()])
await ctx.with(
'reindex',
{},
async (ctx) => {
await manager.withIndexer(ctx, decoded.workspace, token, true, async (indexer) => {
indexer.lastUpdate = Date.now()
if (request?.onlyDrop ?? false) {
await manager.fulltextProducer.send(ctx, decoded.workspace, [workspaceEvents.clearIndex()])
} else {
await manager.fulltextProducer.send(ctx, decoded.workspace, [workspaceEvents.fullReindex()])
}
})
},
{},
{
span: 'inherit'
}
})
)
} catch (err: any) {
Analytics.handleError(err)
console.error(err)
+2 -2
View File
@@ -84,7 +84,7 @@ async function handleCreateDocTx (
}
const msg: VideoTranscodeRequest = { workspaceUuid, blobId, contentType, source }
ctx.info('transcode request', { workspaceUuid, msg })
await producer.send(workspaceUuid, [msg])
await producer.send(ctx, workspaceUuid, [msg])
}
}
@@ -111,7 +111,7 @@ async function handleCommunicationTx (
source
}))
if (messages.length > 0) {
await producer.send(workspaceUuid, messages)
await producer.send(ctx, workspaceUuid, messages)
}
}
}
+5 -13
View File
@@ -63,21 +63,13 @@ async function main (): Promise<void> {
const transcodeProducer = queue.getProducer<VideoTranscodeRequest>(ctx, topicTranscodeRequest)
queue.createConsumer<VideoTranscodeResult>(ctx, topicTranscodeResult, application, async (msgs) => {
for (const msg of msgs) {
for (const res of msg.value) {
await handleTranscodeResult(ctx, msg.workspace, res)
}
}
queue.createConsumer<VideoTranscodeResult>(ctx, topicTranscodeResult, application, async (ctx, msg) => {
await handleTranscodeResult(ctx, msg.workspace, msg.value)
})
queue.createConsumer<TxCUD<Doc>>(ctx, QueueTopic.Tx, queue.getClientId(), async (msgs) => {
for (const msg of msgs) {
const workspaceUuid = msg.workspace
for (const tx of msg.value) {
await handleTx(ctx, workspaceUuid, tx, transcodeProducer)
}
}
queue.createConsumer<TxCUD<Doc>>(ctx, QueueTopic.Tx, queue.getClientId(), async (ctx, msg) => {
const workspaceUuid = msg.workspace
await handleTx(ctx, workspaceUuid, msg.value, transcodeProducer)
})
const shutdownAsync = async (): Promise<void> => {
@@ -351,7 +351,7 @@ async function putEventToQueue (
)
try {
await producer.send(control.workspace.uuid, [{ action, event, modifiedBy, changes }])
await producer.send(control.ctx, control.workspace.uuid, [{ action, event, modifiedBy, changes }])
} catch (err) {
control.ctx.error('Could not queue calendar event', { err, action, event })
}
@@ -75,7 +75,7 @@ async function putEventToQueue (value: Omit<ProcessMessage, 'account'>, control:
const producer = control.queue.getProducer<ProcessMessage>(control.ctx.newChild('queue', {}), QueueTopic.Process)
try {
await producer.send(control.workspace.uuid, [
await producer.send(control.ctx, control.workspace.uuid, [
{
...value,
account: control.txFactory.account
@@ -354,7 +354,7 @@ async function processNotification (
link
}
await producer.send(control.workspace.uuid, [record])
await producer.send(control.ctx, control.workspace.uuid, [record])
} catch (err) {
control.ctx.error('Could not send telegram notification', {
err,
@@ -376,7 +376,7 @@ async function updateWorkspaceSubscription (
if (account == null) {
return
}
await producer.send(control.workspace.uuid, [
await producer.send(control.ctx, control.workspace.uuid, [
{
type: TelegramQueueMessageType.WorkspaceSubscription,
account,
+3 -2
View File
@@ -5,7 +5,7 @@ import { type ConsumerHandle, type PlatformQueue, type PlatformQueueProducer, ty
* A dummy implementation of PlatformQueueProducer for testing and development
*/
class DummyQueueProducer<T> implements PlatformQueueProducer<T> {
async send (id: WorkspaceUuid | string, msgs: T[]): Promise<void> {
async send (ctx: MeasureContext, id: WorkspaceUuid | string, msgs: T[]): Promise<void> {
await Promise.resolve()
}
@@ -39,7 +39,8 @@ export class DummyQueue implements PlatformQueue {
topic: QueueTopic | string,
groupId: string,
onMessage: (
msg: { workspace: WorkspaceUuid, value: T }[],
ctx: MeasureContext,
msg: { workspace: WorkspaceUuid, value: T },
queue: {
pause: () => void
heartbeat: () => Promise<void>
+3 -3
View File
@@ -29,7 +29,7 @@ export interface ConsumerHandle {
export interface ConsumerMessage<T> {
workspace: WorkspaceUuid
value: T[]
value: T
}
export interface ConsumerControl {
@@ -48,7 +48,7 @@ export interface PlatformQueue {
ctx: MeasureContext,
topic: QueueTopic | string,
groupId: string,
onMessage: (msg: ConsumerMessage<T>[], queue: ConsumerControl) => Promise<void>,
onMessage: (ctx: MeasureContext, msg: ConsumerMessage<T>, queue: ConsumerControl) => Promise<void>,
options?: {
fromBegining?: boolean
}
@@ -71,7 +71,7 @@ export interface PlatformQueue {
* Create a producer for a topic.
*/
export interface PlatformQueueProducer<T> {
send: (workspace: WorkspaceUuid, msgs: T[], partitionKey?: string) => Promise<void>
send: (ctx: MeasureContext, workspace: WorkspaceUuid, msgs: T[], partitionKey?: string) => Promise<void>
close: () => Promise<void>
getQueue: () => PlatformQueue
+5 -5
View File
@@ -14,8 +14,8 @@ describe('queue', () => {
const to = setTimeout(() => {
reject(new Error(`Timeout waiting for messages:${msgCount}`))
}, 100000)
queue.createConsumer<string>(testCtx, 'qtest', genId, async (msg) => {
msgCount += msg.length
queue.createConsumer<string>(testCtx, 'qtest', genId, async (ctx, msg) => {
msgCount += 1
console.log('msgCount', msgCount)
if (msgCount === docsCount) {
clearTimeout(to)
@@ -26,7 +26,7 @@ describe('queue', () => {
const producer = queue.getProducer<string>(testCtx, 'qtest')
for (let i = 0; i < docsCount; i++) {
await producer.send(genId as any as WorkspaceUuid, ['msg' + i])
await producer.send(testCtx, genId as any as WorkspaceUuid, ['msg' + i])
}
await p1
@@ -45,7 +45,7 @@ describe('queue', () => {
try {
let counter = 2
const p = new Promise<void>((resolve, reject) => {
queue.createConsumer<string>(testCtx, 'test', genId, async (msg) => {
queue.createConsumer<string>(testCtx, 'test', genId, async (ctx, msg) => {
counter--
if (counter > 0) {
throw new Error('Processing Error')
@@ -55,7 +55,7 @@ describe('queue', () => {
})
const producer = queue.getProducer<string>(testCtx, 'test')
await producer.send(genId as any as WorkspaceUuid, ['msg'])
await producer.send(testCtx, genId as any as WorkspaceUuid, ['msg'])
await p
} finally {
+19 -38
View File
@@ -98,7 +98,7 @@ class PlatformQueueImpl implements PlatformQueue {
ctx: MeasureContext,
topic: QueueTopic | string,
groupId: string,
onMessage: (msg: ConsumerMessage<T>[], queue: ConsumerControl) => Promise<void>,
onMessage: (ctx: MeasureContext, msg: ConsumerMessage<T>, queue: ConsumerControl) => Promise<void>,
options?: {
fromBegining?: boolean
}
@@ -184,7 +184,7 @@ class PlatformQueueProducerImpl implements PlatformQueueProducer<any> {
return this.queue
}
async send (workspace: WorkspaceUuid, msgs: any[], partitionKey?: string): Promise<void> {
async send (ctx: MeasureContext, workspace: WorkspaceUuid, msgs: any[], partitionKey?: string): Promise<void> {
if (this.connected !== undefined) {
await this.connected
this.connected = undefined
@@ -196,7 +196,8 @@ class PlatformQueueProducerImpl implements PlatformQueueProducer<any> {
key: Buffer.from(`${partitionKey ?? workspace}`),
value: Buffer.from(JSON.stringify(m)),
headers: {
workspace
workspace,
meta: JSON.stringify(ctx.extractMeta())
}
}))
})
@@ -222,7 +223,11 @@ class PlatformQueueConsumerImpl implements ConsumerHandle {
readonly config: QueueConfig,
private readonly topic: QueueTopic | string,
groupId: string,
private readonly onMessage: (msg: ConsumerMessage<any>[], queue: ConsumerControl) => Promise<void>,
private readonly onMessage: (
ctx: MeasureContext,
msg: ConsumerMessage<any>,
queue: ConsumerControl
) => Promise<void>,
private readonly options?: {
fromBegining?: boolean
}
@@ -245,12 +250,21 @@ class PlatformQueueConsumerImpl implements ConsumerHandle {
eachMessage: async ({ topic, message, pause, heartbeat }) => {
const msgKey = message.key?.toString() ?? ''
const msgData = JSON.parse(message.value?.toString() ?? '{}')
const meta = JSON.parse(message.headers?.meta?.toString() ?? '{}')
const workspace = (message.headers?.workspace?.toString() ?? msgKey) as WorkspaceUuid
let to = 1
while (true) {
try {
await this.onMessage([{ workspace, value: [msgData] }], { heartbeat, pause })
await this.ctx.with(
'handle-msg',
{},
(ctx) => this.onMessage(ctx, { workspace, value: msgData }, { heartbeat, pause }),
{},
{
meta
}
)
break
} catch (err: any) {
this.ctx.error('failed to process message', { err, msgKey, msgData, workspace })
@@ -262,39 +276,6 @@ class PlatformQueueConsumerImpl implements ConsumerHandle {
}
}
}
// , // TODO: Finish testinf
// eachBatch: async ({ batch, pause, heartbeat, resolveOffset }) => {
// const queueInfo = {
// pause,
// heartbeat
// }
// const batchMessages = batch.messages
// const currentMsg: ConsumerMessage<any> = {
// id: '',
// value: []
// }
// let lastOffset: string = ' '
// const sendLast = async (): Promise<void> => {
// await this.onMessage([currentMsg], queueInfo)
// // Mark last offset as cusomed
// resolveOffset(lastOffset)
// await heartbeat()
// }
// for (const v of batchMessages) {
// const id = v.key?.toString() ?? ''
// if (currentMsg.id !== id && currentMsg.value.length > 0) {
// await sendLast() // Send last message
// currentMsg.id = id
// currentMsg.value = []
// }
// lastOffset = v.offset
// currentMsg.value.push(JSON.parse(v.value?.toString() ?? '{}'))
// }
// await sendLast()
// }
})
}
+6 -1
View File
@@ -49,11 +49,16 @@ export class QueueMiddleware extends BaseMiddleware {
this.connected = undefined
}
const meta = ctx.extractMeta()
await Promise.all([
this.provideBroadcast(ctx),
this.txProducer.send(
ctx,
this.context.workspace.uuid,
ctx.contextData.broadcast.txes.concat(ctx.contextData.broadcast.queue)
ctx.contextData.broadcast.txes
.concat(ctx.contextData.broadcast.queue)
.map((tx) => ({ ...tx, meta: { ...(tx.meta ?? {}), ...meta } }))
)
])
}
+2 -1
View File
@@ -58,7 +58,8 @@ export class TxMiddleware extends BaseMiddleware implements Middleware {
objectClass !== core.class.BenchmarkDoc &&
this.context.hierarchy.findDomain(objectClass) !== DOMAIN_TRANSIENT
) {
txToStore.push(tx)
const { meta, ...txData } = tx
txToStore.push(txData)
}
}
}
+13 -16
View File
@@ -151,18 +151,15 @@ export class TSessionManager implements SessionManager {
ctx.newChild('ws-queue-consume', {}, { span: false }),
QueueTopic.Workspace,
generateId(),
async (messages) => {
for (const msg of messages) {
for (const m of msg.value) {
if (
m.type === QueueWorkspaceEvent.Upgraded ||
m.type === QueueWorkspaceEvent.Restored ||
m.type === QueueWorkspaceEvent.Deleted
) {
// Handle workspace messages
this.workspaceInfoCache.delete(msg.workspace)
}
}
async (ctx, msg) => {
const m = msg.value
if (
m.type === QueueWorkspaceEvent.Upgraded ||
m.type === QueueWorkspaceEvent.Restored ||
m.type === QueueWorkspaceEvent.Deleted
) {
// Handle workspace messages
this.workspaceInfoCache.delete(msg.workspace)
}
}
)
@@ -502,7 +499,7 @@ export class TSessionManager implements SessionManager {
})
workspace = this.createWorkspace(ctx.parent ?? ctx, ctx, token, workspaceInfo.url, workspaceInfo.dataId, branding)
await this.workspaceProducer.send(workspaceUuid, [workspaceEvents.open()])
await this.workspaceProducer.send(ctx, workspaceUuid, [workspaceEvents.open()])
}
if (token.extra?.model === 'upgrade') {
@@ -629,7 +626,7 @@ export class TSessionManager implements SessionManager {
const accountUuid = account.account
if (accountUuid !== systemAccountUuid && accountUuid !== guestAccount) {
await this.usersProducer.send(workspace.wsId.uuid, [
await this.usersProducer.send(ctx, workspace.wsId.uuid, [
userEvents.login({
user: accountUuid,
sessions: this.countUserSessions(workspace, accountUuid),
@@ -933,7 +930,7 @@ export class TSessionManager implements SessionManager {
workspace.sessions.delete(sessionRef.session.sessionId)
const userUuid = sessionRef.session.getUser()
await this.usersProducer.send(workspaceUuid, [
await this.usersProducer.send(ctx, workspaceUuid, [
userEvents.logout({
user: userUuid,
sessions: this.countUserSessions(workspace, userUuid),
@@ -1092,7 +1089,7 @@ export class TSessionManager implements SessionManager {
this.ctx.warn('Closed workspace', logParams)
}
await this.workspaceProducer.send(workspace.wsId.uuid, [workspaceEvents.down()])
await this.workspaceProducer.send(this.ctx, workspace.wsId.uuid, [workspaceEvents.down()])
}
} catch (err: any) {
Analytics.handleError(err)
+2 -2
View File
@@ -128,10 +128,10 @@ export class MigrateClientImpl implements MigrationClient {
}
async fullReindex (): Promise<void> {
await this.queue.send(this.wsIds.uuid, [workspaceEvents.fullReindex()])
await this.queue.send(this.ctx, this.wsIds.uuid, [workspaceEvents.fullReindex()])
}
async reindex (domain: Domain, classes: Ref<Class<Doc>>[]): Promise<void> {
await this.queue.send(this.wsIds.uuid, [workspaceEvents.reindex(domain, classes)])
await this.queue.send(this.ctx, this.wsIds.uuid, [workspaceEvents.reindex(domain, classes)])
}
}
+8 -8
View File
@@ -291,7 +291,7 @@ export class WorkspaceWorker {
time: Date.now() - t
})
await this.workspaceQueue.send(ws.uuid, [workspaceEvents.created()])
await this.workspaceQueue.send(ctx, ws.uuid, [workspaceEvents.created()])
} catch (err: any) {
void opt.errorHandler(ws, err)
@@ -307,7 +307,7 @@ export class WorkspaceWorker {
region: this.region,
time: Date.now() - t
})
await this.workspaceQueue.send(ws.uuid, [workspaceEvents.createFailed()])
await this.workspaceQueue.send(ctx, ws.uuid, [workspaceEvents.createFailed()])
} finally {
if (!opt.console) {
;(logger as FileModelLogger).close()
@@ -390,7 +390,7 @@ export class WorkspaceWorker {
region: this.region,
time: Date.now() - t
})
await this.workspaceQueue.send(ws.uuid, [workspaceEvents.upgraded()])
await this.workspaceQueue.send(ctx, ws.uuid, [workspaceEvents.upgraded()])
} catch (err: any) {
void opt.errorHandler(ws, err)
@@ -407,7 +407,7 @@ export class WorkspaceWorker {
region: this.region,
time: Date.now() - t
})
await this.workspaceQueue.send(ws.uuid, [workspaceEvents.upgradeFailed()])
await this.workspaceQueue.send(ctx, ws.uuid, [workspaceEvents.upgradeFailed()])
} finally {
if (!opt.console) {
;(logger as FileModelLogger).close()
@@ -423,7 +423,7 @@ export class WorkspaceWorker {
const adapter = getWorkspaceDestroyAdapter(dbUrl)
await adapter.deleteWorkspace(ctx, workspace.uuid, workspace.dataId)
await this.workspaceQueue.send(workspace.uuid, [workspaceEvents.clearIndex()])
await this.workspaceQueue.send(ctx, workspace.uuid, [workspaceEvents.clearIndex()])
}
async sendTransactorMaitenance (token: string, ws: WorkspaceUuid): Promise<void> {
@@ -492,7 +492,7 @@ export class WorkspaceWorker {
return
}
await sendEvent('archiving-clean-done', 100)
await this.workspaceQueue.send(workspace.uuid, [workspaceEvents.archived()])
await this.workspaceQueue.send(ctx, workspace.uuid, [workspaceEvents.archived()])
break
}
case 'pending-deletion':
@@ -507,7 +507,7 @@ export class WorkspaceWorker {
return
}
await sendEvent('delete-done', 100)
await this.workspaceQueue.send(workspace.uuid, [workspaceEvents.deleted()])
await this.workspaceQueue.send(ctx, workspace.uuid, [workspaceEvents.deleted()])
break
}
@@ -545,7 +545,7 @@ export class WorkspaceWorker {
workspace.mode = 'active'
await this._upgradeWorkspace(ctx, workspace, opt)
await this.workspaceQueue.send(workspace.uuid, [workspaceEvents.restored()])
await this.workspaceQueue.send(ctx, workspace.uuid, [workspaceEvents.restored()])
}
break
default:
@@ -51,41 +51,38 @@ async function main (): Promise<void> {
ctx,
QueueTopic.CalendarEventCUD,
queue.getClientId(),
async (messages) => {
for (const message of messages) {
const ws = message.workspace
const records = message.value
for (const record of records) {
ctx.info('Processing event', {
ws,
action: record.action,
eventId: record.event.eventId,
objectId: record.event._id,
modifiedBy: record.modifiedBy
})
try {
let skipReason
switch (record.action) {
case 'create':
skipReason = await eventCreated(ctx, ws, record)
break
case 'update':
skipReason = await eventUpdated(ctx, ws, record)
break
case 'delete':
skipReason = await eventDeleted(ctx, ws, record)
break
case 'mixin':
skipReason = await eventMixin(ctx, ws, record)
break
}
if (skipReason !== undefined) {
ctx.info('Notification skipped', { reason: skipReason, objectId: record.event._id })
}
} catch (error) {
ctx.error('Error processing event', { error, ws, record })
}
async (ctx, message) => {
const ws = message.workspace
const record = message.value
ctx.info('Processing event', {
ws,
action: record.action,
eventId: record.event.eventId,
objectId: record.event._id,
modifiedBy: record.modifiedBy
})
try {
let skipReason
switch (record.action) {
case 'create':
skipReason = await eventCreated(ctx, ws, record)
break
case 'update':
skipReason = await eventUpdated(ctx, ws, record)
break
case 'delete':
skipReason = await eventDeleted(ctx, ws, record)
break
case 'mixin':
skipReason = await eventMixin(ctx, ws, record)
break
}
if (skipReason !== undefined) {
ctx.info('Notification skipped', { reason: skipReason, objectId: record.event._id })
}
} catch (error) {
ctx.error('Error processing event', { error, ws, record })
}
}
)
@@ -144,7 +144,7 @@ export class DatalakeImpl implements Datalake {
try {
const events = Array.isArray(name) ? name.map((n) => blobEvents.deleted(n)) : [blobEvents.deleted(name)]
await this.producer.send(workspace, events)
await this.producer.send(ctx, workspace, events)
} catch (err) {
ctx.error('failed to send blob deleted event', { workspace, name, err })
}
@@ -185,7 +185,7 @@ export class DatalakeImpl implements Datalake {
blob != null
? blobEvents.updated(name, { contentType, lastModified, size, etag })
: blobEvents.created(name, { contentType, lastModified, size, etag })
await this.producer.send(workspace, [event])
await this.producer.send(ctx, workspace, [event])
} catch (err) {
ctx.error('failed to send blob created event', { workspace, name, err })
}
@@ -246,7 +246,7 @@ export class DatalakeImpl implements Datalake {
blob != null
? blobEvents.updated(name, { contentType, lastModified, size, etag })
: blobEvents.created(name, { contentType, lastModified, size, etag })
await this.producer.send(workspace, [event])
await this.producer.send(ctx, workspace, [event])
} catch (err) {
this.cache.delete(hash)
ctx.error('failed to send blob created event', { workspace, name, err })
@@ -286,7 +286,7 @@ export class DatalakeImpl implements Datalake {
data != null
? blobEvents.updated(name, { contentType, lastModified, size, etag: hash })
: blobEvents.created(name, { contentType, lastModified, size, etag: hash })
await this.producer.send(workspace, [event])
await this.producer.send(ctx, workspace, [event])
} catch (err) {
ctx.error('failed to send blob created event', { workspace, name, err })
}
@@ -140,15 +140,12 @@ export class GmailController {
this.ctx,
QueueTopic.Tx,
this.queue.getClientId(),
async (msgs) => {
for (const msg of msgs) {
const workspaceUuid = msg.workspace
for (const tx of msg.value) {
const messageEvent = toMessageEvent(tx)
if (messageEvent !== undefined) {
await this.handleNewMessage(workspaceUuid, messageEvent)
}
}
async (ctx, msg) => {
const workspaceUuid = msg.workspace
const messageEvent = toMessageEvent(msg.value)
if (messageEvent !== undefined) {
await this.handleNewMessage(workspaceUuid, messageEvent)
}
},
{
+13 -15
View File
@@ -123,21 +123,19 @@ export class MailWorker {
this.ctx,
QueueTopic.Tx,
this.queue.getClientId(),
async (msgs) => {
for (const msg of msgs) {
const workspaceUuid = msg.workspace
for (const tx of msg.value) {
// Check for new channel creation
if (isNewChannelTx(tx)) {
await this.handleNewChannelTx(workspaceUuid, tx)
continue
}
// Check for message events
const messageEvent = toMessageEvent(tx)
if (messageEvent !== undefined) {
await this.handleNewMessage(workspaceUuid, messageEvent)
}
}
async (ctx, msg) => {
const workspaceUuid = msg.workspace
const tx = msg.value
// Check for new channel creation
if (isNewChannelTx(tx)) {
await this.handleNewChannelTx(workspaceUuid, tx)
return
}
// Check for message events
const messageEvent = toMessageEvent(tx)
if (messageEvent !== undefined) {
await this.handleNewMessage(workspaceUuid, messageEvent)
}
},
{
+4 -8
View File
@@ -54,14 +54,10 @@ async function main (): Promise<void> {
ctx,
QueueTopic.Process,
queue.getClientId(),
async (messages) => {
for (const message of messages) {
const ws = message.workspace
const records = message.value
for (const record of records) {
void messageHandler(record, ws, ctx)
}
}
async (ct, message) => {
const ws = message.workspace
const record = message.value
await messageHandler(record, ws, ctx)
}
)
@@ -91,20 +91,16 @@ export const start = async (): Promise<void> => {
ctx,
QueueTopic.TelegramBot,
queue.getClientId(),
async (messages) => {
for (const message of messages) {
const workspace = message.workspace
const records = message.value
for (const record of records) {
switch (record.type) {
case TelegramQueueMessageType.Notification:
await worker.processNotification(workspace, record, bot)
break
case TelegramQueueMessageType.WorkspaceSubscription:
await worker.processWorkspaceSubscription(workspace, record)
break
}
}
async (ctx, message) => {
const workspace = message.workspace
const record = message.value
switch (record.type) {
case TelegramQueueMessageType.Notification:
await worker.processNotification(workspace, record, bot)
break
case TelegramQueueMessageType.WorkspaceSubscription:
await worker.processWorkspaceSubscription(workspace, record)
break
}
}
)