Skip to content

Commit d3f2830

Browse files
authored
improvement(execution): cut the time between a workflow trigger and its first block (#8434)
* improvement(execution): cut the time between a workflow trigger and its first block - Trigger.dev getJob only retrieves real run ids; caller-chosen ids (schedule_…, workflow-execution:…) go straight to the tag lookup instead of a ~10 s 404 - Execution-log start resolves the snapshot id (cached once referenced) and inserts with ON CONFLICT instead of select + full-state upsert RETURNING * + insert - Preprocessing starts payer attribution and the ban/usage/subscription gate reads as soon as their inputs are known; subscription is only read when needed - Execution core prefetches env + PII policy alongside custom blocks and reuses the webhook job's already-loaded environment - Webhook lookups answer trigger-block deployment from the join they already do; versioned deployment loads serve the materialized cache first - getHighestPrioritySubscription, custom-block rows, and env suspension lookups drop sequential round trips * fix(execution): scope the snapshot FK retry, keep a suspended identity's access read out of the run, tidy tests * improvement(execution): encapsulate early admission reads, plain snapshot identity, explicit suspension paths * improvement(environment): reuse the actor's resolved access when withholding a suspended identity
1 parent 16c7e1a commit d3f2830

29 files changed

Lines changed: 1543 additions & 706 deletions

‎apps/sim/app/api/webhooks/tiktok/route.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -96,11 +96,12 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
9696
const webhooks = await findWebhooksByRoutingKey(envelope.user_openid, requestId, 'tiktok')
9797
let dispatched = 0
9898
let failed = 0
99-
for (const { webhook, workflow } of webhooks) {
99+
for (const { webhook, workflow, triggerBlockDeployed } of webhooks) {
100100
const result = await dispatchResolvedWebhookTarget(webhook, workflow, envelope, request, {
101101
requestId,
102102
receivedAt,
103103
triggerTimestampMs: envelope.create_time * 1000,
104+
triggerBlockDeployed,
104105
})
105106
if (result.outcome === 'queued') dispatched += 1
106107
if (result.outcome === 'failed') failed += 1

‎apps/sim/app/api/webhooks/trigger/[path]/route.ts‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -275,7 +275,11 @@ async function handleWebhookDelivery(
275275
}
276276
const dispatchTargetCount = directWebhooksForPath.length + legacySlackDispatchResults.length
277277

278-
for (const { webhook: foundWebhook, workflow: foundWorkflow } of directWebhooksForPath) {
278+
for (const {
279+
webhook: foundWebhook,
280+
workflow: foundWorkflow,
281+
triggerBlockDeployed,
282+
} of directWebhooksForPath) {
279283
const provider = foundWebhook.provider
280284
if (!provider) {
281285
const missingProviderResponse = NextResponse.json(
@@ -321,6 +325,7 @@ async function handleWebhookDelivery(
321325
path,
322326
receivedAt,
323327
triggerTimestampMs: Number.isFinite(triggerTimestampMs) ? triggerTimestampMs : undefined,
328+
triggerBlockDeployed,
324329
}
325330
)
326331

‎apps/sim/background/quickbooks-webhook-ingress.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -64,13 +64,14 @@ export async function executeQuickBooksWebhookIngress(
6464
const targets = await findWebhooksByRoutingKey(routingKey, payload.requestId, 'quickbooks')
6565
targetCount += targets.length
6666

67-
for (const { webhook, workflow } of targets) {
67+
for (const { webhook, workflow, triggerBlockDeployed } of targets) {
6868
try {
6969
const result = await dispatchResolvedWebhookTarget(webhook, workflow, event, request, {
7070
requestId: payload.requestId,
7171
path: webhook.path ?? undefined,
7272
receivedAt: payload.receivedAt,
7373
triggerTimestampMs: Date.parse(event.time),
74+
triggerBlockDeployed,
7475
})
7576
if (result.outcome === 'queued') processed += 1
7677
else if (result.outcome === 'ignored') ignored += 1

‎apps/sim/background/webhook-execution.ts‎

Lines changed: 23 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ import { SlackExecutionStreamController } from '@/lib/webhooks/slack-execution-s
7272
import { readSlackStreamResponseConfig } from '@/lib/webhooks/slack-stream-config'
7373
import {
7474
executeWorkflowCore,
75+
type PreloadedExecutionEnvironment,
7576
wasExecutionFinalizedByCore,
7677
} from '@/lib/workflows/executor/execution-core'
7778
import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence'
@@ -665,19 +666,21 @@ export async function resolveWebhookExecutionProviderConfig<
665666
options?: WebhookEnvResolutionOptions & {
666667
onEnvironmentSnapshot?: (snapshot: EnvironmentResolutionSnapshot) => void | Promise<void>
667668
actorUserId?: string
669+
/** The same environment load, already started by a caller that had its identities early. */
670+
environment?: Promise<EnvironmentResolutionSnapshot>
668671
}
669672
): Promise<T & { providerConfig: Record<string, unknown> }> {
670673
try {
671674
if (!options) {
672675
return await resolveWebhookRecordProviderConfig(webhookRecord, userId, workspaceId)
673676
}
674677

675-
const { onEnvironmentSnapshot, actorUserId, ...resolutionOptions } = options
678+
const { onEnvironmentSnapshot, actorUserId, environment, ...resolutionOptions } = options
676679
if (onEnvironmentSnapshot && resolutionOptions.envVars === undefined) {
677-
const snapshot =
678-
actorUserId && workspaceId
679-
? await getExecutionEnvironment(userId, actorUserId, workspaceId)
680-
: await getEffectiveEnvironmentSnapshot(userId, workspaceId)
680+
const snapshot = await (environment ??
681+
(actorUserId && workspaceId
682+
? getExecutionEnvironment(userId, actorUserId, workspaceId)
683+
: getEffectiveEnvironmentSnapshot(userId, workspaceId)))
681684
await onEnvironmentSnapshot(snapshot)
682685
resolutionOptions.envVars = {
683686
...snapshot.personalDecrypted,
@@ -838,6 +841,12 @@ async function executeWebhookJobInternal(
838841

839842
try {
840843
return await withResourceOutboundScope({ workspaceId }, async () => {
844+
/**
845+
* The run's environment depends only on identities preprocessing already
846+
* settled, so it loads alongside the workflow state rather than after it.
847+
*/
848+
const environment = getExecutionEnvironment(workflowRecord.userId, actorUserId, workspaceId)
849+
environment.catch(() => {})
841850
const workflowStatePromise = payload.deploymentVersionId
842851
? loadWorkflowDeploymentVersionState(
843852
payload.workflowId,
@@ -883,6 +892,7 @@ async function executeWebhookJobInternal(
883892

884893
const secretScope = { userId: workflowRecord.userId, workspaceId }
885894
let resolvedSecretTraceRegistry = createIncompleteResolvedSecretTraceRegistry(secretScope)
895+
let preloadedEnvironment: PreloadedExecutionEnvironment | undefined
886896
const resolvedWebhookRecord = await resolveWebhookExecutionProviderConfig(
887897
webhookRecord,
888898
payload.provider,
@@ -896,7 +906,14 @@ async function executeWebhookJobInternal(
896906
* selection derived from the workflow owner.
897907
*/
898908
actorUserId,
909+
environment,
899910
onEnvironmentSnapshot: async (secretEnvironment) => {
911+
preloadedEnvironment = {
912+
personalUserId: workflowRecord.userId,
913+
workspaceUserId: actorUserId,
914+
workspaceId,
915+
snapshot: secretEnvironment,
916+
}
900917
try {
901918
resolvedSecretTraceRegistry = await createResolvedSecretTraceRegistry({
902919
personalEncrypted: secretEnvironment.personalEncrypted,
@@ -1151,6 +1168,7 @@ async function executeWebhookJobInternal(
11511168
loggingSession,
11521169
trustedInitialResolvedSecretTraceProvenance:
11531170
resolvedSecretTraceRegistry.exportProvenanceForValue(triggerInput),
1171+
preloadedEnvironment,
11541172
includeFileBase64: false,
11551173
base64MaxBytes: undefined,
11561174
abortSignal: timeoutController.signal,

‎apps/sim/background/workflow-column-execution.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -866,6 +866,7 @@ async function runWorkflowAndWriteTerminal(
866866
triggerType: 'workflow',
867867
checkDeployment: false,
868868
checkRateLimit: false,
869+
includeActorSubscription: true,
869870
skipConcurrencyReservation: true,
870871
logPreprocessingErrors: false,
871872
billingAttribution,

‎apps/sim/executor/handlers/workflow/workflow-handler.test.ts‎

Lines changed: 34 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,7 @@ authInternalMockFns.mockGenerateInternalToken.mockResolvedValue('test-token')
5252

5353
const {
5454
mockExecutorExecute,
55-
mockCreateSnapshot,
55+
mockResolveSnapshot,
5656
mockAdmitCustomBlockChildExecution,
5757
mockTrackChildRun,
5858
mockBuildTraceSpans,
@@ -61,7 +61,7 @@ const {
6161
executorOptions,
6262
} = vi.hoisted(() => ({
6363
mockExecutorExecute: vi.fn(),
64-
mockCreateSnapshot: vi.fn(),
64+
mockResolveSnapshot: vi.fn(),
6565
mockAdmitCustomBlockChildExecution: vi.fn(),
6666
mockTrackChildRun: vi.fn(),
6767
mockBuildTraceSpans: vi.fn(),
@@ -160,7 +160,7 @@ afterAll(() => {
160160
})
161161

162162
vi.mock('@/lib/logs/execution/snapshot/service', () => ({
163-
snapshotService: { createSnapshotWithDeduplication: mockCreateSnapshot },
163+
snapshotService: { resolveSnapshot: mockResolveSnapshot },
164164
}))
165165

166166
vi.mock('@/lib/auth/internal', () => authInternalMock)
@@ -356,7 +356,7 @@ describe('WorkflowBlockHandler', () => {
356356
await expect(handler.execute(ctx, mockBlock, inputs)).rejects.toThrow(
357357
'Child workflow child-workflow-id belongs to a different workspace and cannot be executed'
358358
)
359-
expect(mockCreateSnapshot).not.toHaveBeenCalled()
359+
expect(mockResolveSnapshot).not.toHaveBeenCalled()
360360
expect(mockExecutorExecute).not.toHaveBeenCalled()
361361
expect(mockReadWorkflowDefinitionAsExecutor).toHaveBeenCalledWith(
362362
expect.objectContaining({
@@ -395,7 +395,11 @@ describe('WorkflowBlockHandler', () => {
395395
},
396396
}),
397397
})
398-
mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } })
398+
mockResolveSnapshot.mockResolvedValue({
399+
id: 'snapshot-1',
400+
workflowId: 'workflow-1',
401+
stateHash: 'hash',
402+
})
399403
mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } })
400404

401405
await handler.execute(ctx, mockBlock, inputs)
@@ -463,7 +467,11 @@ describe('WorkflowBlockHandler', () => {
463467
}),
464468
}
465469
})
466-
mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } })
470+
mockResolveSnapshot.mockResolvedValue({
471+
id: 'snapshot-1',
472+
workflowId: 'workflow-1',
473+
stateHash: 'hash',
474+
})
467475
mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } })
468476

469477
await handler.execute(ctx, customBlock, {})
@@ -557,7 +565,11 @@ describe('WorkflowBlockHandler', () => {
557565
}),
558566
}
559567
})
560-
mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } })
568+
mockResolveSnapshot.mockResolvedValue({
569+
id: 'snapshot-1',
570+
workflowId: 'workflow-1',
571+
stateHash: 'hash',
572+
})
561573
mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } })
562574

563575
await handler.execute(ctx, customBlock, {})
@@ -642,7 +654,11 @@ describe('WorkflowBlockHandler', () => {
642654
}),
643655
}
644656
})
645-
mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } })
657+
mockResolveSnapshot.mockResolvedValue({
658+
id: 'snapshot-1',
659+
workflowId: 'workflow-1',
660+
stateHash: 'hash',
661+
})
646662
mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } })
647663

648664
await handler.execute(ctx, customBlock, {})
@@ -706,7 +722,11 @@ describe('WorkflowBlockHandler', () => {
706722
},
707723
}),
708724
})
709-
mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } })
725+
mockResolveSnapshot.mockResolvedValue({
726+
id: 'snapshot-1',
727+
workflowId: 'workflow-1',
728+
stateHash: 'hash',
729+
})
710730
mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } })
711731

712732
await handler.execute(ctx, mockBlock, inputs)
@@ -818,7 +838,11 @@ describe('WorkflowBlockHandler', () => {
818838
}),
819839
}
820840
})
821-
mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } })
841+
mockResolveSnapshot.mockResolvedValue({
842+
id: 'snapshot-1',
843+
workflowId: 'workflow-1',
844+
stateHash: 'hash',
845+
})
822846
mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } })
823847
})
824848

‎apps/sim/executor/handlers/workflow/workflow-handler.ts‎

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -472,11 +472,9 @@ export class WorkflowBlockHandler implements BlockHandler {
472472
}
473473
}
474474

475-
const childSnapshotResult = await snapshotService.createSnapshotWithDeduplication(
476-
workflowId,
477-
childWorkflow.workflowState
478-
)
479-
childWorkflowSnapshotId = childSnapshotResult.snapshot.id
475+
childWorkflowSnapshotId = (
476+
await snapshotService.resolveSnapshot(workflowId, childWorkflow.workflowState)
477+
).id
480478

481479
const childDepth = (ctx.childWorkflowContext?.depth ?? 0) + 1
482480
const withinSseChildDepth = childDepth <= DEFAULTS.MAX_SSE_CHILD_DEPTH
Lines changed: 109 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,109 @@
1+
/** Subscription selection against real PostgreSQL: tier priority, scope tie-break, and entitlement. */
2+
import { db } from '@sim/db'
3+
import { member, organization, subscription, user } from '@sim/db/schema'
4+
import { generateId } from '@sim/utils/id'
5+
import { inArray } from 'drizzle-orm'
6+
import { afterAll, describe, expect, it } from 'vitest'
7+
import { getHighestPrioritySubscription } from '@/lib/billing/core/plan'
8+
9+
const userIds: string[] = []
10+
const organizationIds: string[] = []
11+
12+
async function createUser(): Promise<string> {
13+
const id = `plan-user-${generateId()}`
14+
const now = new Date()
15+
await db.insert(user).values({
16+
id,
17+
name: 'Plan Test',
18+
email: `${id}@plan.test`,
19+
emailVerified: true,
20+
createdAt: now,
21+
updatedAt: now,
22+
})
23+
userIds.push(id)
24+
return id
25+
}
26+
27+
async function createOrganizationWithMember(userId: string): Promise<string> {
28+
const id = `plan-org-${generateId()}`
29+
await db.insert(organization).values({ id, name: 'Plan Org', slug: id })
30+
await db.insert(member).values({ id: generateId(), userId, organizationId: id, role: 'member' })
31+
organizationIds.push(id)
32+
return id
33+
}
34+
35+
async function createSubscription(
36+
referenceId: string,
37+
plan: 'pro' | 'team' | 'enterprise',
38+
status = 'active'
39+
): Promise<string> {
40+
const id = generateId()
41+
await db.insert(subscription).values({
42+
id,
43+
plan,
44+
referenceId,
45+
status,
46+
...(plan === 'enterprise' ? { metadata: { workspaces: 'unlimited' } } : {}),
47+
})
48+
return id
49+
}
50+
51+
afterAll(async () => {
52+
await db
53+
.delete(subscription)
54+
.where(inArray(subscription.referenceId, [...userIds, ...organizationIds]))
55+
if (organizationIds.length > 0) {
56+
await db.delete(organization).where(inArray(organization.id, organizationIds))
57+
}
58+
if (userIds.length > 0) await db.delete(user).where(inArray(user.id, userIds))
59+
})
60+
61+
describe('getHighestPrioritySubscription', () => {
62+
it('returns null when the user has no entitled subscription', async () => {
63+
const userId = await createUser()
64+
await createOrganizationWithMember(userId)
65+
66+
expect(await getHighestPrioritySubscription(userId)).toBeNull()
67+
})
68+
69+
it('prefers the higher tier regardless of scope', async () => {
70+
const userId = await createUser()
71+
const organizationId = await createOrganizationWithMember(userId)
72+
await createSubscription(userId, 'pro')
73+
const enterpriseId = await createSubscription(organizationId, 'enterprise')
74+
75+
expect((await getHighestPrioritySubscription(userId))?.id).toBe(enterpriseId)
76+
77+
const personalTeamUserId = await createUser()
78+
const proOrganizationId = await createOrganizationWithMember(personalTeamUserId)
79+
await createSubscription(proOrganizationId, 'pro')
80+
const personalTeamId = await createSubscription(personalTeamUserId, 'team')
81+
82+
expect((await getHighestPrioritySubscription(personalTeamUserId))?.id).toBe(personalTeamId)
83+
})
84+
85+
it('prefers the organization subscription over a personal one of the same tier', async () => {
86+
const userId = await createUser()
87+
const organizationId = await createOrganizationWithMember(userId)
88+
await createSubscription(userId, 'team')
89+
const organizationTeamId = await createSubscription(organizationId, 'team')
90+
91+
expect((await getHighestPrioritySubscription(userId))?.id).toBe(organizationTeamId)
92+
})
93+
94+
it('ignores subscriptions that are not in an entitled status', async () => {
95+
const userId = await createUser()
96+
const organizationId = await createOrganizationWithMember(userId)
97+
await createSubscription(organizationId, 'enterprise', 'canceled')
98+
const personalProId = await createSubscription(userId, 'pro', 'past_due')
99+
100+
expect((await getHighestPrioritySubscription(userId))?.id).toBe(personalProId)
101+
})
102+
103+
it('reads a personal subscription for a user with no organization', async () => {
104+
const userId = await createUser()
105+
const personalProId = await createSubscription(userId, 'pro')
106+
107+
expect((await getHighestPrioritySubscription(userId))?.id).toBe(personalProId)
108+
})
109+
})

0 commit comments

Comments
 (0)