diff --git a/apps/sim/app/api/webhooks/tiktok/route.ts b/apps/sim/app/api/webhooks/tiktok/route.ts index 4a70ec8cf45..e3732bfceb5 100644 --- a/apps/sim/app/api/webhooks/tiktok/route.ts +++ b/apps/sim/app/api/webhooks/tiktok/route.ts @@ -96,11 +96,12 @@ export const POST = withRouteHandler(async (request: NextRequest) => { const webhooks = await findWebhooksByRoutingKey(envelope.user_openid, requestId, 'tiktok') let dispatched = 0 let failed = 0 - for (const { webhook, workflow } of webhooks) { + for (const { webhook, workflow, triggerBlockDeployed } of webhooks) { const result = await dispatchResolvedWebhookTarget(webhook, workflow, envelope, request, { requestId, receivedAt, triggerTimestampMs: envelope.create_time * 1000, + triggerBlockDeployed, }) if (result.outcome === 'queued') dispatched += 1 if (result.outcome === 'failed') failed += 1 diff --git a/apps/sim/app/api/webhooks/trigger/[path]/route.ts b/apps/sim/app/api/webhooks/trigger/[path]/route.ts index 41d0d4b8179..32082e4466f 100644 --- a/apps/sim/app/api/webhooks/trigger/[path]/route.ts +++ b/apps/sim/app/api/webhooks/trigger/[path]/route.ts @@ -275,7 +275,11 @@ async function handleWebhookDelivery( } const dispatchTargetCount = directWebhooksForPath.length + legacySlackDispatchResults.length - for (const { webhook: foundWebhook, workflow: foundWorkflow } of directWebhooksForPath) { + for (const { + webhook: foundWebhook, + workflow: foundWorkflow, + triggerBlockDeployed, + } of directWebhooksForPath) { const provider = foundWebhook.provider if (!provider) { const missingProviderResponse = NextResponse.json( @@ -321,6 +325,7 @@ async function handleWebhookDelivery( path, receivedAt, triggerTimestampMs: Number.isFinite(triggerTimestampMs) ? triggerTimestampMs : undefined, + triggerBlockDeployed, } ) diff --git a/apps/sim/background/quickbooks-webhook-ingress.ts b/apps/sim/background/quickbooks-webhook-ingress.ts index 3ab7ed54623..7e7436c49ba 100644 --- a/apps/sim/background/quickbooks-webhook-ingress.ts +++ b/apps/sim/background/quickbooks-webhook-ingress.ts @@ -64,13 +64,14 @@ export async function executeQuickBooksWebhookIngress( const targets = await findWebhooksByRoutingKey(routingKey, payload.requestId, 'quickbooks') targetCount += targets.length - for (const { webhook, workflow } of targets) { + for (const { webhook, workflow, triggerBlockDeployed } of targets) { try { const result = await dispatchResolvedWebhookTarget(webhook, workflow, event, request, { requestId: payload.requestId, path: webhook.path ?? undefined, receivedAt: payload.receivedAt, triggerTimestampMs: Date.parse(event.time), + triggerBlockDeployed, }) if (result.outcome === 'queued') processed += 1 else if (result.outcome === 'ignored') ignored += 1 diff --git a/apps/sim/background/webhook-execution.ts b/apps/sim/background/webhook-execution.ts index 95f7b1ba9a8..827cd55e684 100644 --- a/apps/sim/background/webhook-execution.ts +++ b/apps/sim/background/webhook-execution.ts @@ -72,6 +72,7 @@ import { SlackExecutionStreamController } from '@/lib/webhooks/slack-execution-s import { readSlackStreamResponseConfig } from '@/lib/webhooks/slack-stream-config' import { executeWorkflowCore, + type PreloadedExecutionEnvironment, wasExecutionFinalizedByCore, } from '@/lib/workflows/executor/execution-core' import { handlePostExecutionPauseState } from '@/lib/workflows/executor/pause-persistence' @@ -665,6 +666,8 @@ export async function resolveWebhookExecutionProviderConfig< options?: WebhookEnvResolutionOptions & { onEnvironmentSnapshot?: (snapshot: EnvironmentResolutionSnapshot) => void | Promise actorUserId?: string + /** The same environment load, already started by a caller that had its identities early. */ + environment?: Promise } ): Promise }> { try { @@ -672,12 +675,12 @@ export async function resolveWebhookExecutionProviderConfig< return await resolveWebhookRecordProviderConfig(webhookRecord, userId, workspaceId) } - const { onEnvironmentSnapshot, actorUserId, ...resolutionOptions } = options + const { onEnvironmentSnapshot, actorUserId, environment, ...resolutionOptions } = options if (onEnvironmentSnapshot && resolutionOptions.envVars === undefined) { - const snapshot = - actorUserId && workspaceId - ? await getExecutionEnvironment(userId, actorUserId, workspaceId) - : await getEffectiveEnvironmentSnapshot(userId, workspaceId) + const snapshot = await (environment ?? + (actorUserId && workspaceId + ? getExecutionEnvironment(userId, actorUserId, workspaceId) + : getEffectiveEnvironmentSnapshot(userId, workspaceId))) await onEnvironmentSnapshot(snapshot) resolutionOptions.envVars = { ...snapshot.personalDecrypted, @@ -838,6 +841,12 @@ async function executeWebhookJobInternal( try { return await withResourceOutboundScope({ workspaceId }, async () => { + /** + * The run's environment depends only on identities preprocessing already + * settled, so it loads alongside the workflow state rather than after it. + */ + const environment = getExecutionEnvironment(workflowRecord.userId, actorUserId, workspaceId) + environment.catch(() => {}) const workflowStatePromise = payload.deploymentVersionId ? loadWorkflowDeploymentVersionState( payload.workflowId, @@ -883,6 +892,7 @@ async function executeWebhookJobInternal( const secretScope = { userId: workflowRecord.userId, workspaceId } let resolvedSecretTraceRegistry = createIncompleteResolvedSecretTraceRegistry(secretScope) + let preloadedEnvironment: PreloadedExecutionEnvironment | undefined const resolvedWebhookRecord = await resolveWebhookExecutionProviderConfig( webhookRecord, payload.provider, @@ -896,7 +906,14 @@ async function executeWebhookJobInternal( * selection derived from the workflow owner. */ actorUserId, + environment, onEnvironmentSnapshot: async (secretEnvironment) => { + preloadedEnvironment = { + personalUserId: workflowRecord.userId, + workspaceUserId: actorUserId, + workspaceId, + snapshot: secretEnvironment, + } try { resolvedSecretTraceRegistry = await createResolvedSecretTraceRegistry({ personalEncrypted: secretEnvironment.personalEncrypted, @@ -1151,6 +1168,7 @@ async function executeWebhookJobInternal( loggingSession, trustedInitialResolvedSecretTraceProvenance: resolvedSecretTraceRegistry.exportProvenanceForValue(triggerInput), + preloadedEnvironment, includeFileBase64: false, base64MaxBytes: undefined, abortSignal: timeoutController.signal, diff --git a/apps/sim/background/workflow-column-execution.ts b/apps/sim/background/workflow-column-execution.ts index d9601d2f223..cbd43a3b95f 100644 --- a/apps/sim/background/workflow-column-execution.ts +++ b/apps/sim/background/workflow-column-execution.ts @@ -866,6 +866,7 @@ async function runWorkflowAndWriteTerminal( triggerType: 'workflow', checkDeployment: false, checkRateLimit: false, + includeActorSubscription: true, skipConcurrencyReservation: true, logPreprocessingErrors: false, billingAttribution, diff --git a/apps/sim/executor/handlers/workflow/workflow-handler.test.ts b/apps/sim/executor/handlers/workflow/workflow-handler.test.ts index c9507937368..665bd926eee 100644 --- a/apps/sim/executor/handlers/workflow/workflow-handler.test.ts +++ b/apps/sim/executor/handlers/workflow/workflow-handler.test.ts @@ -52,7 +52,7 @@ authInternalMockFns.mockGenerateInternalToken.mockResolvedValue('test-token') const { mockExecutorExecute, - mockCreateSnapshot, + mockResolveSnapshot, mockAdmitCustomBlockChildExecution, mockTrackChildRun, mockBuildTraceSpans, @@ -61,7 +61,7 @@ const { executorOptions, } = vi.hoisted(() => ({ mockExecutorExecute: vi.fn(), - mockCreateSnapshot: vi.fn(), + mockResolveSnapshot: vi.fn(), mockAdmitCustomBlockChildExecution: vi.fn(), mockTrackChildRun: vi.fn(), mockBuildTraceSpans: vi.fn(), @@ -160,7 +160,7 @@ afterAll(() => { }) vi.mock('@/lib/logs/execution/snapshot/service', () => ({ - snapshotService: { createSnapshotWithDeduplication: mockCreateSnapshot }, + snapshotService: { resolveSnapshot: mockResolveSnapshot }, })) vi.mock('@/lib/auth/internal', () => authInternalMock) @@ -356,7 +356,7 @@ describe('WorkflowBlockHandler', () => { await expect(handler.execute(ctx, mockBlock, inputs)).rejects.toThrow( 'Child workflow child-workflow-id belongs to a different workspace and cannot be executed' ) - expect(mockCreateSnapshot).not.toHaveBeenCalled() + expect(mockResolveSnapshot).not.toHaveBeenCalled() expect(mockExecutorExecute).not.toHaveBeenCalled() expect(mockReadWorkflowDefinitionAsExecutor).toHaveBeenCalledWith( expect.objectContaining({ @@ -395,7 +395,11 @@ describe('WorkflowBlockHandler', () => { }, }), }) - mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } }) + mockResolveSnapshot.mockResolvedValue({ + id: 'snapshot-1', + workflowId: 'workflow-1', + stateHash: 'hash', + }) mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } }) await handler.execute(ctx, mockBlock, inputs) @@ -463,7 +467,11 @@ describe('WorkflowBlockHandler', () => { }), } }) - mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } }) + mockResolveSnapshot.mockResolvedValue({ + id: 'snapshot-1', + workflowId: 'workflow-1', + stateHash: 'hash', + }) mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } }) await handler.execute(ctx, customBlock, {}) @@ -557,7 +565,11 @@ describe('WorkflowBlockHandler', () => { }), } }) - mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } }) + mockResolveSnapshot.mockResolvedValue({ + id: 'snapshot-1', + workflowId: 'workflow-1', + stateHash: 'hash', + }) mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } }) await handler.execute(ctx, customBlock, {}) @@ -642,7 +654,11 @@ describe('WorkflowBlockHandler', () => { }), } }) - mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } }) + mockResolveSnapshot.mockResolvedValue({ + id: 'snapshot-1', + workflowId: 'workflow-1', + stateHash: 'hash', + }) mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } }) await handler.execute(ctx, customBlock, {}) @@ -706,7 +722,11 @@ describe('WorkflowBlockHandler', () => { }, }), }) - mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } }) + mockResolveSnapshot.mockResolvedValue({ + id: 'snapshot-1', + workflowId: 'workflow-1', + stateHash: 'hash', + }) mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } }) await handler.execute(ctx, mockBlock, inputs) @@ -818,7 +838,11 @@ describe('WorkflowBlockHandler', () => { }), } }) - mockCreateSnapshot.mockResolvedValue({ snapshot: { id: 'snapshot-1' } }) + mockResolveSnapshot.mockResolvedValue({ + id: 'snapshot-1', + workflowId: 'workflow-1', + stateHash: 'hash', + }) mockExecutorExecute.mockResolvedValue({ success: true, output: { data: 'ok' } }) }) diff --git a/apps/sim/executor/handlers/workflow/workflow-handler.ts b/apps/sim/executor/handlers/workflow/workflow-handler.ts index 961606ac8a1..51fdc0f93e7 100644 --- a/apps/sim/executor/handlers/workflow/workflow-handler.ts +++ b/apps/sim/executor/handlers/workflow/workflow-handler.ts @@ -472,11 +472,9 @@ export class WorkflowBlockHandler implements BlockHandler { } } - const childSnapshotResult = await snapshotService.createSnapshotWithDeduplication( - workflowId, - childWorkflow.workflowState - ) - childWorkflowSnapshotId = childSnapshotResult.snapshot.id + childWorkflowSnapshotId = ( + await snapshotService.resolveSnapshot(workflowId, childWorkflow.workflowState) + ).id const childDepth = (ctx.childWorkflowContext?.depth ?? 0) + 1 const withinSseChildDepth = childDepth <= DEFAULTS.MAX_SSE_CHILD_DEPTH diff --git a/apps/sim/lib/billing/core/plan.integration.ts b/apps/sim/lib/billing/core/plan.integration.ts new file mode 100644 index 00000000000..7df0202b7c5 --- /dev/null +++ b/apps/sim/lib/billing/core/plan.integration.ts @@ -0,0 +1,109 @@ +/** Subscription selection against real PostgreSQL: tier priority, scope tie-break, and entitlement. */ +import { db } from '@sim/db' +import { member, organization, subscription, user } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { inArray } from 'drizzle-orm' +import { afterAll, describe, expect, it } from 'vitest' +import { getHighestPrioritySubscription } from '@/lib/billing/core/plan' + +const userIds: string[] = [] +const organizationIds: string[] = [] + +async function createUser(): Promise { + const id = `plan-user-${generateId()}` + const now = new Date() + await db.insert(user).values({ + id, + name: 'Plan Test', + email: `${id}@plan.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + userIds.push(id) + return id +} + +async function createOrganizationWithMember(userId: string): Promise { + const id = `plan-org-${generateId()}` + await db.insert(organization).values({ id, name: 'Plan Org', slug: id }) + await db.insert(member).values({ id: generateId(), userId, organizationId: id, role: 'member' }) + organizationIds.push(id) + return id +} + +async function createSubscription( + referenceId: string, + plan: 'pro' | 'team' | 'enterprise', + status = 'active' +): Promise { + const id = generateId() + await db.insert(subscription).values({ + id, + plan, + referenceId, + status, + ...(plan === 'enterprise' ? { metadata: { workspaces: 'unlimited' } } : {}), + }) + return id +} + +afterAll(async () => { + await db + .delete(subscription) + .where(inArray(subscription.referenceId, [...userIds, ...organizationIds])) + if (organizationIds.length > 0) { + await db.delete(organization).where(inArray(organization.id, organizationIds)) + } + if (userIds.length > 0) await db.delete(user).where(inArray(user.id, userIds)) +}) + +describe('getHighestPrioritySubscription', () => { + it('returns null when the user has no entitled subscription', async () => { + const userId = await createUser() + await createOrganizationWithMember(userId) + + expect(await getHighestPrioritySubscription(userId)).toBeNull() + }) + + it('prefers the higher tier regardless of scope', async () => { + const userId = await createUser() + const organizationId = await createOrganizationWithMember(userId) + await createSubscription(userId, 'pro') + const enterpriseId = await createSubscription(organizationId, 'enterprise') + + expect((await getHighestPrioritySubscription(userId))?.id).toBe(enterpriseId) + + const personalTeamUserId = await createUser() + const proOrganizationId = await createOrganizationWithMember(personalTeamUserId) + await createSubscription(proOrganizationId, 'pro') + const personalTeamId = await createSubscription(personalTeamUserId, 'team') + + expect((await getHighestPrioritySubscription(personalTeamUserId))?.id).toBe(personalTeamId) + }) + + it('prefers the organization subscription over a personal one of the same tier', async () => { + const userId = await createUser() + const organizationId = await createOrganizationWithMember(userId) + await createSubscription(userId, 'team') + const organizationTeamId = await createSubscription(organizationId, 'team') + + expect((await getHighestPrioritySubscription(userId))?.id).toBe(organizationTeamId) + }) + + it('ignores subscriptions that are not in an entitled status', async () => { + const userId = await createUser() + const organizationId = await createOrganizationWithMember(userId) + await createSubscription(organizationId, 'enterprise', 'canceled') + const personalProId = await createSubscription(userId, 'pro', 'past_due') + + expect((await getHighestPrioritySubscription(userId))?.id).toBe(personalProId) + }) + + it('reads a personal subscription for a user with no organization', async () => { + const userId = await createUser() + const personalProId = await createSubscription(userId, 'pro') + + expect((await getHighestPrioritySubscription(userId))?.id).toBe(personalProId) + }) +}) diff --git a/apps/sim/lib/billing/core/plan.test.ts b/apps/sim/lib/billing/core/plan.test.ts deleted file mode 100644 index b26f3a66239..00000000000 --- a/apps/sim/lib/billing/core/plan.test.ts +++ /dev/null @@ -1,112 +0,0 @@ -import { member, organization, subscription } from '@sim/db/schema' -import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing' -import { billingSubscriptionUtilsMock } from '@sim/testing/mocks/billing-subscription-utils.mock' -import { afterAll, beforeEach, describe, expect, it, vi } from 'vitest' - -vi.mock('@/lib/billing/subscriptions/utils', () => billingSubscriptionUtilsMock) - -import { getHighestPrioritySubscription } from '@/lib/billing/core/plan' - -/** - * `getHighestPrioritySubscription` issues up to four queries keyed by table: - * - `subscription` for the user's personal subs (parallelized with members) - * - `member` for the user's org memberships (parallelized with subs) - * - `organization` for the org-existence follow-up - * - `subscription` again for the org-scoped subs follow-up - * - * Results are routed by the table object passed to `.from()` via - * `queueTableRows` (FIFO per table: first `subscription` read = personal, - * second = org). `dbChainMockFns.from` call args record which tables were - * queried so we can assert the parallelized pair both run and that follow-ups - * are skipped when appropriate. - */ -function fromTables(): unknown[] { - return dbChainMockFns.from.mock.calls.map(([table]) => table) -} - -interface SubRow { - id: string - referenceId: string - plan: string - status: string -} - -function personalPro(userId: string): SubRow { - return { id: 'sub-personal-pro', referenceId: userId, plan: 'pro', status: 'active' } -} - -function orgEnterprise(orgId: string): SubRow { - return { id: 'sub-org-enterprise', referenceId: orgId, plan: 'enterprise', status: 'active' } -} - -describe('getHighestPrioritySubscription', () => { - beforeEach(() => { - resetDbChainMock() - }) - - afterAll(() => { - resetDbChainMock() - }) - - it('picks the org Enterprise sub over a personal Pro sub (priority order)', async () => { - queueTableRows(subscription, [personalPro('user-1')]) // personalSubs query - queueTableRows(member, [{ organizationId: 'org-1' }]) // memberships query - queueTableRows(organization, [{ id: 'org-1' }]) // org-existence query - queueTableRows(subscription, [orgEnterprise('org-1')]) // org-subscriptions query - - const result = await getHighestPrioritySubscription('user-1') - - expect(result).not.toBeNull() - expect(result?.id).toBe('sub-org-enterprise') - expect(result?.plan).toBe('enterprise') - }) - - it('selection is deterministic regardless of which parallelized query resolves first', async () => { - queueTableRows(subscription, [personalPro('user-1')]) - queueTableRows(member, [{ organizationId: 'org-1' }]) - queueTableRows(organization, [{ id: 'org-1' }]) - queueTableRows(subscription, [orgEnterprise('org-1')]) - - const result = await getHighestPrioritySubscription('user-1') - - expect(result?.id).toBe('sub-org-enterprise') - }) - - it('returns the personal sub and skips org follow-ups when there are no memberships', async () => { - queueTableRows(subscription, [personalPro('user-1')]) - queueTableRows(member, []) - - const result = await getHighestPrioritySubscription('user-1') - - expect(result?.id).toBe('sub-personal-pro') - expect(result?.plan).toBe('pro') - // org-existence + org-subscription follow-ups are NOT issued. - expect(fromTables()).not.toContain(organization) - expect(fromTables().filter((t) => t === subscription)).toHaveLength(1) - }) - - it('excludes orphaned org memberships whose organization row no longer exists', async () => { - queueTableRows(subscription, []) - queueTableRows(member, [{ organizationId: 'ghost-org' }]) // membership points at a deleted org - queueTableRows(organization, []) - - const result = await getHighestPrioritySubscription('user-1') - - // Org subs are never fetched (no valid org ids) -> falls back to null. - expect(result).toBeNull() - expect(fromTables()).toContain(organization) - // Only the initial personal-subs read on `subscription`; org-subs query skipped. - expect(fromTables().filter((t) => t === subscription)).toHaveLength(1) - }) - - it('falls back to the personal sub when the only org is orphaned', async () => { - queueTableRows(subscription, [personalPro('user-1')]) - queueTableRows(member, [{ organizationId: 'ghost-org' }]) - queueTableRows(organization, []) - - const result = await getHighestPrioritySubscription('user-1') - - expect(result?.id).toBe('sub-personal-pro') - expect(fromTables().filter((t) => t === subscription)).toHaveLength(1) - }) -}) diff --git a/apps/sim/lib/billing/core/plan.ts b/apps/sim/lib/billing/core/plan.ts index 991b1ff822f..c2ec8357284 100644 --- a/apps/sim/lib/billing/core/plan.ts +++ b/apps/sim/lib/billing/core/plan.ts @@ -1,7 +1,7 @@ import { db } from '@sim/db' import { member, organization, subscription } from '@sim/db/schema' import { createLogger } from '@sim/logger' -import { and, eq, inArray } from 'drizzle-orm' +import { and, eq, getTableColumns, inArray } from 'drizzle-orm' import { checkEnterprisePlan, checkProPlan, @@ -82,47 +82,21 @@ export async function getHighestPrioritySubscription( ) { const { onError = 'return-null', executor = db } = options try { - const [personalSubs, memberships] = await Promise.all([ + const entitled = inArray(subscription.status, ENTITLED_SUBSCRIPTION_STATUSES) + const [personalSubs, orgSubs] = await Promise.all([ executor .select() .from(subscription) - .where( - and( - eq(subscription.referenceId, userId), - inArray(subscription.status, ENTITLED_SUBSCRIPTION_STATUSES) - ) - ), + .where(and(eq(subscription.referenceId, userId), entitled)), + // The `organization` join keeps orphaned memberships from contributing a subscription. executor - .select({ organizationId: member.organizationId }) + .select(getTableColumns(subscription)) .from(member) - .where(eq(member.userId, userId)), + .innerJoin(organization, eq(organization.id, member.organizationId)) + .innerJoin(subscription, eq(subscription.referenceId, organization.id)) + .where(and(eq(member.userId, userId), entitled)), ]) - const orgIds = memberships.map((m: { organizationId: string }) => m.organizationId) - - let orgSubs: typeof personalSubs = [] - if (orgIds.length > 0) { - // Verify orgs exist to filter out orphaned subscriptions - const existingOrgs = await executor - .select({ id: organization.id }) - .from(organization) - .where(inArray(organization.id, orgIds)) - - const validOrgIds = existingOrgs.map((o) => o.id) - - if (validOrgIds.length > 0) { - orgSubs = await executor - .select() - .from(subscription) - .where( - and( - inArray(subscription.referenceId, validOrgIds), - inArray(subscription.status, ENTITLED_SUBSCRIPTION_STATUSES) - ) - ) - } - } - if (personalSubs.length === 0 && orgSubs.length === 0) return null return pickHighestPrioritySubscription( diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts index 1d6508b258c..70247fdefd8 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.test.ts @@ -227,19 +227,19 @@ describe('TriggerDevJobQueue status mapping', () => { ['WAITING', 'processing'], ])('maps active Trigger.dev status %s to %s', async (triggerStatus, jobStatus) => { mockRetrieve.mockResolvedValueOnce({ - id: 'run-1', + id: 'run_1', payload: {}, status: triggerStatus, taskIdentifier: 'workflow-execution', }) const queue = new TriggerDevJobQueue() - await expect(queue.getJob('run-1')).resolves.toMatchObject({ status: jobStatus }) + await expect(queue.getJob('run_1')).resolves.toMatchObject({ status: jobStatus }) }) it('dates a run cancelled before it was dequeued by its last transition', async () => { mockRetrieve.mockResolvedValueOnce({ - id: 'run-1', + id: 'run_1', payload: {}, status: 'CANCELED', taskIdentifier: 'workflow-execution', @@ -248,7 +248,7 @@ describe('TriggerDevJobQueue status mapping', () => { }) const queue = new TriggerDevJobQueue() - await expect(queue.getJob('run-1')).resolves.toMatchObject({ + await expect(queue.getJob('run_1')).resolves.toMatchObject({ status: 'cancelled', startedAt: undefined, completedAt: new Date('2026-08-05T12:00:02.000Z'), @@ -257,7 +257,7 @@ describe('TriggerDevJobQueue status mapping', () => { it('prefers the reported finish over the last transition once the run has drained', async () => { mockRetrieve.mockResolvedValueOnce({ - id: 'run-1', + id: 'run_1', payload: {}, status: 'COMPLETED', taskIdentifier: 'workflow-execution', @@ -268,11 +268,55 @@ describe('TriggerDevJobQueue status mapping', () => { }) const queue = new TriggerDevJobQueue() - await expect(queue.getJob('run-1')).resolves.toMatchObject({ + await expect(queue.getJob('run_1')).resolves.toMatchObject({ status: 'completed', completedAt: new Date('2026-08-05T12:00:04.000Z'), }) }) + + /** Trigger.dev can only retrieve its own run ids; anything else fails, and slowly. */ + function retrieveOnlyRunIds() { + mockRetrieve.mockImplementation(async (id: string) => { + if (!id.startsWith('run_')) throw new Error(`Retrieved a caller-chosen job id: ${id}`) + return { id, payload: {}, status: 'QUEUED', taskIdentifier: 'schedule-execution' } + }) + } + + it('resolves a caller-chosen job id through its tag', async () => { + retrieveOnlyRunIds() + mockList.mockReturnValueOnce(createListPage([{ id: 'run_1', tags: ['jobId:schedule_abc'] }])) + const queue = new TriggerDevJobQueue() + + await expect(queue.getJob('schedule_abc')).resolves.toMatchObject({ + id: 'run_1', + status: 'pending', + }) + }) + + it('returns null for a caller-chosen job id with no tagged run', async () => { + retrieveOnlyRunIds() + mockList.mockReturnValueOnce(createListPage([])) + const queue = new TriggerDevJobQueue() + + await expect(queue.getJob('schedule_abc')).resolves.toBeNull() + }) + + it('falls back to the tag lookup when a run id is not found', async () => { + mockRetrieve.mockRejectedValueOnce(new MockApiError(404, 'Not found')) + mockList.mockReturnValueOnce(createListPage([{ id: 'run_2', tags: ['jobId:run_missing'] }])) + mockRetrieve.mockResolvedValueOnce({ + id: 'run_2', + payload: {}, + status: 'EXECUTING', + taskIdentifier: 'workflow-execution', + }) + const queue = new TriggerDevJobQueue() + + await expect(queue.getJob('run_missing')).resolves.toMatchObject({ + id: 'run_2', + status: 'processing', + }) + }) }) describe('TriggerDevJobQueue cancellation', () => { diff --git a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts index b0de23c96fb..e65f04f9efb 100644 --- a/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts +++ b/apps/sim/lib/core/async-jobs/backends/trigger-dev.ts @@ -1,6 +1,7 @@ import { createLogger } from '@sim/logger' import { sha256Hex } from '@sim/security/hash' import { toError } from '@sim/utils/errors' +import { isRecordLike } from '@sim/utils/object' import { taskContext } from '@trigger.dev/core/v3' import { ApiError, runs, type TriggerOptions, tasks } from '@trigger.dev/sdk' import { resolveTriggerRegion } from '@/lib/core/async-jobs/region' @@ -49,6 +50,42 @@ function classifyTriggerEnqueueError(error: unknown): AsyncJobEnqueueError { }) } +/** Trigger's friendly run ids — the only ids `runs.retrieve` can resolve. */ +const TRIGGER_RUN_ID_PREFIX = 'run_' + +type TriggerRun = Awaited> + +function isTriggerNotFoundError(error: unknown): boolean { + return ( + (error instanceof Error && error.message.toLowerCase().includes('not found')) || + (isRecordLike(error) && error.status === 404) + ) +} + +/** + * Retrieves a run by its Trigger run id, or null when `jobId` is not one or no + * such run exists. A caller-chosen job id (`schedule_…`, `workflow-execution:…`) + * can never resolve here, and Trigger takes ~10 s to answer that 404, so it is + * not sent at all. + */ +async function retrieveRunById(jobId: string): Promise { + if (!jobId.startsWith(TRIGGER_RUN_ID_PREFIX)) return null + try { + return await runs.retrieve(jobId) + } catch (error) { + if (isTriggerNotFoundError(error)) return null + throw error + } +} + +/** Resolves a caller-chosen job id through the `jobId:` tag set at enqueue. */ +async function retrieveRunByJobIdTag(jobId: string): Promise { + for await (const candidate of runs.list({ tag: `jobId:${jobId}`, limit: 1 })) { + return runs.retrieve(candidate.id) + } + return null +} + function buildExecutionTag(executionId: string): string { const tag = `executionId:${executionId}` return tag.length <= 128 ? tag : `executionIdHash:${sha256Hex(executionId)}` @@ -380,25 +417,10 @@ export class TriggerDevJobQueue implements JobQueueBackend { async getJob(jobId: string): Promise { try { - let run: Awaited> - try { - run = await runs.retrieve(jobId) - } catch (error) { - const isNotFound = - (error instanceof Error && error.message.toLowerCase().includes('not found')) || - (error && typeof error === 'object' && 'status' in error && error.status === 404) - if (!isNotFound) throw error - - let runId: string | undefined - for await (const candidate of runs.list({ tag: `jobId:${jobId}`, limit: 1 })) { - runId = candidate.id - break - } - if (!runId) { - logger.debug('Job not found in trigger.dev', { jobId }) - return null - } - run = await runs.retrieve(runId) + const run = (await retrieveRunById(jobId)) ?? (await retrieveRunByJobIdTag(jobId)) + if (!run) { + logger.debug('Job not found in trigger.dev', { jobId }) + return null } const payload = run.payload as Record @@ -428,11 +450,7 @@ export class TriggerDevJobQueue implements JobQueueBackend { metadata, } } catch (error) { - const isNotFound = - (error instanceof Error && error.message.toLowerCase().includes('not found')) || - (error && typeof error === 'object' && 'status' in error && error.status === 404) - - if (isNotFound) { + if (isTriggerNotFoundError(error)) { logger.debug('Job not found in trigger.dev', { jobId }) return null } diff --git a/apps/sim/lib/environment/execution-environment.integration.ts b/apps/sim/lib/environment/execution-environment.integration.ts new file mode 100644 index 00000000000..b38e792f4a6 --- /dev/null +++ b/apps/sim/lib/environment/execution-environment.integration.ts @@ -0,0 +1,122 @@ +/** + * Execution environment resolution against real PostgreSQL: which identity lends the + * personal slice and which one authorizes the workspace slice, including suspension + * and lapsed workspace access. + */ +import { db } from '@sim/db' +import { environment, permissions, user, workspace, workspaceEnvironment } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, inArray } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { encryptSecret } from '@/lib/core/security/encryption' +import { getExecutionEnvironment } from '@/lib/environment/utils' + +const workspaceId = generateId() +const owner = `env-owner-${generateId()}` +const actor = `env-actor-${generateId()}` +const suspended = `env-suspended-${generateId()}` +const outsider = `env-outsider-${generateId()}` +const userIds = [owner, actor, suspended, outsider] + +async function encryptedVariables(values: Record) { + const entries = await Promise.all( + Object.entries(values).map(async ([key, value]) => [ + key, + (await encryptSecret(value)).encrypted, + ]) + ) + return Object.fromEntries(entries) +} + +beforeAll(async () => { + const now = new Date() + await db.insert(user).values( + userIds.map((id) => ({ + id, + name: id, + email: `${id}@environment.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + ...(id === suspended ? { banned: true } : {}), + })) + ) + await db.insert(workspace).values({ + id: workspaceId, + name: 'Environment', + ownerId: owner, + billedAccountUserId: actor, + }) + await db.insert(permissions).values( + [owner, actor, suspended].map((userId) => ({ + id: generateId(), + userId, + entityType: 'workspace', + entityId: workspaceId, + permissionType: 'admin' as const, + })) + ) + await db.insert(workspaceEnvironment).values({ + id: generateId(), + workspaceId, + variables: await encryptedVariables({ SHARED: 'workspace-value' }), + }) + await db.insert(environment).values( + await Promise.all( + userIds.map(async (userId) => ({ + id: userId, + userId, + variables: await encryptedVariables({ PERSONAL: `${userId}-personal` }), + })) + ) + ) +}) + +afterAll(async () => { + await db.delete(workspace).where(eq(workspace.id, workspaceId)) + await db.delete(user).where(inArray(user.id, userIds)) +}) + +describe('getExecutionEnvironment', () => { + it('lends the personal slice of an identity that is both owner and actor', async () => { + const env = await getExecutionEnvironment(owner, owner, workspaceId) + + expect(env.personalDecrypted).toEqual({ PERSONAL: `${owner}-personal` }) + expect(env.workspaceDecrypted).toEqual({ SHARED: 'workspace-value' }) + }) + + it('keeps the owner personal slice and the actor workspace slice when they differ', async () => { + const env = await getExecutionEnvironment(owner, actor, workspaceId) + + expect(env.personalDecrypted).toEqual({ PERSONAL: `${owner}-personal` }) + expect(env.workspaceDecrypted).toEqual({ SHARED: 'workspace-value' }) + }) + + it('withholds a suspended identity personal slice when it is also the actor', async () => { + const env = await getExecutionEnvironment(suspended, suspended, workspaceId) + + expect(env.personalDecrypted).toEqual({}) + expect(env.personalEncrypted).toEqual({}) + expect(env.workspaceDecrypted).toEqual({ SHARED: 'workspace-value' }) + }) + + it('withholds a suspended owner personal slice and resolves the actor workspace slice', async () => { + const env = await getExecutionEnvironment(suspended, actor, workspaceId) + + expect(env.personalDecrypted).toEqual({}) + expect(env.workspaceDecrypted).toEqual({ SHARED: 'workspace-value' }) + }) + + it('lends no personal slice from an identity that cannot reach the workspace', async () => { + const env = await getExecutionEnvironment(outsider, actor, workspaceId) + + expect(env.personalDecrypted).toEqual({}) + expect(env.workspaceDecrypted).toEqual({ SHARED: 'workspace-value' }) + }) + + it('refuses when neither identity can reach the workspace', async () => { + await expect(getExecutionEnvironment(outsider, outsider, workspaceId)).rejects.toThrow( + /Access denied/ + ) + }) +}) diff --git a/apps/sim/lib/environment/utils.ts b/apps/sim/lib/environment/utils.ts index 3e43c5f48ac..c1f6c38cee4 100644 --- a/apps/sim/lib/environment/utils.ts +++ b/apps/sim/lib/environment/utils.ts @@ -493,37 +493,66 @@ export async function getExecutionEnvironment( /** * A suspended account lends nothing, from any path. * - * Checked before the single-identity shortcut below rather than alongside the - * access lookups, because "the caller already cleared this identity" does not - * hold everywhere: a custom-block child is admitted by - * `admitCustomBlockChildExecution`, which checks usage limits and nothing - * else, and a provider URL-validation challenge resolves with no admission at - * all. Behind the shortcut, a publisher who is also their workspace's billing - * account made both identities equal and skipped the gate entirely — the one - * arrangement where suspension was silently ignored. + * Applied on every path, including the single-identity shortcut below, because + * "the caller already cleared this identity" does not hold everywhere: a + * custom-block child is admitted by `admitCustomBlockChildExecution`, which checks + * usage limits and nothing else, and a provider URL-validation challenge resolves + * with no admission at all. Behind the shortcut, a publisher who is also their + * workspace's billing account made both identities equal and skipped the gate + * entirely — the one arrangement where suspension was silently ignored. * * Only the personal namespace is withheld. Workspace variables belong to the * workspace rather than to a person, so they keep resolving and the runs a * suspended member's teammates depend on keep working — which is the whole * reason admission stopped blocking on this identity in the first place. + * + * The lookup depends on nothing the reads below produce, so it runs alongside + * them, and a suspended identity's personal slice is dropped from their result. + * The one exception is two distinct identities with no workspace: the personal + * read there is not needed once suspension is known, so it waits for the answer. */ - if ((await getActivelyBannedUserIds([personalUserId])).length > 0) { + const personalIdentitySuspended = getActivelyBannedUserIds([personalUserId]).then( + (bannedUserIds) => bannedUserIds.length > 0 + ) + const withholdPersonalSlice = (snapshot: EnvironmentResolutionSnapshot) => { logger.error('Personal-environment identity is suspended; resolving workspace variables only', { personalUserId, workspaceUserId, workspaceId, }) - return toWorkspaceOnlySnapshot(await getPersonalAndWorkspaceEnv(workspaceUserId, workspaceId)) + return toWorkspaceOnlySnapshot(snapshot) } - if (!workspaceId || workspaceUserId === personalUserId) { + if (workspaceUserId === personalUserId) { + const [suspended, snapshot] = await Promise.all([ + personalIdentitySuspended, + getPersonalAndWorkspaceEnv(personalUserId, workspaceId), + ]) + return suspended ? withholdPersonalSlice(snapshot) : snapshot + } + + if (!workspaceId) { + if (await personalIdentitySuspended) { + return withholdPersonalSlice(await getPersonalAndWorkspaceEnv(workspaceUserId, workspaceId)) + } return getPersonalAndWorkspaceEnv(personalUserId, workspaceId) } - const [actorAccess, personalAccess] = await Promise.all([ + const personalAccessRead = checkWorkspaceAccess(workspaceId, personalUserId) + personalAccessRead.catch(() => {}) + const [suspended, actorAccess] = await Promise.all([ + personalIdentitySuspended, checkWorkspaceAccess(workspaceId, workspaceUserId), - checkWorkspaceAccess(workspaceId, personalUserId), ]) + // A suspended identity's access is never consulted, so its read cannot fail the run. + if (suspended) { + return withholdPersonalSlice( + await getPersonalAndWorkspaceEnv(workspaceUserId, workspaceId, { + workspaceAccess: actorAccess, + }) + ) + } + const personalAccess = await personalAccessRead /** * A workspace that no longer exists and one an identity may not read are diff --git a/apps/sim/lib/execution/preprocessing.test.ts b/apps/sim/lib/execution/preprocessing.test.ts index f1a6c5ef9d9..e5c6f5ce1fd 100644 --- a/apps/sim/lib/execution/preprocessing.test.ts +++ b/apps/sim/lib/execution/preprocessing.test.ts @@ -372,6 +372,40 @@ describe('preprocessExecution ban gate', () => { expect(loggingSession.safeStart).toHaveBeenCalled() }) + it('keeps the 404 for a missing workflow when the gates it overlaps would reject', async () => { + workflowAuthzMockFns.mockGetActiveWorkflowRecord.mockResolvedValueOnce(null) + mockGetActivelyBannedUserIds.mockResolvedValue(['actor-1']) + mockCheckAttributedUsageLimits.mockResolvedValue({ + isExceeded: true, + payerUsage: { currentUsage: 11, limit: 10 }, + }) + + const result = await preprocessExecution({ + ...baseOptions, + workflowRecord: undefined, + billingAttribution: ORGANIZATION_ATTRIBUTION, + }) + + expect(result).toMatchObject({ + success: false, + error: { statusCode: 404, message: 'Workflow not found' }, + }) + }) + + it('reports a serialized attribution for another workspace before any gate result', async () => { + mockGetActivelyBannedUserIds.mockResolvedValue(['actor-1']) + + const result = await preprocessExecution({ + ...baseOptions, + billingAttribution: { ...ORGANIZATION_ATTRIBUTION, workspaceId: 'workspace-2' }, + }) + + expect(result).toMatchObject({ + success: false, + error: { statusCode: 500, message: 'Error resolving billing account' }, + }) + }) + it('returns 403 (ban precedence) when ban, usage, and rate limit all fail simultaneously', async () => { mockGetActivelyBannedUserIds.mockResolvedValue(['billed-account-1']) mockCheckAttributedUsageLimits.mockResolvedValue({ diff --git a/apps/sim/lib/execution/preprocessing.ts b/apps/sim/lib/execution/preprocessing.ts index 4f367235995..1c8d56f1cce 100644 --- a/apps/sim/lib/execution/preprocessing.ts +++ b/apps/sim/lib/execution/preprocessing.ts @@ -72,6 +72,11 @@ export interface PreprocessExecutionOptions { * for every surface, including async-queued v1 runs. */ rateLimitCounter?: 'sync' | 'async' + /** + * Also resolve the actor's highest-priority subscription for a caller that paces + * its own rate limiting. The rate-limit gate resolves it whenever it runs. + */ + includeActorSubscription?: boolean checkDeployment?: boolean skipUsageLimits?: boolean /** @@ -163,7 +168,8 @@ export interface PreprocessExecutionSuccess { success: true actorUserId: string workflowRecord: WorkflowRecord - actorSubscription: SubscriptionInfo + /** Resolved when the rate-limit gate ran or the caller asked for it. */ + actorSubscription?: SubscriptionInfo billingAttribution: BillingAttributionSnapshot executionTimeout: { sync: number @@ -181,6 +187,152 @@ export type PreprocessExecutionResult = PreprocessExecutionSuccess | PreprocessE type WorkflowRecord = typeof workflow.$inferSelect type SubscriptionInfo = HighestPrioritySubscription +/** The admission gates' lookups for one acting identity and payer. */ +interface GateReads { + actorUserId: string + billingAttribution: BillingAttributionSnapshot + bannedUserIds: Promise + actorSubscription: Promise | undefined + usage: Promise>> | undefined +} + +/** + * Blocks when an identity this run actually acts as has an active ban or + * blocked email domain. + * + * `userId` is a candidate unless the caller declares it a stored reference. + * The default is deliberately the blocking one: callers overload the + * parameter, and only the caller knows which kind it passed, so a call site + * that forgets to say must fail closed rather than silently admit a + * suspended account. + * + * `useAuthenticatedUserAsActor` cannot stand in for that declaration, which + * an earlier revision of this gate assumed. Resume passes the live + * authenticated resumer as `userId` and leaves that flag false on purpose — + * attribution is captured before the pause and must not move — so keying on + * it excluded exactly the person who just acted. + * + * A stored reference being banned must not take down work their teammates + * still depend on — but it must not lend that person's credentials either, + * which is why {@link getExecutionEnvironment} drops a suspended identity's + * personal namespace rather than this gate blocking the whole run. + */ +function banCandidateIds( + actorUserId: string, + userId: string, + userIdIsStoredReference: boolean +): string[] { + return !userIdIsStoredReference && userId && userId !== 'unknown' && userId !== actorUserId + ? [actorUserId, userId] + : [actorUserId] +} + +/** + * The payer and admission-gate reads of one preprocessing pass. None of them + * depends on the workflow read, so each starts as soon as its own inputs are + * known — the payer when the caller already knows the workspace, the gates once + * the actor and payer are known — and runs alongside the workflow read. Callers + * still consume the results in preprocessing's fixed order, and a read started + * for inputs the run does not end up using is simply resolved again. + */ +function startAdmissionReads(params: { + userId: string + useAuthenticatedUserAsActor: boolean + userIdIsStoredReference: boolean + providedBillingAttribution: BillingAttributionSnapshot | undefined + knownWorkspaceId: string | undefined + includeSubscription: boolean + includeUsage: boolean +}) { + const { userId, useAuthenticatedUserAsActor, providedBillingAttribution, knownWorkspaceId } = + params + + const resolveAttribution = (workspaceId: string) => + useAuthenticatedUserAsActor && userId + ? withDatabaseReadRetry( + () => resolveBillingAttribution({ actorUserId: userId, workspaceId }), + { + label: 'resolveBillingAttribution', + } + ) + : withDatabaseReadRetry(() => resolveSystemBillingAttribution(workspaceId), { + label: 'resolveSystemBillingAttribution', + }) + + const startGateReads = ( + actorUserId: string, + billingAttribution: BillingAttributionSnapshot + ): GateReads => { + const ids = banCandidateIds(actorUserId, userId, params.userIdIsStoredReference) + const bannedUserIds = withDatabaseReadRetry(() => getActivelyBannedUserIds(ids), { + label: 'getActivelyBannedUserIds', + }) + const actorSubscription = params.includeSubscription + ? getHighestPrioritySubscription(actorUserId) + : undefined + const usage = params.includeUsage + ? withDatabaseReadRetry(() => checkExecutionUsageLimits(billingAttribution), { + label: 'checkExecutionUsageLimits', + }) + : undefined + // The gates await these inside their own error handling; a run rejected before + // it reaches them must not leave them unhandled. + for (const read of [bannedUserIds, actorSubscription, usage]) read?.catch(() => {}) + return { actorUserId, billingAttribution, bannedUserIds, actorSubscription, usage } + } + + /** A serialized attribution is validated once; a failure surfaces where the payer is read. */ + let providedAttribution: { snapshot: BillingAttributionSnapshot } | { error: unknown } | undefined + if (providedBillingAttribution) { + try { + providedAttribution = { + snapshot: assertBillingAttributionSnapshot(providedBillingAttribution), + } + } catch (error) { + providedAttribution = { error } + } + } + + const earlyAttribution = + !providedBillingAttribution && knownWorkspaceId + ? { workspaceId: knownWorkspaceId, attribution: resolveAttribution(knownWorkspaceId) } + : undefined + earlyAttribution?.attribution.catch(() => {}) + + const earlyGateReads: Promise = ( + providedAttribution && 'snapshot' in providedAttribution + ? Promise.resolve( + startGateReads(providedAttribution.snapshot.actorUserId, providedAttribution.snapshot) + ) + : earlyAttribution + ? earlyAttribution.attribution.then((attribution) => + attribution.actorUserId + ? startGateReads(attribution.actorUserId, attribution) + : undefined + ) + : Promise.resolve(undefined) + ).catch(() => undefined) + + return { + providedAttribution, + /** The workspace payer, reusing the early read when it was for this workspace. */ + attributionFor: (workspaceId: string) => + earlyAttribution?.workspaceId === workspaceId + ? earlyAttribution.attribution + : resolveAttribution(workspaceId), + /** The gate reads, reusing the early ones when they were for this actor and payer. */ + gateReadsFor: async ( + actorUserId: string, + billingAttribution: BillingAttributionSnapshot + ): Promise => { + const early = await earlyGateReads + return early?.actorUserId === actorUserId && early.billingAttribution === billingAttribution + ? early + : startGateReads(actorUserId, billingAttribution) + }, + } +} + export async function preprocessExecution( options: PreprocessExecutionOptions ): Promise { @@ -193,6 +345,7 @@ export async function preprocessExecution( requestId, checkRateLimit = triggerType !== 'manual' && triggerType !== 'chat', rateLimitCounter = 'sync', + includeActorSubscription = false, checkDeployment = triggerType !== 'manual', skipUsageLimits = false, skipConcurrencyReservation = false, @@ -234,6 +387,17 @@ export async function preprocessExecution( `Prefetched workflow record ID mismatch: expected ${workflowId}, got ${prefetchedWorkflowRecord.id}` ) } + + const admission = startAdmissionReads({ + userId, + useAuthenticatedUserAsActor, + userIdIsStoredReference, + providedBillingAttribution, + knownWorkspaceId: prefetchedWorkflowRecord?.workspaceId || providedWorkspaceId, + includeSubscription: checkRateLimit || includeActorSubscription, + includeUsage: !skipUsageLimits, + }) + let workflowRecord: WorkflowRecord | null = prefetchedWorkflowRecord ?? null if (!workflowRecord) { try { @@ -349,8 +513,10 @@ export async function preprocessExecution( let billingAttribution: BillingAttributionSnapshot | null = null try { - if (providedBillingAttribution) { - const validatedAttribution = assertBillingAttributionSnapshot(providedBillingAttribution) + const { providedAttribution } = admission + if (providedAttribution) { + if ('error' in providedAttribution) throw providedAttribution.error + const validatedAttribution = providedAttribution.snapshot if (validatedAttribution.workspaceId !== workspaceId) { throw new Error( `Billing attribution workspace mismatch: expected ${workspaceId}, received ${validatedAttribution.workspaceId}` @@ -370,10 +536,7 @@ export async function preprocessExecution( } if (!actorUserId) { - billingAttribution = await withDatabaseReadRetry( - () => resolveSystemBillingAttribution(workspaceId), - { label: 'resolveSystemBillingAttribution' } - ) + billingAttribution = await admission.attributionFor(workspaceId) actorUserId = billingAttribution.actorUserId logger.info(`[${requestId}] Using atomically resolved system actor and payer`, { actorUserId, @@ -410,11 +573,7 @@ export async function preprocessExecution( } if (!billingAttribution) { - const attributionInput = { actorUserId, workspaceId } - billingAttribution = await withDatabaseReadRetry( - () => resolveBillingAttribution(attributionInput), - { label: 'resolveBillingAttribution' } - ) + billingAttribution = await admission.attributionFor(workspaceId) } } catch (error) { logger.error(`[${requestId}] Error resolving billing attribution`, { error, workflowId }) @@ -472,37 +631,11 @@ export async function preprocessExecution( } } + const gateReads = await admission.gateReadsFor(actorUserId, billingAttribution) + const banCheck = (async (): Promise => { - /** - * Blocks when an identity this run actually acts as has an active ban or - * blocked email domain. - * - * `userId` is a candidate unless the caller declares it a stored reference. - * The default is deliberately the blocking one: callers overload the - * parameter, and only the caller knows which kind it passed, so a call site - * that forgets to say must fail closed rather than silently admit a - * suspended account. - * - * `useAuthenticatedUserAsActor` cannot stand in for that declaration, which - * an earlier revision of this gate assumed. Resume passes the live - * authenticated resumer as `userId` and leaves that flag false on purpose — - * attribution is captured before the pause and must not move — so keying on - * it excluded exactly the person who just acted. - * - * A stored reference being banned must not take down work their teammates - * still depend on — but it must not lend that person's credentials either, - * which is why {@link getExecutionEnvironment} drops a suspended identity's - * personal namespace rather than this gate blocking the whole run. - */ - const banCandidateIds = [actorUserId] - if (!userIdIsStoredReference && userId && userId !== 'unknown' && userId !== actorUserId) { - banCandidateIds.push(userId) - } try { - const bannedUserIds = await withDatabaseReadRetry( - () => getActivelyBannedUserIds(banCandidateIds), - { label: 'getActivelyBannedUserIds' } - ) + const bannedUserIds = await gateReads.bannedUserIds if (bannedUserIds.length > 0) { logger.warn(`[${requestId}] Execution blocked: banned account`, { workflowId, @@ -560,8 +693,6 @@ export async function preprocessExecution( } })() - const subscriptionFetch = getHighestPrioritySubscription(actorUserId) - /** * Returns the usage failure and reservation snapshot together so concurrent * read gates do not communicate through mutable outer state. @@ -570,13 +701,10 @@ export async function preprocessExecution( failure: GateFailure | null snapshot: UsageSnapshot | null }> => { - if (skipUsageLimits) return { failure: null, snapshot: null } + if (!gateReads.usage) return { failure: null, snapshot: null } let snapshot: UsageSnapshot | null = null try { - const usageCheck = await withDatabaseReadRetry( - () => checkExecutionUsageLimits(billingAttribution), - { label: 'checkExecutionUsageLimits' } - ) + const usageCheck = await gateReads.usage snapshot = usageCheck.payerUsage ? { ...usageCheck.payerUsage, @@ -670,7 +798,7 @@ export async function preprocessExecution( */ const [banFailure, actorSubscription, usageResult] = await Promise.all([ banCheck, - subscriptionFetch, + gateReads.actorSubscription, usageCheckTask, ]) @@ -684,7 +812,7 @@ export async function preprocessExecution( const rateLimiter = new RateLimiter() const info = await rateLimiter.checkRateLimitWithSubscription( actorUserId, - actorSubscription, + actorSubscription ?? null, triggerType, rateLimitCounter === 'async' ) diff --git a/apps/sim/lib/logs/execution/logger.test.ts b/apps/sim/lib/logs/execution/logger.test.ts index 5c1d0c10134..cbd20a79913 100644 --- a/apps/sim/lib/logs/execution/logger.test.ts +++ b/apps/sim/lib/logs/execution/logger.test.ts @@ -127,27 +127,10 @@ vi.mock('@/lib/logs/execution/progress-markers', () => ({ // Mock snapshot service vi.mock('@/lib/logs/execution/snapshot/service', () => ({ snapshotService: { - createSnapshotWithDeduplication: vi.fn(() => - Promise.resolve({ - snapshot: { - id: 'snapshot-123', - workflowId: 'workflow-123', - stateHash: 'hash-123', - stateData: { blocks: {}, edges: [], loops: {}, parallels: {} }, - createdAt: '2024-01-01T00:00:00.000Z', - }, - isNew: true, - }) - ), - getSnapshot: vi.fn(() => - Promise.resolve({ - id: 'snapshot-123', - workflowId: 'workflow-123', - stateHash: 'hash-123', - stateData: { blocks: {}, edges: [], loops: {}, parallels: {} }, - createdAt: '2024-01-01T00:00:00.000Z', - }) + resolveSnapshot: vi.fn(() => + Promise.resolve({ id: 'snapshot-123', workflowId: 'workflow-123', stateHash: 'hash' }) ), + rememberReferencedSnapshot: vi.fn(), }, })) @@ -161,7 +144,6 @@ describe('ExecutionLogger', () => { describe('interface implementation', () => { test('marks new execution rows as contract-aware before any provenance is available', async () => { - dbChainMockFns.limit.mockResolvedValueOnce([]) dbChainMockFns.returning.mockResolvedValueOnce([ { id: 'log-1', diff --git a/apps/sim/lib/logs/execution/logger.ts b/apps/sim/lib/logs/execution/logger.ts index b6ca886ba74..650481bc68d 100644 --- a/apps/sim/lib/logs/execution/logger.ts +++ b/apps/sim/lib/logs/execution/logger.ts @@ -8,7 +8,12 @@ import { workspace, } from '@sim/db/schema' import { createLogger } from '@sim/logger' -import { describeError, getErrorMessage } from '@sim/utils/errors' +import { + describeError, + getErrorMessage, + getPostgresConstraintName, + getPostgresErrorCode, +} from '@sim/utils/errors' import { generateId } from '@sim/utils/id' import { and, eq, inArray, sql } from 'drizzle-orm' import { checkUsageStatus as checkResolvedUsageStatus } from '@/lib/billing/calculations/usage-monitor' @@ -72,7 +77,6 @@ import type { ExecutionLoggerService as IExecutionLoggerService, TraceSpan, WorkflowExecutionLog, - WorkflowExecutionSnapshot, WorkflowState, } from '@/lib/logs/types' import { emitExecutionCompletedEvent } from '@/lib/workspace-events/emitter' @@ -95,6 +99,9 @@ const EXECUTION_LOG_IDLE_TIMEOUT_MS = 5_000 // Bounds the wait for the per-execution usage-reconcile advisory lock. Generous // (favor waiting over dropping a charge); only trips on a pathological lock hold. const USAGE_RECONCILE_LOCK_TIMEOUT_MS = 10_000 +const FOREIGN_KEY_VIOLATION = '23503' +/** The log's snapshot foreign key as Postgres names it (identifiers are cut at 63 bytes). */ +const STATE_SNAPSHOT_FOREIGN_KEY = 'workflow_execution_logs_state_snapshot_id_workflow_execution_sn' type ExecutionData = WorkflowExecutionLog['executionData'] @@ -598,6 +605,11 @@ export class ExecutionLogger implements IExecutionLoggerService { } } + /** + * Creates the execution's `running` log row, pointing it at the (deduplicated) + * workflow-state snapshot. Idempotent per execution id: when the row already + * exists — a retried job or a resumed start — only its deadline is refreshed. + */ async startWorkflowExecution(params: { workflowId: string workspaceId: string @@ -609,10 +621,7 @@ export class ExecutionLogger implements IExecutionLoggerService { workflowState: WorkflowState deploymentVersionId?: string executionDeadlineAt?: Date - }): Promise<{ - workflowLog: WorkflowExecutionLog - snapshot: WorkflowExecutionSnapshot - }> { + }): Promise { const { workflowId, workspaceId, @@ -628,99 +637,74 @@ export class ExecutionLogger implements IExecutionLoggerService { execLog.debug('Starting workflow execution') - // Check if execution log already exists (idempotency check) - const existingLog = await execDb - .select() - .from(workflowExecutionLogs) - .where(eq(workflowExecutionLogs.executionId, executionId)) - .limit(1) + const insertRunningLog = async (stateSnapshotId: string) => { + const [inserted] = await execDb + .insert(workflowExecutionLogs) + .values({ + id: generateId(), + workflowId, + workspaceId, + executionId, + stateSnapshotId, + deploymentVersionId: deploymentVersionId ?? null, + level: 'info', + status: 'running', + trigger: trigger.type, + startedAt: new Date(), + endedAt: null, + totalDurationMs: null, + executionDeadlineAt: executionDeadlineAt ?? null, + executionData: { + secretProjectionVersion: SECRET_PROJECTION_VERSION, + environment, + trigger, + ...(billingAttribution ? { billingAttribution } : {}), + ...(trigger.data?.correlation ? { correlation: trigger.data.correlation } : {}), + hasTraceSpans: false, + traceSpanCount: 0, + }, + }) + .onConflictDoNothing({ target: workflowExecutionLogs.executionId }) + .returning({ id: workflowExecutionLogs.id }) + return inserted + } - if (existingLog.length > 0) { - execLog.debug('Execution log already exists, skipping duplicate INSERT (idempotent)') - await execDb - .update(workflowExecutionLogs) - .set({ executionDeadlineAt: executionDeadlineAt ?? null }) - .where( - and( - eq(workflowExecutionLogs.executionId, executionId), - sql`${workflowExecutionLogs.status} IN ('pending', 'running')` - ) - ) - const snapshot = await snapshotService.getSnapshot(existingLog[0].stateSnapshotId) - if (!snapshot) { - throw new Error(`Snapshot ${existingLog[0].stateSnapshotId} not found for existing log`) - } - return { - workflowLog: { - id: existingLog[0].id, - workflowId: existingLog[0].workflowId, - executionId: existingLog[0].executionId, - stateSnapshotId: existingLog[0].stateSnapshotId, - level: existingLog[0].level as 'info' | 'error', - trigger: existingLog[0].trigger as ExecutionTrigger['type'], - startedAt: existingLog[0].startedAt.toISOString(), - endedAt: existingLog[0].endedAt?.toISOString() || existingLog[0].startedAt.toISOString(), - totalDurationMs: existingLog[0].totalDurationMs || 0, - executionData: existingLog[0].executionData as WorkflowExecutionLog['executionData'], - createdAt: existingLog[0].createdAt.toISOString(), - }, - snapshot, + let snapshot = await snapshotService.resolveSnapshot(workflowId, workflowState) + let inserted: { id: string } | undefined + try { + inserted = await insertRunningLog(snapshot.id) + } catch (error) { + if ( + getPostgresErrorCode(error) !== FOREIGN_KEY_VIOLATION || + getPostgresConstraintName(error) !== STATE_SNAPSHOT_FOREIGN_KEY + ) { + throw error } + /** + * A snapshot resolved before the insert can be deleted underneath it by + * orphan cleanup when no log references it yet. Resolve it again from the + * database, which recreates the row if it is gone. + */ + snapshot = await snapshotService.resolveSnapshot(workflowId, workflowState, { fresh: true }) + inserted = await insertRunningLog(snapshot.id) } - const snapshotResult = await snapshotService.createSnapshotWithDeduplication( - workflowId, - workflowState - ) - - const startTime = new Date() - - const [workflowLog] = await execDb - .insert(workflowExecutionLogs) - .values({ - id: generateId(), - workflowId, - workspaceId, - executionId, - stateSnapshotId: snapshotResult.snapshot.id, - deploymentVersionId: deploymentVersionId ?? null, - level: 'info', - status: 'running', - trigger: trigger.type, - startedAt: startTime, - endedAt: null, - totalDurationMs: null, - executionDeadlineAt: executionDeadlineAt ?? null, - executionData: { - secretProjectionVersion: SECRET_PROJECTION_VERSION, - environment, - trigger, - ...(billingAttribution ? { billingAttribution } : {}), - ...(trigger.data?.correlation ? { correlation: trigger.data.correlation } : {}), - hasTraceSpans: false, - traceSpanCount: 0, - }, - }) - .returning() - - execLog.debug('Created workflow log', { logId: workflowLog.id }) - - return { - workflowLog: { - id: workflowLog.id, - workflowId: workflowLog.workflowId, - executionId: workflowLog.executionId, - stateSnapshotId: workflowLog.stateSnapshotId, - level: workflowLog.level as 'info' | 'error', - trigger: workflowLog.trigger as ExecutionTrigger['type'], - startedAt: workflowLog.startedAt.toISOString(), - endedAt: workflowLog.endedAt?.toISOString() || workflowLog.startedAt.toISOString(), - totalDurationMs: workflowLog.totalDurationMs || 0, - executionData: workflowLog.executionData as WorkflowExecutionLog['executionData'], - createdAt: workflowLog.createdAt.toISOString(), - }, - snapshot: snapshotResult.snapshot, + if (inserted) { + snapshotService.rememberReferencedSnapshot(snapshot) + execLog.debug('Created workflow log', { logId: inserted.id }) + return } + + execLog.debug('Execution log already exists, skipping duplicate INSERT (idempotent)') + await execDb + .update(workflowExecutionLogs) + .set({ executionDeadlineAt: executionDeadlineAt ?? null }) + .where( + and( + eq(workflowExecutionLogs.executionId, executionId), + sql`${workflowExecutionLogs.status} IN ('pending', 'running')` + ) + ) } /** diff --git a/apps/sim/lib/logs/execution/snapshot/service.test.ts b/apps/sim/lib/logs/execution/snapshot/service.test.ts index cec122498e7..da67f7c73c3 100644 --- a/apps/sim/lib/logs/execution/snapshot/service.test.ts +++ b/apps/sim/lib/logs/execution/snapshot/service.test.ts @@ -1,35 +1,8 @@ import { databaseMock } from '@sim/testing' -import { idMock, idMockFns } from '@sim/testing/mocks/id.mock' import { describe, expect, it, vi } from 'vitest' - -vi.mock('@sim/utils/id', () => idMock) - import { SnapshotService } from '@/lib/logs/execution/snapshot/service' import type { WorkflowState } from '@/lib/logs/types' -idMockFns.mockGenerateId.mockImplementation(() => 'generated-uuid-1') -idMockFns.mockGenerateShortId.mockImplementation(() => 'generated-short-1') - -const mockState: WorkflowState = { - blocks: { - block1: { - id: 'block1', - name: 'Test Agent', - type: 'agent', - position: { x: 100, y: 200 }, - subBlocks: {}, - outputs: {}, - enabled: true, - horizontalHandles: true, - advancedMode: false, - height: 0, - }, - }, - edges: [{ id: 'edge1', source: 'block1', target: 'block2' }], - loops: {}, - parallels: {}, -} - describe('SnapshotService', () => { describe('computeStateHash', () => { it.concurrent('should ignore position changes', () => { @@ -190,108 +163,6 @@ describe('SnapshotService', () => { }) }) - describe('createSnapshotWithDeduplication', () => { - type SnapshotRow = { - id: string - workflowId: string - stateHash: string - stateData: WorkflowState - createdAt: Date - } - - /** Mock the insert → values → onConflictDoUpdate → returning chain. */ - function mockUpsertReturning(rows: SnapshotRow[]) { - let capturedConflictConfig: Record | undefined - const onConflictDoUpdate = vi.fn().mockImplementation((config: Record) => { - capturedConflictConfig = config - return { returning: vi.fn().mockResolvedValue(rows) } - }) - const values = vi.fn().mockReturnValue({ onConflictDoUpdate }) - databaseMock.db.insert = vi.fn().mockReturnValue({ values }) - databaseMock.db.select = vi.fn() - return { values, onConflictDoUpdate, getConflictConfig: () => capturedConflictConfig } - } - - it('reuses the existing snapshot atomically when the returned id differs', async () => { - const service = new SnapshotService() - const workflowId = 'wf-123' - - mockUpsertReturning([ - { - id: 'existing-snapshot-id', - workflowId, - stateHash: 'abc123', - stateData: mockState, - createdAt: new Date('2026-02-19T00:00:00Z'), - }, - ]) - - const result = await service.createSnapshotWithDeduplication(workflowId, mockState) - - expect(result.snapshot.id).toBe('existing-snapshot-id') - expect(result.isNew).toBe(false) - expect(databaseMock.db.select).not.toHaveBeenCalled() - }) - - it('SET targets only state_hash on conflict, never the large state_data', async () => { - const service = new SnapshotService() - const workflowId = 'wf-123' - - const { onConflictDoUpdate, getConflictConfig } = mockUpsertReturning([ - { - id: 'generated-uuid-1', - workflowId, - stateHash: 'abc123', - stateData: mockState, - createdAt: new Date('2026-02-19T00:00:00Z'), - }, - ]) - - await service.createSnapshotWithDeduplication(workflowId, mockState) - - expect(onConflictDoUpdate).toHaveBeenCalledTimes(1) - const config = getConflictConfig() - expect(config?.target).toBeDefined() - // The crux of this change: the SET touches state_hash only, so the unchanged - // TOASTed state_data jsonb is never rewritten. - expect(config?.set).toHaveProperty('stateHash') - expect(config?.set).not.toHaveProperty('stateData') - }) - - it('does not throw on concurrent inserts with the same hash', async () => { - const service = new SnapshotService() - const workflowId = 'wf-123' - - const newRow: SnapshotRow = { - id: 'generated-uuid-1', - workflowId, - stateHash: 'abc123', - stateData: mockState, - createdAt: new Date('2026-02-19T00:00:00Z'), - } - const existingRow: SnapshotRow = { ...newRow, id: 'existing-snapshot-id' } - - let upsertCall = 0 - databaseMock.db.insert = vi.fn().mockImplementation(() => ({ - values: vi.fn().mockReturnValue({ - onConflictDoUpdate: vi.fn().mockReturnValue({ - returning: vi.fn().mockResolvedValue(upsertCall++ === 0 ? [newRow] : [existingRow]), - }), - }), - })) - - const [result1, result2] = await Promise.all([ - service.createSnapshotWithDeduplication(workflowId, mockState), - service.createSnapshotWithDeduplication(workflowId, mockState), - ]) - - expect(result1.snapshot.id).toBe('generated-uuid-1') - expect(result1.isNew).toBe(true) - expect(result2.snapshot.id).toBe('existing-snapshot-id') - expect(result2.isNew).toBe(false) - }) - }) - describe('cleanupOrphanedSnapshots', () => { function setupCleanupMocks(selectBatches: Array>) { const limitFn = vi.fn() diff --git a/apps/sim/lib/logs/execution/snapshot/service.ts b/apps/sim/lib/logs/execution/snapshot/service.ts index 33db22dba97..5bc96a9465c 100644 --- a/apps/sim/lib/logs/execution/snapshot/service.ts +++ b/apps/sim/lib/logs/execution/snapshot/service.ts @@ -4,97 +4,102 @@ import { createLogger } from '@sim/logger' import { sha256Hex } from '@sim/security/hash' import { generateId } from '@sim/utils/id' import { and, eq, inArray, lt, notExists, sql } from 'drizzle-orm' +import { LRUCache } from 'lru-cache' import { consumeRowBudget, type RowBudget } from '@/lib/cleanup/batch-delete' -import type { - SnapshotService as ISnapshotService, - SnapshotCreationResult, - WorkflowExecutionSnapshot, - WorkflowExecutionSnapshotInsert, - WorkflowState, -} from '@/lib/logs/types' +import type { SnapshotService as ISnapshotService, WorkflowState } from '@/lib/logs/types' import { normalizedStringify, normalizeWorkflowState } from '@/lib/workflows/comparison' const logger = createLogger('SnapshotService') +const SNAPSHOT_ID_CACHE_MAX_ENTRIES = 1000 +const SNAPSHOT_ID_CACHE_TTL_MS = 5 * 60 * 1000 + +/** + * Remembers which row holds a workflow's snapshot for a state hash, so repeat runs + * of an unchanged workflow neither look it up nor ship its (often hundreds of KB) + * state to the database. Only ids a log row was just inserted against are + * remembered: orphan cleanup deletes only unreferenced snapshots, so a referenced id + * stays valid for the entry's lifetime. A caller that still hits a missing row + * resolves again with `{ fresh: true }`. + */ +const snapshotIdCache = new LRUCache({ + max: SNAPSHOT_ID_CACHE_MAX_ENTRIES, + ttl: SNAPSHOT_ID_CACHE_TTL_MS, +}) + +/** The snapshot row holding one workflow's state, identified by the state's hash. */ +export interface ResolvedSnapshot { + id: string + workflowId: string + stateHash: string +} + +const snapshotCacheKey = ({ workflowId, stateHash }: Omit) => + `${workflowId}:${stateHash}` + export class SnapshotService implements ISnapshotService { - async createSnapshot( + /** + * Resolves the snapshot row holding `state` for `workflowId`, creating the row + * only when no identical state (same normalized hash) is stored yet. + */ + async resolveSnapshot( workflowId: string, - state: WorkflowState - ): Promise { - const result = await this.createSnapshotWithDeduplication(workflowId, state) - return result.snapshot + state: WorkflowState, + options: { fresh?: boolean } = {} + ): Promise { + const stateHash = this.computeStateHash(state) + if (!options.fresh) { + const cachedId = snapshotIdCache.get(snapshotCacheKey({ workflowId, stateHash })) + if (cachedId) return { id: cachedId, workflowId, stateHash } + } + + const [existing] = await dbFor('exec') + .select({ id: workflowExecutionSnapshots.id }) + .from(workflowExecutionSnapshots) + .where( + and( + eq(workflowExecutionSnapshots.workflowId, workflowId), + eq(workflowExecutionSnapshots.stateHash, stateHash) + ) + ) + .limit(1) + + const id = existing?.id ?? (await this.insertSnapshot(workflowId, stateHash, state)) + return { id, workflowId, stateHash } } - async createSnapshotWithDeduplication( + /** Remembers a snapshot once a log row referencing it has been inserted. */ + rememberReferencedSnapshot(snapshot: ResolvedSnapshot): void { + snapshotIdCache.set(snapshotCacheKey(snapshot), snapshot.id) + } + + /** + * Inserts the snapshot, or — when a concurrent run stored the identical + * (workflowId, stateHash) row first — returns that row's id without rewriting it. + * + * The hash is a sha256 of the normalized state, so an existing row's stateData is + * byte-identical; there is nothing to update. SET touches only the small + * state_hash column so the upsert still RETURNs the row's id, while the unchanged, + * TOASTed stateData keeps its existing out-of-line storage. + */ + private async insertSnapshot( workflowId: string, + stateHash: string, state: WorkflowState - ): Promise { - const stateHash = this.computeStateHash(state) - - const snapshotData: WorkflowExecutionSnapshotInsert = { - id: generateId(), - workflowId, - stateHash, - stateData: state, - } - - /** - * Insert the snapshot, or — when an identical (workflowId, stateHash) row - * already exists — return it without rewriting the large stateData jsonb. - * - * The hash is a sha256 of the normalized state, so an existing row's stateData - * is byte-identical; there is nothing to update. The previous implementation - * SET state_data on conflict, which rewrote the full (tens-of-KB) jsonb every - * run. We keep a single atomic upsert — so RETURNING always yields the row and - * there is no race with snapshot cleanup (unlike DO NOTHING + a follow-up - * select) — but SET only the small state_hash column to itself. Under Postgres - * MVCC the unchanged, TOASTed stateData is not rewritten: its existing - * out-of-line storage is reused, so the per-execution write drops from the - * full blob to a tiny heap tuple. - */ - const [upsertedSnapshot] = await dbFor('exec') + ): Promise { + const [row] = await dbFor('exec') .insert(workflowExecutionSnapshots) - .values(snapshotData) + .values({ id: generateId(), workflowId, stateHash, stateData: state }) .onConflictDoUpdate({ target: [workflowExecutionSnapshots.workflowId, workflowExecutionSnapshots.stateHash], - set: { - stateHash: sql`excluded.state_hash`, - }, + set: { stateHash: sql`excluded.state_hash` }, }) - .returning() - - const isNew = upsertedSnapshot.id === snapshotData.id + .returning({ id: workflowExecutionSnapshots.id }) logger.info( - isNew - ? `Created new snapshot for workflow ${workflowId} (hash: ${stateHash.slice(0, 12)}..., blocks: ${Object.keys(state.blocks || {}).length})` - : `Reusing existing snapshot for workflow ${workflowId} (hash: ${stateHash.slice(0, 12)}...)` + `Stored snapshot for workflow ${workflowId} (hash: ${stateHash.slice(0, 12)}..., blocks: ${Object.keys(state.blocks || {}).length})` ) - - return { - snapshot: { - ...upsertedSnapshot, - stateData: upsertedSnapshot.stateData as WorkflowState, - createdAt: upsertedSnapshot.createdAt.toISOString(), - }, - isNew, - } - } - - async getSnapshot(id: string): Promise { - const [snapshot] = await dbFor('exec') - .select() - .from(workflowExecutionSnapshots) - .where(eq(workflowExecutionSnapshots.id, id)) - .limit(1) - - if (!snapshot) return null - - return { - ...snapshot, - stateData: snapshot.stateData as WorkflowState, - createdAt: snapshot.createdAt.toISOString(), - } + return row.id } computeStateHash(state: WorkflowState): string { diff --git a/apps/sim/lib/logs/execution/start-execution.integration.ts b/apps/sim/lib/logs/execution/start-execution.integration.ts new file mode 100644 index 00000000000..ea835d3d966 --- /dev/null +++ b/apps/sim/lib/logs/execution/start-execution.integration.ts @@ -0,0 +1,190 @@ +/** + * Execution-log start against real PostgreSQL: snapshot deduplication, the per-execution + * idempotent start, and recovery when orphan cleanup removes a snapshot a run resolved. + */ +import { db } from '@sim/db' +import { + user, + workflow, + workflowExecutionLogs, + workflowExecutionSnapshots, + workspace, +} from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq, inArray } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { executionLogger } from '@/lib/logs/execution/logger' +import type { WorkflowState } from '@/lib/logs/types' + +const ids = { + owner: `start-log-owner-${generateId()}`, + workspace: generateId(), + workflow: generateId(), +} + +function stateWith(label: string): WorkflowState { + return { + blocks: { + start: { + id: 'start', + type: 'starter', + name: label, + position: { x: 0, y: 0 }, + subBlocks: {}, + outputs: {}, + enabled: true, + }, + }, + edges: [], + loops: {}, + parallels: {}, + } +} + +async function startExecution(executionId: string, state: WorkflowState, deadline?: Date) { + await executionLogger.startWorkflowExecution({ + workflowId: ids.workflow, + workspaceId: ids.workspace, + executionId, + trigger: { type: 'api', source: 'api', timestamp: new Date().toISOString() }, + environment: { + variables: {}, + workflowId: ids.workflow, + executionId, + userId: ids.owner, + workspaceId: ids.workspace, + }, + workflowState: state, + executionDeadlineAt: deadline, + }) +} + +async function logRow(executionId: string) { + const [row] = await db + .select() + .from(workflowExecutionLogs) + .where(eq(workflowExecutionLogs.executionId, executionId)) + return row +} + +async function snapshotRows() { + return db + .select({ id: workflowExecutionSnapshots.id, stateHash: workflowExecutionSnapshots.stateHash }) + .from(workflowExecutionSnapshots) + .where(eq(workflowExecutionSnapshots.workflowId, ids.workflow)) +} + +beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: ids.owner, + name: 'Start Log', + email: `${ids.owner}@start-log.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(workspace).values({ + id: ids.workspace, + name: 'Start Log', + ownerId: ids.owner, + billedAccountUserId: ids.owner, + }) + await db.insert(workflow).values({ + id: ids.workflow, + userId: ids.owner, + workspaceId: ids.workspace, + name: 'Start Log', + lastSynced: now, + createdAt: now, + updatedAt: now, + }) +}) + +afterAll(async () => { + // Snapshots outlive their workflow (workflow_id is set null), so remove them first. + await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.workflowId, ids.workflow)) + await db + .delete(workflowExecutionSnapshots) + .where(eq(workflowExecutionSnapshots.workflowId, ids.workflow)) + await db.delete(workspace).where(eq(workspace.id, ids.workspace)) + await db.delete(user).where(eq(user.id, ids.owner)) +}) + +describe('startWorkflowExecution', () => { + it('stores one snapshot per distinct state and points every run at it', async () => { + const state = stateWith('dedup') + const [first, second] = [generateId(), generateId()] + await startExecution(first, state) + await startExecution(second, state) + + const [firstLog, secondLog] = [await logRow(first), await logRow(second)] + expect(firstLog.status).toBe('running') + expect(secondLog.stateSnapshotId).toBe(firstLog.stateSnapshotId) + + const changed = generateId() + await startExecution(changed, stateWith('dedup changed')) + expect((await logRow(changed)).stateSnapshotId).not.toBe(firstLog.stateSnapshotId) + + const hashes = (await snapshotRows()).map((row) => row.stateHash) + expect(new Set(hashes).size).toBe(hashes.length) + }) + + it('keeps one row for a repeated start and refreshes the deadline only while it runs', async () => { + const executionId = generateId() + const state = stateWith('idempotent') + await startExecution(executionId, state, new Date('2030-01-01T00:00:00Z')) + await startExecution(executionId, state, new Date('2030-01-02T00:00:00Z')) + + const rows = await db + .select() + .from(workflowExecutionLogs) + .where(eq(workflowExecutionLogs.executionId, executionId)) + expect(rows).toHaveLength(1) + expect(rows[0].executionDeadlineAt?.toISOString()).toBe('2030-01-02T00:00:00.000Z') + + await db + .update(workflowExecutionLogs) + .set({ status: 'completed' }) + .where(eq(workflowExecutionLogs.executionId, executionId)) + await startExecution(executionId, state, new Date('2030-01-03T00:00:00Z')) + expect((await logRow(executionId)).executionDeadlineAt?.toISOString()).toBe( + '2030-01-02T00:00:00.000Z' + ) + }) + + it('recovers when orphan cleanup deleted the snapshot a previous run resolved', async () => { + const state = stateWith('cleaned up') + const first = generateId() + await startExecution(first, state) + const { stateSnapshotId } = await logRow(first) + + await db.delete(workflowExecutionLogs).where(eq(workflowExecutionLogs.executionId, first)) + await db + .delete(workflowExecutionSnapshots) + .where(eq(workflowExecutionSnapshots.id, stateSnapshotId)) + + const second = generateId() + await startExecution(second, state) + const secondLog = await logRow(second) + const [snapshot] = await db + .select({ id: workflowExecutionSnapshots.id }) + .from(workflowExecutionSnapshots) + .where(eq(workflowExecutionSnapshots.id, secondLog.stateSnapshotId)) + expect(secondLog.status).toBe('running') + expect(snapshot).toBeDefined() + }) + + it('stores a single snapshot when concurrent runs start a new state together', async () => { + const state = stateWith('concurrent') + const executionIds = Array.from({ length: 6 }, () => generateId()) + await Promise.all(executionIds.map((executionId) => startExecution(executionId, state))) + + const logs = await db + .select({ stateSnapshotId: workflowExecutionLogs.stateSnapshotId }) + .from(workflowExecutionLogs) + .where(inArray(workflowExecutionLogs.executionId, executionIds)) + expect(logs).toHaveLength(executionIds.length) + expect(new Set(logs.map((log) => log.stateSnapshotId)).size).toBe(1) + }) +}) diff --git a/apps/sim/lib/logs/types.ts b/apps/sim/lib/logs/types.ts index 2117cd1881a..a592c1c2d43 100644 --- a/apps/sim/lib/logs/types.ts +++ b/apps/sim/lib/logs/types.ts @@ -101,17 +101,6 @@ export interface ExecutionLastCompletedBlock { success: boolean } -export interface WorkflowExecutionSnapshot { - id: string - workflowId: string | null - stateHash: string - stateData: WorkflowState - createdAt: string -} - -export type WorkflowExecutionSnapshotInsert = Omit -export type WorkflowExecutionSnapshotSelect = WorkflowExecutionSnapshot - export interface WorkflowExecutionLog { id: string workflowId: string | null @@ -515,17 +504,16 @@ export interface BatchInsertResult { } export interface SnapshotService { - createSnapshot(workflowId: string, state: WorkflowState): Promise - getSnapshot(id: string): Promise + resolveSnapshot( + workflowId: string, + state: WorkflowState, + options?: { fresh?: boolean } + ): Promise<{ id: string; workflowId: string; stateHash: string }> + rememberReferencedSnapshot(snapshot: { id: string; workflowId: string; stateHash: string }): void computeStateHash(state: WorkflowState): string cleanupOrphanedSnapshots(olderThanDays: number): Promise } -export interface SnapshotCreationResult { - snapshot: WorkflowExecutionSnapshot - isNew: boolean -} - export interface ExecutionLoggerService { loadTraceSpansForProjection(params: { executionId: string @@ -552,10 +540,7 @@ export interface ExecutionLoggerService { actorUserId?: string | null billingAttribution?: BillingAttributionSnapshot workflowState: WorkflowState - }): Promise<{ - workflowLog: WorkflowExecutionLog - snapshot: WorkflowExecutionSnapshot - }> + }): Promise completeWorkflowExecution(params: { executionId: string diff --git a/apps/sim/lib/webhooks/processor.ts b/apps/sim/lib/webhooks/processor.ts index e618d917621..ebb2ef2b60c 100644 --- a/apps/sim/lib/webhooks/processor.ts +++ b/apps/sim/lib/webhooks/processor.ts @@ -5,7 +5,7 @@ import { toError } from '@sim/utils/errors' import { generateId } from '@sim/utils/id' import { isRecordLike } from '@sim/utils/object' import { truncate } from '@sim/utils/string' -import { and, eq, isNull, or } from 'drizzle-orm' +import { and, eq, isNull, or, sql } from 'drizzle-orm' import { type NextRequest, NextResponse } from 'next/server' import { releaseExecutionSlot } from '@/lib/billing/calculations/usage-reservation' import type { BillingAttributionSnapshot } from '@/lib/billing/core/billing-attribution' @@ -46,7 +46,18 @@ const logger = createLogger('WebhookProcessor') type WebhookRecord = typeof webhook.$inferSelect type WorkflowRecord = typeof workflow.$inferSelect -type WebhookTarget = { webhook: WebhookRecord; workflow: WorkflowRecord } +type WebhookTarget = { + webhook: WebhookRecord + workflow: WorkflowRecord + /** Whether the webhook's trigger block exists in the workflow's active deployment. */ + triggerBlockDeployed: boolean +} + +/** + * Answers {@link blockExistsInDeployment} from the active deployment a lookup + * already joins, so delivery does not read it a second time. + */ +const triggerBlockDeployedColumn = sql`coalesce(json_typeof(${workflowDeploymentVersion.state} -> 'blocks' -> ${webhook.blockId}) = 'object', false)` type ResolvedWebhookRecord = Omit & { provider: string providerConfig: Record @@ -67,6 +78,8 @@ export interface WebhookProcessorOptions { triggerTimestampMs?: number /** Provider-authenticated external actor. Never derived from workflow input. */ subject?: ExternalUserSubject + /** The lookup's answer for this webhook's trigger block; read fresh when absent. */ + triggerBlockDeployed?: boolean } export interface WebhookPreprocessingResult { @@ -363,6 +376,7 @@ export async function findAllWebhooksForPath( .select({ webhook: webhook, workflow: workflow, + triggerBlockDeployed: triggerBlockDeployedColumn, }) .from(webhook) .innerJoin(workflow, eq(webhook.workflowId, workflow.id)) @@ -467,6 +481,7 @@ export async function findWebhooksByRoutingKey( .select({ webhook: webhook, workflow: workflow, + triggerBlockDeployed: triggerBlockDeployedColumn, }) .from(webhook) .innerJoin(workflow, eq(webhook.workflowId, workflow.id)) @@ -896,7 +911,9 @@ export async function dispatchResolvedWebhookTarget( } if (webhookRecord.blockId) { - const blockExists = await blockExistsInDeployment(foundWorkflow.id, webhookRecord.blockId) + const blockExists = + options.triggerBlockDeployed ?? + (await blockExistsInDeployment(foundWorkflow.id, webhookRecord.blockId)) if (!blockExists) { const verificationResponse = handlePreDeploymentVerification(webhookRecord, options.requestId) return { diff --git a/apps/sim/lib/webhooks/slack-dispatch.ts b/apps/sim/lib/webhooks/slack-dispatch.ts index 56e64c33967..18c83b8430b 100644 --- a/apps/sim/lib/webhooks/slack-dispatch.ts +++ b/apps/sim/lib/webhooks/slack-dispatch.ts @@ -108,7 +108,7 @@ export async function dispatchSlackWebhooks( return mapWithConcurrency( webhooks, SLACK_WEBHOOK_DISPATCH_CONCURRENCY, - async ({ webhook: foundWebhook, workflow: foundWorkflow }) => { + async ({ webhook: foundWebhook, workflow: foundWorkflow, triggerBlockDeployed }) => { const result = await dispatchResolvedWebhookTarget( foundWebhook, foundWorkflow, @@ -119,6 +119,7 @@ export async function dispatchSlackWebhooks( receivedAt, triggerTimestampMs, subject, + triggerBlockDeployed, } ) diff --git a/apps/sim/lib/webhooks/trigger-block-deployment.integration.ts b/apps/sim/lib/webhooks/trigger-block-deployment.integration.ts new file mode 100644 index 00000000000..4b94dd6a7ca --- /dev/null +++ b/apps/sim/lib/webhooks/trigger-block-deployment.integration.ts @@ -0,0 +1,148 @@ +/** + * Webhook delivery's deployment reads against real PostgreSQL: the lookup's + * trigger-block answer must match the standalone check, and a cached deployment + * version is only ever served to the workflow it belongs to. + */ +import { db } from '@sim/db' +import { user, webhook, workflow, workflowDeploymentVersion, workspace } from '@sim/db/schema' +import { generateId } from '@sim/utils/id' +import { eq } from 'drizzle-orm' +import { afterAll, beforeAll, describe, expect, it } from 'vitest' +import { findAllWebhooksForPath } from '@/lib/webhooks/processor' +import { + blockExistsInDeployment, + loadWorkflowDeploymentVersionState, +} from '@/lib/workflows/persistence/utils' + +const owner = `trigger-block-owner-${generateId()}` +const workspaceId = generateId() +const deployedWorkflow = generateId() +const undeployedWorkflow = generateId() +const deploymentVersion = generateId() +const paths = { + deployedBlock: `deployed-${generateId()}`, + missingBlock: `missing-${generateId()}`, + noDeployment: `legacy-${generateId()}`, +} + +beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: owner, + name: 'Trigger Block', + email: `${owner}@trigger-block.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db + .insert(workspace) + .values({ id: workspaceId, name: 'Trigger Block', ownerId: owner, billedAccountUserId: owner }) + await db.insert(workflow).values( + [deployedWorkflow, undeployedWorkflow].map((id) => ({ + id, + userId: owner, + workspaceId, + name: id, + lastSynced: now, + createdAt: now, + updatedAt: now, + isDeployed: id === deployedWorkflow, + })) + ) + await db.insert(workflowDeploymentVersion).values({ + id: deploymentVersion, + workflowId: deployedWorkflow, + version: 1, + isActive: true, + state: { + blocks: { + trigger: { + id: 'trigger', + type: 'generic_webhook', + name: 'Webhook', + position: { x: 0, y: 0 }, + subBlocks: {}, + outputs: {}, + enabled: true, + }, + }, + edges: [], + loops: {}, + parallels: {}, + }, + }) + await db.insert(webhook).values([ + { + id: generateId(), + workflowId: deployedWorkflow, + deploymentVersionId: deploymentVersion, + blockId: 'trigger', + path: paths.deployedBlock, + provider: 'generic', + }, + { + id: generateId(), + workflowId: deployedWorkflow, + deploymentVersionId: deploymentVersion, + blockId: 'removed-trigger', + path: paths.missingBlock, + provider: 'generic', + }, + { + id: generateId(), + workflowId: undeployedWorkflow, + blockId: 'trigger', + path: paths.noDeployment, + provider: 'generic', + }, + ]) +}) + +afterAll(async () => { + await db.delete(workspace).where(eq(workspace.id, workspaceId)) + await db.delete(user).where(eq(user.id, owner)) +}) + +async function lookup(path: string) { + const targets = await findAllWebhooksForPath({ requestId: 'trigger-block-test', path }) + expect(targets).toHaveLength(1) + return targets[0] +} + +describe('trigger block deployment', () => { + it('reports a trigger block present in the active deployment', async () => { + const target = await lookup(paths.deployedBlock) + + expect(target.triggerBlockDeployed).toBe(true) + expect(await blockExistsInDeployment(deployedWorkflow, 'trigger')).toBe(true) + }) + + it('reports a trigger block the active deployment no longer contains', async () => { + const target = await lookup(paths.missingBlock) + + expect(target.triggerBlockDeployed).toBe(false) + expect(await blockExistsInDeployment(deployedWorkflow, 'removed-trigger')).toBe(false) + }) + + it('reports no trigger block for a workflow without an active deployment', async () => { + const target = await lookup(paths.noDeployment) + + expect(target.triggerBlockDeployed).toBe(false) + expect(await blockExistsInDeployment(undeployedWorkflow, 'trigger')).toBe(false) + }) + + it('serves a deployment version only to the workflow it belongs to', async () => { + const state = await loadWorkflowDeploymentVersionState( + deployedWorkflow, + deploymentVersion, + workspaceId + ) + expect(Object.keys(state.blocks)).toEqual(['trigger']) + expect(state.deploymentVersionId).toBe(deploymentVersion) + + await expect( + loadWorkflowDeploymentVersionState(undeployedWorkflow, deploymentVersion, workspaceId) + ).rejects.toThrow(/was not found/) + }) +}) diff --git a/apps/sim/lib/workflows/custom-blocks/operations.ts b/apps/sim/lib/workflows/custom-blocks/operations.ts index e6dedafaefd..ffafec56483 100644 --- a/apps/sim/lib/workflows/custom-blocks/operations.ts +++ b/apps/sim/lib/workflows/custom-blocks/operations.ts @@ -8,6 +8,7 @@ import { } from '@sim/db/schema' import { createLogger } from '@sim/logger' import { generateId, generateShortId } from '@sim/utils/id' +import { omit } from '@sim/utils/object' import { and, eq, isNull, ne, sql } from 'drizzle-orm' import { isOrganizationFeatureEntitled } from '@/lib/billing/core/subscription' import { acquireOrganizationMutationLock } from '@/lib/billing/organizations/membership' @@ -126,6 +127,29 @@ function applyInputPlaceholders( }) } +const CUSTOM_BLOCK_ROW_COLUMNS = { + type: customBlock.type, + name: customBlock.name, + description: customBlock.description, + workflowId: customBlock.workflowId, + outputs: customBlock.outputs, + enabled: customBlock.enabled, +} + +function toCustomBlockRow({ + outputs, + ...row +}: { + type: string + name: string + description: string + workflowId: string + outputs: CustomBlockOutput[] | null + enabled: boolean +}): CustomBlockRow & { enabled: boolean } { + return { ...row, exposedOutputs: outputs ?? [] } +} + /** * The org's custom blocks for the server overlay (`withCustomBlockOverlay`). * Includes DISABLED rows (carrying `enabled`) so a still-placed disabled block @@ -139,31 +163,33 @@ export async function getCustomBlockRowsForOrg( organizationId: string ): Promise> { const rows = await db - .select({ - type: customBlock.type, - name: customBlock.name, - description: customBlock.description, - workflowId: customBlock.workflowId, - outputs: customBlock.outputs, - enabled: customBlock.enabled, - }) + .select(CUSTOM_BLOCK_ROW_COLUMNS) .from(customBlock) .where(eq(customBlock.organizationId, organizationId)) - return rows.map(({ outputs, ...r }) => ({ ...r, exposedOutputs: outputs ?? [] })) + return rows.map(toCustomBlockRow) } /** * The custom-block rows in scope for a workspace's organization, for wrapping an * execution in `withCustomBlockOverlay`. Returns `[]` when the workspace has no - * organization (nothing to resolve). + * organization or the organization is not entitled to custom blocks. + * + * Every execution starts here, so the rows are read first, joined through the + * workspace: most organizations have none, which settles the answer in one query + * without the entitlement check. */ export async function getCustomBlockRowsForWorkspace( workspaceId: string ): Promise { - const organizationId = await eligibleOrgForWorkspace(workspaceId) - if (!organizationId) return [] - return getCustomBlockRowsForOrg(organizationId) + const rows = await db + .select({ ...CUSTOM_BLOCK_ROW_COLUMNS, organizationId: customBlock.organizationId }) + .from(customBlock) + .innerJoin(workspace, eq(workspace.organizationId, customBlock.organizationId)) + .where(eq(workspace.id, workspaceId)) + if (rows.length === 0) return [] + if (!(await isCustomBlocksEligibleForOrganization(rows[0].organizationId))) return [] + return rows.map((row) => toCustomBlockRow(omit(row, ['organizationId']))) } /** diff --git a/apps/sim/lib/workflows/custom-blocks/workspace-rows.integration.ts b/apps/sim/lib/workflows/custom-blocks/workspace-rows.integration.ts new file mode 100644 index 00000000000..a40c76c63f2 --- /dev/null +++ b/apps/sim/lib/workflows/custom-blocks/workspace-rows.integration.ts @@ -0,0 +1,125 @@ +/** The custom-block rows an execution overlays, read against real PostgreSQL. */ +import { db } from '@sim/db' +import { customBlock, organization, user, workflow, workspace } from '@sim/db/schema' +import { + billingSubscriptionMock, + billingSubscriptionMockFns, +} from '@sim/testing/mocks/billing-subscription.mock' +import { generateId } from '@sim/utils/id' +import { eq, inArray } from 'drizzle-orm' +import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest' + +vi.mock('@/lib/billing/core/subscription', () => billingSubscriptionMock) + +import { getCustomBlockRowsForWorkspace } from '@/lib/workflows/custom-blocks/operations' + +const mockIsOrganizationFeatureEntitled = + billingSubscriptionMockFns.mockIsOrganizationFeatureEntitled +const owner = `custom-rows-owner-${generateId()}` +const organizations = { withBlocks: `org-${generateId()}`, withoutBlocks: `org-${generateId()}` } +const workspaces = { + withBlocks: generateId(), + withoutBlocks: generateId(), + personal: generateId(), +} +const sourceWorkflow = generateId() + +beforeAll(async () => { + const now = new Date() + await db.insert(user).values({ + id: owner, + name: 'Custom Rows', + email: `${owner}@custom-rows.test`, + emailVerified: true, + createdAt: now, + updatedAt: now, + }) + await db.insert(organization).values([ + { id: organizations.withBlocks, name: 'With Blocks', slug: organizations.withBlocks }, + { id: organizations.withoutBlocks, name: 'Without Blocks', slug: organizations.withoutBlocks }, + ]) + await db.insert(workspace).values([ + { + id: workspaces.withBlocks, + name: 'With Blocks', + ownerId: owner, + billedAccountUserId: owner, + organizationId: organizations.withBlocks, + }, + { + id: workspaces.withoutBlocks, + name: 'Without Blocks', + ownerId: owner, + billedAccountUserId: owner, + organizationId: organizations.withoutBlocks, + }, + { id: workspaces.personal, name: 'Personal', ownerId: owner, billedAccountUserId: owner }, + ]) + await db.insert(workflow).values({ + id: sourceWorkflow, + userId: owner, + workspaceId: workspaces.withBlocks, + name: 'Source', + lastSynced: now, + createdAt: now, + updatedAt: now, + }) + await db.insert(customBlock).values([ + { + id: generateId(), + organizationId: organizations.withBlocks, + workflowId: sourceWorkflow, + type: `custom_block_${generateId()}`, + name: 'Enabled', + outputs: [{ blockId: 'b1', path: 'result', name: 'result' }], + }, + { + id: generateId(), + organizationId: organizations.withBlocks, + workflowId: sourceWorkflow, + type: `custom_block_${generateId()}`, + name: 'Disabled', + enabled: false, + }, + ]) +}) + +beforeEach(() => { + mockIsOrganizationFeatureEntitled.mockReset() + mockIsOrganizationFeatureEntitled.mockResolvedValue(true) +}) + +afterAll(async () => { + await db.delete(workspace).where(inArray(workspace.id, Object.values(workspaces))) + await db.delete(organization).where(inArray(organization.id, Object.values(organizations))) + await db.delete(user).where(eq(user.id, owner)) +}) + +describe('getCustomBlockRowsForWorkspace', () => { + it('returns every block of an entitled organization, disabled ones included', async () => { + const rows = await getCustomBlockRowsForWorkspace(workspaces.withBlocks) + + expect(rows.map((row) => row.name).sort()).toEqual(['Disabled', 'Enabled']) + expect(rows.find((row) => row.name === 'Enabled')).toMatchObject({ + workflowId: sourceWorkflow, + enabled: true, + exposedOutputs: [{ blockId: 'b1', path: 'result', name: 'result' }], + }) + expect(rows.find((row) => row.name === 'Disabled')).toMatchObject({ + enabled: false, + exposedOutputs: [], + }) + }) + + it('returns nothing for an organization that is not entitled', async () => { + mockIsOrganizationFeatureEntitled.mockResolvedValue(false) + + expect(await getCustomBlockRowsForWorkspace(workspaces.withBlocks)).toEqual([]) + }) + + it('returns nothing for an organization without blocks or a workspace without one', async () => { + expect(await getCustomBlockRowsForWorkspace(workspaces.withoutBlocks)).toEqual([]) + expect(await getCustomBlockRowsForWorkspace(workspaces.personal)).toEqual([]) + expect(await getCustomBlockRowsForWorkspace(generateId())).toEqual([]) + }) +}) diff --git a/apps/sim/lib/workflows/executor/execution-core.ts b/apps/sim/lib/workflows/executor/execution-core.ts index fa337c3cef1..12cbbe5c8f0 100644 --- a/apps/sim/lib/workflows/executor/execution-core.ts +++ b/apps/sim/lib/workflows/executor/execution-core.ts @@ -22,7 +22,10 @@ import { import { isOutboundRoutingEnabled } from '@/lib/core/network/config.server' import { runWithOutboundOrganization } from '@/lib/core/network/context.server' import { withDatabaseReadRetry } from '@/lib/db/read-retry' -import { getExecutionEnvironment } from '@/lib/environment/utils' +import { + type EnvironmentResolutionSnapshot, + getExecutionEnvironment, +} from '@/lib/environment/utils' import { clearExecutionCancellation } from '@/lib/execution/cancellation' import { connectExecutionSignalHub } from '@/lib/execution/execution-signal' import { getStoredFileReferenceScope, processInputFileFields } from '@/lib/execution/files' @@ -114,6 +117,14 @@ function describeErrorCause(error: unknown): Record | undefined } } +/** An execution environment together with the identities it was resolved for. */ +export interface PreloadedExecutionEnvironment { + personalUserId: string | undefined + workspaceUserId: string + workspaceId: string + snapshot: EnvironmentResolutionSnapshot +} + export interface ExecuteWorkflowCoreOptions { snapshot: ExecutionSnapshot callbacks: ExecutionCallbacks @@ -129,6 +140,12 @@ export interface ExecuteWorkflowCoreOptions { trustedInitialResolvedSecretTraceProvenance?: ResolvedSecretTraceProvenanceV1 /** Immutable deployment admitted by the durable parent log for a resumed execution. */ resumeDeploymentVersionId?: string + /** + * Environment the caller already resolved for this run, reused instead of loading + * and decrypting it again. Used only when it was resolved for exactly the + * identities and workspace this run resolves its environment for. + */ + preloadedEnvironment?: PreloadedExecutionEnvironment /** Run-from-block mode: execute starting from a specific block using cached upstream outputs */ runFromBlock?: { startBlockId: string @@ -432,6 +449,126 @@ async function finalizeExecutionError(params: { return finalized } +interface ExecutionEnvironmentIdentities { + /** Undefined for an anonymous public-API run, which lends no personal namespace. */ + personalEnvUserId: string | undefined + workspaceEnvUserId: string +} + +/** Whose personal and workspace variables this run resolves; throws on incomplete metadata. */ +function resolveExecutionEnvironmentIdentities( + metadata: ExecutionSnapshot['metadata'] +): ExecutionEnvironmentIdentities { + /** + * Personal variables belong to whoever is running, whenever that is knowable. + * `enforceCredentialAccess` is the principal layer's own answer to "is there + * an identifiable caller": it is set from `principal.kind !== 'workspace_api_key'`, + * so a session, personal API key, or delegated run reads its own personal + * variables rather than borrowing the workflow owner's. + * + * The workflow owner remains the fallback for a workspace API key, schedule, + * or webhook. Someone in the workspace configured each of those, and a + * deployed workflow is routinely authored against its owner's personal keys. + * + * An anonymous public-API run resolves no personal variables at all. Anyone + * can call that endpoint, so there is no caller to read as and no person whose + * private namespace it would be reasonable to lend — such a workflow runs on + * workspace secrets alone. + */ + const identifiedCallerUserId = + (metadata.isClientSession && metadata.sessionUserId) || + (metadata.enforceCredentialAccess ? metadata.userId : undefined) + + const personalEnvUserId = metadata.isPublicApiAccess + ? undefined + : identifiedCallerUserId || metadata.workflowUserId + + if (!metadata.isPublicApiAccess && !personalEnvUserId) { + throw new Error('Missing workflowUserId in execution metadata') + } + + /** + * The actor already carries the identity each trigger kind should authorize + * workspace secrets against: the caller for a session, personal API key, or + * delegated principal, and the workspace billing account for a workspace API + * key, schedule, webhook, or anonymous public-API call, where no caller is + * identifiable. Deriving it again here would only risk disagreeing with the + * principal layer. + */ + const workspaceEnvUserId = metadata.userId || personalEnvUserId + if (!workspaceEnvUserId) { + throw new Error('Missing execution actor in execution metadata') + } + + return { personalEnvUserId, workspaceEnvUserId } +} + +function readPiiPolicyRow(workspaceId: string) { + return db + .select({ orgSettings: organization.dataRetentionSettings }) + .from(workspace) + .leftJoin(organization, eq(organization.id, workspace.organizationId)) + .where(eq(workspace.id, workspaceId)) + .limit(1) + .then(([row]) => row) +} + +/** + * Reads that depend only on who the run is and which workspace it runs in, so they + * can start before the custom-block overlay is resolved. + */ +interface ExecutionReads { + environment: Promise + /** + * The org/workspace PII redaction policy, resolved once for both the input stage + * and the block-outputs stage. Stored rules are the source of truth; absence + * yields the disabled default. + */ + piiPolicyRow: Promise>> +} + +function startExecutionReads( + options: ExecuteWorkflowCoreOptions, + workspaceId: string, + { personalEnvUserId, workspaceEnvUserId }: ExecutionEnvironmentIdentities +): ExecutionReads { + const { preloadedEnvironment } = options + const environment = + preloadedEnvironment && + preloadedEnvironment.personalUserId === personalEnvUserId && + preloadedEnvironment.workspaceUserId === workspaceEnvUserId && + preloadedEnvironment.workspaceId === workspaceId + ? Promise.resolve(preloadedEnvironment.snapshot) + : withDatabaseReadRetry( + () => getExecutionEnvironment(personalEnvUserId, workspaceEnvUserId, workspaceId), + { label: 'getExecutionEnvironment' } + ) + const piiPolicyRow = withDatabaseReadRetry(() => readPiiPolicyRow(workspaceId), { + label: 'resolvePiiRedactionPolicy', + }) + // Awaited later by the run; a run that fails first must not leave them unhandled. + environment.catch(() => {}) + piiPolicyRow.catch(() => {}) + return { environment, piiPolicyRow } +} + +/** + * Starts the run's identity-scoped reads ahead of the overlay. Metadata that cannot + * name those identities starts nothing here: the run itself rejects it, inside its + * own error handling. + */ +function prefetchExecutionReads(options: ExecuteWorkflowCoreOptions): ExecutionReads | undefined { + const { metadata } = options.snapshot + if (!metadata.workspaceId) return undefined + let identities: ExecutionEnvironmentIdentities + try { + identities = resolveExecutionEnvironmentIdentities(metadata) + } catch { + return undefined + } + return startExecutionReads(options, metadata.workspaceId, identities) +} + /** * Establish the custom-block registry overlay for the execution's organization, * then run the core. Wrapping here — the shared choke point for the sync route and @@ -452,22 +589,29 @@ export async function executeWorkflowCore( ): Promise { connectExecutionSignalHub() const workspaceId = options.snapshot.metadata.workspaceId - const rows = workspaceId - ? await withDatabaseReadRetry(() => getCustomBlockRowsForWorkspace(workspaceId), { - label: 'getCustomBlockRowsForWorkspace', - }) - : [] - const execute = () => withCustomBlockOverlay(rows, () => executeWorkflowCoreImpl(options)) - if (!isOutboundRoutingEnabled()) return execute() - const context = await resolveActiveWorkflowApplicationContext({ - workflowId: options.snapshot.metadata.workflowId, - assertedWorkspaceId: workspaceId, - }) - return runWithOutboundOrganization(context.workspaceOrganizationId, execute) + const prefetchedReads = prefetchExecutionReads(options) + const [rows, outboundContext] = await Promise.all([ + workspaceId + ? withDatabaseReadRetry(() => getCustomBlockRowsForWorkspace(workspaceId), { + label: 'getCustomBlockRowsForWorkspace', + }) + : [], + isOutboundRoutingEnabled() + ? resolveActiveWorkflowApplicationContext({ + workflowId: options.snapshot.metadata.workflowId, + assertedWorkspaceId: workspaceId, + }) + : undefined, + ]) + const execute = () => + withCustomBlockOverlay(rows, () => executeWorkflowCoreImpl(options, prefetchedReads)) + if (!outboundContext) return execute() + return runWithOutboundOrganization(outboundContext.workspaceOrganizationId, execute) } async function executeWorkflowCoreImpl( - options: ExecuteWorkflowCoreOptions + options: ExecuteWorkflowCoreOptions, + prefetchedReads: ExecutionReads | undefined ): Promise { const { snapshot, @@ -529,46 +673,9 @@ async function executeWorkflowCoreImpl( } try { - /** - * Personal variables belong to whoever is running, whenever that is knowable. - * `enforceCredentialAccess` is the principal layer's own answer to "is there - * an identifiable caller": it is set from `principal.kind !== 'workspace_api_key'`, - * so a session, personal API key, or delegated run reads its own personal - * variables rather than borrowing the workflow owner's. - * - * The workflow owner remains the fallback for a workspace API key, schedule, - * or webhook. Someone in the workspace configured each of those, and a - * deployed workflow is routinely authored against its owner's personal keys. - * - * An anonymous public-API run resolves no personal variables at all. Anyone - * can call that endpoint, so there is no caller to read as and no person whose - * private namespace it would be reasonable to lend — such a workflow runs on - * workspace secrets alone. - */ - const identifiedCallerUserId = - (metadata.isClientSession && metadata.sessionUserId) || - (metadata.enforceCredentialAccess ? metadata.userId : undefined) - - const personalEnvUserId = metadata.isPublicApiAccess - ? undefined - : identifiedCallerUserId || metadata.workflowUserId - - if (!metadata.isPublicApiAccess && !personalEnvUserId) { - throw new Error('Missing workflowUserId in execution metadata') - } - - /** - * The actor already carries the identity each trigger kind should authorize - * workspace secrets against: the caller for a session, personal API key, or - * delegated principal, and the workspace billing account for a workspace API - * key, schedule, webhook, or anonymous public-API call, where no caller is - * identifiable. Deriving it again here would only risk disagreeing with the - * principal layer. - */ - const workspaceEnvUserId = metadata.userId || personalEnvUserId - if (!workspaceEnvUserId) { - throw new Error('Missing execution actor in execution metadata') - } + const identities = resolveExecutionEnvironmentIdentities(metadata) + const reads = prefetchedReads ?? startExecutionReads(options, providedWorkspaceId, identities) + const { personalEnvUserId, workspaceEnvUserId } = identities /** * Resolves the workflow state from the override, the draft tables, or the @@ -640,12 +747,10 @@ async function executeWorkflowCoreImpl( } } - const [workflowState, env] = await Promise.all([ + const [workflowState, env, piiPolicyRow] = await Promise.all([ withDatabaseReadRetry(loadWorkflowState, { label: 'loadWorkflowState' }), - withDatabaseReadRetry( - () => getExecutionEnvironment(personalEnvUserId, workspaceEnvUserId, providedWorkspaceId), - { label: 'getExecutionEnvironment' } - ), + reads.environment, + reads.piiPolicyRow, ]) const { blocks, loops, parallels } = workflowState @@ -954,22 +1059,8 @@ async function executeWorkflowCoreImpl( allowLargeValueWorkflowScope, }) - // Resolve the org/workspace PII redaction policy once; serves both the input - // stage (below) and the block-outputs stage (threaded into the executor). - // Stored rules are the source of truth; absence yields the disabled default - // with one indexed lookup and no masking cost for non-PII organizations. - const [row] = await withDatabaseReadRetry( - () => - db - .select({ orgSettings: organization.dataRetentionSettings }) - .from(workspace) - .leftJoin(organization, eq(organization.id, workspace.organizationId)) - .where(eq(workspace.id, providedWorkspaceId)) - .limit(1), - { label: 'resolvePiiRedactionPolicy' } - ) const piiRedaction: EffectivePiiRedaction = resolveEffectivePiiRedaction({ - orgSettings: row?.orgSettings, + orgSettings: piiPolicyRow?.orgSettings, workspaceId: providedWorkspaceId, }) diff --git a/apps/sim/lib/workflows/persistence/utils.ts b/apps/sim/lib/workflows/persistence/utils.ts index aeece86a6b7..b3d9fc0d979 100644 --- a/apps/sim/lib/workflows/persistence/utils.ts +++ b/apps/sim/lib/workflows/persistence/utils.ts @@ -99,13 +99,19 @@ export interface DeployedWorkflowData extends NormalizedWorkflowData { variables?: Record } +/** + * Whether the active deployment of `workflowId` contains `blockId`. Answered by + * the database so the (often hundreds of KB) deployed state never leaves it. + */ export async function blockExistsInDeployment( workflowId: string, blockId: string ): Promise { try { const [result] = await db - .select({ state: workflowDeploymentVersion.state }) + .select({ + exists: sql`json_typeof(${workflowDeploymentVersion.state} -> 'blocks' -> ${blockId}) = 'object'`, + }) .from(workflowDeploymentVersion) .where( and( @@ -115,12 +121,7 @@ export async function blockExistsInDeployment( ) .limit(1) - if (!result?.state) { - return false - } - - const state = result.state as WorkflowState - return !!state.blocks?.[blockId] + return result?.exists === true } catch (error) { logger.error(`Error checking block ${blockId} in deployment for workflow ${workflowId}:`, error) return false @@ -136,11 +137,22 @@ const DEPLOYED_STATE_CACHE_TTL_MS = 5 * 60 * 1000 * absolute on purpose — it bounds the one non-immutable part, the live credential * remap in `applyBlockMigrations` — so credential changes still propagate. */ -const deployedStateCache = new LRUCache({ +const deployedStateCache = new LRUCache< + string, + { workflowId: string; state: DeployedWorkflowData } +>({ max: DEPLOYED_STATE_CACHE_MAX_ENTRIES, ttl: DEPLOYED_STATE_CACHE_TTL_MS, }) +function getCachedDeploymentState( + workflowId: string, + deploymentVersionId: string +): DeployedWorkflowData | undefined { + const cached = deployedStateCache.get(deploymentVersionId) + return cached?.workflowId === workflowId ? structuredClone(cached.state) : undefined +} + /** Evicts one deployed-state entry, or clears the cache when no id is given. */ export function invalidateDeployedStateCache(deploymentVersionId?: string): void { if (deploymentVersionId) { @@ -196,10 +208,9 @@ export async function materializeDeploymentState( executor?: DbOrTx, options: { cache?: boolean } = {} ): Promise { - const cached = options.cache === false ? undefined : deployedStateCache.get(version.id) - if (cached) { - return structuredClone(cached) - } + const cached = + options.cache === false ? undefined : getCachedDeploymentState(workflowId, version.id) + if (cached) return cached const state = version.state as WorkflowState & { variables?: Record } @@ -247,7 +258,9 @@ export async function materializeDeploymentState( deploymentVersionId: version.id, } - if (options.cache !== false) deployedStateCache.set(version.id, deployedState) + if (options.cache !== false) { + deployedStateCache.set(version.id, { workflowId, state: deployedState }) + } return structuredClone(deployedState) } @@ -289,12 +302,17 @@ export async function loadDeployedWorkflowState( /** * Loads an immutable deployment snapshot by ID for work admitted before a later cutover. + * A cached materialization of this workflow's version (the same entry + * {@link materializeDeploymentState} serves) is returned without reading the row again. */ export async function loadWorkflowDeploymentVersionState( workflowId: string, deploymentVersionId: string, providedWorkspaceId?: string ): Promise { + const cached = getCachedDeploymentState(workflowId, deploymentVersionId) + if (cached) return cached + const [version] = await db .select({ id: workflowDeploymentVersion.id,