Skip to content

Commit 055ff9c

Browse files
authored
improvement(outbox): load handler modules only for the event types a run will process (#8731)
* improvement(outbox): load handler modules only for the event types a run will process * refactor(outbox): take only lazy handler groups; resolve eligible types in one place * fix(outbox): fail the run after maintenance when a handler module could not load * improvement(outbox): import every registry event type from a light module; prove failed loads in Postgres * improvement(outbox): source enterprise and document-processing event types from dependency-free modules
1 parent 1741ff9 commit 055ff9c

110 files changed

Lines changed: 841 additions & 351 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎apps/sim/app/api/cron/reconcile-billing-seats/route.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@ import { toError } from '@sim/utils/errors'
33
import { type NextRequest, NextResponse } from 'next/server'
44
import { verifyCronAuth } from '@/lib/auth/internal'
55
import { reconcileTeamSeatDrift } from '@/lib/billing/organizations/seat-drift'
6-
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-handlers'
6+
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
77
import { findDeadLetteredEvents } from '@/lib/core/outbox/service'
88
import { generateRequestId } from '@/lib/core/utils/request'
99
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'

‎apps/sim/app/api/v1/admin/outbox/[id]/requeue/route.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,11 +6,11 @@ import { and, eq, sql } from 'drizzle-orm'
66
import { NextResponse } from 'next/server'
77
import { adminV1RequeueOutboxEventContract } from '@/lib/api/contracts/v1/admin'
88
import { getValidationErrorMessage, parseRequest } from '@/lib/api/server'
9+
import { enterpriseMetadataSyncPayloadSchema } from '@/lib/billing/enterprise-outbox'
910
import {
1011
ENTERPRISE_METADATA_SYNC_EVENT_TYPE,
1112
ENTERPRISE_PROVISION_EVENT_TYPE,
12-
enterpriseMetadataSyncPayloadSchema,
13-
} from '@/lib/billing/enterprise-outbox'
13+
} from '@/lib/billing/enterprise-outbox-events'
1414
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
1515
import { withAdminAuthParams } from '@/app/api/v1/admin/middleware'
1616

‎apps/sim/app/api/v1/admin/subscriptions/[id]/route.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ import {
3333
} from '@/lib/api/contracts/v1/admin'
3434
import { parseRequest } from '@/lib/api/server'
3535
import { requireStripeClient } from '@/lib/billing/stripe-client'
36-
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-handlers'
36+
import { OUTBOX_EVENT_TYPES } from '@/lib/billing/webhooks/outbox-events'
3737
import { enqueueOutboxEvent } from '@/lib/core/outbox/service'
3838
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
3939
import { withAdminAuthParams } from '@/app/api/v1/admin/middleware'
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,2 +1,3 @@
11
export const PERMISSION_ACCESS_REQUEST_CREATED_EVENT = 'permission-access-request.created'
22
export const PERMISSION_ACCESS_REQUEST_DECIDED_EVENT = 'permission-access-request.decided'
3+
export const PERMISSION_ACCESS_REQUEST_NOTIFY_ADMIN_EVENT = 'permission-access-request.notify-admin'

‎apps/sim/ee/access-requests/lib/notifications.ts‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,12 @@ import { getMyAccessRequestHref } from '@/ee/access-requests/lib/navigation'
1818
import {
1919
PERMISSION_ACCESS_REQUEST_CREATED_EVENT,
2020
PERMISSION_ACCESS_REQUEST_DECIDED_EVENT,
21+
PERMISSION_ACCESS_REQUEST_NOTIFY_ADMIN_EVENT,
2122
} from '@/ee/access-requests/lib/notification-events'
2223
import { isAccessRequestEnabled } from '@/ee/access-requests/lib/settings'
2324

2425
const logger = createLogger('PermissionAccessRequestNotifications')
2526
const ADMIN_RECIPIENT_PAGE_SIZE = 50
26-
const ADMIN_NOTIFICATION_EVENT = 'permission-access-request.notify-admin'
2727
const notificationPayloadSchema = z.object({ requestId: z.string().min(1).max(256) }).strict()
2828
const createdPayloadSchema = notificationPayloadSchema.extend({
2929
afterMemberId: z.string().min(1).max(256).optional(),
@@ -176,8 +176,8 @@ export const permissionAccessRequestOutboxHandlers = {
176176
.insert(outboxEvent)
177177
.values(
178178
recipients.map((recipient) => ({
179-
id: `${ADMIN_NOTIFICATION_EVENT}:${requestId}:${recipient.userId}`,
180-
eventType: ADMIN_NOTIFICATION_EVENT,
179+
id: `${PERMISSION_ACCESS_REQUEST_NOTIFY_ADMIN_EVENT}:${requestId}:${recipient.userId}`,
180+
eventType: PERMISSION_ACCESS_REQUEST_NOTIFY_ADMIN_EVENT,
181181
payload: { requestId, recipientUserId: recipient.userId },
182182
}))
183183
)
@@ -187,7 +187,7 @@ export const permissionAccessRequestOutboxHandlers = {
187187
return continueOutboxHandler('Continue access request administrator notifications')
188188
}
189189
},
190-
[ADMIN_NOTIFICATION_EVENT]: async (rawPayload, context) => {
190+
[PERMISSION_ACCESS_REQUEST_NOTIFY_ADMIN_EVENT]: async (rawPayload, context) => {
191191
const { requestId, recipientUserId } = adminPayloadSchema.parse(rawPayload)
192192
const request = await loadRequest(requestId)
193193
if (

‎apps/sim/ee/workspace-forking/application/admit-sync.ts‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,11 @@ import { generateId } from '@sim/utils/id'
22
import { truncate } from '@sim/utils/string'
33
import { enqueueOutboxEvent } from '@/lib/core/outbox/service'
44
import type { DbOrTx } from '@/lib/db/types'
5+
import {
6+
WORKSPACE_MCP_CHANGED_EVENT,
7+
WORKSPACE_OPERATION_OBSERVE_EVENT,
8+
WORKSPACE_WORKFLOWS_CHANGED_EVENT,
9+
} from '@/lib/workspaces/operations/outbox-events'
510
import {
611
insertWorkspaceOperationReceipt,
712
type WorkspaceOperationReport,
@@ -100,7 +105,7 @@ export async function admitForkSync(
100105
}
101106
if (params.mcpAttachmentServerIds.length)
102107
report.effectEventIds!.push(
103-
await enqueueOutboxEvent(tx, 'workspace.mcp.changed', {
108+
await enqueueOutboxEvent(tx, WORKSPACE_MCP_CHANGED_EVENT, {
104109
serverIds: params.mcpAttachmentServerIds,
105110
})
106111
)
@@ -117,10 +122,10 @@ export async function admitForkSync(
117122
? 'completed_with_warnings'
118123
: 'completed'
119124
}
120-
await enqueueOutboxEvent(tx, 'workspace.workflows.changed', {
125+
await enqueueOutboxEvent(tx, WORKSPACE_WORKFLOWS_CHANGED_EVENT, {
121126
workspaceId: params.targetWorkspaceId,
122127
})
123-
await enqueueOutboxEvent(tx, 'workspace.operation.observe', {
128+
await enqueueOutboxEvent(tx, WORKSPACE_OPERATION_OBSERVE_EVENT, {
124129
workspaceId: report.workspaceId,
125130
operationId: report.operationId,
126131
})
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
export const FORK_CONTENT_COPY_EVENT = 'workspace.fork.content.copy'

‎apps/sim/ee/workspace-forking/application/content-outbox.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import {
1111
} from '@/lib/core/outbox/service'
1212
import type { DbOrTx } from '@/lib/db/types'
1313
import type { WorkspaceOperationReport } from '@/lib/workspaces/operations/receipts'
14+
import { FORK_CONTENT_COPY_EVENT } from '@/ee/workspace-forking/application/content-outbox-event'
1415
import {
1516
type ForkContentCopyPayload,
1617
runForkContentCopy,
@@ -153,11 +154,11 @@ export async function enqueueDurableForkContent(
153154
})
154155
if (Buffer.byteLength(JSON.stringify(payload)) > 8 * 1024 * 1024)
155156
throw new OrchestrationError('payload_too_large', 'Fork background work exceeds 8 MiB')
156-
return enqueueOutboxEvent(tx, 'workspace.fork.content.copy', payload)
157+
return enqueueOutboxEvent(tx, FORK_CONTENT_COPY_EVENT, payload)
157158
}
158159

159160
export const forkContentOutboxHandlers = {
160-
'workspace.fork.content.copy': withOutboxHandlerTimeout(async (raw, context) => {
161+
[FORK_CONTENT_COPY_EVENT]: withOutboxHandlerTimeout(async (raw, context) => {
161162
const payload = contentPayloadSchema.parse(raw)
162163
const [receipt] = await db
163164
.select({ report: workspaceOperationReceipt.report })

‎apps/sim/ee/workspace-forking/lib/create-fork.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import {
1616
collectReferencedFileFolderPaths,
1717
} from '@/lib/workflows/references/reference-scan'
1818
import type { ForkRemapKind } from '@/lib/workflows/references/remap-references'
19+
import { WORKSPACE_OPERATION_OBSERVE_EVENT } from '@/lib/workspaces/operations/outbox-events'
1920
import {
2021
findWorkspaceOperationReceipt,
2122
insertWorkspaceOperationReceipt,
@@ -625,7 +626,7 @@ export async function createFork(params: CreateForkParams): Promise<CreateForkRe
625626
requestId: admission.requestId,
626627
})
627628
}
628-
await enqueueOutboxEvent(tx, 'workspace.operation.observe', {
629+
await enqueueOutboxEvent(tx, WORKSPACE_OPERATION_OBSERVE_EVENT, {
629630
workspaceId: report.workspaceId,
630631
operationId: report.operationId,
631632
})

‎apps/sim/lib/admin/dashboard-credit-grant.test.ts‎

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,6 @@ vi.mock('@/lib/billing/enterprise-provisioning', () => ({
5050
getLatestEnterpriseProvisionings: vi.fn(async () => new Map()),
5151
}))
5252
vi.mock('@/lib/billing/enterprise-outbox', () => ({
53-
ENTERPRISE_METADATA_SYNC_EVENT_TYPE: 'stripe.sync-enterprise-metadata',
5453
resolveEnterpriseMetadataIntent: vi.fn(),
5554
}))
5655
vi.mock('@/lib/core/outbox/service', () => outboxServiceMock)

0 commit comments

Comments
 (0)