Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion apps/sim/app/api/webhooks/tiktok/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 6 additions & 1 deletion apps/sim/app/api/webhooks/trigger/[path]/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -321,6 +325,7 @@ async function handleWebhookDelivery(
path,
receivedAt,
triggerTimestampMs: Number.isFinite(triggerTimestampMs) ? triggerTimestampMs : undefined,
triggerBlockDeployed,
}
)

Expand Down
3 changes: 2 additions & 1 deletion apps/sim/background/quickbooks-webhook-ingress.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
28 changes: 23 additions & 5 deletions apps/sim/background/webhook-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down Expand Up @@ -665,19 +666,21 @@ export async function resolveWebhookExecutionProviderConfig<
options?: WebhookEnvResolutionOptions & {
onEnvironmentSnapshot?: (snapshot: EnvironmentResolutionSnapshot) => void | Promise<void>
actorUserId?: string
/** The same environment load, already started by a caller that had its identities early. */
environment?: Promise<EnvironmentResolutionSnapshot>
}
): Promise<T & { providerConfig: Record<string, unknown> }> {
try {
if (!options) {
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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -1151,6 +1168,7 @@ async function executeWebhookJobInternal(
loggingSession,
trustedInitialResolvedSecretTraceProvenance:
resolvedSecretTraceRegistry.exportProvenanceForValue(triggerInput),
preloadedEnvironment,
includeFileBase64: false,
base64MaxBytes: undefined,
abortSignal: timeoutController.signal,
Expand Down
1 change: 1 addition & 0 deletions apps/sim/background/workflow-column-execution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -866,6 +866,7 @@ async function runWorkflowAndWriteTerminal(
triggerType: 'workflow',
checkDeployment: false,
checkRateLimit: false,
includeActorSubscription: true,
skipConcurrencyReservation: true,
logPreprocessingErrors: false,
billingAttribution,
Expand Down
44 changes: 34 additions & 10 deletions apps/sim/executor/handlers/workflow/workflow-handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ authInternalMockFns.mockGenerateInternalToken.mockResolvedValue('test-token')

const {
mockExecutorExecute,
mockCreateSnapshot,
mockResolveSnapshot,
mockAdmitCustomBlockChildExecution,
mockTrackChildRun,
mockBuildTraceSpans,
Expand All @@ -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(),
Expand Down Expand Up @@ -160,7 +160,7 @@ afterAll(() => {
})

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

vi.mock('@/lib/auth/internal', () => authInternalMock)
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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, {})
Expand Down Expand Up @@ -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, {})
Expand Down Expand Up @@ -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, {})
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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' } })
})

Expand Down
8 changes: 3 additions & 5 deletions apps/sim/executor/handlers/workflow/workflow-handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
109 changes: 109 additions & 0 deletions apps/sim/lib/billing/core/plan.integration.ts
Original file line number Diff line number Diff line change
@@ -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<string> {
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<string> {
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<string> {
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)
})
})
Loading
Loading